diff --git a/packages/core/src/bus.ts b/packages/core/src/bus.ts index 92934de9610..fc4bf9c1cda 100644 --- a/packages/core/src/bus.ts +++ b/packages/core/src/bus.ts @@ -152,16 +152,19 @@ export interface Interface { export class Service extends Context.Service()("@opencode/Bus") {} -export interface LayerOptions { +interface Options { readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect /** Maximum durable rows read per page while replaying or tailing an aggregate log. */ readonly logReadPageSize?: number + /** Retain durable event payloads for historical log reads and replay. */ + readonly persist?: boolean } -export const layerWith = (options?: LayerOptions) => - Layer.effect( - Service, - Effect.gen(function* () { +export function configured(options?: Options) { + return makeGlobalNode({ + service: Service, + deps: [Database.node], + layer: Layer.effect(Service, Effect.gen(function* () { const pubsub = { live: yield* PubSub.unbounded(), durable: new Map>>(), @@ -171,6 +174,7 @@ export const layerWith = (options?: LayerOptions) => const listeners = new Array() const { db } = yield* Database.Service const logReadPageSize = options?.logReadPageSize ?? 512 + const persist = options?.persist ?? false const getOrCreate = (definition: Event.Definition) => Effect.gen(function* () { @@ -251,6 +255,7 @@ export const layerWith = (options?: LayerOptions) => ) } if (input && input.seq <= latest) { + if (!persist) return const stored = yield* db .select() .from(EventTable) @@ -292,19 +297,21 @@ export const layerWith = (options?: LayerOptions) => }), ) } - const stored = yield* db - .select({ aggregateID: EventTable.aggregate_id, seq: EventTable.seq }) - .from(EventTable) - .where(eq(EventTable.id, event.id)) - .get() - .pipe(Effect.orDie) - if (stored) - yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`, - }), - ) + if (persist) { + const stored = yield* db + .select({ aggregateID: EventTable.aggregate_id, seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.id, event.id)) + .get() + .pipe(Effect.orDie) + if (stored) + yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`, + }), + ) + } const committed = { ...event, durable: { aggregateID, seq, version: durable.version }, @@ -325,20 +332,21 @@ export const layerWith = (options?: LayerOptions) => }) .run() .pipe(Effect.orDie) - yield* db - .insert(EventTable) - .values([ - { - id: event.id, - aggregate_id: aggregateID, - seq, - created: DateTime.toEpochMillis(event.created ?? DateTime.makeUnsafe(0)), - type: versionedType(definition.type, durable.version), - data: encoded, - }, - ]) - .run() - .pipe(Effect.orDie) + if (persist) + yield* db + .insert(EventTable) + .values([ + { + id: event.id, + aggregate_id: aggregateID, + seq, + created: DateTime.toEpochMillis(event.created ?? DateTime.makeUnsafe(0)), + type: versionedType(definition.type, durable.version), + data: encoded, + }, + ]) + .run() + .pipe(Effect.orDie) return { aggregateID, seq } }), { behavior: "immediate" }, @@ -683,8 +691,8 @@ export const layerWith = (options?: LayerOptions) => remove, claim, }) - }), - ) + })), + }) +} -export const layer = layerWith() -export const node = makeGlobalNode({ service: Service, layer: layer, deps: [Database.node] }) +export const node = configured() diff --git a/packages/core/test/bus.test.ts b/packages/core/test/bus.test.ts index 2b7d049ff60..4f5515d41b7 100644 --- a/packages/core/test/bus.test.ts +++ b/packages/core/test/bus.test.ts @@ -100,9 +100,19 @@ const tail = (bus: Bus.Interface, input: { aggregateID: string; after?: number } bus.log({ ...input, follow: true }).pipe(Stream.filter((item): item is Event.Payload => !Bus.isSynced(item))) const it = testEffect( - AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, Location.node]), [[Location.node, locationLayer]]), + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, Location.node]), [ + [Location.node, locationLayer], + [Bus.node, Bus.configured({ persist: true })], + ]), +) +const itWithoutLocation = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]), [ + [Bus.node, Bus.configured({ persist: true })], + ]), +) +const itWithoutPersistence = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node])), ) -const itWithoutLocation = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]))) describe("Bus", () => { it.effect("subscribes to multiple event definitions with a discriminated payload union", () => @@ -254,6 +264,27 @@ describe("Bus", () => { }), ) + itWithoutPersistence.effect("projects durable events without retaining their payloads", () => + Effect.gen(function* () { + const bus = yield* Bus.Service + const { db } = yield* Database.Service + const aggregateID = Event.ID.create() + yield* db.run("CREATE TABLE IF NOT EXISTS event_commit_probe (value text NOT NULL)") + yield* bus.project(SyncMessage, () => + db.run("INSERT INTO event_commit_probe (value) VALUES ('projected')").pipe(Effect.orDie, Effect.asVoid), + ) + + const event = yield* bus.publish(SyncMessage, { id: aggregateID, text: "hello" }) + + expect(event.durable?.seq).toBe(Event.Seq.make(0)) + expect(yield* db.all("SELECT value FROM event_commit_probe")).toEqual([{ value: "projected" }]) + expect(yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).all()).toEqual([]) + expect( + yield* db.select().from(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).all(), + ).toEqual([{ aggregate_id: aggregateID, seq: 0, owner_id: null }]) + }), + ) + it.effect("rejects local commit hooks on live-only events", () => Effect.gen(function* () { const bus = yield* Bus.Service @@ -472,12 +503,18 @@ describe("Bus", () => { const readStarted = yield* Deferred.make() const continueRead = yield* Deferred.make() let pause = true - const eventLayer = Bus.layerWith({ - beforeAggregateRead: () => - pause - ? Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(continueRead))) - : Effect.void, - }).pipe(Layer.provide(LayerNode.compile(Database.node))) + const eventLayer = AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]), [ + [ + Bus.node, + Bus.configured({ + persist: true, + beforeAggregateRead: () => + pause + ? Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(continueRead))) + : Effect.void, + }), + ], + ]) yield* Effect.gen(function* () { const bus = yield* Bus.Service @@ -492,7 +529,7 @@ describe("Bus", () => { expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([ [0, durableData(aggregateID, "during handoff")], ]) - }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer))) + }).pipe(Effect.provide(eventLayer)) }), ) @@ -1235,7 +1272,9 @@ describe("Bus", () => { it.effect("log replays across configured read pages", () => Effect.gen(function* () { - const eventLayer = Bus.layerWith({ logReadPageSize: 2 }).pipe(Layer.provide(LayerNode.compile(Database.node))) + const eventLayer = AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]), [ + [Bus.node, Bus.configured({ persist: true, logReadPageSize: 2 })], + ]) yield* Effect.gen(function* () { const bus = yield* Bus.Service @@ -1257,7 +1296,7 @@ describe("Bus", () => { "log.synced", ]) expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID, seq: Event.Seq.make(4) }) - }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer))) + }).pipe(Effect.provide(eventLayer)) }), ) @@ -1266,15 +1305,21 @@ describe("Bus", () => { const readStarted = yield* Deferred.make() const releaseRead = yield* Deferred.make() const firstRead = yield* Ref.make(true) - const eventLayer = Bus.layerWith({ - beforeAggregateRead: () => - Ref.getAndSet(firstRead, false).pipe( - Effect.flatMap((shouldBlock) => { - if (!shouldBlock) return Effect.void - return Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(releaseRead))) - }), - ), - }).pipe(Layer.provide(LayerNode.compile(Database.node))) + const eventLayer = AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]), [ + [ + Bus.node, + Bus.configured({ + persist: true, + beforeAggregateRead: () => + Ref.getAndSet(firstRead, false).pipe( + Effect.flatMap((shouldBlock) => { + if (!shouldBlock) return Effect.void + return Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(releaseRead))) + }), + ), + }), + ], + ]) yield* Effect.gen(function* () { const bus = yield* Bus.Service @@ -1294,7 +1339,7 @@ describe("Bus", () => { { type: "log.synced", aggregateID, seq: Event.Seq.make(0) }, Event.Seq.make(1), ]) - }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer))) + }).pipe(Effect.provide(eventLayer)) }), ) diff --git a/packages/core/test/instruction-state.test.ts b/packages/core/test/instruction-state.test.ts index d63ea574973..fdb4bd56baf 100644 --- a/packages/core/test/instruction-state.test.ts +++ b/packages/core/test/instruction-state.test.ts @@ -17,7 +17,11 @@ import { SessionSchema } from "@opencode-ai/core/session/schema" import { InstructionBlobTable, InstructionStateTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql" import { testEffect } from "./lib/effect" -const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node]))) +const it = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node]), [ + [Bus.node, Bus.configured({ persist: true })], + ]), +) const source = (name: string, read: Effect.Effect) => Instructions.make({ diff --git a/packages/core/test/session-compaction.test.ts b/packages/core/test/session-compaction.test.ts index 689acdd6250..b5ce5baec7e 100644 --- a/packages/core/test/session-compaction.test.ts +++ b/packages/core/test/session-compaction.test.ts @@ -81,6 +81,7 @@ const it = testEffect( AppNodeBuilder.build( LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, SessionCompaction.node]), [ + [Bus.node, Bus.configured({ persist: true })], [llmClient, client], [Config.node, config], [SessionRunnerModel.node, models], diff --git a/packages/core/test/session-create.test.ts b/packages/core/test/session-create.test.ts index c3d9f2be1f2..a1837f00632 100644 --- a/packages/core/test/session-create.test.ts +++ b/packages/core/test/session-create.test.ts @@ -40,6 +40,7 @@ const it = testEffect( AppNodeBuilder.build( LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]), [ + [Bus.node, Bus.configured({ persist: true })], [Project.node, projects], [SessionExecution.node, SessionExecution.noopLayer], ], @@ -562,7 +563,10 @@ describe("Session.create", () => { const targetDatabase = Database.layer({ path: path.join(tmp.path, "target.sqlite") }) const targetLayer = AppNodeBuilder.build( LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node]), - [[Database.node, targetDatabase]], + [ + [Database.node, targetDatabase], + [Bus.node, Bus.configured({ persist: true })], + ], ) yield* Effect.gen(function* () { diff --git a/packages/core/test/session-generate.test.ts b/packages/core/test/session-generate.test.ts index a3f416ce303..3b6de779e1a 100644 --- a/packages/core/test/session-generate.test.ts +++ b/packages/core/test/session-generate.test.ts @@ -133,6 +133,7 @@ const it = testEffect( SessionGenerateNode.node, ]), [ + [Bus.node, Bus.configured({ persist: true })], [llmClient, client], [SessionRunnerModel.node, models], [InstructionBuiltIns.node, builtins], diff --git a/packages/core/test/session-log.test.ts b/packages/core/test/session-log.test.ts index 4cfea987de2..e72ff388e0d 100644 --- a/packages/core/test/session-log.test.ts +++ b/packages/core/test/session-log.test.ts @@ -29,6 +29,7 @@ const it = testEffect( AppNodeBuilder.build( LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]), [ + [Bus.node, Bus.configured({ persist: true })], [Project.node, projects], [SessionExecution.node, SessionExecution.noopLayer], ], @@ -45,9 +46,7 @@ describe("Session.log", () => { const items = Array.from(yield* Stream.runCollect(session.log({ sessionID: created.id }))) - // Session creation commits a non-public durable event, so the marker's - // seq covers more of the aggregate than the public events emitted. - expect(items.map((item) => item.type)).toEqual(["session.renamed", "log.synced"]) + expect(items.map((item) => item.type)).toEqual(["session.created", "session.renamed", "log.synced"]) expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID: created.id, seq: Event.Seq.make(1) }) }), ) @@ -57,7 +56,7 @@ describe("Session.log", () => { const session = yield* Session.Service const created = yield* session.create({ location }) const fiber = yield* session - .log({ sessionID: created.id, follow: true }) + .log({ sessionID: created.id, after: Event.Seq.make(0), follow: true }) .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow diff --git a/packages/core/test/session-projector.test.ts b/packages/core/test/session-projector.test.ts index 1f870fec62f..c13b4195d32 100644 --- a/packages/core/test/session-projector.test.ts +++ b/packages/core/test/session-projector.test.ts @@ -31,7 +31,11 @@ import { import { testEffect } from "./lib/effect" import { Snapshot } from "@opencode-ai/core/snapshot" -const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node]))) +const it = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node]), [ + [Bus.node, Bus.configured({ persist: true })], + ]), +) const sessionsLayer = AppNodeBuilder.build(Session.node, [[SessionExecution.node, SessionExecution.noopLayer]]) const sessionID = Session.ID.make("ses_projector_test") const created = DateTime.makeUnsafe(0) diff --git a/packages/core/test/session-prompt.test.ts b/packages/core/test/session-prompt.test.ts index a735cd82e95..457175c35c4 100644 --- a/packages/core/test/session-prompt.test.ts +++ b/packages/core/test/session-prompt.test.ts @@ -68,6 +68,7 @@ const it = testEffect( AppNodeBuilder.build( LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]), [ + [Bus.node, Bus.configured({ persist: true })], [SessionExecution.node, execution], [LocationServiceMap.node, locations], ], diff --git a/packages/core/test/session-runner-recorded.test.ts b/packages/core/test/session-runner-recorded.test.ts index 3d6b3de38fd..56bc52da1ac 100644 --- a/packages/core/test/session-runner-recorded.test.ts +++ b/packages/core/test/session-runner-recorded.test.ts @@ -159,6 +159,7 @@ const testLayer = (llmClient: Layer.Layer) => Session.node, ]), [ + [Bus.node, Bus.configured({ persist: true })], [LayerNodePlatform.llmClient, llmClient], [Permission.node, permission], [Catalog.node, promptCatalog], diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 5ac96a6def1..6092c10edfe 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -423,6 +423,7 @@ const it = testEffect( Session.node, ]), [ + [Bus.node, Bus.configured({ persist: true })], [LayerNodePlatform.llmClient, client], [Permission.node, permission], [Catalog.node, promptCatalog], diff --git a/packages/core/test/session-tool-progress.test.ts b/packages/core/test/session-tool-progress.test.ts index 74c312ebec2..2cda24108d5 100644 --- a/packages/core/test/session-tool-progress.test.ts +++ b/packages/core/test/session-tool-progress.test.ts @@ -2,6 +2,7 @@ import { describe, expect } from "bun:test" import { asc, eq } from "drizzle-orm" import { DateTime, Effect, Schema } from "effect" import { Database } from "@opencode-ai/core/database/database" +import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder" import { LayerNode } from "@opencode-ai/util/effect/layer-node" import { Bus } from "@opencode-ai/core/bus" import { Agent } from "@opencode-ai/core/agent" @@ -18,7 +19,11 @@ import { SessionProjector } from "@opencode-ai/core/session/projector" import { SessionTable, SessionMessageTable } from "@opencode-ai/core/session/sql" import { testEffect } from "./lib/effect" -const it = testEffect(LayerNode.compile(LayerNode.group([Database.node, Bus.node, SessionProjector.node]))) +const it = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node]), [ + [Bus.node, Bus.configured({ persist: true })], + ]), +) const timestamp = DateTime.makeUnsafe(1) const model = { id: Model.ID.make("model"), providerID: Provider.ID.make("provider") } diff --git a/packages/sdk-next/test/embedded.test.ts b/packages/sdk-next/test/embedded.test.ts index 2f229c4efeb..3b3f823b169 100644 --- a/packages/sdk-next/test/embedded.test.ts +++ b/packages/sdk-next/test/embedded.test.ts @@ -207,7 +207,7 @@ it.live( () => withEmbedded("opencode-embedded-", (fixture) => Effect.gen(function* () { - const opencode = yield* fixture.sdk.OpenCode.create() + const opencode = yield* fixture.sdk.OpenCode.create({ events: { persist: true } }) const id = sessionID(fixture) const model = fixture.sdk.Model.Ref.make({ id: fixture.sdk.Model.ID.make("embedded"), @@ -266,7 +266,7 @@ it.live( const wakeContext = yield* opencode.sessions.context({ sessionID: id }) const pendingAfterPromote = yield* opencode.sessions.pending.list({ sessionID: id }) const event = yield* opencode.sessions.log({ sessionID: id }).pipe( - Stream.filter((item) => item.type !== "log.synced"), + Stream.filter((item) => item.type === "session.model.selected"), Stream.take(1), Stream.runHead, Effect.map(Option.getOrUndefined), diff --git a/packages/server/src/options.ts b/packages/server/src/options.ts index b4f8eccd5b1..cd88c825a6b 100644 --- a/packages/server/src/options.ts +++ b/packages/server/src/options.ts @@ -18,6 +18,11 @@ export const ServerOptions = Schema.Struct({ password: Schema.optional(Schema.String), simulation: Schema.optional(Schema.Boolean), database: Schema.optional(Database.Options), + events: Schema.optional( + Schema.Struct({ + persist: Schema.optional(Schema.Boolean), + }), + ), models: Schema.optional(ModelsDev.Options), observability: Schema.optional(Observability.Options), config: Schema.optional( diff --git a/packages/server/src/routes.ts b/packages/server/src/routes.ts index f2d0d68df73..f75c01cf092 100644 --- a/packages/server/src/routes.ts +++ b/packages/server/src/routes.ts @@ -83,6 +83,7 @@ function makeRoutes( const pluginRuntimeCell = PluginRuntime.makeCell() const replacements: LayerNode.Replacements = [ [Database.node, Database.configured(options.database)], + [Bus.node, Bus.configured({ persist: options.events?.persist })], [App.node, App.configured(options.app)], [ModelsDev.node, ModelsDev.configured(options.models)], [Watcher.node, Watcher.configured({ enabled: options.fs?.filewatcher })], diff --git a/packages/server/test/options.test.ts b/packages/server/test/options.test.ts index 2cae02d393f..5e528072f84 100644 --- a/packages/server/test/options.test.ts +++ b/packages/server/test/options.test.ts @@ -18,3 +18,7 @@ test("accepts optional app metadata", () => { Option.getOrThrow(decode({ app: { name: "sdk", version: "1.2.3", channel: "beta" } })).app, ).toEqual({ name: "sdk", version: "1.2.3", channel: "beta" }) }) + +test("accepts durable event persistence configuration", () => { + expect(Option.getOrThrow(decode({ events: { persist: true } })).events).toEqual({ persist: true }) +})