From a00822b06e2fbc78855523ff7773fe1efe85d7f6 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 17 Aug 2026 15:09:35 -0400 Subject: [PATCH] feat(core): Session.snapshot --- packages/core/src/session.ts | 59 ++++++++++++++ packages/core/test/session-snapshot.test.ts | 87 +++++++++++++++++++++ 2 files changed, 146 insertions(+) create mode 100644 packages/core/test/session-snapshot.test.ts diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index a7df3882500..a90148281fd 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -53,6 +53,7 @@ import { Shell as ShellSchema } from "@opencode-ai/schema/shell" import { KeyedMutex } from "./effect/keyed-mutex.js" import { fileURLToPath } from "url" import { SessionEnvironment } from "./session/environment.js" +import { EventSequenceTable } from "./event/sql.js" // get project -> project.locations // @@ -194,6 +195,19 @@ export interface Interface { * unhandled compaction barriers. */ readonly inbox: (sessionID: SessionSchema.ID) => Effect.Effect + readonly snapshot: (input: { + sessionID: SessionSchema.ID + recent?: number + }) => Effect.Effect< + { + readonly session: SessionSchema.Info + readonly children: SessionSchema.Info[] + readonly inbox: SessionInbox.Info[] + readonly messages: SessionMessage.Info[] + readonly seq: Event.Seq + }, + NotFoundError | MessageDecodeError + > readonly cancelInbox: (input: InboxItemRef) => Effect.Effect readonly steerInbox: (input: InboxItemRef) => Effect.Effect readonly queueInbox: (input: InboxItemRef) => Effect.Effect @@ -557,6 +571,51 @@ const layer = Layer.effect( yield* result.get(sessionID) return yield* SessionInbox.list(db, sessionID) }), + snapshot: Effect.fn("Session.snapshot")(function* (input) { + return yield* db + .transaction(() => + Effect.gen(function* () { + const row = yield* db + .select() + .from(SessionTable) + .where(eq(SessionTable.id, input.sessionID)) + .get() + .pipe(Effect.orDie) + if (!row) return yield* new NotFoundError({ sessionID: input.sessionID }) + const children = yield* db + .select() + .from(SessionTable) + .where(eq(SessionTable.parent_id, input.sessionID)) + .orderBy(desc(SessionTable.time_updated), desc(SessionTable.id)) + .all() + .pipe(Effect.orDie) + const inbox = yield* SessionInbox.list(db, input.sessionID) + const messages = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.session_id, input.sessionID)) + .orderBy(desc(SessionMessageTable.seq)) + .limit(input.recent ?? 200) + .all() + .pipe(Effect.orDie) + const sequence = yield* db + .select({ seq: EventSequenceTable.seq }) + .from(EventSequenceTable) + .where(eq(EventSequenceTable.aggregate_id, input.sessionID)) + .get() + .pipe(Effect.orDie) + if (!sequence) return yield* Effect.die(new Error(`Session ${input.sessionID} has no event sequence`)) + return { + session: fromRow(row), + children: children.map(fromRow), + inbox, + messages: yield* Effect.forEach(messages.toReversed(), decode), + seq: Event.Seq.make(sequence.seq), + } + }), + ) + .pipe(Effect.catchTag("SqlError", Effect.die)) + }), cancelInbox: Effect.fn("Session.cancelInbox")((input) => mutatePending(input, SessionInbox.cancel)), steerInbox: Effect.fn("Session.steerInbox")((input) => mutatePending(input, SessionInbox.steer, true)), queueInbox: Effect.fn("Session.queueInbox")((input) => mutatePending(input, SessionInbox.queue)), diff --git a/packages/core/test/session-snapshot.test.ts b/packages/core/test/session-snapshot.test.ts new file mode 100644 index 00000000000..6c5c1ddf5bc --- /dev/null +++ b/packages/core/test/session-snapshot.test.ts @@ -0,0 +1,87 @@ +import { describe, expect } from "bun:test" +import { Effect } 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 { Location } from "@opencode-ai/core/location" +import { Project } from "@opencode-ai/core/project" +import { AbsolutePath } from "@opencode-ai/core/schema" +import { Session } from "@opencode-ai/core/session" +import { SessionProjector } from "@opencode-ai/core/session/projector" +import { SessionExecution } from "@opencode-ai/core/session/execution" +import { SessionStore } from "@opencode-ai/core/session/store" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { Event } from "@opencode-ai/schema/event" +import { testEffect } from "./lib/effect" +import { globalProjectLayer } from "./lib/project" + +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, globalProjectLayer], + [SessionExecution.node, SessionExecution.noopLayer], + ], + ), +) +const location = Location.Ref.make({ directory: AbsolutePath.make("/project") }) + +describe("Session.snapshot", () => { + it.effect("returns an empty projected session at its aggregate watermark", () => + Effect.gen(function* () { + const sessions = yield* Session.Service + const created = yield* sessions.create({ location }) + + expect(yield* sessions.snapshot({ sessionID: created.id })).toEqual({ + session: created, + children: [], + inbox: [], + messages: [], + seq: Event.Seq.make(0), + }) + }), + ) + + it.effect("returns the most recent messages in aggregate order", () => + Effect.gen(function* () { + const sessions = yield* Session.Service + const bus = yield* Bus.Service + const created = yield* sessions.create({ location }) + yield* Effect.forEach(["first", "second", "third"], (text) => + bus.publish(SessionEvent.Synthetic, { sessionID: created.id, text }), + ) + + const snapshot = yield* sessions.snapshot({ sessionID: created.id, recent: 2 }) + + expect(snapshot.messages.map((message) => (message.type === "synthetic" ? message.text : message.type))).toEqual([ + "second", + "third", + ]) + expect(snapshot.seq).toBe(Event.Seq.make(3)) + }), + ) + + it.effect("keeps rows and watermark consistent during concurrent publication", () => + Effect.gen(function* () { + const sessions = yield* Session.Service + const bus = yield* Bus.Service + const created = yield* sessions.create({ location }) + const publish = Effect.forEach( + Array.from({ length: 40 }, (_, index) => index + 1), + (index) => bus.publish(SessionEvent.Synthetic, { sessionID: created.id, text: String(index) }), + ) + const read = Effect.forEach(Array.from({ length: 40 }), () => sessions.snapshot({ sessionID: created.id })) + + const [, snapshots] = yield* Effect.all([publish, read], { concurrency: "unbounded" }) + + snapshots.forEach((snapshot) => { + expect(snapshot.messages).toHaveLength(snapshot.seq) + expect( + snapshot.messages.map((message) => (message.type === "synthetic" ? Number(message.text) : -1)), + ).toEqual(Array.from({ length: snapshot.seq }, (_, index) => index + 1)) + }) + }), + ) +})