From 087d356d418fbee11f641240af025e96be34ddce Mon Sep 17 00:00:00 2001 From: Dax Raad Date: Mon, 25 May 2026 14:25:23 -0400 Subject: [PATCH] chore(core): format event service --- packages/core/src/event.ts | 30 +++++++++++++++++++++++++----- 1 file changed, 25 insertions(+), 5 deletions(-) diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index 2ff7313d4..777226e87 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -87,7 +87,11 @@ export function define= existing.sync.version) { registry.set(input.type, definition) } - if (input.sync) syncRegistry.set(versionedType(input.type, input.sync.version), definition as Definition & { readonly sync: NonNullable }) + if (input.sync) + syncRegistry.set( + versionedType(input.type, input.sync.version), + definition as Definition & { readonly sync: NonNullable }, + ) return definition as Schema.Schema>>> & Definition> } @@ -110,7 +114,10 @@ export interface Interface { readonly subscribe: (definition: D) => Stream.Stream> readonly all: () => Stream.Stream readonly project: (definition: D, projector: Projector) => Effect.Effect - readonly replay: (event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string }) => Effect.Effect + readonly replay: ( + event: SerializedEvent, + options?: { readonly publish?: boolean; readonly ownerID?: string }, + ) => Effect.Effect readonly replayAll: ( events: SerializedEvent[], options?: { readonly publish?: boolean; readonly ownerID?: string }, @@ -145,7 +152,10 @@ export const layer = Layer.effect( }), ) - function commitSyncEvent(event: Payload, input?: { readonly seq: number; readonly aggregateID: string; readonly ownerID?: string }) { + function commitSyncEvent( + event: Payload, + input?: { readonly seq: number; readonly aggregateID: string; readonly ownerID?: string }, + ) { return Effect.gen(function* () { const definition = registry.get(event.type) const sync = definition?.sync @@ -272,13 +282,23 @@ export const layer = Layer.effect( const source = events[0]?.aggregateID if (!source) return undefined if (events.some((event) => event.aggregateID !== source)) { - yield* Effect.die(new InvalidSyncEventError({ type: events[0]?.type ?? "unknown", message: "Replay events must belong to the same aggregate" })) + yield* Effect.die( + new InvalidSyncEventError({ + type: events[0]?.type ?? "unknown", + message: "Replay events must belong to the same aggregate", + }), + ) } const start = events[0]?.seq ?? 0 for (const [index, event] of events.entries()) { const seq = start + index if (event.seq !== seq) { - yield* Effect.die(new InvalidSyncEventError({ type: event.type, message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}` })) + yield* Effect.die( + new InvalidSyncEventError({ + type: event.type, + message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`, + }), + ) } } for (const event of events) {