From da57b272771e1cbde3079ee9d92dfd9492d5bdc6 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Fri, 28 Aug 2026 14:39:43 -0400 Subject: [PATCH] refactor(core): flatten durable commit validation (#45662) --- packages/core/src/bus.ts | 300 +++++++++++++++++++-------------------- 1 file changed, 147 insertions(+), 153 deletions(-) diff --git a/packages/core/src/bus.ts b/packages/core/src/bus.ts index 6ca5da6bb4d..efd00328f27 100644 --- a/packages/core/src/bus.ts +++ b/packages/core/src/bus.ts @@ -294,162 +294,156 @@ export function configured(options?: Options) { ) { return Effect.gen(function* () { const durable = definition.durable - if (durable) { - const aggregateID = (event.data as Record)[durable.aggregate] - if (typeof aggregateID !== "string") { - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Expected string aggregate field ${durable.aggregate}`, - }), - ) - } else { - if (input && input.aggregateID !== aggregateID) { - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`, - }), - ) - } - const list = projectors.get(versionedType(definition.type, durable.version)) ?? [] - return yield* Effect.uninterruptible( - Effect.gen(function* () { - const committed = yield* db - .transaction( - () => - Effect.gen(function* () { - const row = yield* db - .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id }) - .from(EventSequenceTable) - .where(eq(EventSequenceTable.aggregate_id, aggregateID)) - .get() - .pipe(Effect.orDie) - const latest = row?.seq ?? -1 - const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record< - string, - unknown - > - if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) { - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`, - }), - ) - } - if (input && input.seq <= latest) { - if (!persist) return - const stored = yield* db - .select() - .from(EventTable) - .where(and(eq(EventTable.aggregate_id, aggregateID), eq(EventTable.seq, input.seq))) - .get() - .pipe(Effect.orDie) - if ( - stored?.id === event.id && - stored.type === versionedType(definition.type, durable.version) && - stored.created === (event.created ?? 0) && - isDeepStrictEqual(stored.data, encoded) - ) { - if (input.ownerID && row?.ownerID == null) { - yield* db - .update(EventSequenceTable) - .set({ owner_id: input.ownerID }) - .where(eq(EventSequenceTable.aggregate_id, aggregateID)) - .run() - .pipe(Effect.orDie) - } - return - } - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Replay diverged at aggregate ${aggregateID} sequence ${input.seq}`, - }), - ) - } - if (input && row?.ownerID && row.ownerID !== input.ownerID) { - return - } - const seq = input?.seq ?? latest + 1 - if (input && seq !== latest + 1) { - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`, - }), - ) - } - if (persist) { - const stored = yield* db - .select({ aggregateID: EventTable.aggregate_id, seq: EventTable.seq }) - .from(EventTable) - .where(eq(EventTable.id, event.id)) - .get() - .pipe(Effect.orDie) - if (stored) - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`, - }), - ) - } - const committed = { - ...event, - durable: { aggregateID, seq, version: durable.version }, - } as Event.Payload - const route = yield* prepareRoutes([committed]) - for (const projector of list) { - yield* projector(committed) - } - if (commit) yield* commit(seq) - yield* db - .insert(EventSequenceTable) - .values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }]) - .onConflictDoUpdate({ - target: EventSequenceTable.aggregate_id, - set: { - seq: sql`max(${EventSequenceTable.seq}, ${seq})`, - ...(input?.ownerID && row?.ownerID == null ? { owner_id: input.ownerID } : {}), - }, - }) - .run() - .pipe(Effect.orDie) - if (persist) + if (!durable) return yield* Effect.void + const aggregateID = (event.data as Record)[durable.aggregate] + if (typeof aggregateID !== "string") + return yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Expected string aggregate field ${durable.aggregate}`, + }), + ) + if (input && input.aggregateID !== aggregateID) { + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`, + }), + ) + } + const list = projectors.get(versionedType(definition.type, durable.version)) ?? [] + return yield* Effect.uninterruptible( + Effect.gen(function* () { + const committed = yield* db + .transaction( + () => + Effect.gen(function* () { + const row = yield* db + .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id }) + .from(EventSequenceTable) + .where(eq(EventSequenceTable.aggregate_id, aggregateID)) + .get() + .pipe(Effect.orDie) + const latest = row?.seq ?? -1 + const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record + if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) { + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`, + }), + ) + } + if (input && input.seq <= latest) { + if (!persist) return + const stored = yield* db + .select() + .from(EventTable) + .where(and(eq(EventTable.aggregate_id, aggregateID), eq(EventTable.seq, input.seq))) + .get() + .pipe(Effect.orDie) + if ( + stored?.id === event.id && + stored.type === versionedType(definition.type, durable.version) && + stored.created === (event.created ?? 0) && + isDeepStrictEqual(stored.data, encoded) + ) { + if (input.ownerID && row?.ownerID == null) { yield* db - .insert(EventTable) - .values([ - { - id: event.id, - aggregate_id: aggregateID, - seq, - created: event.created ?? 0, - type: versionedType(definition.type, durable.version), - data: encoded, - }, - ]) + .update(EventSequenceTable) + .set({ owner_id: input.ownerID }) + .where(eq(EventSequenceTable.aggregate_id, aggregateID)) .run() .pipe(Effect.orDie) - return { aggregateID, seq, event: committed, route } - }), - { behavior: "immediate" }, - ) - .pipe(Effect.orDie) - if (committed) { - committed.route() - yield* Effect.forEach( - pubsub.durable.get(committed.aggregateID) ?? [], - (wake) => PubSub.publish(wake, undefined), - { discard: true }, - ) - } - return committed - }), - ) - } - } + } + return + } + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Replay diverged at aggregate ${aggregateID} sequence ${input.seq}`, + }), + ) + } + if (input && row?.ownerID && row.ownerID !== input.ownerID) { + return + } + const seq = input?.seq ?? latest + 1 + if (input && seq !== latest + 1) { + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`, + }), + ) + } + if (persist) { + const stored = yield* db + .select({ aggregateID: EventTable.aggregate_id, seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.id, event.id)) + .get() + .pipe(Effect.orDie) + if (stored) + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`, + }), + ) + } + const committed = { + ...event, + durable: { aggregateID, seq, version: durable.version }, + } as Event.Payload + const route = yield* prepareRoutes([committed]) + for (const projector of list) { + yield* projector(committed) + } + if (commit) yield* commit(seq) + yield* db + .insert(EventSequenceTable) + .values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }]) + .onConflictDoUpdate({ + target: EventSequenceTable.aggregate_id, + set: { + seq: sql`max(${EventSequenceTable.seq}, ${seq})`, + ...(input?.ownerID && row?.ownerID == null ? { owner_id: input.ownerID } : {}), + }, + }) + .run() + .pipe(Effect.orDie) + if (persist) + yield* db + .insert(EventTable) + .values([ + { + id: event.id, + aggregate_id: aggregateID, + seq, + created: event.created ?? 0, + type: versionedType(definition.type, durable.version), + data: encoded, + }, + ]) + .run() + .pipe(Effect.orDie) + return { aggregateID, seq, event: committed, route } + }), + { behavior: "immediate" }, + ) + .pipe(Effect.orDie) + if (committed) { + committed.route() + yield* Effect.forEach( + pubsub.durable.get(committed.aggregateID) ?? [], + (wake) => PubSub.publish(wake, undefined), + { discard: true }, + ) + } + return committed + }), + ) }) }