export * as SessionProjector from "./projector" import { and, asc, desc, eq, gt, inArray, lt, or, sql } from "drizzle-orm" import { DateTime, Effect, Layer, Schema } from "effect" import { Database } from "../database/database" import { EventV2 } from "../event" import { makeGlobalNode } from "../effect/app-node" import { SessionEvent } from "./event" import { SessionV1 } from "../v1/session" import { WorkspaceTable } from "../control-plane/workspace.sql" import { SessionMessage } from "./message" import { SessionMessageUpdater } from "./message-updater" import { SessionInput } from "./input" import { WorkspaceV2 } from "../workspace" import { SessionContextCheckpoint } from "./context-checkpoint" import { MessageTable, PartTable, SessionContextCheckpointTable, SessionInputTable, SessionMessageTable, SessionTable, } from "./sql" import type { DeepMutable } from "../schema" import { Slug } from "../util/slug" type DatabaseService = Database.Interface["db"] type MessageEvent = Exclude const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Message) const encodeMessage = Schema.encodeSync(SessionMessage.Message) export class SessionAlreadyProjected extends Error {} type Usage = { cost: number tokens: { input: number output: number reasoning: number cache: { read: number; write: number } } } const ForkBatchSize = 500 const emptyUsage = (): Usage => ({ cost: 0, tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, }) const forkTitle = (value: string) => { const match = value.match(/^(.+) \(fork #(\d+)\)$/) if (match) return `${match[1]} (fork #${Number.parseInt(match[2], 10) + 1})` return `${value} (fork #1)` } function usage(part: (typeof SessionV1.Event.PartUpdated.Type)["data"]["part"] | unknown): Usage | undefined { if (typeof part !== "object" || part === null) return undefined const value = part as Record if (value.type !== "step-finish") return undefined if (!("cost" in value) || !("tokens" in value)) return undefined return { cost: value.cost as Usage["cost"], tokens: value.tokens as Usage["tokens"] } } function addUsage(target: Usage, value: Usage) { target.cost += value.cost target.tokens.input += value.tokens.input target.tokens.output += value.tokens.output target.tokens.reasoning += value.tokens.reasoning target.tokens.cache.read += value.tokens.cache.read target.tokens.cache.write += value.tokens.cache.write } function messageUsage(row: typeof SessionMessageTable.$inferSelect): Usage | undefined { if (row.type !== "assistant") return undefined const message = decodeMessage({ ...row.data, id: row.id, type: row.type }) if (message.type !== "assistant" || message.cost === undefined || message.tokens === undefined) return undefined return { cost: message.cost, tokens: message.tokens } } function sessionRow(info: SessionV1.SessionInfo): typeof SessionTable.$inferInsert { return { id: info.id, project_id: info.projectID, workspace_id: info.workspaceID ?? null, parent_id: info.parentID, slug: info.slug, directory: info.directory, path: info.path, title: info.title, agent: info.agent, model: info.model, version: info.version, share_url: info.share?.url, summary_additions: info.summary?.additions, summary_deletions: info.summary?.deletions, summary_files: info.summary?.files, summary_diffs: info.summary?.diffs ? [...info.summary.diffs] : undefined, metadata: info.metadata, cost: info.cost ?? 0, tokens_input: (info.tokens ?? { input: 0 }).input, tokens_output: (info.tokens ?? { output: 0 }).output, tokens_reasoning: (info.tokens ?? { reasoning: 0 }).reasoning, tokens_cache_read: (info.tokens ?? { cache: { read: 0 } }).cache.read, tokens_cache_write: (info.tokens ?? { cache: { write: 0 } }).cache.write, revert: info.revert ? { ...info.revert, messageID: SessionMessage.ID.make(info.revert.messageID) } : null, permission: info.permission ? [...info.permission] : undefined, time_created: info.time.created, time_updated: info.time.updated, time_compacting: info.time.compacting, time_archived: info.time.archived, } } function messageData( info: (typeof SessionV1.Event.MessageUpdated.Type)["data"]["info"], ): typeof MessageTable.$inferInsert.data { const { id: _, sessionID: __, ...rest } = info return rest as DeepMutable } function partData(part: (typeof SessionV1.Event.PartUpdated.Type)["data"]["part"]): typeof PartTable.$inferInsert.data { const { id: _, messageID: __, sessionID: ___, ...rest } = part return rest as DeepMutable } function applyUsage( db: DatabaseService, sessionID: (typeof SessionV1.Event.MessageUpdated.Type)["data"]["sessionID"], value: Usage, sign = 1, ) { return db .update(SessionTable) .set({ cost: sql`${SessionTable.cost} + ${value.cost * sign}`, tokens_input: sql`${SessionTable.tokens_input} + ${value.tokens.input * sign}`, tokens_output: sql`${SessionTable.tokens_output} + ${value.tokens.output * sign}`, tokens_reasoning: sql`${SessionTable.tokens_reasoning} + ${value.tokens.reasoning * sign}`, tokens_cache_read: sql`${SessionTable.tokens_cache_read} + ${value.tokens.cache.read * sign}`, tokens_cache_write: sql`${SessionTable.tokens_cache_write} + ${value.tokens.cache.write * sign}`, time_updated: sql`${SessionTable.time_updated}`, }) .where(eq(SessionTable.id, sessionID)) .run() .pipe(Effect.orDie) } const projectFork = Effect.fn("SessionProjector.projectFork")(function* ( db: DatabaseService, event: typeof SessionEvent.Forked.Type, ) { const parent = yield* db .select() .from(SessionTable) .where(eq(SessionTable.id, event.data.parentID)) .get() .pipe(Effect.orDie) if (!parent) return yield* Effect.die(`Fork parent session not found: ${event.data.parentID}`) const boundary = event.data.messageID ? yield* db .select({ seq: SessionMessageTable.seq }) .from(SessionMessageTable) .where( and( eq(SessionMessageTable.session_id, event.data.parentID), eq(SessionMessageTable.id, event.data.messageID), ), ) .get() .pipe(Effect.orDie) : undefined if (event.data.messageID && !boundary) return yield* Effect.die(`Fork boundary message not found: ${event.data.messageID}`) const copied = yield* db .select({ seq: SessionMessageTable.seq }) .from(SessionMessageTable) .where( and( eq(SessionMessageTable.session_id, event.data.parentID), boundary === undefined ? undefined : lt(SessionMessageTable.seq, boundary.seq), ), ) .orderBy(desc(SessionMessageTable.seq)) .limit(1) .get() .pipe(Effect.orDie) const copiedSeq = copied?.seq ?? 0 const stored = yield* db .insert(SessionTable) .values({ id: event.data.sessionID, parent_id: event.data.parentID, project_id: parent.project_id, workspace_id: parent.workspace_id, slug: Slug.create(), directory: parent.directory, path: parent.path, title: forkTitle(parent.title), agent: parent.agent, model: parent.model, version: parent.version, cost: 0, tokens_input: 0, tokens_output: 0, tokens_reasoning: 0, tokens_cache_read: 0, tokens_cache_write: 0, time_created: DateTime.toEpochMillis(event.data.timestamp), time_updated: DateTime.toEpochMillis(event.data.timestamp), }) .onConflictDoNothing() .returning({ sessionID: SessionTable.id }) .get() .pipe(Effect.orDie) if (!stored) return yield* Effect.die(new SessionAlreadyProjected()) // The fork inherits the parent's transcript, so it inherits the context // checkpoint that transcript was built against: copied message seqs keep // folding at the same baseline horizon. const checkpoint = yield* db .select() .from(SessionContextCheckpointTable) .where(eq(SessionContextCheckpointTable.session_id, event.data.parentID)) .get() .pipe(Effect.orDie) if (checkpoint) { yield* db .insert(SessionContextCheckpointTable) .values({ ...checkpoint, session_id: event.data.sessionID }) .run() .pipe(Effect.orDie) } const usage = emptyUsage() let cursor = -1 while (true) { const rows = yield* db .select() .from(SessionMessageTable) .where( and( eq(SessionMessageTable.session_id, event.data.parentID), gt(SessionMessageTable.seq, cursor), copiedSeq === 0 ? undefined : lt(SessionMessageTable.seq, copiedSeq + 1), ), ) .orderBy(asc(SessionMessageTable.seq)) .limit(ForkBatchSize) .all() .pipe(Effect.orDie) if (rows.length === 0) break const idMap = new Map(rows.map((row) => [row.id, SessionMessage.ID.create()])) yield* db .insert(SessionMessageTable) .values( rows.map((row) => { const id = idMap.get(row.id) if (!id) throw new Error(`Fork message ID mapping missing: ${row.id}`) return { id, session_id: event.data.sessionID, type: row.type, seq: row.seq, time_created: row.time_created, time_updated: row.time_updated, data: row.type === "synthetic" ? { ...row.data, sessionID: event.data.sessionID } : row.data, } }), ) .run() .pipe(Effect.orDie) const inputRows = yield* db .select() .from(SessionInputTable) .where( and( eq(SessionInputTable.session_id, event.data.parentID), inArray( SessionInputTable.id, rows.map((row) => row.id), ), ), ) .all() .pipe(Effect.orDie) if (inputRows.length > 0) { yield* db .insert(SessionInputTable) .values( inputRows.flatMap((row) => { const id = idMap.get(row.id) return id ? [ { id, session_id: event.data.sessionID, prompt: row.prompt, delivery: row.delivery, admitted_seq: row.admitted_seq, promoted_seq: row.promoted_seq, time_created: row.time_created, }, ] : [] }), ) .run() .pipe(Effect.orDie) } for (const row of rows) { const value = messageUsage(row) if (value) addUsage(usage, value) } cursor = rows.at(-1)!.seq } yield* db .update(SessionTable) .set({ cost: usage.cost, tokens_input: usage.tokens.input, tokens_output: usage.tokens.output, tokens_reasoning: usage.tokens.reasoning, tokens_cache_read: usage.tokens.cache.read, tokens_cache_write: usage.tokens.cache.write, }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie) if (copiedSeq > 0) yield* EventV2.reserveSequence(db, event.data.sessionID, copiedSeq) }) 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.Message) => { if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence") 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.Message) => insertMessage(db, event, message) const adapter: SessionMessageUpdater.Adapter = { getCurrentAssistant() { return Effect.gen(function* () { // A newer turn 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 }) }, getCurrentShell(callID) { return Effect.gen(function* () { const rows = yield* db .select() .from(SessionMessageTable) .where(and(eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.type, "shell"))) .orderBy(desc(SessionMessageTable.seq)) .all() .pipe(Effect.orDie) return rows .map(decodeRow) .find((message): message is SessionMessage.Shell => message.type === "shell" && message.callID === callID) }) }, updateAssistant: updateMessage, updateShell: updateMessage, appendMessage, } yield* SessionMessageUpdater.update(adapter, event) }) } function insertMessage(db: DatabaseService, event: SessionEvent.Event, message: SessionMessage.Message) { if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence") 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) } const layer = Layer.effectDiscard( Effect.gen(function* () { const events = yield* EventV2.Service const { db } = yield* Database.Service yield* events.project(SessionV1.Event.Created, (event) => Effect.gen(function* () { const stored = yield* db .insert(SessionTable) .values(sessionRow(event.data.info)) .onConflictDoNothing() .returning({ sessionID: SessionTable.id }) .get() .pipe(Effect.orDie) if (!stored) return yield* Effect.die(new SessionAlreadyProjected()) if (event.data.info.workspaceID) { yield* db .update(WorkspaceTable) .set({ time_used: Date.now() }) .where(eq(WorkspaceTable.id, event.data.info.workspaceID)) .run() .pipe(Effect.orDie) } }), ) yield* events.project(SessionV1.Event.Updated, (event) => db .update(SessionTable) .set(sessionRow(event.data.info)) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie), ) yield* events.project(SessionEvent.Moved, (event) => Effect.gen(function* () { yield* db .update(SessionTable) .set({ directory: event.data.location.directory, path: event.data.subdirectory, workspace_id: event.data.location.workspaceID ? WorkspaceV2.ID.make(event.data.location.workspaceID) : null, time_updated: DateTime.toEpochMillis(event.data.timestamp), }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie) yield* SessionContextCheckpoint.reset(db, event.data.sessionID) }), ) yield* events.project(SessionV1.Event.Deleted, (event) => db.delete(SessionTable).where(eq(SessionTable.id, event.data.sessionID)).run().pipe(Effect.orDie), ) yield* events.project(SessionV1.Event.MessageUpdated, (event) => Effect.gen(function* () { const time_created = event.data.info.time.created const id = event.data.info.id const sessionID = event.data.info.sessionID const data = messageData(event.data.info) yield* db .insert(MessageTable) .values({ id, session_id: sessionID, time_created, data }) .onConflictDoUpdate({ target: MessageTable.id, set: { data } }) .run() .pipe(Effect.orDie) }), ) yield* events.project(SessionV1.Event.MessageRemoved, (event) => Effect.gen(function* () { const rows = yield* db .select() .from(PartTable) .where(and(eq(PartTable.message_id, event.data.messageID), eq(PartTable.session_id, event.data.sessionID))) .all() .pipe(Effect.orDie) for (const row of rows) { const previous = usage(row.data) if (previous) yield* applyUsage(db, event.data.sessionID, previous, -1) } yield* db .delete(MessageTable) .where(and(eq(MessageTable.id, event.data.messageID), eq(MessageTable.session_id, event.data.sessionID))) .run() .pipe(Effect.orDie) }), ) yield* events.project(SessionV1.Event.PartRemoved, (event) => Effect.gen(function* () { const row = yield* db .select() .from(PartTable) .where(and(eq(PartTable.id, event.data.partID), eq(PartTable.session_id, event.data.sessionID))) .get() .pipe(Effect.orDie) const previous = row && usage(row.data) if (previous) yield* applyUsage(db, event.data.sessionID, previous, -1) yield* db .delete(PartTable) .where(and(eq(PartTable.id, event.data.partID), eq(PartTable.session_id, event.data.sessionID))) .run() .pipe(Effect.orDie) }), ) yield* events.project(SessionV1.Event.PartUpdated, (event) => Effect.gen(function* () { const id = event.data.part.id const messageID = event.data.part.messageID const sessionID = event.data.part.sessionID const data = partData(event.data.part) const row = yield* db.select().from(PartTable).where(eq(PartTable.id, id)).get().pipe(Effect.orDie) yield* db .insert(PartTable) .values({ id, message_id: messageID, session_id: sessionID, time_created: event.data.time, data }) .onConflictDoUpdate({ target: PartTable.id, set: { data } }) .run() .pipe(Effect.orDie) const previous = row && usage(row.data) const next = usage(event.data.part) if (previous) yield* applyUsage(db, row.session_id, previous, -1) if (next) yield* applyUsage(db, sessionID, next) }), ) yield* events.project(SessionEvent.AgentSwitched, (event) => 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, Effect.andThen(run(db, event))), ) yield* events.project(SessionEvent.ModelSwitched, (event) => Effect.gen(function* () { 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* run(db, event) }), ) yield* events.project(SessionEvent.Renamed, (event) => db .update(SessionTable) .set({ title: event.data.title, time_updated: DateTime.toEpochMillis(event.data.timestamp) }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie), ) yield* events.project(SessionEvent.Forked, (event) => projectFork(db, event)) yield* events.project(SessionEvent.Prompted, (event) => Effect.gen(function* () { if (event.durable === undefined) return yield* Effect.die("Durable Session event is missing aggregate sequence") yield* SessionInput.projectPrompted(db, { id: event.data.messageID, sessionID: event.data.sessionID, prompt: event.data.prompt, delivery: event.data.delivery, timeCreated: event.data.timestamp, promotedSeq: event.durable.seq, }) yield* run(db, event) }), ) yield* events.project(SessionEvent.PromptAdmitted, (event) => Effect.gen(function* () { if (event.durable === undefined) return yield* Effect.die("Durable Session event is missing aggregate sequence") yield* SessionInput.projectAdmitted(db, { admittedSeq: event.durable.seq, id: event.data.messageID, sessionID: event.data.sessionID, prompt: event.data.prompt, delivery: event.data.delivery, timeCreated: event.data.timestamp, }) }), ) yield* events.project(SessionEvent.ContextUpdated, (event) => run(db, event)) yield* events.project(SessionEvent.Synthetic, (event) => run(db, event)) yield* events.project(SessionEvent.Skill.Activated, (event) => insertMessage(db, event, { id: event.data.messageID, type: "skill", name: event.data.name, text: event.data.text, time: { created: event.data.timestamp }, }), ) 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.Progress, (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.Ended, (event) => run(db, event)) yield* events.project(SessionEvent.RevertEvent.Staged, (event) => db .update(SessionTable) .set({ revert: { ...event.data.revert, files: event.data.revert.files ? [...event.data.revert.files] : undefined }, time_updated: DateTime.toEpochMillis(event.data.timestamp), }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie, Effect.asVoid), ) yield* events.project(SessionEvent.RevertEvent.Cleared, (event) => db .update(SessionTable) .set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie, Effect.asVoid), ) yield* events.project(SessionEvent.RevertEvent.Committed, (event) => Effect.gen(function* () { const boundary = yield* db .select({ seq: SessionMessageTable.seq }) .from(SessionMessageTable) .where( and( eq(SessionMessageTable.session_id, event.data.sessionID), eq(SessionMessageTable.id, event.data.messageID), ), ) .get() .pipe(Effect.orDie) if (!boundary) return yield* Effect.die(`Revert boundary message not found: ${event.data.messageID}`) yield* db .delete(SessionMessageTable) .where( and(eq(SessionMessageTable.session_id, event.data.sessionID), gt(SessionMessageTable.seq, boundary.seq)), ) .run() .pipe(Effect.orDie) yield* db .delete(SessionInputTable) .where( and( eq(SessionInputTable.session_id, event.data.sessionID), or(gt(SessionInputTable.admitted_seq, boundary.seq), gt(SessionInputTable.promoted_seq, boundary.seq)), ), ) .run() .pipe(Effect.orDie) yield* db .update(SessionTable) .set({ revert: null, time_updated: DateTime.toEpochMillis(event.data.timestamp) }) .where(eq(SessionTable.id, event.data.sessionID)) .run() .pipe(Effect.orDie) yield* SessionContextCheckpoint.reset(db, event.data.sessionID) }), ) }), ) export const node = makeGlobalNode({ name: "session-projector", layer, deps: [EventV2.node, Database.node] })