feat(core): add sync event replay controls

This commit is contained in:
Dax Raad
2026-05-25 00:28:30 -04:00
parent 05dca9f698
commit 06bd34c953
2 changed files with 204 additions and 9 deletions
+102 -9
View File
@@ -37,6 +37,14 @@ export type Payload<D extends Definition = Definition> = {
export type Projector<D extends Definition = Definition> = (event: Payload<D>) => Effect.Effect<void>
type AnyProjector = (event: Payload) => Effect.Effect<void>
export type SerializedEvent = {
readonly id: ID
readonly type: string
readonly seq: number
readonly aggregateID: string
readonly data: Record<string, unknown>
}
export class InvalidSyncEventError extends Schema.TaggedErrorClass<InvalidSyncEventError>()(
"EventV2.InvalidSyncEvent",
{
@@ -50,6 +58,7 @@ export function versionedType(type: string, version: number) {
}
export const registry = new Map<string, Definition>()
const syncRegistry = new Map<string, Definition & { readonly sync: NonNullable<Definition["sync"]> }>()
export function define<const Type extends string, Fields extends Schema.Struct.Fields>(input: {
readonly type: Type
@@ -78,6 +87,7 @@ export function define<const Type extends string, Fields extends Schema.Struct.F
if (input.sync === undefined || existing?.sync === undefined || input.sync.version >= 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<Definition["sync"]> })
return definition as Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> &
Definition<Type, Schema.Struct<Fields>>
}
@@ -100,6 +110,13 @@ export interface Interface {
readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
readonly all: () => Stream.Stream<Payload>
readonly project: <D extends Definition>(definition: D, projector: Projector<D>) => Effect.Effect<void>
readonly replay: (event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string }) => Effect.Effect<void>
readonly replayAll: (
events: SerializedEvent[],
options?: { readonly publish?: boolean; readonly ownerID?: string },
) => Effect.Effect<string | undefined>
readonly remove: (aggregateID: string) => Effect.Effect<void>
readonly claim: (aggregateID: string, ownerID: string) => Effect.Effect<void>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/Event") {}
@@ -128,7 +145,7 @@ export const layer = Layer.effect(
}),
)
function runProjectors<D extends Definition>(event: Payload<D>) {
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
@@ -155,19 +172,30 @@ export const layer = Layer.effect(
.transaction(
() =>
Effect.gen(function* () {
for (const projector of list) {
yield* projector(event as Payload)
}
const row = yield* db
.select({ seq: EventSequenceTable.seq })
.select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
.from(EventSequenceTable)
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.get()
.pipe(Effect.orDie)
const seq = row?.seq != null ? row.seq + 1 : 0
const latest = row?.seq ?? -1
if (input && input.seq <= latest) return
if (input && row?.ownerID && row.ownerID !== input.ownerID) return
const seq = input?.seq ?? latest + 1
if (input && seq !== latest + 1) {
yield* Effect.die(
new InvalidSyncEventError({
type: event.type,
message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
}),
)
}
for (const projector of list) {
yield* projector(event as Payload)
}
yield* db
.insert(EventSequenceTable)
.values([{ aggregate_id: aggregateID, seq }])
.values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }])
.onConflictDoUpdate({
target: EventSequenceTable.aggregate_id,
set: { seq },
@@ -207,7 +235,7 @@ export const layer = Layer.effect(
...(location ? { location } : {}),
data,
} as Payload<D>
yield* runProjectors(event)
yield* commitSyncEvent(event as Payload)
const pubsub = typed.get(event.type)
if (pubsub) yield* PubSub.publish(pubsub, event as Payload)
yield* PubSub.publish(all, event as Payload)
@@ -215,6 +243,71 @@ export const layer = Layer.effect(
})
}
function replay(event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string }) {
return Effect.gen(function* () {
const definition = syncRegistry.get(event.type)
if (!definition) {
yield* Effect.die(
new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` }),
)
} else {
const payload = {
id: event.id,
type: definition.type,
version: definition.sync.version,
data: event.data,
} as Payload
yield* commitSyncEvent(payload, { seq: event.seq, aggregateID: event.aggregateID, ownerID: options?.ownerID })
if (options?.publish) {
const pubsub = typed.get(payload.type)
if (pubsub) yield* PubSub.publish(pubsub, payload)
yield* PubSub.publish(all, payload)
}
}
})
}
function replayAll(events: SerializedEvent[], options?: { readonly publish?: boolean; readonly ownerID?: string }) {
return Effect.gen(function* () {
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" }))
}
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}` }))
}
}
for (const event of events) {
yield* replay(event, options)
}
return source
})
}
function remove(aggregateID: string) {
return db
.transaction(() =>
Effect.gen(function* () {
yield* db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
}),
)
.pipe(Effect.orDie)
}
function claim(aggregateID: string, ownerID: string) {
return db
.update(EventSequenceTable)
.set({ owner_id: ownerID })
.where(eq(EventSequenceTable.aggregate_id, aggregateID))
.run()
.pipe(Effect.orDie)
}
const subscribe = <D extends Definition>(definition: D): Stream.Stream<Payload<D>> =>
Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe(
Stream.map((event) => event as Payload<D>),
@@ -229,7 +322,7 @@ export const layer = Layer.effect(
projectors.set(definition.type, list)
})
return Service.of({ publish, subscribe, all: streamAll, project })
return Service.of({ publish, subscribe, all: streamAll, project, replay, replayAll, remove, claim })
}),
)
+102
View File
@@ -163,4 +163,106 @@ describe("EventV2", () => {
expect(received).toEqual([SyncMessage.type, "stream"])
}),
)
it.effect("replays sync events through projectors", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
const aggregateID = EventV2.ID.create()
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "hello" },
})
expect(received[0]?.type).toBe(SyncMessage.type)
expect(received[0]?.data).toEqual({ id: aggregateID, text: "hello" })
}),
)
it.effect("replayAll validates contiguous aggregate events", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const aggregateID = EventV2.ID.create()
const source = yield* events.replayAll([
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "one" },
},
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "two" },
},
])
expect(source).toBe(aggregateID)
}),
)
it.effect("claim fences replay owners", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
const aggregateID = EventV2.ID.create()
yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
yield* events.claim(aggregateID, "owner-a")
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
yield* events.replay(
{
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 1,
aggregateID,
data: { id: aggregateID, text: "ignored" },
},
{ ownerID: "owner-b" },
)
expect(received).toHaveLength(0)
}),
)
it.effect("remove clears sync event sequence", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const received = new Array<EventV2.Payload>()
const aggregateID = EventV2.ID.create()
yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
yield* events.remove(aggregateID)
yield* events.project(SyncMessage, (event) =>
Effect.sync(() => {
received.push(event)
}),
)
yield* events.replay({
id: EventV2.ID.create(),
type: EventV2.versionedType(SyncMessage.type, 1),
seq: 0,
aggregateID,
data: { id: aggregateID, text: "replayed" },
})
expect(received[0]?.data).toEqual({ id: aggregateID, text: "replayed" })
}),
)
})