@@ -215,154 +215,175 @@ export const layerWith = (options?: LayerOptions) =>
215215 ) {
216216 return Effect . gen ( function * ( ) {
217217 const durable = definition ?. durable
218- if ( durable ) {
219- const aggregateID = ( event . data as Record < string , unknown > ) [ durable . aggregate ]
220- if ( typeof aggregateID !== "string" ) {
221- yield * Effect . die (
222- new InvalidDurableEventError ( {
223- type : event . type ,
224- message : `Expected string aggregate field ${ durable . aggregate } ` ,
225- } ) ,
226- )
227- } else {
228- if ( input && input . aggregateID !== aggregateID ) {
229- yield * Effect . die (
230- new InvalidDurableEventError ( {
231- type : event . type ,
232- message : `Aggregate mismatch: expected ${ input . aggregateID } , got ${ aggregateID } ` ,
233- } ) ,
218+ if ( ! durable ) return
219+ const aggregateID = ( event . data as Record < string , unknown > ) [ durable . aggregate ]
220+ if ( typeof aggregateID !== "string" )
221+ return yield * Effect . die (
222+ new InvalidDurableEventError ( {
223+ type : event . type ,
224+ message : `Expected string aggregate field ${ durable . aggregate } ` ,
225+ } ) ,
226+ )
227+ if ( input && input . aggregateID !== aggregateID )
228+ return yield * Effect . die (
229+ new InvalidDurableEventError ( {
230+ type : event . type ,
231+ message : `Aggregate mismatch: expected ${ input . aggregateID } , got ${ aggregateID } ` ,
232+ } ) ,
233+ )
234+ const list = projectors . get ( event . type ) ?? [ ]
235+ return yield * Effect . uninterruptible (
236+ Effect . gen ( function * ( ) {
237+ const committed = yield * db
238+ . transaction ( ( ) => commitTransaction ( definition , durable , event , aggregateID , list , input , commit ) , {
239+ behavior : "immediate" ,
240+ } )
241+ . pipe ( Effect . orDie )
242+ if ( committed ) {
243+ yield * Effect . forEach (
244+ pubsub . durable . get ( committed . aggregateID ) ?? [ ] ,
245+ ( wake ) => PubSub . publish ( wake , undefined ) ,
246+ { discard : true } ,
234247 )
235248 }
236- const list = projectors . get ( event . type ) ?? [ ]
237- return yield * Effect . uninterruptible (
238- Effect . gen ( function * ( ) {
239- const committed = yield * db
240- . transaction (
241- ( ) =>
242- Effect . gen ( function * ( ) {
243- const row = yield * db
244- . select ( { seq : EventSequenceTable . seq , ownerID : EventSequenceTable . owner_id } )
245- . from ( EventSequenceTable )
246- . where ( eq ( EventSequenceTable . aggregate_id , aggregateID ) )
247- . get ( )
248- . pipe ( Effect . orDie )
249- const latest = row ?. seq ?? - 1
250- const encoded = Schema . encodeUnknownSync ( definition . data ) ( event . data ) as Record <
251- string ,
252- unknown
253- >
254- if ( input ?. strictOwner && row ?. ownerID && row . ownerID !== input . ownerID ) {
255- yield * Effect . die (
256- new InvalidDurableEventError ( {
257- type : event . type ,
258- message : `Replay owner mismatch for aggregate ${ aggregateID } : expected ${ row . ownerID } , got ${ input . ownerID ?? "none" } ` ,
259- } ) ,
260- )
261- }
262- if ( input && input . seq <= latest ) {
263- const stored = yield * db
264- . select ( )
265- . from ( EventTable )
266- . where ( and ( eq ( EventTable . aggregate_id , aggregateID ) , eq ( EventTable . seq , input . seq ) ) )
267- . get ( )
268- . pipe ( Effect . orDie )
269- if (
270- stored ?. id === event . id &&
271- stored . type === versionedType ( definition . type , durable . version ) &&
272- isDeepStrictEqual ( stored . data , encoded )
273- ) {
274- if ( input . ownerID && row ?. ownerID == null ) {
275- yield * db
276- . update ( EventSequenceTable )
277- . set ( { owner_id : input . ownerID } )
278- . where ( eq ( EventSequenceTable . aggregate_id , aggregateID ) )
279- . run ( )
280- . pipe ( Effect . orDie )
281- }
282- return
283- }
284- yield * Effect . die (
285- new InvalidDurableEventError ( {
286- type : event . type ,
287- message : `Replay diverged at aggregate ${ aggregateID } sequence ${ input . seq } ` ,
288- } ) ,
289- )
290- }
291- if ( input && row ?. ownerID && row . ownerID !== input . ownerID ) {
292- return
293- }
294- const seq = input ?. seq ?? latest + 1
295- if ( input && seq !== latest + 1 ) {
296- yield * Effect . die (
297- new InvalidDurableEventError ( {
298- type : event . type ,
299- message : `Sequence mismatch for aggregate ${ aggregateID } : expected ${ latest + 1 } , got ${ seq } ` ,
300- } ) ,
301- )
302- }
303- const stored = yield * db
304- . select ( { aggregateID : EventTable . aggregate_id , seq : EventTable . seq } )
305- . from ( EventTable )
306- . where ( eq ( EventTable . id , event . id ) )
307- . get ( )
308- . pipe ( Effect . orDie )
309- if ( stored )
310- yield * Effect . die (
311- new InvalidDurableEventError ( {
312- type : event . type ,
313- message : `Event ${ event . id } already exists at aggregate ${ stored . aggregateID } sequence ${ stored . seq } ` ,
314- } ) ,
315- )
316- const committed = {
317- ...event ,
318- durable : { aggregateID, seq, version : durable . version } ,
319- } as Payload
320- for ( const projector of list ) {
321- yield * projector ( committed )
322- }
323- if ( commit ) yield * commit ( seq )
324- yield * db
325- . insert ( EventSequenceTable )
326- . values ( [ { aggregate_id : aggregateID , seq, owner_id : input ?. ownerID } ] )
327- . onConflictDoUpdate ( {
328- target : EventSequenceTable . aggregate_id ,
329- set : {
330- seq,
331- ...( input ?. ownerID && row ?. ownerID == null ? { owner_id : input . ownerID } : { } ) ,
332- } ,
333- } )
334- . run ( )
335- . pipe ( Effect . orDie )
336- yield * db
337- . insert ( EventTable )
338- . values ( [
339- {
340- id : event . id ,
341- aggregate_id : aggregateID ,
342- seq,
343- type : versionedType ( definition . type , durable . version ) ,
344- data : encoded ,
345- } ,
346- ] )
347- . run ( )
348- . pipe ( Effect . orDie )
349- return { aggregateID, seq }
350- } ) ,
351- { behavior : "immediate" } ,
352- )
353- . pipe ( Effect . orDie )
354- if ( committed ) {
355- yield * Effect . forEach (
356- pubsub . durable . get ( committed . aggregateID ) ?? [ ] ,
357- ( wake ) => PubSub . publish ( wake , undefined ) ,
358- { discard : true } ,
359- )
360- }
361- return committed
362- } ) ,
363- )
249+ return committed
250+ } ) ,
251+ )
252+ } )
253+ }
254+
255+ function commitTransaction (
256+ definition : Definition ,
257+ durable : NonNullable < Definition [ "durable" ] > ,
258+ event : Payload ,
259+ aggregateID : string ,
260+ list : Subscriber [ ] ,
261+ input ?: {
262+ readonly seq : number
263+ readonly aggregateID : string
264+ readonly ownerID ?: string
265+ readonly strictOwner ?: boolean
266+ } ,
267+ commit ?: ( seq : number ) => Effect . Effect < void > ,
268+ ) {
269+ return Effect . gen ( function * ( ) {
270+ const row = yield * db
271+ . select ( { seq : EventSequenceTable . seq , ownerID : EventSequenceTable . owner_id } )
272+ . from ( EventSequenceTable )
273+ . where ( eq ( EventSequenceTable . aggregate_id , aggregateID ) )
274+ . get ( )
275+ . pipe ( Effect . orDie )
276+ const latest = row ?. seq ?? - 1
277+ const encoded = Schema . encodeUnknownSync ( definition . data ) ( event . data ) as Record < string , unknown >
278+ if ( input ?. strictOwner && row ?. ownerID && row . ownerID !== input . ownerID ) {
279+ yield * Effect . die (
280+ new InvalidDurableEventError ( {
281+ type : event . type ,
282+ message : `Replay owner mismatch for aggregate ${ aggregateID } : expected ${ row . ownerID } , got ${ input . ownerID ?? "none" } ` ,
283+ } ) ,
284+ )
285+ }
286+ if ( input && input . seq <= latest )
287+ return yield * reconcileReplay ( definition , durable , event , aggregateID , input , encoded , row ?. ownerID )
288+ if ( input && row ?. ownerID && row . ownerID !== input . ownerID ) {
289+ return
290+ }
291+ const seq = input ?. seq ?? latest + 1
292+ if ( input && seq !== latest + 1 ) {
293+ yield * Effect . die (
294+ new InvalidDurableEventError ( {
295+ type : event . type ,
296+ message : `Sequence mismatch for aggregate ${ aggregateID } : expected ${ latest + 1 } , got ${ seq } ` ,
297+ } ) ,
298+ )
299+ }
300+ const stored = yield * db
301+ . select ( { aggregateID : EventTable . aggregate_id , seq : EventTable . seq } )
302+ . from ( EventTable )
303+ . where ( eq ( EventTable . id , event . id ) )
304+ . get ( )
305+ . pipe ( Effect . orDie )
306+ if ( stored )
307+ yield * Effect . die (
308+ new InvalidDurableEventError ( {
309+ type : event . type ,
310+ message : `Event ${ event . id } already exists at aggregate ${ stored . aggregateID } sequence ${ stored . seq } ` ,
311+ } ) ,
312+ )
313+ const committed = {
314+ ...event ,
315+ durable : { aggregateID, seq, version : durable . version } ,
316+ } as Payload
317+ for ( const projector of list ) {
318+ yield * projector ( committed )
319+ }
320+ if ( commit ) yield * commit ( seq )
321+ yield * db
322+ . insert ( EventSequenceTable )
323+ . values ( [ { aggregate_id : aggregateID , seq, owner_id : input ?. ownerID } ] )
324+ . onConflictDoUpdate ( {
325+ target : EventSequenceTable . aggregate_id ,
326+ set : {
327+ seq,
328+ ...( input ?. ownerID && row ?. ownerID == null ? { owner_id : input . ownerID } : { } ) ,
329+ } ,
330+ } )
331+ . run ( )
332+ . pipe ( Effect . orDie )
333+ yield * db
334+ . insert ( EventTable )
335+ . values ( [
336+ {
337+ id : event . id ,
338+ aggregate_id : aggregateID ,
339+ seq,
340+ type : versionedType ( definition . type , durable . version ) ,
341+ data : encoded ,
342+ } ,
343+ ] )
344+ . run ( )
345+ . pipe ( Effect . orDie )
346+ return { aggregateID, seq }
347+ } )
348+ }
349+
350+ function reconcileReplay (
351+ definition : Definition ,
352+ durable : NonNullable < Definition [ "durable" ] > ,
353+ event : Payload ,
354+ aggregateID : string ,
355+ input : { readonly seq : number ; readonly ownerID ?: string } ,
356+ encoded : Record < string , unknown > ,
357+ ownerID : string | null | undefined ,
358+ ) {
359+ return Effect . gen ( function * ( ) {
360+ const stored = yield * db
361+ . select ( )
362+ . from ( EventTable )
363+ . where ( and ( eq ( EventTable . aggregate_id , aggregateID ) , eq ( EventTable . seq , input . seq ) ) )
364+ . get ( )
365+ . pipe ( Effect . orDie )
366+ if (
367+ stored ?. id === event . id &&
368+ stored . type === versionedType ( definition . type , durable . version ) &&
369+ isDeepStrictEqual ( stored . data , encoded )
370+ ) {
371+ if ( input . ownerID && ownerID == null ) {
372+ yield * db
373+ . update ( EventSequenceTable )
374+ . set ( { owner_id : input . ownerID } )
375+ . where ( eq ( EventSequenceTable . aggregate_id , aggregateID ) )
376+ . run ( )
377+ . pipe ( Effect . orDie )
364378 }
379+ return
365380 }
381+ yield * Effect . die (
382+ new InvalidDurableEventError ( {
383+ type : event . type ,
384+ message : `Replay diverged at aggregate ${ aggregateID } sequence ${ input . seq } ` ,
385+ } ) ,
386+ )
366387 } )
367388 }
368389
0 commit comments