From f8f505d508e842c64a11e8c929e6a3e5644d7f60 Mon Sep 17 00:00:00 2001 From: Dante Date: Sun, 4 Oct 2026 17:57:17 +0800 Subject: [PATCH] refactor(core): extract transcript projection Adapt the transcript-only extraction previously proposed in #36811 to current v2 agent, model and location lookups. Keep transaction and registration ordering unchanged. --- .../core/src/session/projection/transcript.ts | 201 ++++++++++++++ packages/core/src/session/projector.ts | 245 +++--------------- .../core/test/transcript-projection.test.ts | 216 +++++++++++++++ 3 files changed, 447 insertions(+), 215 deletions(-) create mode 100644 packages/core/src/session/projection/transcript.ts create mode 100644 packages/core/test/transcript-projection.test.ts diff --git a/packages/core/src/session/projection/transcript.ts b/packages/core/src/session/projection/transcript.ts new file mode 100644 index 000000000000..307e157a291b --- /dev/null +++ b/packages/core/src/session/projection/transcript.ts @@ -0,0 +1,201 @@ +export * as TranscriptProjection from "./transcript.js" + +import { and, desc, eq, sql } from "drizzle-orm" +import { DateTime, Effect, Schema } from "effect" +import { Database } from "../../database/database.js" +import { Agent } from "@opencode/schema/agent" +import { Model } from "@opencode/schema/model" +import { Workspace } from "@opencode/schema/workspace" +import { AbsolutePath, RelativePath } from "../../schema.js" +import { SessionEvent } from "../event.js" +import { SessionMessage } from "../message.js" +import { SessionMessageUpdater } from "../message-updater.js" +import { SessionMessageTable, SessionTable } from "../sql.js" + +type DatabaseService = Database.Interface["db"] +type MessageEvent = Exclude< + SessionEvent.DurableEvent, + typeof SessionEvent.Forked.Type | typeof SessionEvent.Deleted.Type +> + +const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info) +const encodeMessage = Schema.encodeSync(SessionMessage.Info) + +export function project(db: DatabaseService, event: MessageEvent) { + return Effect.gen(function* () { + const decodeRow = (row: typeof SessionMessageTable.$inferSelect) => + decodeMessage({ ...row.data, id: row.id, type: row.type }) + const updateMessage = (message: SessionMessage.Info) => { + const encoded = encodeMessage(message) + const { id, type, ...data } = encoded + return db + .update(SessionMessageTable) + .set({ type, time_created: DateTime.toEpochMillis(message.time.created), data }) + .where( + and( + eq(SessionMessageTable.id, SessionMessage.ID.make(id)), + eq(SessionMessageTable.session_id, event.data.sessionID), + ), + ) + .run() + .pipe(Effect.orDie) + } + const appendMessage = (message: SessionMessage.Info) => appendAtEventSequence(db, event, message) + const adapter: SessionMessageUpdater.Adapter = { + getAgent() { + return db + .select({ agent: SessionTable.agent }) + .from(SessionTable) + .where(eq(SessionTable.id, event.data.sessionID)) + .get() + .pipe( + Effect.orDie, + Effect.map((row) => (row?.agent ? Agent.ID.make(row.agent) : undefined)), + ) + }, + getModel() { + return db + .select({ model: SessionTable.model }) + .from(SessionTable) + .where(eq(SessionTable.id, event.data.sessionID)) + .get() + .pipe( + Effect.orDie, + Effect.map((row) => (row?.model ? Schema.decodeUnknownSync(Model.Ref)(row.model) : undefined)), + ) + }, + getLocation() { + return db + .select({ + directory: SessionTable.directory, + workspaceID: SessionTable.workspace_id, + projectID: SessionTable.project_id, + subpath: SessionTable.path, + }) + .from(SessionTable) + .where(eq(SessionTable.id, event.data.sessionID)) + .get() + .pipe( + Effect.orDie, + Effect.map((row) => + row + ? { + location: { + directory: AbsolutePath.make(row.directory), + workspaceID: row.workspaceID ? Workspace.ID.make(row.workspaceID) : undefined, + }, + projectID: row.projectID, + subpath: row.subpath === null ? undefined : RelativePath.make(row.subpath), + } + : undefined, + ), + ) + }, + getCurrentAssistant() { + return Effect.gen(function* () { + // A newer step supersedes stale incomplete rows; never resume an older assistant projection. + const row = yield* db + .select() + .from(SessionMessageTable) + .where( + and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "assistant")), + ) + .orderBy(desc(SessionMessageTable.seq)) + .limit(1) + .get() + .pipe(Effect.orDie) + if (!row) return + const message = decodeRow(row) + return message.type === "assistant" && !message.time.completed ? message : undefined + }) + }, + getAssistant(messageID) { + return Effect.gen(function* () { + const row = yield* db + .select() + .from(SessionMessageTable) + .where( + and( + eq(SessionMessageTable.id, messageID), + eq(SessionMessageTable.session_id, event.data.sessionID), + eq(SessionMessageTable.type, "assistant"), + ), + ) + .get() + .pipe(Effect.orDie) + if (!row) return + const message = decodeRow(row) + return message.type === "assistant" ? message : undefined + }) + }, + getShell(shellID) { + return Effect.gen(function* () { + const row = yield* db + .select() + .from(SessionMessageTable) + .where( + and( + eq(SessionMessageTable.session_id, event.data.sessionID), + eq(SessionMessageTable.type, "shell"), + sql`json_extract(${SessionMessageTable.data}, '$.shellID') = ${shellID}`, + ), + ) + .orderBy(desc(SessionMessageTable.seq)) + .limit(1) + .get() + .pipe(Effect.orDie) + if (!row) return + const message = decodeRow(row) + return message.type === "shell" ? message : undefined + }) + }, + getCompaction() { + return Effect.gen(function* () { + const row = yield* db + .select() + .from(SessionMessageTable) + .where( + and( + eq(SessionMessageTable.session_id, event.data.sessionID), + eq(SessionMessageTable.type, "compaction"), + sql`json_extract(${SessionMessageTable.data}, '$.status') = 'running'`, + ), + ) + .orderBy(desc(SessionMessageTable.seq)) + .limit(1) + .get() + .pipe(Effect.orDie) + if (!row) return + const message = decodeRow(row) + return message.type === "compaction" ? message : undefined + }) + }, + updateAssistant: updateMessage, + updateShell: updateMessage, + updateCompaction: updateMessage, + appendMessage, + } + yield* SessionMessageUpdater.update(adapter, event) + }) +} + +export function appendAtEventSequence( + db: DatabaseService, + event: SessionEvent.DurableEvent, + message: SessionMessage.Info, +) { + const encoded = encodeMessage(message) + const { id, type, ...data } = encoded + return db + .insert(SessionMessageTable) + .values({ + id: SessionMessage.ID.make(id), + session_id: event.data.sessionID, + type, + seq: event.durable.seq, + time_created: DateTime.toEpochMillis(message.time.created), + data, + }) + .run() + .pipe(Effect.orDie) +} diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index f3b5ec439906..d5be8bf2d386 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -1,16 +1,14 @@ export * as SessionProjector from "./projector.js" import { and, asc, desc, eq, gt, gte, inArray, isNull, lt, lte, or, sql } from "drizzle-orm" -import { DateTime, Effect, Layer, Schema, Stream } from "effect" +import { DateTime, Effect, Layer, Stream } from "effect" import path from "path" import { Database } from "../database/database.js" import { Bus } from "../bus.js" import { makeGlobalNode } from "@opencode/util/effect/app-node" -import { Agent } from "@opencode/schema/agent" -import { Model } from "@opencode/schema/model" import { SessionEvent } from "./event.js" import { SessionMessage } from "./message.js" -import { SessionMessageUpdater } from "./message-updater.js" +import { TranscriptProjection } from "./projection/transcript.js" import { SessionInbox } from "./inbox.js" import { Workspace } from "@opencode/schema/workspace" import { InstructionState } from "./instruction-state.js" @@ -26,14 +24,6 @@ import type { SessionSchema } from "./schema.js" import { ProjectTable } from "../project/sql.js" type DatabaseService = Database.Interface["db"] -type MessageEvent = Exclude< - SessionEvent.DurableEvent, - typeof SessionEvent.Forked.Type | typeof SessionEvent.Deleted.Type -> - -const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info) -const encodeMessage = Schema.encodeSync(SessionMessage.Info) - export class SessionAlreadyProjected extends Error {} type Usage = { @@ -225,181 +215,6 @@ const projectFork = Effect.fn("SessionProjector.projectFork")(function* ( yield* InstructionState.initialize(db, event.data.sessionID, event.durable.seq, event.data.instructions) }) -function run(db: DatabaseService, event: MessageEvent) { - return Effect.gen(function* () { - const decodeRow = (row: typeof SessionMessageTable.$inferSelect) => - decodeMessage({ ...row.data, id: row.id, type: row.type }) - const updateMessage = (message: SessionMessage.Info) => { - const encoded = encodeMessage(message) - const { id, type, ...data } = encoded - return db - .update(SessionMessageTable) - .set({ type, time_created: DateTime.toEpochMillis(message.time.created), data }) - .where( - and( - eq(SessionMessageTable.id, SessionMessage.ID.make(id)), - eq(SessionMessageTable.session_id, event.data.sessionID), - ), - ) - .run() - .pipe(Effect.orDie) - } - const appendMessage = (message: SessionMessage.Info) => insertMessage(db, event, message) - const adapter: SessionMessageUpdater.Adapter = { - getAgent() { - return db - .select({ agent: SessionTable.agent }) - .from(SessionTable) - .where(eq(SessionTable.id, event.data.sessionID)) - .get() - .pipe( - Effect.orDie, - Effect.map((row) => (row?.agent ? Agent.ID.make(row.agent) : undefined)), - ) - }, - getModel() { - return db - .select({ model: SessionTable.model }) - .from(SessionTable) - .where(eq(SessionTable.id, event.data.sessionID)) - .get() - .pipe( - Effect.orDie, - Effect.map((row) => (row?.model ? Schema.decodeUnknownSync(Model.Ref)(row.model) : undefined)), - ) - }, - getLocation() { - return db - .select({ - directory: SessionTable.directory, - workspaceID: SessionTable.workspace_id, - projectID: SessionTable.project_id, - subpath: SessionTable.path, - }) - .from(SessionTable) - .where(eq(SessionTable.id, event.data.sessionID)) - .get() - .pipe( - Effect.orDie, - Effect.map((row) => - row - ? { - location: { - directory: AbsolutePath.make(row.directory), - workspaceID: row.workspaceID ? Workspace.ID.make(row.workspaceID) : undefined, - }, - projectID: row.projectID, - subpath: row.subpath === null ? undefined : RelativePath.make(row.subpath), - } - : undefined, - ), - ) - }, - getCurrentAssistant() { - return Effect.gen(function* () { - // A newer step supersedes stale incomplete rows; never resume an older assistant projection. - const row = yield* db - .select() - .from(SessionMessageTable) - .where( - and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "assistant")), - ) - .orderBy(desc(SessionMessageTable.seq)) - .limit(1) - .get() - .pipe(Effect.orDie) - if (!row) return - const message = decodeRow(row) - return message.type === "assistant" && !message.time.completed ? message : undefined - }) - }, - getAssistant(messageID) { - return Effect.gen(function* () { - const row = yield* db - .select() - .from(SessionMessageTable) - .where( - and( - eq(SessionMessageTable.id, messageID), - eq(SessionMessageTable.session_id, event.data.sessionID), - eq(SessionMessageTable.type, "assistant"), - ), - ) - .get() - .pipe(Effect.orDie) - if (!row) return - const message = decodeRow(row) - return message.type === "assistant" ? message : undefined - }) - }, - getShell(shellID) { - return Effect.gen(function* () { - const row = yield* db - .select() - .from(SessionMessageTable) - .where( - and( - eq(SessionMessageTable.session_id, event.data.sessionID), - eq(SessionMessageTable.type, "shell"), - sql`json_extract(${SessionMessageTable.data}, '$.shellID') = ${shellID}`, - ), - ) - .orderBy(desc(SessionMessageTable.seq)) - .limit(1) - .get() - .pipe(Effect.orDie) - if (!row) return - const message = decodeRow(row) - return message.type === "shell" ? message : undefined - }) - }, - getCompaction() { - return Effect.gen(function* () { - const row = yield* db - .select() - .from(SessionMessageTable) - .where( - and( - eq(SessionMessageTable.session_id, event.data.sessionID), - eq(SessionMessageTable.type, "compaction"), - sql`json_extract(${SessionMessageTable.data}, '$.status') = 'running'`, - ), - ) - .orderBy(desc(SessionMessageTable.seq)) - .limit(1) - .get() - .pipe(Effect.orDie) - if (!row) return - const message = decodeRow(row) - return message.type === "compaction" ? message : undefined - }) - }, - updateAssistant: updateMessage, - updateShell: updateMessage, - updateCompaction: updateMessage, - appendMessage, - } - yield* SessionMessageUpdater.update(adapter, event) - }) -} - -function insertMessage(db: DatabaseService, event: SessionEvent.DurableEvent, message: SessionMessage.Info) { - const encoded = encodeMessage(message) - const { id, type, ...data } = encoded - return db - .insert(SessionMessageTable) - .values({ - id: SessionMessage.ID.make(id), - session_id: event.data.sessionID, - type, - seq: event.durable.seq, - time_created: DateTime.toEpochMillis(message.time.created), - data, - }) - .run() - .pipe(Effect.orDie) -} - function projectIdle( db: DatabaseService, event: @@ -408,7 +223,7 @@ function projectIdle( | typeof SessionEvent.Execution.Interrupted.Type, ) { return Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) if (event.type === SessionEvent.Execution.Interrupted.type && event.data.reason === "shutdown") return const time = event.created const outcome = @@ -465,7 +280,7 @@ const layer = Layer.effectDiscard( ) yield* bus.project(SessionEvent.Moved, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* db .update(SessionTable) .set({ @@ -545,7 +360,7 @@ const layer = Layer.effectDiscard( ) yield* bus.project(SessionEvent.AgentSelected, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* db .update(SessionTable) .set({ agent: event.data.agent, time_updated: event.created }) @@ -556,7 +371,7 @@ const layer = Layer.effectDiscard( ) yield* bus.project(SessionEvent.ModelSelected, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* db .update(SessionTable) .set({ model: event.data.model, time_updated: event.created }) @@ -603,7 +418,7 @@ const layer = Layer.effectDiscard( .run() .pipe(Effect.orDie) }) - yield* bus.project(SessionEvent.MessageContentUpdated, (event) => run(db, event)) + yield* bus.project(SessionEvent.MessageContentUpdated, (event) => TranscriptProjection.project(db, event)) yield* bus.project(SessionEvent.UsageRecorded, (event) => applyUsage(db, event.data.sessionID, event.data)) yield* bus.project(SessionEvent.Forked, (event) => projectFork(db, event)) yield* bus.project(SessionEvent.InboxDelivered, (event) => @@ -613,7 +428,7 @@ const layer = Layer.effectDiscard( sessionID: event.data.sessionID, }) if (input.type === "compaction" || input.type === "move") return - yield* insertMessage( + yield* TranscriptProjection.appendAtEventSequence( db, event, input.type === "user" @@ -673,47 +488,47 @@ const layer = Layer.effectDiscard( yield* bus.project(SessionEvent.Execution.Interrupted, (event) => projectIdle(db, event)) yield* bus.project(SessionEvent.InstructionsUpdated, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* InstructionState.apply(db, event.data.sessionID, event.durable.seq, event.data.delta) }), ) - yield* bus.project(SessionEvent.Synthetic, (event) => run(db, event)) - yield* bus.project(SessionEvent.Skill.Activated, (event) => run(db, event)) - yield* bus.project(SessionEvent.Shell.Started, (event) => run(db, event)) - yield* bus.project(SessionEvent.Shell.Ended, (event) => run(db, event)) - yield* bus.project(SessionEvent.Step.Started, (event) => run(db, event)) - yield* bus.project(SessionEvent.Step.Streamed, (event) => run(db, event)) + yield* bus.project(SessionEvent.Synthetic, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Skill.Activated, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Shell.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Shell.Ended, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Step.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Step.Streamed, (event) => TranscriptProjection.project(db, event)) yield* bus.project(SessionEvent.Step.Ended, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* applyUsage(db, event.data.sessionID, event.data) }), ) yield* bus.project(SessionEvent.Step.Failed, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) if (event.data.cost !== undefined && event.data.tokens !== undefined) yield* applyUsage(db, event.data.sessionID, { cost: event.data.cost, tokens: event.data.tokens }) }), ) - yield* bus.project(SessionEvent.Text.Started, (event) => run(db, event)) - yield* bus.project(SessionEvent.Text.Ended, (event) => run(db, event)) - yield* bus.project(SessionEvent.Tool.Input.Started, (event) => run(db, event)) - yield* bus.project(SessionEvent.Tool.Input.Ended, (event) => run(db, event)) - yield* bus.project(SessionEvent.Tool.Called, (event) => run(db, event)) - yield* bus.project(SessionEvent.Tool.Success, (event) => run(db, event)) - yield* bus.project(SessionEvent.Tool.Failed, (event) => run(db, event)) - yield* bus.project(SessionEvent.Reasoning.Started, (event) => run(db, event)) - yield* bus.project(SessionEvent.Reasoning.Ended, (event) => run(db, event)) - yield* bus.project(SessionEvent.RetryScheduled, (event) => run(db, event)) - yield* bus.project(SessionEvent.Compaction.Started, (event) => run(db, event)) + yield* bus.project(SessionEvent.Text.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Text.Ended, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Tool.Input.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Tool.Input.Ended, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Tool.Called, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Tool.Success, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Tool.Failed, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Reasoning.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Reasoning.Ended, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.RetryScheduled, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Compaction.Started, (event) => TranscriptProjection.project(db, event)) yield* bus.project(SessionEvent.Compaction.Ended, (event) => Effect.gen(function* () { - yield* run(db, event) + yield* TranscriptProjection.project(db, event) yield* InstructionState.advanceEpoch(db, event.data.sessionID, event.durable.seq) }), ) - yield* bus.project(SessionEvent.Compaction.Failed, (event) => run(db, event)) + yield* bus.project(SessionEvent.Compaction.Failed, (event) => TranscriptProjection.project(db, event)) yield* bus.project(SessionEvent.RevertEvent.Staged, (event) => Effect.gen(function* () { const revert = event.data.revert diff --git a/packages/core/test/transcript-projection.test.ts b/packages/core/test/transcript-projection.test.ts new file mode 100644 index 000000000000..87f78e4bb7e9 --- /dev/null +++ b/packages/core/test/transcript-projection.test.ts @@ -0,0 +1,216 @@ +import { describe, expect } from "bun:test" +import { DateTime, Effect, Schema } from "effect" +import { eq } from "drizzle-orm" +import { Database } from "@opencode/core/database/database" +import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder" +import { LayerNode } from "@opencode/util/effect/layer-node" +import { Bus } from "@opencode/core/bus" +import { ProjectTable } from "@opencode/core/project/sql" +import { AbsolutePath } from "@opencode/core/schema" +import { SessionEvent } from "@opencode/core/session/event" +import { SessionMessage } from "@opencode/core/session/message" +import { TranscriptProjection } from "@opencode/core/session/projection/transcript" +import { SessionMessageTable, SessionTable } from "@opencode/core/session/sql" +import { Agent } from "@opencode/core/agent" +import { Model } from "@opencode/schema/model" +import { Provider } from "@opencode/schema/provider" +import { Project } from "@opencode/schema/project" +import { Session } from "@opencode/schema/session" +import { Shell } from "@opencode/schema/shell" +import { testEffect } from "./lib/effect" + +const it = testEffect( + AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node]), [ + Bus.node.replace(Bus.configured({ persist: true })), + ]), +) +const sessionID = Session.ID.make("ses_transcript") +const foreignID = Session.ID.make("ses_foreign") +const model = { id: Model.ID.make("model"), providerID: Provider.ID.make("provider") } +const created = DateTime.makeUnsafe(0) +const encode = Schema.encodeSync(SessionMessage.Info) +const decode = Schema.decodeUnknownSync(SessionMessage.Info) + +const seed = Effect.gen(function* () { + const db = (yield* Database.Service).db + const bus = yield* Bus.Service + yield* db + .insert(ProjectTable) + .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] }) + .run() + yield* db + .insert(SessionTable) + .values( + [sessionID, foreignID].map((id) => ({ + id, + project_id: Project.ID.global, + slug: id, + directory: "/project", + title: id, + version: "test", + })), + ) + .run() + yield* bus.project(SessionEvent.Step.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Text.Started, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Shell.Ended, (event) => TranscriptProjection.project(db, event)) + yield* bus.project(SessionEvent.Compaction.Failed, (event) => TranscriptProjection.project(db, event)) + return { db, bus } +}) + +const row = (message: SessionMessage.Info, seq: number, session = sessionID) => { + const { id, type, ...data } = encode(message) + return { + id: SessionMessage.ID.make(id), + session_id: session, + type, + seq, + time_created: DateTime.toEpochMillis(message.time.created), + data, + } +} +const assistant = (id: string, completed?: DateTime.Utc) => + SessionMessage.Assistant.make({ + id: SessionMessage.ID.make(id), + type: "assistant", + agent: Agent.defaultID, + model, + content: [], + time: { created, completed }, + }) +const messages = (db: Database.Interface["db"], session = sessionID) => + db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.session_id, session)) + .all() + .pipe(Effect.map((rows) => rows.map((row) => decode({ ...row.data, id: row.id, type: row.type })))) + +describe("TranscriptProjection", () => { + it.effect("settles only the newest incomplete assistant by sequence", () => + Effect.gen(function* () { + const { db, bus } = yield* seed + yield* db + .insert(SessionMessageTable) + .values([row(assistant("msg_older"), 1), row(assistant("msg_newer"), 2)]) + .run() + const event = yield* bus.publish(SessionEvent.Step.Started, { + sessionID, + assistantMessageID: SessionMessage.ID.make("msg_next"), + agent: Agent.defaultID, + model, + started: 0, + }) + const stored = (yield* messages(db)).filter( + (message): message is SessionMessage.Assistant => message.type === "assistant", + ) + expect(stored.find((message) => message.id === "msg_older")?.time.completed).toBeUndefined() + expect(stored.find((message) => message.id === "msg_newer")?.time.completed).toEqual( + DateTime.makeUnsafe(event.created), + ) + const appended = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.id, SessionMessage.ID.make("msg_next"))) + .get() + expect(appended?.seq).toBe(event.durable.seq) + }), + ) + + it.effect("does not fall back to an older incomplete assistant after a completed one", () => + Effect.gen(function* () { + const { db, bus } = yield* seed + yield* db + .insert(SessionMessageTable) + .values([row(assistant("msg_older"), 1), row(assistant("msg_completed", DateTime.makeUnsafe(1)), 2)]) + .run() + yield* bus.publish(SessionEvent.Step.Started, { + sessionID, + assistantMessageID: SessionMessage.ID.make("msg_next"), + agent: Agent.defaultID, + model, + started: 0, + }) + const stored = (yield* messages(db)).filter( + (message): message is SessionMessage.Assistant => message.type === "assistant", + ) + expect(stored.find((message) => message.id === "msg_older")?.time.completed).toBeUndefined() + expect(stored.find((message) => message.id === "msg_completed")?.time.completed).toEqual(DateTime.makeUnsafe(1)) + }), + ) + + it.effect("does not update an assistant belonging to another session", () => + Effect.gen(function* () { + const { db, bus } = yield* seed + const foreign = assistant("msg_foreign") + yield* db + .insert(SessionMessageTable) + .values(row(foreign, 1, foreignID)) + .run() + yield* bus.publish(SessionEvent.Text.Started, { sessionID, assistantMessageID: foreign.id, ordinal: 0 }) + expect(yield* messages(db, foreignID)).toEqual([foreign]) + expect(yield* messages(db)).toEqual([]) + }), + ) + + it.effect("keeps shell lookup scoped to the event session", () => + Effect.gen(function* () { + const { db, bus } = yield* seed + const foreign = SessionMessage.Shell.make({ + id: SessionMessage.ID.make("msg_shell"), + type: "shell", + shellID: Shell.ID.make("sh_shared"), + command: "pwd", + status: "running", + time: { created }, + }) + yield* db + .insert(SessionMessageTable) + .values(row(foreign, 1, foreignID)) + .run() + yield* bus.publish(SessionEvent.Shell.Ended, { + sessionID, + shell: Shell.Info.make({ + id: foreign.shellID, + status: "exited", + command: "pwd", + cwd: "/project", + shell: "/bin/sh", + file: "/tmp/sh_shared.out", + exit: 0, + metadata: {}, + time: { started: 0, completed: 1 }, + }), + output: { output: "done", cursor: 4, size: 4, truncated: false }, + }) + expect(yield* messages(db, foreignID)).toEqual([foreign]) + expect(yield* messages(db)).toEqual([]) + }), + ) + + it.effect("does not finish a running compaction belonging to another session", () => + Effect.gen(function* () { + const { db, bus } = yield* seed + const foreign = SessionMessage.CompactionRunning.make({ + id: SessionMessage.ID.make("msg_compaction"), + type: "compaction", + status: "running", + reason: "auto", + summary: "", + recent: "", + time: { created }, + }) + yield* db + .insert(SessionMessageTable) + .values(row(foreign, 1, foreignID)) + .run() + yield* bus.publish(SessionEvent.Compaction.Failed, { + sessionID, + reason: "auto", + error: { type: "compaction.failed", message: "Failed" }, + }) + expect(yield* messages(db, foreignID)).toEqual([foreign]) + expect(yield* messages(db)).toMatchObject([{ type: "compaction", status: "failed" }]) + }), + ) +})