From 05dca9f6989f2c78abf61c326c4741c226a140b2 Mon Sep 17 00:00:00 2001 From: Dax Raad Date: Mon, 25 May 2026 00:13:17 -0400 Subject: [PATCH] feat(core): persist sync event projections --- packages/core/src/event.ts | 119 +++++++++++++++++++---- packages/core/src/session/event.ts | 6 +- packages/core/src/session/projector.ts | 48 ++++----- packages/core/test/event.test.ts | 56 ++++++++--- packages/opencode/src/event-v2-bridge.ts | 4 +- packages/opencode/src/sync/index.ts | 12 +-- 6 files changed, 181 insertions(+), 64 deletions(-) diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index a07c97b69..f8ddad439 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -1,6 +1,9 @@ export * as EventV2 from "./event" import { Context, Effect, Layer, Option, PubSub, Schema, Stream } from "effect" +import { eq } from "drizzle-orm" +import { Database } from "./database/database" +import { EventSequenceTable, EventTable } from "./event/sql" import { Location } from "./location" import { withStatics } from "./schema" import { Identifier } from "./util/identifier" @@ -13,8 +16,10 @@ export type ID = typeof ID.Type export type Definition = { readonly type: Type - readonly version?: number - readonly aggregate?: string + readonly sync?: { + readonly version: number + readonly aggregate: string + } readonly data: DataSchema } @@ -32,12 +37,26 @@ export type Payload = { export type Projector = (event: Payload) => Effect.Effect type AnyProjector = (event: Payload) => Effect.Effect +export class InvalidSyncEventError extends Schema.TaggedErrorClass()( + "EventV2.InvalidSyncEvent", + { + type: Schema.String, + message: Schema.String, + }, +) {} + +export function versionedType(type: string, version: number) { + return `${type}.${version}` +} + export const registry = new Map() export function define(input: { readonly type: Type - readonly version?: number - readonly aggregate?: string + readonly sync?: { + readonly version: number + readonly aggregate: string + } readonly schema: Fields }): Schema.Schema>>> & Definition> { const Data = Schema.Struct(input.schema) @@ -52,11 +71,13 @@ export function define= existing.sync.version) { + registry.set(input.type, definition) + } return definition as Schema.Schema>>> & Definition> } @@ -76,7 +97,6 @@ export interface Interface { data: Data, options?: PublishOptions, ) => Effect.Effect> - readonly publishEvent: (event: Payload) => Effect.Effect> readonly subscribe: (definition: D) => Stream.Stream> readonly all: () => Stream.Stream readonly project: (definition: D, projector: Projector) => Effect.Effect @@ -90,6 +110,7 @@ export const layer = Layer.effect( const all = yield* PubSub.unbounded() const typed = new Map>() const projectors = new Map() + const { db } = yield* Database.Service const getOrCreate = (definition: Definition) => Effect.gen(function* () { @@ -107,15 +128,71 @@ export const layer = Layer.effect( }), ) - function publishEvent(event: Payload) { + function runProjectors(event: Payload) { return Effect.gen(function* () { - for (const projector of projectors.get(event.type) ?? []) { - yield* projector(event as Payload) + const definition = registry.get(event.type) + const sync = definition?.sync + if (sync) { + if (event.version !== sync.version) { + yield* Effect.die( + new InvalidSyncEventError({ + type: event.type, + message: `Expected event version ${sync.version}, got ${event.version}`, + }), + ) + } + const aggregateID = (event.data as Record)[sync.aggregate] + if (typeof aggregateID !== "string") { + yield* Effect.die( + new InvalidSyncEventError({ + type: event.type, + message: `Expected string aggregate field ${sync.aggregate}`, + }), + ) + } else { + const list = projectors.get(event.type) ?? [] + yield* db + .transaction( + () => + Effect.gen(function* () { + for (const projector of list) { + yield* projector(event as Payload) + } + const row = yield* db + .select({ seq: EventSequenceTable.seq }) + .from(EventSequenceTable) + .where(eq(EventSequenceTable.aggregate_id, aggregateID)) + .get() + .pipe(Effect.orDie) + const seq = row?.seq != null ? row.seq + 1 : 0 + yield* db + .insert(EventSequenceTable) + .values([{ aggregate_id: aggregateID, seq }]) + .onConflictDoUpdate({ + target: EventSequenceTable.aggregate_id, + set: { seq }, + }) + .run() + .pipe(Effect.orDie) + yield* db + .insert(EventTable) + .values([ + { + id: event.id, + aggregate_id: aggregateID, + seq, + type: versionedType(definition.type, sync.version), + data: event.data as Record, + }, + ]) + .run() + .pipe(Effect.orDie) + }), + { behavior: "immediate" }, + ) + .pipe(Effect.orDie) + } } - const pubsub = typed.get(event.type) - if (pubsub) yield* PubSub.publish(pubsub, event as Payload) - yield* PubSub.publish(all, event as Payload) - return event }) } @@ -126,11 +203,15 @@ export const layer = Layer.effect( id: options?.id ?? ID.create(), ...(options?.metadata ? { metadata: options.metadata } : {}), type: definition.type, - ...(definition.version === undefined ? {} : { version: definition.version }), + ...(definition.sync === undefined ? {} : { version: definition.sync.version }), ...(location ? { location } : {}), data, } as Payload - return yield* publishEvent(event) + yield* runProjectors(event) + const pubsub = typed.get(event.type) + if (pubsub) yield* PubSub.publish(pubsub, event as Payload) + yield* PubSub.publish(all, event as Payload) + return event }) } @@ -148,8 +229,8 @@ export const layer = Layer.effect( projectors.set(definition.type, list) }) - return Service.of({ publish, publishEvent, subscribe, all: streamAll, project }) + return Service.of({ publish, subscribe, all: streamAll, project }) }), ) -export const defaultLayer = layer +export const defaultLayer = layer.pipe(Layer.provide(Database.defaultLayer)) diff --git a/packages/core/src/session/event.ts b/packages/core/src/session/event.ts index 27119fcf9..825c49025 100644 --- a/packages/core/src/session/event.ts +++ b/packages/core/src/session/event.ts @@ -24,8 +24,10 @@ const Base = { } const options = { - aggregate: "sessionID", - version: 1, + sync: { + aggregate: "sessionID", + version: 1, + }, } as const export const UnknownError = Schema.Struct({ diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index b961e00db..72d90b4f4 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -174,7 +174,7 @@ function run(db: DatabaseService, event: SessionEvent.Event) { export const layer = Layer.effectDiscard( Effect.gen(function* () { const events = yield* EventV2.Service - const database = yield* Database.Service + const { db } = yield* Database.Service yield* events.project(SessionEvent.AgentSwitched, (event) => Effect.gen(function* () { const message = Schema.encodeSync(SessionMessage.AgentSwitched)( @@ -187,13 +187,13 @@ export const layer = Layer.effectDiscard( }), ) const data = { metadata: message.metadata, agent: message.agent, time: message.time } - yield* database.db + yield* db .update(SessionTable) .set({ agent: event.data.agent, time_updated: DateTime.toEpochMillis(event.data.timestamp) }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie) - yield* database.db + yield* db .insert(SessionMessageTable) .values([ { @@ -220,13 +220,13 @@ export const layer = Layer.effectDiscard( }), ) const data = { metadata: message.metadata, model: message.model, time: message.time } - yield* database.db + yield* db .update(SessionTable) .set({ model: event.data.model, time_updated: DateTime.toEpochMillis(event.data.timestamp) }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie) - yield* database.db + yield* db .insert(SessionMessageTable) .values([ { @@ -241,25 +241,25 @@ export const layer = Layer.effectDiscard( .pipe(Effect.orDie) }), ) - yield* events.project(SessionEvent.Prompted, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Synthetic, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Shell.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Shell.Ended, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Step.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Step.Ended, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Step.Failed, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Text.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Text.Ended, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Tool.Input.Ended, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Tool.Called, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Tool.Success, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Tool.Failed, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Reasoning.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Reasoning.Ended, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Retried, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Compaction.Started, (event) => run(database.db, event)) - yield* events.project(SessionEvent.Compaction.Ended, (event) => run(database.db, event)) + yield* events.project(SessionEvent.Prompted, (event) => run(db, event)) + yield* events.project(SessionEvent.Synthetic, (event) => run(db, event)) + yield* events.project(SessionEvent.Shell.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Shell.Ended, (event) => run(db, event)) + yield* events.project(SessionEvent.Step.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Step.Ended, (event) => run(db, event)) + yield* events.project(SessionEvent.Step.Failed, (event) => run(db, event)) + yield* events.project(SessionEvent.Text.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Text.Ended, (event) => run(db, event)) + yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Tool.Input.Ended, (event) => run(db, event)) + yield* events.project(SessionEvent.Tool.Called, (event) => run(db, event)) + yield* events.project(SessionEvent.Tool.Success, (event) => run(db, event)) + yield* events.project(SessionEvent.Tool.Failed, (event) => run(db, event)) + yield* events.project(SessionEvent.Reasoning.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Reasoning.Ended, (event) => run(db, event)) + yield* events.project(SessionEvent.Retried, (event) => run(db, event)) + yield* events.project(SessionEvent.Compaction.Started, (event) => run(db, event)) + yield* events.project(SessionEvent.Compaction.Ended, (event) => run(db, event)) }), ) diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index 821193688..24b0df0a2 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -9,8 +9,8 @@ const locationLayer = Layer.succeed( Location.Service, Location.Service.of({ directory: AbsolutePath.make("project"), workspaceID: "workspace" }), ) -const it = testEffect(EventV2.layer.pipe(Layer.provideMerge(locationLayer))) -const itWithoutLocation = testEffect(EventV2.layer) +const it = testEffect(EventV2.defaultLayer.pipe(Layer.provideMerge(locationLayer))) +const itWithoutLocation = testEffect(EventV2.defaultLayer) const Message = EventV2.define({ type: "test.message", @@ -19,6 +19,18 @@ const Message = EventV2.define({ }, }) +const SyncMessage = EventV2.define({ + type: "test.sync", + sync: { + version: 1, + aggregate: "id", + }, + schema: { + id: Schema.String, + text: Schema.String, + }, +}) + const GlobalMessage = EventV2.define({ type: "test.global", schema: { @@ -28,8 +40,12 @@ const GlobalMessage = EventV2.define({ const VersionedMessage = EventV2.define({ type: "test.versioned", - version: 2, + sync: { + version: 2, + aggregate: "id", + }, schema: { + id: Schema.String, text: Schema.String, }, }) @@ -64,7 +80,7 @@ describe("EventV2", () => { it.effect("publishes definition version", () => Effect.gen(function* () { const events = yield* EventV2.Service - const event = yield* events.publish(VersionedMessage, { text: "hello" }) + const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" }) expect(event.type).toBe("test.versioned") expect(event.version).toBe(2) @@ -77,6 +93,23 @@ describe("EventV2", () => { }), ) + it.effect("keeps the latest sync definition in the registry", () => + Effect.sync(() => { + const latest = EventV2.define({ + type: "test.out-of-order", + sync: { version: 2, aggregate: "id" }, + schema: { id: Schema.String }, + }) + EventV2.define({ + type: "test.out-of-order", + sync: { version: 1, aggregate: "id" }, + schema: { id: Schema.String }, + }) + + expect(EventV2.registry.get("test.out-of-order")).toBe(latest) + }), + ) + it.effect("publishes to typed and wildcard subscriptions", () => Effect.gen(function* () { const events = yield* EventV2.Service @@ -94,16 +127,17 @@ describe("EventV2", () => { Effect.gen(function* () { const events = yield* EventV2.Service const received = new Array() - yield* events.project(Message, (event) => + yield* events.project(SyncMessage, (event) => Effect.sync(() => { received.push(event) }), ) - const event = yield* events.publish(Message, { text: "hello" }) - yield* events.publish(Message, { text: "after unsubscribe" }) + const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" }) + yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" }) - expect(received).toEqual([event, expect.objectContaining({ data: { text: "after unsubscribe" } })]) + expect(received[0]).toEqual(event) + expect(received[1]?.data).toEqual({ id: "one", text: "after unsubscribe" }) }), ) @@ -116,17 +150,17 @@ describe("EventV2", () => { Stream.runForEach(() => Effect.sync(() => received.push("stream"))), Effect.forkScoped, ) - yield* events.project(Message, (event) => + yield* events.project(SyncMessage, (event) => Effect.sync(() => { received.push(event.type) }), ) yield* Effect.yieldNow - yield* events.publish(Message, { text: "hello" }) + yield* events.publish(SyncMessage, { id: "one", text: "hello" }) yield* Fiber.join(fiber) - expect(received).toEqual([Message.type, "stream"]) + expect(received).toEqual([SyncMessage.type, "stream"]) }), ) }) diff --git a/packages/opencode/src/event-v2-bridge.ts b/packages/opencode/src/event-v2-bridge.ts index 668d701d0..23923c504 100644 --- a/packages/opencode/src/event-v2-bridge.ts +++ b/packages/opencode/src/event-v2-bridge.ts @@ -16,8 +16,8 @@ import { Context, Effect, Layer, Option, Stream } from "effect" export function toSyncDefinition(definition: D) { const result = { type: definition.type, - version: definition.version, - aggregate: definition.aggregate, + version: definition.sync?.version, + aggregate: definition.sync?.aggregate, schema: definition.data, properties: definition.data, } diff --git a/packages/opencode/src/sync/index.ts b/packages/opencode/src/sync/index.ts index 82cbfb169..25cae5b9a 100644 --- a/packages/opencode/src/sync/index.ts +++ b/packages/opencode/src/sync/index.ts @@ -218,11 +218,11 @@ export function reset() { export function init(input: { projectors: Array<[Definition, ProjectorFunc]>; convertEvent?: ConvertEvent }) { projectors = new Map(input.projectors.map(([def, func]) => [versionedType(def.type, def.version), func])) for (let entry of EventV2.registry.values()) { - if (!entry.version || !entry.aggregate) continue + if (!entry.sync) continue register({ type: entry.type, - version: entry.version, - aggregate: entry.aggregate, + version: entry.sync.version, + aggregate: entry.sync.aggregate, properties: entry.data, schema: entry.data, }) @@ -392,15 +392,15 @@ export function effectPayloads() { .values() .filter( (definition) => - definition.version !== undefined && !registry.has(versionedType(definition.type, definition.version)), + definition.sync !== undefined && !registry.has(versionedType(definition.type, definition.sync.version)), ) .map((definition) => EffectSchema.Struct({ type: EffectSchema.Literal("sync"), - name: EffectSchema.Literal(versionedType(definition.type, definition.version!)), + name: EffectSchema.Literal(versionedType(definition.type, definition.sync!.version)), id: EffectSchema.String, seq: EffectSchema.Finite, - aggregateID: EffectSchema.Literal(definition.aggregate!), + aggregateID: EffectSchema.Literal(definition.sync!.aggregate), data: definition.data, }).annotate({ identifier: `SyncEvent.${definition.type}` }), )