feat(core): persist sync event projections

This commit is contained in:
Dax Raad
2026-05-25 00:13:17 -04:00
parent dedcb9ba91
commit 05dca9f698
6 changed files with 181 additions and 64 deletions
+100 -19
View File
@@ -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<Type extends string = string, DataSchema extends Schema.Top = Schema.Top> = {
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<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 class InvalidSyncEventError extends Schema.TaggedErrorClass<InvalidSyncEventError>()(
"EventV2.InvalidSyncEvent",
{
type: Schema.String,
message: Schema.String,
},
) {}
export function versionedType(type: string, version: number) {
return `${type}.${version}`
}
export const registry = new Map<string, Definition>()
export function define<const Type extends string, Fields extends Schema.Struct.Fields>(input: {
readonly type: Type
readonly version?: number
readonly aggregate?: string
readonly sync?: {
readonly version: number
readonly aggregate: string
}
readonly schema: Fields
}): Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> & Definition<Type, Schema.Struct<Fields>> {
const Data = Schema.Struct(input.schema)
@@ -52,11 +71,13 @@ export function define<const Type extends string, Fields extends Schema.Struct.F
const definition = Object.assign(Payload, {
type: input.type,
...(input.version === undefined ? {} : { version: input.version }),
...(input.aggregate === undefined ? {} : { aggregate: input.aggregate }),
...(input.sync === undefined ? {} : { sync: input.sync }),
data: Data,
})
registry.set(input.type, definition)
const existing = registry.get(input.type)
if (input.sync === undefined || existing?.sync === undefined || input.sync.version >= existing.sync.version) {
registry.set(input.type, definition)
}
return definition as Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> &
Definition<Type, Schema.Struct<Fields>>
}
@@ -76,7 +97,6 @@ export interface Interface {
data: Data<D>,
options?: PublishOptions,
) => Effect.Effect<Payload<D>>
readonly publishEvent: <D extends Definition>(event: Payload<D>) => Effect.Effect<Payload<D>>
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>
@@ -90,6 +110,7 @@ export const layer = Layer.effect(
const all = yield* PubSub.unbounded<Payload>()
const typed = new Map<string, PubSub.PubSub<Payload>>()
const projectors = new Map<string, AnyProjector[]>()
const { db } = yield* Database.Service
const getOrCreate = (definition: Definition) =>
Effect.gen(function* () {
@@ -107,15 +128,71 @@ export const layer = Layer.effect(
}),
)
function publishEvent<D extends Definition>(event: Payload<D>) {
function runProjectors<D extends Definition>(event: Payload<D>) {
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<string, unknown>)[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<string, unknown>,
},
])
.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<D>
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))
+4 -2
View File
@@ -24,8 +24,10 @@ const Base = {
}
const options = {
aggregate: "sessionID",
version: 1,
sync: {
aggregate: "sessionID",
version: 1,
},
} as const
export const UnknownError = Schema.Struct({
+24 -24
View File
@@ -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))
}),
)
+45 -11
View File
@@ -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<EventV2.Payload>()
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"])
}),
)
})
+2 -2
View File
@@ -16,8 +16,8 @@ import { Context, Effect, Layer, Option, Stream } from "effect"
export function toSyncDefinition<D extends EventV2.Definition>(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,
}
+6 -6
View File
@@ -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}` }),
)