diff --git a/packages/opencode/src/server/projectors.ts b/packages/opencode/src/server/projectors.ts index 296151b30b..87ed254fff 100644 --- a/packages/opencode/src/server/projectors.ts +++ b/packages/opencode/src/server/projectors.ts @@ -2,7 +2,8 @@ import z from "zod" import sessionProjectors from "../session/projectors" import { SyncEvent } from "@/sync" import { Session } from "@/session" -import { SessionID } from "@/session/schema" +import { SessionTable } from "@/session/session.sql" +import { Database, eq } from "@/storage/db" let initialized = false @@ -16,10 +17,14 @@ export function initProjectors() { projectors: sessionProjectors, convertEvent: (type, data) => { if (type === "session.updated") { - const sessionID = (data as z.infer).sessionID + const id = (data as z.infer).sessionID + const row = Database.use((db) => db.select().from(SessionTable).where(eq(SessionTable.id, id)).get()) + + if (!row) return data + return { - sessionID: SessionID.zod, - info: Session.get(sessionID), + sessionID: id, + info: Session.fromRow(row), } } return data diff --git a/packages/opencode/src/session/index.ts b/packages/opencode/src/session/index.ts index d0652f6199..a379cd228c 100644 --- a/packages/opencode/src/session/index.ts +++ b/packages/opencode/src/session/index.ts @@ -16,6 +16,7 @@ import { SessionTable } from "./session.sql" import { ProjectTable } from "../project/project.sql" import { Storage } from "@/storage/storage" import { Log } from "../util/log" +import { updateSchema } from "../util/update-schema" import { MessageV2 } from "./message-v2" import { Instance } from "../project/instance" import { SessionPrompt } from "./prompt" @@ -198,9 +199,9 @@ export namespace Session { aggregate: "sessionID", schema: z.object({ sessionID: SessionID.zod, - info: Info.partial().extend({ - share: Info.shape.share.unwrap().partial().optional(), - time: Info.shape.time.partial().optional(), + info: updateSchema(Info).extend({ + share: updateSchema(Info.shape.share.unwrap()).optional(), + time: updateSchema(Info.shape.time).optional(), }), }), busSchema: z.object({ @@ -378,7 +379,7 @@ export namespace Session { const { ShareNext } = await import("@/share/share-next") await ShareNext.remove(id) - SyncEvent.run(Event.Updated, { sessionID: id, info: { share: { url: undefined } } }) + SyncEvent.run(Event.Updated, { sessionID: id, info: { share: { url: null } } }) }) export const setTitle = fn( @@ -437,7 +438,7 @@ export namespace Session { sessionID, info: { time: { updated: Date.now() }, - revert: undefined, + revert: null, }, }) }) @@ -615,6 +616,9 @@ export namespace Session { await unshare(sessionID).catch(() => {}) SyncEvent.run(Event.Deleted, { sessionID, info: session }) + + // Eagerly remove event sourcing data to free up space + SyncEvent.remove(sessionID) } catch (e) { log.error(e) } diff --git a/packages/opencode/src/session/projectors.ts b/packages/opencode/src/session/projectors.ts index 09d178431a..88c53d9ef8 100644 --- a/packages/opencode/src/session/projectors.ts +++ b/packages/opencode/src/session/projectors.ts @@ -5,7 +5,7 @@ import { MessageV2 } from "./message-v2" import { SessionTable, MessageTable, PartTable } from "./session.sql" import { ProjectTable } from "../project/project.sql" -export type DeepPartial = T extends object ? { [K in keyof T]?: DeepPartial } : T +export type DeepPartial = T extends object ? { [K in keyof T]?: DeepPartial | null } : T function grab( obj: T, @@ -18,7 +18,12 @@ function grab( if (val && typeof val === "object" && cb) { return cb(val) } - return (val === undefined ? null : val) as X | undefined + if (val === undefined) { + throw new Error( + "Session update failure: pass `null` to clear a field instead of `undefined`: " + JSON.stringify(obj), + ) + } + return val as X | undefined } export function toPartialRow(info: DeepPartial) { diff --git a/packages/opencode/src/sync/index.ts b/packages/opencode/src/sync/index.ts index e9f1ff9369..bfff84031f 100644 --- a/packages/opencode/src/sync/index.ts +++ b/packages/opencode/src/sync/index.ts @@ -35,14 +35,11 @@ export namespace SyncEvent { let projectors: Map | undefined const versions = new Map() let frozen = false - let convertEvent: ((type: string, event: Event["data"]) => Record) | undefined + let convertEvent: (type: string, event: Event["data"]) => Promise> | Record const Bus = new EventEmitter<{ event: [{ def: Definition; event: Event }] }>() - export function init(input: { - projectors: Array<[Definition, ProjectorFunc]> - convertEvent?: Exclude - }) { + export function init(input: { projectors: Array<[Definition, ProjectorFunc]>; convertEvent?: typeof convertEvent }) { projectors = new Map(input.projectors) // Install all the latest event defs to the bus. We only ever emit @@ -58,7 +55,7 @@ export namespace SyncEvent { // Freeze the system so it clearly errors if events are defined // after `init` which would cause bugs frozen = true - convertEvent = input.convertEvent + convertEvent = input.convertEvent || ((_, data) => data) } export function versionedType(type: A): A @@ -99,7 +96,7 @@ export namespace SyncEvent { return [def, func as ProjectorFunc] } - function process(def: Def, input: Event) { + function process(def: Def, event: Event, options: { publish: boolean }) { if (projectors == null) { throw new Error("No projectors available. Call `SyncEvent.init` to install projectors") } @@ -112,29 +109,47 @@ export namespace SyncEvent { // idempotent: need to ignore any events already logged Database.transaction((tx) => { - projector(tx, input.data) + projector(tx, event.data) if (Flag.OPENCODE_EXPERIMENTAL_WORKSPACES) { tx.insert(EventSequenceTable) .values({ - aggregate_id: input.aggregateID, - seq: input.seq, + aggregate_id: event.aggregateID, + seq: event.seq, }) .onConflictDoUpdate({ target: EventSequenceTable.aggregate_id, - set: { seq: input.seq }, + set: { seq: event.seq }, }) .run() tx.insert(EventTable) .values({ - id: input.id, - seq: input.seq, - aggregate_id: input.aggregateID, - name: def.type, - data: input.data as Record, + id: event.id, + seq: event.seq, + aggregate_id: event.aggregateID, + name: versionedType(def.type, def.version), + data: event.data as Record, }) .run() } + + Database.effect(() => { + Bus.emit("event", { + def, + event, + }) + + if (options?.publish) { + const result = convertEvent(def.type, event.data) + if (result instanceof Promise) { + result.then((data) => { + ProjectBus.publish({ type: def.type, properties: def.schema }, data) + }) + } else { + ProjectBus.publish({ type: def.type, properties: def.schema }, result) + } + } + }) }) } @@ -144,7 +159,7 @@ export namespace SyncEvent { // and it validets all the sequence ids // * when loading events from db, apply zod validation to ensure shape - export function replay(event: SerializedEvent) { + export function replay(event: SerializedEvent, options?: { republish: boolean }) { const def = registry.get(event.type) if (!def) { throw new Error(`Unknown event type: ${event.type}`) @@ -158,12 +173,17 @@ export namespace SyncEvent { .get(), ) - const expected = row ? row.seq + 1 : 0 + const latest = row?.seq ?? -1 + if (event.seq <= latest) { + return + } + + const expected = latest + 1 if (event.seq !== expected) { throw new Error(`Sequence mismatch for aggregate "${event.aggregateID}": expected ${expected}, got ${event.seq}`) } - process(def, event) + process(def, event, { publish: !!options?.republish }) } export function run(def: Def, data: Event["data"]) { @@ -192,22 +212,7 @@ export namespace SyncEvent { const seq = row?.seq != null ? row.seq + 1 : 0 const event = { id, seq, aggregateID: agg, data } - process(def, event) - - Database.effect(() => { - Bus.emit("event", { - def, - event, - }) - - ProjectBus.publish( - { - type: def.type, - properties: def.schema, - }, - convertEvent ? convertEvent(def.type, event.data) : event.data, - ) - }) + process(def, event, { publish: true }) }, { behavior: "immediate", @@ -215,6 +220,13 @@ export namespace SyncEvent { ) } + export function remove(aggregateID: string) { + Database.transaction((tx) => { + tx.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run() + tx.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run() + }) + } + export function subscribeAll(handler: (event: { def: Definition; event: Event }) => void) { Bus.on("event", handler) return () => Bus.off("event", handler) diff --git a/packages/opencode/src/util/update-schema.ts b/packages/opencode/src/util/update-schema.ts new file mode 100644 index 0000000000..f2246ece33 --- /dev/null +++ b/packages/opencode/src/util/update-schema.ts @@ -0,0 +1,13 @@ +import z from "zod" + +export function updateSchema(schema: z.ZodObject) { + const next = {} as { + [K in keyof T]: z.ZodOptional> + } + + for (const [k, v] of Object.entries(schema.required().shape) as [keyof T & string, z.ZodTypeAny][]) { + next[k] = v.nullable() as unknown as (typeof next)[typeof k] + } + + return z.object(next) +}