diff --git a/docs/adr/0019-session-surface.md b/docs/adr/0019-session-surface.md new file mode 100644 index 0000000000..159db911b1 --- /dev/null +++ b/docs/adr/0019-session-surface.md @@ -0,0 +1,63 @@ +# ADR 0019: Session surface — a linked list over the event log for LLM message derivation + +Status: accepted (2026-06-17) + +## Context + +The `Session` event log is the single source of truth ([ADR 0003](0003-event-sourced-sessions.md)), but the only view over it was `deriveMessages()` — a linear scan that filtered and transformed raw events into `Message[]`. This creates problems for session-history-manipulating plugins (compaction, tool-call result pruning, etc.). Without a central mechanism, each plugin would need to wrap `agent/request` to rewrite the message list — a pattern that suffers from listener-ordering fragility, provides no durable record of what was changed, and forces repeated changes to the core `deriveMessages()` whenever a new manipulation is added. A central hub in the `session` package, with a provenance-recording mechanism and enough flexibility for future plugins to manipulate session history through a stable API, lays a solid foundation for plugin development. + +## Decision + +Add a **surface** — a derived, cached linked list of "surface nodes" (the subset of events that produce LLM messages) — maintained by `surfaceOp` markers in the event log. + +### Two new top-level fields on `SessionEvent` + +Every `SessionEvent` gains two optional fields (structural metadata, like `seq`/`time`): + +- **`sourceEventSeqs?: number[]`** — seq numbers of events that are provenance sources (e.g., the `assistant/chunk` seqs that built an `assistant/message`, or the surface nodes shadowed by a compaction marker). Provenance is a core design principle; without it, the replace-range operation cannot be validated on replay. +- **`surfaceOp?: SurfaceOp`** — how this event entered the surface. Absent for non-surface events. + +### SurfaceOp: two operations + +```ts +export type SurfaceOp = + | 'append' // normal tail append + | { op: 'replace'; start: number; end: number } // shadow [start, end] inclusive +``` + +1. **Append** — add a new node to the tail. Used by `user/message`, `assistant/message`, `tool/result`, `context/message`, `steering/message`. The loop passes `surfaceOp: 'append'` on all such appends, and `sourceEventSeqs` where applicable (e.g., `assistant/message` records its `assistant/chunk` sources; `tool/result` records its `tool/call` source). + +2. **Replace** — remove nodes from `start` through `end` (both inclusive) and insert a new node in their place. Both `start` and `end` must be valid surface node seqs in the current surface; `start === end` replaces a single node. The node's `sourceEventSeqs` must contain every shadowed surface node. The shadowed events remain in the log but are no longer on the surface. + +The both-ends-inclusive design was chosen over half-open `[start, endExclusive)` because the surface is a doubly-linked list — both ends are naturally named by node seqs, and single-node replacement (`start === end`) is a common case that reads naturally with inclusive semantics. + +### SurfaceManager: delta-based, not full rebuild + +A `SurfaceManager` class (private to `Session`) maintains the cached linked list. It tracks `_lastProcessedSeq` and processes only the **delta** (new events since the last access) rather than rescanning the entire log. Because the log is append-only, prior events never change — full rebuild is only needed after a wholesale log replacement (e.g., seeding). + +Why delta processing? The naive approach (a dirty flag + full rebuild on every access) would be O(N²) over a session's lifetime — every single-event append triggers a complete scan of all prior events. Delta processing is O(1) when no new events and O(new events) when new events arrive. + +`deriveMessages()` uses the surface when surface markers exist, falling back to the existing linear scan for sessions without markers (backward compatibility). + +### Persistence + +The new fields are serialized as top-level JSON properties. The JSONL backend requires zero changes — `JSON.stringify`/`JSON.parse` preserve everything transparently. The SQLite backend adds two nullable TEXT columns (`source_event_seqs`, `surface_op`) with an `ALTER TABLE` migration (SCHEMA_VERSION 1 → 2). The session format `version` stays at 1 — the new fields are optional and backward-compatible. + +### Crash recovery + +The `repair.ts` module synthesizes `tool/result` closers for orphaned tool calls after a crash. These closers carry `surfaceOp: 'append'` and `sourceEventSeqs` pointing to the orphaned `tool/call` event, so the rehydrated surface is valid. + +### Invariants + +The dev-mode invariants plugin validates: `sourceEventSeqs` references (non-empty, no duplicates, references earlier events, references known seqs) and `surfaceOp` (replace start ≤ end). + +## Consequences + +- **`packages/session`**: New `surface.ts` (`SurfaceManager`), new types (`SurfaceOp`, `SurfaceAppendOpts`), new fields on `SessionEvent`, modified `append()` (third optional `SurfaceAppendOpts` param), refactored `deriveMessages()` (surface path + legacy fallback), surface-aware `repair.ts`. +- **`packages/agent-loop`**: All surface-capable appends pass surface opts. Chunk seqs are collected for `assistant/message` provenance; `tool/call` seqs are captured for `tool/result` provenance. +- **`packages/session-persistence-sqlite`**: Schema migration v1 → v2 (two new nullable TEXT columns). +- **`packages/invariants`**: Surface-related validation rules. +- **`packages/session-persistence-jsonl`**: No changes required. +- **`packages/session-persistence`**: Abstract interface unchanged. + +The surface is the foundation for future compaction: a compaction plugin appends a new event (e.g., `compaction/marker`, added to `SessionEventMap` via declaration merging) with `surfaceOp: { op: 'replace', start, end }` and `sourceEventSeqs` covering the shadowed nodes. Replay preserves the compaction decision deterministically. diff --git a/docs/adr/README.md b/docs/adr/README.md index d9671c6652..8f3180979e 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -30,3 +30,4 @@ Do NOT write an ADR for: a mechanical or local choice (a variable name, a one-fi | [0016](0016-pnpm-over-yarn.md) | pnpm as the package manager instead of Yarn 4 | accepted | | [0017](0017-turn-enclosure-invariant.md) | Every session event is enclosed in a turn | accepted | | [0018](0018-session-persistence.md) | Session persistence as an abstract service over `SessionEvent` | accepted | +| [0019](0019-session-surface.md) | Session surface — a linked list over the event log for LLM message derivation | accepted | diff --git a/packages/agent-loop/src/agent.ts b/packages/agent-loop/src/agent.ts index 64576186c3..0735784616 100644 --- a/packages/agent-loop/src/agent.ts +++ b/packages/agent-loop/src/agent.ts @@ -78,7 +78,7 @@ export class LoopAgent implements Agent { // A turn is open in the LOG (decided from the log, not agent status — // status can be `running` with no turn open): the context/message is // turn-enclosed by that turn, so append it directly. - this.session.append('context/message', { content, source }) + this.session.append('context/message', { content, source }, { surfaceOp: 'append' }) return } // No turn open: wrap the injection in a one-shot turn so every event stays @@ -95,7 +95,7 @@ export class LoopAgent implements Agent { // can't happen for our fixed trigger — no turn was opened and none is owed.) try { this.session.append('turn/start', { turn, trigger: { kind: 'injection', source } }) - this.session.append('context/message', { content, source }) + this.session.append('context/message', { content, source }, { surfaceOp: 'append' }) } finally { // Close the turn if turn/start made it into the log. Contain a throwing // turn/end listener: Session.append pushes before notifying, so a throw diff --git a/packages/agent-loop/src/loop.ts b/packages/agent-loop/src/loop.ts index d2b1ed270b..f1bb879626 100644 --- a/packages/agent-loop/src/loop.ts +++ b/packages/agent-loop/src/loop.ts @@ -276,7 +276,7 @@ async function runTurn(ctx: Context, agent: LoopAgent, handle: LoopHandle, turn: // every event in the log is turn-enclosed. turn/end is now owed, so a throw // while appending these is caught below and the turn is still closed. for (const message of queued) { - session.append('user/message', { content: message.content, source: message.source }) + session.append('user/message', { content: message.content, source: message.source }, { surfaceOp: 'append' }) } ctx.emit('agent/turn-start', agent, turn) @@ -410,7 +410,7 @@ async function runTurn(ctx: Context, agent: LoopAgent, handle: LoopHandle, turn: function drainSteering(ctx: Context, agent: LoopAgent, turn: number): boolean { const messages = agent.inbox.drainSteering() for (const message of messages) { - agent.session.append('steering/message', { turn, content: message.content, source: message.source }) + agent.session.append('steering/message', { turn, content: message.content, source: message.source }, { surfaceOp: 'append' }) ctx.emit('agent/steering', agent, turn, message.content, message.source) } return messages.length > 0 @@ -446,10 +446,12 @@ async function runStep( // --- Model call (streaming-first; raw chunks are the replay record) --- const assembler = new BlockAssembler() + const chunkSeqs: number[] = [] for await (const chunk of ctx.llm.stream(request)) { /* v8 ignore next -- signal.reason always set by agent.abort() which provides a default */ if (signal.aborted) throw new Error(String(signal.reason ?? 'aborted')) - session.append('assistant/chunk', { turn, step, chunk }) + const chunkEvent = session.append('assistant/chunk', { turn, step, chunk }) + chunkSeqs.push(chunkEvent.seq) ctx.emit('agent/stream-chunk', agent, turn, step, chunk) assembler.push(chunk) } @@ -468,7 +470,7 @@ async function runStep( let message: Message = assembler.message() message = await ctx.waterfall('agent/step-result', agent, turn, step, message, () => Promise.resolve(message)) - session.append('assistant/message', { turn, step, content: message.content }) + session.append('assistant/message', { turn, step, content: message.content }, { surfaceOp: 'append', sourceEventSeqs: chunkSeqs }) if (assembler.usage) { session.append('usage', { turn, step, usage: assembler.usage }) } @@ -480,7 +482,7 @@ async function runStep( for (const call of toolCalls) { /* v8 ignore next -- signal.reason always set by agent.abort() which provides a default */ if (signal.aborted) throw new Error(String(signal.reason ?? 'aborted')) - session.append('tool/call', { turn, step, callId: call.id, name: call.name, arguments: call.arguments }) + const callEvent = session.append('tool/call', { turn, step, callId: call.id, name: call.name, arguments: call.arguments }) let parsedArguments: unknown try { parsedArguments = call.arguments ? JSON.parse(call.arguments) : {} @@ -506,7 +508,7 @@ async function runStep( content: result.content, isError: result.isError, ...result.error ? { error: result.error } : {}, - }) + }, { surfaceOp: 'append', sourceEventSeqs: [callEvent.seq] }) // signal CAN flip during the await above (abort() inside a tool); // the analyzer can't see through the await boundary. // signal can flip during the await above (abort() inside a tool); diff --git a/packages/invariants/src/index.ts b/packages/invariants/src/index.ts index ebb4decfd8..ac536eeceb 100644 --- a/packages/invariants/src/index.ts +++ b/packages/invariants/src/index.ts @@ -62,6 +62,8 @@ interface SessionTrace { * `step/end` — a result must arrive in the same step as its call. */ pendingCalls: Set + /** Every seq seen so far — validates `sourceEventSeqs` references. */ + knownSeqs: Set } /** @@ -105,6 +107,30 @@ function checkEvent(trace: SessionTrace, event: SessionEvent): void { } trace.lastSeq = event.seq + // --- Surface invariants --- + if (event.sourceEventSeqs !== undefined) { + if (event.sourceEventSeqs.length === 0) { + throw new InvariantError('sourceEventSeqs must not be empty when present') + } + const unique = new Set(event.sourceEventSeqs) + if (unique.size !== event.sourceEventSeqs.length) { + throw new InvariantError('sourceEventSeqs must not contain duplicates') + } + for (const ref of event.sourceEventSeqs) { + if (ref >= event.seq) { + throw new InvariantError(`sourceEventSeqs must reference earlier events: ${ref} >= current seq ${event.seq}`) + } + if (!trace.knownSeqs.has(ref)) { + throw new InvariantError(`sourceEventSeqs references unknown seq ${ref}`) + } + } + } + if (event.surfaceOp !== undefined && typeof event.surfaceOp !== 'string') { + if (event.surfaceOp.start > event.surfaceOp.end) { + throw new InvariantError(`surface replace: start ${event.surfaceOp.start} must be <= end ${event.surfaceOp.end}`) + } + } + // Boundary/step-scoped events have explicit cases; every OTHER event type — // including plugin-added (merge-extensible) SessionEventMap keys — is caught // by the `default` and must be turn-enclosed (ADR 0017). No assertNever: an @@ -185,6 +211,8 @@ function checkEvent(trace: SessionTrace, event: SessionEvent): void { break } } + // Track every seq seen — used above to validate sourceEventSeqs references. + trace.knownSeqs.add(event.seq) } /** Legal agent status transitions (the only state machine the loop guarantees). */ @@ -216,7 +244,7 @@ export function apply(ctx: Context, config: Config = {}): void { // (re-)apply seeds the baseline, so a reload never produces a false positive. const lastStatus = new WeakMap() - const freshTrace = (): SessionTrace => ({ lastSeq: -1, openTurn: null, openStep: null, pendingCalls: new Set() }) + const freshTrace = (): SessionTrace => ({ lastSeq: -1, openTurn: null, openStep: null, pendingCalls: new Set(), knownSeqs: new Set() }) /** Build (or rebuild) a session's trace by replaying its whole log; freeze it. */ const seedSession = (session: Session): SessionTrace => { diff --git a/packages/invariants/tests/invariants.spec.ts b/packages/invariants/tests/invariants.spec.ts index b96d172c22..69d63aed4e 100644 --- a/packages/invariants/tests/invariants.spec.ts +++ b/packages/invariants/tests/invariants.spec.ts @@ -387,3 +387,115 @@ describe('HMR safety', () => { expect(Object.isFrozen(session.events[0])).toBe(false) }) }) + +describe('surface invariants', () => { + it('accepts well-formed surface metadata', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + // Events must be turn-enclosed and step-scoped events need an open step. + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('step/start', { turn: 1, step: 1 }) + expect(() => { + session.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [1] }) + }).not.toThrow() + }) + + it('accepts replace surface op', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('step/start', { turn: 1, step: 1 }) + session.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: { op: 'replace', start: 2, end: 2 }, sourceEventSeqs: [2] }) + // no throw — well-formed replace op + }) + + it('rejects empty sourceEventSeqs', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [] }) + }).toThrow(InvariantError) + }) + + it('rejects duplicate sourceEventSeqs', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [1, 1] }) + }).toThrow(/must not contain duplicates/) + }) + + it('rejects sourceEventSeqs referencing the event itself (self-reference)', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) // seq 0 + // The next event is seq 1. Referencing its own seq fails on "must reference + // earlier events" (the check order is: earlier first, then unknown). + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [1] }) + }).toThrow(/must reference earlier/) + }) + + it('accepts sourceEventSeqs referencing a valid earlier event', async () => { + // Positive test: ref < current seq and ref is in knownSeqs → passes. + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('step/start', { turn: 1, step: 1 }) + // seqs so far: 0, 1. The next event at seq 2 references seq 1 → valid. + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [1] }) + }).not.toThrow() + }) + + it('rejects sourceEventSeqs referencing a far-future seq', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [99] }) + }).toThrow(/must reference earlier/) + }) + + it('rejects sourceEventSeqs referencing unknown seq (gap in event log)', async () => { + // The unknown-seq check fires when a ref passes the "earlier" test but is + // not in knownSeqs — only possible with a gap in seqs. We create a gap by + // directly manipulating the private log array to skip a seq. + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('step/start', { turn: 1, step: 1 }) + // Push a fake event at seq 3 into the internal log, creating a gap at seq 2. + // The invariants plugin replays session.events on every append, so it sees + // this gap during trace reconstruction. + ;(session as unknown as { log: unknown[] }).log.push({ + type: 'assistant/chunk', + seq: 3, + time: Date.now(), + data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'x' } }, + }) + // Now the log has seqs 0, 1, 3 (gap at 2). Append at what session believes + // is seq 3 (log.length). Reference seq 2: passes earlier (2 < 3) but not + // in knownSeqs ({0, 1, 3} — gap at 2). + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [2] }) + }).toThrow(/unknown seq 2/) + }) + + it('rejects replace op with start > end', async () => { + const { ctx } = await setup() + const session = ctx.sessions.create() + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('step/start', { turn: 1, step: 1 }) + session.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 2 + // start > end is invalid (reversed order). + expect(() => { + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: { op: 'replace', start: 2, end: 1 }, sourceEventSeqs: [2] }) + }).toThrow(/must be <= end/) + }) +}) diff --git a/packages/session-persistence-sqlite/README.md b/packages/session-persistence-sqlite/README.md index 3a7a3c8163..dace03c529 100644 --- a/packages/session-persistence-sqlite/README.md +++ b/packages/session-persistence-sqlite/README.md @@ -6,7 +6,7 @@ A SQLite durable session-persistence backend — a second `SessionPersistence` i ## Storage model -Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). Out-of-log metadata (`SessionMeta`) lives in a `sessions` row, including the mutable `SessionSummary` fields (`updatedAt`, `title`, `firstPrompt`) that `update()` rewrites without touching the event log. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`has`/`list` report exactly the sessions that have a row), so no separate column is needed. +Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data, source_event_seqs, surface_op)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). The two `TEXT` columns `source_event_seqs` and `surface_op` are nullable; they store the event's optional surface-metadata fields (see [ADR 0019](../../docs/adr/0019-session-surface.md)). The schema migrates from v1 to v2 via `ALTER TABLE ADD COLUMN` — existing rows get NULL for both columns, which is correct for events written before surface support. Out-of-log metadata (`SessionMeta`) lives in a `sessions` row, including the mutable `SessionSummary` fields (`updatedAt`, `title`, `firstPrompt`) that `update()` rewrites without touching the event log. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`has`/`list` report exactly the sessions that have a row), so no separate column is needed. The repo targets Node ≥ 24 (the root `engines` field), which includes the stable `node:sqlite` module. The database opens with `foreign_keys = ON` (so `ON DELETE CASCADE` drops a session's events with its row) and `journal_mode = WAL`. The table-layout version is stored in `PRAGMA user_version` and checked on open: a fresh database is stamped with the current `SCHEMA_VERSION`; a database written by a newer, incompatible build (higher `user_version`) is rejected rather than opened against an unknown layout. diff --git a/packages/session-persistence-sqlite/src/index.ts b/packages/session-persistence-sqlite/src/index.ts index bd4b4d9435..8a6cd8cb09 100644 --- a/packages/session-persistence-sqlite/src/index.ts +++ b/packages/session-persistence-sqlite/src/index.ts @@ -33,6 +33,18 @@ import { export { SCHEMA_VERSION } from './schema.ts' +/** + * Serialize an event's surface-metadata fields for SQL binding. Both fields are + * nullable TEXT columns — null when the event has no surface metadata (non-surface + * events, events written before surface support). + */ +function surfaceBindings(event: SessionEvent): [string | null, string | null] { + return [ + event.sourceEventSeqs ? JSON.stringify(event.sourceEventSeqs) : null, + event.surfaceOp !== undefined ? JSON.stringify(event.surfaceOp) : null, + ] +} + /** Plugin configuration. */ export interface Config { /** @@ -180,13 +192,14 @@ export class SessionPersistenceSqlite extends SessionPersistence { // durably closes the interrupted turn before returning, so by the time any // append runs the stored log is balanced and contiguous.) const insertEvent = this.db.prepare( - 'INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)', + 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)', ) this.db.exec('BEGIN') try { if (!state.materialized) this.writeRow(state.meta) for (const event of events) { - insertEvent.run(id, event.seq, event.type, event.time, JSON.stringify(event.data)) + const [surfaceSeqs, surfaceOp] = surfaceBindings(event) + insertEvent.run(id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp) } // Bump updatedAt on every append (the mutable summary lives in the row). const updatedAt = Date.now() @@ -220,7 +233,7 @@ export class SessionPersistenceSqlite extends SessionPersistence { // discarded (not unloadable); only a parse error / seq gap in the COMMITTED // region (at or before the last turn/end) throws (genuine corruption). const eventRows = this.db - .prepare('SELECT seq, type, time, data FROM events WHERE session_id = ? ORDER BY seq') + .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op FROM events WHERE session_id = ? ORDER BY seq') .all(id) as unknown as EventRow[] const { preserved, tornFrom } = scanRows(eventRows) @@ -248,9 +261,12 @@ export class SessionPersistenceSqlite extends SessionPersistence { this.db.prepare('DELETE FROM events WHERE session_id = ? AND seq >= ?').run(id, tornFrom) } if (closers.length > 0) { - const insertEvent = this.db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)') + const insertEvent = this.db.prepare( + 'INSERT INTO events (session_id, seq, type, time, data, source_event_seqs, surface_op) VALUES (?, ?, ?, ?, ?, ?, ?)', + ) for (const event of closers) { - insertEvent.run(id, event.seq, event.type, event.time, JSON.stringify(event.data)) + const [surfaceSeqs, surfaceOp] = surfaceBindings(event) + insertEvent.run(id, event.seq, event.type, event.time, JSON.stringify(event.data), surfaceSeqs, surfaceOp) } } this.db.exec('COMMIT') @@ -497,7 +513,7 @@ export class SessionPersistenceSqlite extends SessionPersistence { /** The preserved events for a session id (torn tail excluded, turn NOT yet closed). */ private eventsFor(id: SessionId): SessionEvent[] { const rows = this.db - .prepare('SELECT seq, type, time, data FROM events WHERE session_id = ? ORDER BY seq') + .prepare('SELECT seq, type, time, data, source_event_seqs, surface_op FROM events WHERE session_id = ? ORDER BY seq') .all(id) as unknown as EventRow[] // Scan on seq+type columns, parsing `data` only for the preserved prefix (a // malformed torn tail must not throw here — same as loadCore). Returns the diff --git a/packages/session-persistence-sqlite/src/schema.ts b/packages/session-persistence-sqlite/src/schema.ts index 1dad51698b..84b357937a 100644 --- a/packages/session-persistence-sqlite/src/schema.ts +++ b/packages/session-persistence-sqlite/src/schema.ts @@ -8,14 +8,14 @@ */ import { DatabaseSync } from 'node:sqlite' -import type { SessionEvent, SessionId, SessionMeta } from '@deepseek-ai/dsh-session' +import type { SessionEvent, SessionId, SessionMeta, SurfaceOp } from '@deepseek-ai/dsh-session' /** * The on-disk schema version. Bumped only on a breaking change to the table * layout; orthogonal to a session's own `version` (which versions the EVENT * vocabulary, stored per session in the `sessions` row). */ -export const SCHEMA_VERSION = 1 +export const SCHEMA_VERSION = 2 /** * A row of the `sessions` table — the out-of-log metadata (`SessionMeta`). The @@ -41,6 +41,10 @@ export interface EventRow { type: string time: number data: string + /** JSON-encoded `number[]` — the event's sourceEventSeqs, or null. */ + source_event_seqs: string | null + /** JSON-encoded `SurfaceOp` — how the event entered the surface, or null. */ + surface_op: string | null } /** @@ -72,6 +76,13 @@ export function openDatabase(path: string): DatabaseSync { // constant (SCHEMA_VERSION is a trusted in-code number, not user input). db.exec(`PRAGMA user_version = ${SCHEMA_VERSION}`) } + if (onDisk === 1) { + // Migrate from v1 to v2: add surface-metadata columns (nullable — existing + // rows get NULL, which is correct for events written before surface existed). + db.exec('ALTER TABLE events ADD COLUMN source_event_seqs TEXT') + db.exec('ALTER TABLE events ADD COLUMN surface_op TEXT') + db.exec(`PRAGMA user_version = ${SCHEMA_VERSION}`) + } db.exec(` CREATE TABLE IF NOT EXISTS sessions ( id TEXT PRIMARY KEY, @@ -86,11 +97,13 @@ export function openDatabase(path: string): DatabaseSync { `) db.exec(` CREATE TABLE IF NOT EXISTS events ( - session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, - seq INTEGER NOT NULL, - type TEXT NOT NULL, - time INTEGER NOT NULL, - data TEXT NOT NULL, + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + seq INTEGER NOT NULL, + type TEXT NOT NULL, + time INTEGER NOT NULL, + data TEXT NOT NULL, + source_event_seqs TEXT, + surface_op TEXT, PRIMARY KEY (session_id, seq) ) STRICT `) @@ -113,12 +126,19 @@ export function rowToMeta(row: SessionRow): SessionMeta { /** Reconstruct a {@link SessionEvent} from an `events` row (parses `data`). */ export function rowToEvent(row: EventRow): SessionEvent { - return { - type: row.type, + const event = { + type: row.type as SessionEvent['type'], seq: row.seq, time: row.time, data: JSON.parse(row.data) as SessionEvent['data'], } as SessionEvent + if (row.source_event_seqs !== null) { + event.sourceEventSeqs = JSON.parse(row.source_event_seqs) as number[] + } + if (row.surface_op !== null) { + event.surfaceOp = JSON.parse(row.surface_op) as SurfaceOp + } + return event } /** diff --git a/packages/session-persistence-sqlite/tests/sqlite.spec.ts b/packages/session-persistence-sqlite/tests/sqlite.spec.ts index 44d788ab72..61b166f257 100644 --- a/packages/session-persistence-sqlite/tests/sqlite.spec.ts +++ b/packages/session-persistence-sqlite/tests/sqlite.spec.ts @@ -1,12 +1,13 @@ import { afterEach, describe, expect, it } from 'vitest' import { Context } from 'cordis' +import { DatabaseSync } from 'node:sqlite' import { mkdtemp, rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' import type { Session, SessionEvent, SessionMeta } from '@deepseek-ai/dsh-session' import SessionPersistenceSqlite, { SCHEMA_VERSION } from '@deepseek-ai/dsh-session-persistence-sqlite' -import { openDatabase, scanRows, type EventRow } from '../src/schema.ts' +import { openDatabase, rowToEvent, scanRows, type EventRow } from '../src/schema.ts' import { runPersistenceContract, meta, oneTurnLog } from '../../session-persistence/tests/contract.ts' const dirs: string[] = [] @@ -42,7 +43,7 @@ describe('scanRows', () => { // scanRows works off EventRows (data is a JSON string column); build them from // SessionEvents so the unit tests read in terms of the event vocabulary. const rows = (events: SessionEvent[]): EventRow[] => - events.map(e => ({ seq: e.seq, type: e.type, time: e.time, data: JSON.stringify(e.data) })) + events.map(e => ({ seq: e.seq, type: e.type, time: e.time, data: JSON.stringify(e.data), source_event_seqs: null, surface_op: null })) it('preserves the full log when it ends exactly on a turn/end (no torn tail)', () => { const { preserved, tornFrom } = scanRows(rows(oneTurnLog())) @@ -91,8 +92,8 @@ describe('scanRows', () => { it('throws on an unparsable row inside the committed region', () => { const withCorruptCommitted: EventRow[] = [ - { seq: 0, type: 'turn/start', time: 1, data: '{not json' }, // corrupt, sits before a turn/end - { seq: 1, type: 'turn/end', time: 2, data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }) }, + { seq: 0, type: 'turn/start', time: 1, data: '{not json', source_event_seqs: null, surface_op: null }, // corrupt, sits before a turn/end + { seq: 1, type: 'turn/end', time: 2, data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }), source_event_seqs: null, surface_op: null }, ] expect(() => scanRows(withCorruptCommitted)).toThrow(/unparsable committed event/) }) @@ -100,7 +101,7 @@ describe('scanRows', () => { it('tolerates an unparsable torn-tail row after the last turn/end', () => { const withCorruptTail: EventRow[] = [ ...rows(oneTurnLog()), - { seq: 6, type: 'turn/start', time: 7, data: '{not json' }, // torn fragment, no committed turn/end after + { seq: 6, type: 'turn/start', time: 7, data: '{not json', source_event_seqs: null, surface_op: null }, // torn fragment, no committed turn/end after ] const { preserved, tornFrom } = scanRows(withCorruptTail) expect(preserved).toEqual(oneTurnLog()) @@ -343,7 +344,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => { }) it('exposes the schema version constant', () => { - expect(SCHEMA_VERSION).toBe(1) + expect(SCHEMA_VERSION).toBe(2) }) }) @@ -749,4 +750,121 @@ describe('SessionPersistenceSqlite: edge cases', () => { await expect(ctx.parallel('session/flush', session)).rejects.toThrow(/id collision/) await ctx.fiber.dispose() }) + + it('migrates a v1 database to v2 (adds surface columns)', async () => { + const path = await freshDbPath() + // Manually create a v1 database with the OLD schema (no surface columns). + const db = new DatabaseSync(path) + db.exec('PRAGMA user_version = 1') + db.exec(` + CREATE TABLE sessions ( + id TEXT PRIMARY KEY, + version INTEGER NOT NULL, + created_at INTEGER NOT NULL, + cwd TEXT, + parent_session TEXT, + updated_at INTEGER NOT NULL, + title TEXT, + first_prompt TEXT + ) STRICT + `) + db.exec(` + CREATE TABLE events ( + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + seq INTEGER NOT NULL, + type TEXT NOT NULL, + time INTEGER NOT NULL, + data TEXT NOT NULL, + PRIMARY KEY (session_id, seq) + ) STRICT + `) + db.close() + // Re-open with v2 code: migration adds the surface columns and stamps v2. + const db2 = openDatabase(path) + const version = (db2.prepare('PRAGMA user_version').get() as { user_version: number }).user_version + expect(version).toBe(2) + const info = db2.prepare("PRAGMA table_info('events')").all() as Array<{ name: string }> + const names = info.map(c => c.name) + expect(names).toContain('source_event_seqs') + expect(names).toContain('surface_op') + db2.close() + }) +}) + +describe('surface field round-trip', () => { + it('rowToEvent parses surface fields from EventRow columns', () => { + const row: EventRow = { + seq: 0, type: 'assistant/message', time: 1, + data: JSON.stringify({ turn: 1, step: 1, content: [] }), + source_event_seqs: JSON.stringify([3, 5]), + surface_op: JSON.stringify('append'), + } + const event = rowToEvent(row) + expect(event.sourceEventSeqs).toEqual([3, 5]) + expect(event.surfaceOp).toBe('append') + }) + + it('rowToEvent handles replace surfaceOp object', () => { + const row: EventRow = { + seq: 0, type: 'assistant/message', time: 1, + data: JSON.stringify({ turn: 1, step: 1, content: [] }), + source_event_seqs: JSON.stringify([0, 1]), + surface_op: JSON.stringify({ op: 'replace', start: 0, end: 1 }), + } + const event = rowToEvent(row) + expect(event.sourceEventSeqs).toEqual([0, 1]) + expect(event.surfaceOp).toEqual({ op: 'replace', start: 0, end: 1 }) + }) + + it('scanRows with surface columns reconstructs events with surface fields', () => { + const rows: EventRow[] = [ + { seq: 0, type: 'user/message', time: 1, + data: JSON.stringify({ content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }), + source_event_seqs: null, surface_op: '{"op":"replace","start":0,"end":0}' }, + { seq: 1, type: 'turn/end', time: 2, + data: JSON.stringify({ turn: 1, reason: { kind: 'completed' } }), + source_event_seqs: null, surface_op: null }, + ] + const { preserved } = scanRows(rows) + expect(preserved).toHaveLength(2) + expect(preserved[0]!.surfaceOp).toEqual({ op: 'replace', start: 0, end: 0 }) + expect(preserved[0]!.sourceEventSeqs).toBeUndefined() + expect(preserved[1]!.surfaceOp).toBeUndefined() + }) + + it('append and load round-trips surface fields through SQLite', async () => { + const ctx = new Context() + await ctx.plugin(SessionStore) + const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' }) + const session = ctx.sessions.create('roundtrip-surface') + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + session.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append', sourceEventSeqs: [0] }) + session.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + await ctx.parallel('session/flush', session) + const loaded = await ctx.sessionPersistence.load(SessionId('roundtrip-surface')) + expect(loaded.events).toHaveLength(4) + const um = loaded.events[1]! + expect(um.surfaceOp).toBe('append') + expect(um.sourceEventSeqs).toBeUndefined() + const am = loaded.events[2]! + expect(am.surfaceOp).toBe('append') + expect(am.sourceEventSeqs).toEqual([0]) + await fiber.dispose() + }) + + it('persists events with surfaceOp but no sourceEventSeqs (covers null branch in surfaceBindings)', async () => { + const ctx = new Context() + await ctx.plugin(SessionStore) + const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: ':memory:' }) + const session = ctx.sessions.create('surface-noseq') + session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + session.append('steering/message', { turn: 1, content: [], source: { kind: 'user' } }, { surfaceOp: 'append' }) + session.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + await ctx.parallel('session/flush', session) + const loaded = await ctx.sessionPersistence.load(SessionId('surface-noseq')) + expect(loaded.events[1]!.surfaceOp).toBe('append') + expect(loaded.events[1]!.sourceEventSeqs).toBeUndefined() + await fiber.dispose() + }) }) diff --git a/packages/session/README.md b/packages/session/README.md index 7cf443fa71..b7787a768c 100644 --- a/packages/session/README.md +++ b/packages/session/README.md @@ -1,6 +1,6 @@ # dsh-session -Event-sourced session log and in-memory store. A `Session` is the append-only source of truth for an agent's whole interaction history — the LLM message history is *derived* from it. +Event-sourced session log and in-memory store. A `Session` is the append-only source of truth for an agent's whole interaction history — the LLM message history is *derived* from it. A **surface** layer (a linked list of message-producing events) is maintained on top of the raw log for efficient derivation and compaction. ## Service: `SessionStore` (ctx key: `sessions`) @@ -24,16 +24,17 @@ Creates and holds event-sourced `Session` instances. Persistence is intentionall Plain class (not a Cordis Service). Create via `ctx.sessions.create()`. -- `session.append(type, data): SessionEvent` — synchronous, never blocks on I/O. **Throws** if `data` is not losslessly JSON-serializable (BigInt, function, symbol, undefined, non-finite number, circular ref, or an exotic object like Map/Set/Date) — the event log is the durable source of truth, so this invariant is enforced at the source (exported as `isJsonValue` for backends to reuse on their replay/fork entry points). -- `session.deriveMessages(): Message[]` — derive the LLM message history from the event log. Raw `assistant/chunk` events are skipped; `context/message` and `steering/message` render as tagged synthetic user messages. +- `session.append(type, data, opts?): SessionEvent` — synchronous, never blocks on I/O. **Throws** if `data` is not losslessly JSON-serializable (BigInt, function, symbol, undefined, non-finite number, circular ref, or an exotic object like Map/Set/Date) — the event log is the durable source of truth, so this invariant is enforced at the source (exported as `isJsonValue` for backends to reuse on their replay/fork entry points). An optional third parameter `opts: SurfaceAppendOpts` carries surface metadata: `surfaceOp` controls how the event enters the surface linked list, and `sourceEventSeqs` records provenance (the seq numbers of events this one derives from). +- `session.deriveMessages(): Message[]` — derive the LLM message history. If any event in the log carries `surfaceOp`, derivation walks the surface linked list (skipping non-surface events). Otherwise, falls back to a linear scan of the raw log (legacy sessions without surface markers). +- `session.surface: SurfaceManager` — the derived surface, lazily rebuilt from `surfaceOp` markers in the log. Processes only new events (delta) on each access — the log is append-only, so prior events never change. - `session.events`, `session.seq`, `session.id` - `session.header: SessionHeader` — immutable creation metadata (`version`, `id`, `createdAt`, optional `cwd`/`parentSession`). Kept out of the event log (a storage concern, not replayable state); a minimal v1 header is synthesized for bare `Session` construction. -### Metadata types (`types.ts`) +### Surface types -- `SessionHeader` — immutable, written once: `{ version, id, createdAt, cwd?, parentSession? }`. -- `SessionSummary` — mutable, updateable without touching the log: `{ updatedAt, title?, firstPrompt? }`. -- `SessionMeta = SessionHeader & SessionSummary` — owned here (beside `SessionId`) because `Session.header` is typed by it; persistence backends re-export these rather than own them (which would force a package cycle). +- `SurfaceOp` — how a surface node entered the linked list: `'append'` (normal tail append) or `{ op: 'replace', start, end }` (replace nodes from `start` through `end` inclusive — both must be valid surface node seqs; `start === end` replaces a single node). Used by compaction to shadow old nodes without deleting them. +- `SurfaceAppendOpts` — `{ surfaceOp?: SurfaceOp; sourceEventSeqs?: number[] }`, the optional third parameter to `session.append()`. +- `SurfaceNode` — `{ seq: number; prev: number | null; next: number | null }`, one node in the surface linked list. ### Session event vocabulary (`types.ts`) @@ -43,11 +44,23 @@ Merge-extensible via `SessionEventMap` — a compaction plugin adds `compaction/ Also defines `TurnTriggerMap` and `TurnEndReasonMap` (merge-extensible sum types for typed turn boundaries — `kind`-tagged instead of strings). +Every `SessionEvent` carries two optional top-level fields (structural metadata): + +- `sourceEventSeqs?: number[]` — seq numbers of provenance sources (e.g., the `assistant/chunk` seqs behind an `assistant/message`, or the shadowed nodes behind a compaction marker). +- `surfaceOp?: SurfaceOp` — how this event entered the surface. Absent for non-surface events (boundaries, chunks, usage, errors). + +### Metadata types (`types.ts`) + +- `SessionHeader` — immutable, written once: `{ version, id, createdAt, cwd?, parentSession? }`. +- `SessionSummary` — mutable, updateable without touching the log: `{ updatedAt, title?, firstPrompt? }`. +- `SessionMeta = SessionHeader & SessionSummary` — owned here (beside `SessionId`) because `Session.header` is typed by it; persistence backends re-export these rather than own them (which would force a package cycle). + ### Extension points - Persistence plugins: subscribe to `session/event` (write-behind) and drain on `session/flush` (awaited) and fiber dispose. A durable backend reads the log and reloads it into a live session; the metadata seam (`SessionHeader`/`SessionSummary`/`SessionMeta`, `session.header`) is what such a backend stores beside the log. -- Replay/fork: `ctx.sessions.create(id, { seed })` seeds a new session with an existing event log. +- Replay/fork: `ctx.sessions.create(id, { seed })` seeds a new session with an existing event log. The surface rebuilds deterministically from `surfaceOp` markers in the seeded events. +- Compaction: a future plugin appends a new event with `surfaceOp: { op: 'replace', start, end }` to shadow old surface nodes. ### What is NOT here (TODO) -- **Session branching/tree** (pi-style entry tree) — defered unless needed beyond seed-based forking. +- **Session branching/tree** (pi-style entry tree) — deferred unless needed beyond seed-based forking. diff --git a/packages/session/src/index.ts b/packages/session/src/index.ts index 34eefb24a3..d0d4d44aed 100644 --- a/packages/session/src/index.ts +++ b/packages/session/src/index.ts @@ -10,12 +10,14 @@ import { Context, Service } from 'cordis' import { isAbsolute } from 'node:path' import type { ContentBlock, Message, MessageSource } from '@deepseek-ai/dsh-llm' import { SessionId } from './types.ts' -import type { CreateSessionOptions, SessionEvent, SessionEventMap, SessionEventType, SessionHeader } from './types.ts' +import type { CreateSessionOptions, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SurfaceAppendOpts } from './types.ts' import { isJsonValue } from './json.ts' +import { SurfaceManager } from './surface.ts' export * from './types.ts' export { isJsonValue } from './json.ts' export { interruptedTurnClosers } from './repair.ts' +export type { SurfaceNode } from './surface.ts' declare module 'cordis' { interface Context { @@ -65,6 +67,21 @@ export class Session { /** Set by the store so appends are observable; undefined when detached. */ onAppend: ((event: SessionEvent) => void) | undefined + /** + * Derived surface — a cached linked list of message-producing events. + * Lazily rebuilt from `surfaceOp` markers in the log; processes only new + * events (delta) on each access — the log is append-only, so prior events + * never change. + * `append`. Undefined until first accessed (including after fork/seed). + */ + private _surface: SurfaceManager | undefined + + /** The surface linked list over this session's event log. */ + get surface(): SurfaceManager { + if (!this._surface) this._surface = new SurfaceManager(this.log) + return this._surface + } + /** * Immutable creation metadata (format version, cwd, lineage). Supplied by * the store via `ctx.sessions.create()`. When a `Session` is constructed @@ -117,6 +134,11 @@ export class Session { * `onAppend`. The hot path never blocks on I/O — persistence plugins buffer * asynchronously. * + * @param type - The event type (key of {@link SessionEventMap}). + * @param data - The event payload; must be JSON-serializable. + * @param opts - Optional surface metadata: `surfaceOp` controls how the + * event enters the surface linked list; `sourceEventSeqs` records + * provenance (the seq numbers of events this one derives from). * @throws if `data` is not losslessly JSON-serializable (BigInt, function, * symbol, undefined, non-finite number, circular ref, or an exotic object * like Map/Set/Date). The event log is the durable source of truth, so this @@ -125,7 +147,7 @@ export class Session { * throw surfaces at the buggy caller's append site, not asynchronously in a * backend flush. */ - append(type: T, data: SessionEventMap[T]): SessionEvent { + append(type: T, data: SessionEventMap[T], opts?: SurfaceAppendOpts): SessionEvent { if (!isJsonValue(data)) { throw new Error(`session event "${type}" carries non-JSON-serializable data`) } @@ -138,14 +160,29 @@ export class Session { // validated. structuredClone is safe because serializability was just // checked. The returned event carries the SAME snapshot, so a caller reading // back `event.data` sees the logged value, not its own mutable input. - const event = { type, seq: this.log.length, time: Date.now(), data: structuredClone(data) } as SessionEvent + // + // Surface metadata is snapshot separately: sourceEventSeqs (number[] — + // primitives, so array spread is a complete copy) and surfaceOp (a string + // primitive, or cloned if it's a replace object). + const event = { + type, + seq: this.log.length, + time: Date.now(), + data: structuredClone(data), + ...opts?.sourceEventSeqs !== undefined ? { sourceEventSeqs: [...opts.sourceEventSeqs] } : {}, + ...opts?.surfaceOp !== undefined ? { + surfaceOp: typeof opts.surfaceOp === 'string' ? opts.surfaceOp : structuredClone(opts.surfaceOp), + } : {}, + } as SessionEvent this.log.push(event) this.onAppend?.(event) return event } /** - * Derive the LLM message history from the event log. + * Derive the LLM message history from the session surface (when surface + * markers exist) or from a linear scan of the raw event log (legacy sessions + * without surface markers). * * - `user/message` → user message * - `assistant/message` → assistant message (chunks are skipped — they are @@ -163,43 +200,62 @@ export class Session { * negligible next to a model call. */ deriveMessages(): Message[] { + if (this.surface.hasSurface) { + const messages: Message[] = [] + for (const node of this.surface.nodes) { + // Surface nodes are built from this.log — node.seq is always a valid + // index by construction. The non-null assertion expresses that invariant. + // eslint-disable-next-line @typescript-eslint/no-non-null-assertion + const msg = this._deriveOneMessage(this.log[node.seq]!) + if (msg) messages.push(msg) + } + return messages + } + // Legacy path: linear scan for sessions without surface markers. const messages: Message[] = [] for (const event of this.log) { - // Intentionally non-exhaustive: only message-producing events derive - // history; turn/step boundaries, chunks, usage, and errors are - // trace/replay data. - // eslint-disable-next-line @typescript-eslint/switch-exhaustiveness-check - switch (event.type) { - case 'user/message': { - messages.push({ role: 'user', content: structuredClone(event.data.content) }) - break - } - case 'assistant/message': { - messages.push({ role: 'assistant', content: structuredClone(event.data.content) }) - break - } - case 'tool/result': { - const { callId, content, isError } = event.data - messages.push({ - role: 'user', - content: [{ type: 'tool-result', toolCallId: callId, content: structuredClone(content), isError }], - }) - break - } - case 'context/message': { - const { content, source } = event.data - messages.push({ role: 'user', content: renderTagged('context', structuredClone(content), source) }) - break - } - case 'steering/message': { - const { content, source } = event.data - messages.push({ role: 'user', content: renderTagged('steering', structuredClone(content), source) }) - break - } - } + const msg = this._deriveOneMessage(event) + if (msg) messages.push(msg) } return messages } + + /** + * Derive a single LLM message from one event, or null if the event type + * does not produce a message. Extracted so both the surface path and the + * legacy linear-scan path share the same derivation rules. + */ + private _deriveOneMessage(event: SessionEvent): Message | null { + // Intentionally non-exhaustive: only message-producing events derive + // history; turn/step boundaries, chunks, usage, and errors are + // trace/replay data. + + switch (event.type) { + case 'user/message': { + return { role: 'user', content: structuredClone(event.data.content) } + } + case 'assistant/message': { + return { role: 'assistant', content: structuredClone(event.data.content) } + } + case 'tool/result': { + const { callId, content, isError } = event.data + return { + role: 'user', + content: [{ type: 'tool-result', toolCallId: callId, content: structuredClone(content), isError }], + } + } + case 'context/message': { + const { content, source } = event.data + return { role: 'user', content: renderTagged('context', structuredClone(content), source) } + } + case 'steering/message': { + const { content, source } = event.data + return { role: 'user', content: renderTagged('steering', structuredClone(content), source) } + } + default: + return null + } + } } /** diff --git a/packages/session/src/repair.ts b/packages/session/src/repair.ts index 6a3f60a681..d313af12c3 100644 --- a/packages/session/src/repair.ts +++ b/packages/session/src/repair.ts @@ -61,7 +61,12 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session // call is "pending" until its matching tool/result arrives. Reset at every // turn boundary so a committed earlier turn (already balanced) never leaks a // phantom pending call into the interrupted-turn repair. - const pendingCalls = new Map() + // Track pending tool calls with their callSeq (the seq of the `tool/call` + // event, captured for surface sourceEventSeqs provenance on the synthetic + // result). CallSeq is set from `tool/call` events; the assistant/message + // block scan may register a call first (it appears earlier in the log), and + // the later `tool/call` event fills in the seq. + const pendingCalls = new Map() for (const event of events) { switch (event.type) { case 'turn/start': @@ -87,6 +92,18 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session if (block.type === 'tool-call') pendingCalls.set(block.id, { step: event.data.step }) } break + case 'tool/call': + // Capture the tool/call event seq for surface provenance on the + // synthesized tool/result. The entry may already exist (registered by + // the assistant/message above) or may be new (if the assistant/message + // came from a prior step that was already closed). + { + const entry = pendingCalls.get(event.data.callId) + if (entry) { + entry.callSeq = event.seq + } + } + break case 'tool/result': pendingCalls.delete(event.data.callId) break @@ -112,7 +129,7 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session // crash, so deriveMessages() yields a valid provider transcript on resume (a // dangling assistant tool-call is rejected by every provider). Insertion // order follows the Map (insertion = log order of the assistant messages). - for (const [callId, { step }] of pendingCalls) { + for (const [callId, { step, callSeq }] of pendingCalls) { closers.push({ type: 'tool/result', seq: seq++, @@ -125,6 +142,8 @@ export function interruptedTurnClosers(events: readonly SessionEvent[]): Session isError: true, error: { name: 'InterruptedError', code: 'interrupted' }, }, + surfaceOp: 'append', + ...callSeq !== undefined ? { sourceEventSeqs: [callSeq] } : {}, }) } diff --git a/packages/session/src/surface.ts b/packages/session/src/surface.ts new file mode 100644 index 0000000000..a09dbd001f --- /dev/null +++ b/packages/session/src/surface.ts @@ -0,0 +1,129 @@ +/** + * Surface layer on top of the session event log: a derived, cached linked list + * of events that produce LLM messages. Rebuilt deterministically from + * `surfaceOp` markers in the log — the log is the source of truth; the surface + * is a view. + * + * @module @deepseek-ai/dsh-session/surface + */ + +import type { SessionEvent, SurfaceOp } from './types.ts' + +/** One node in the surface linked list. */ +export interface SurfaceNode { + /** The event seq of this surface node. */ + seq: number + /** The previous surface node's seq, or null if this is the head. */ + prev: number | null + /** The next surface node's seq, or null if this is the tail. */ + next: number | null +} + +/** + * Maintains a cached linked list of surface nodes, rebuilt lazily from + * `surfaceOp` markers in the event log. Because the log is append-only, it + * processes only the delta since the last rebuild — new events are folded + * into the existing surface in O(new events) rather than rescanning the + * whole log. + */ +export class SurfaceManager { + /** Surface nodes in linked-list order (head to tail). Empty until first access. */ + private _nodes: SurfaceNode[] = [] + /** Map from event seq → node for O(1) lookup during replacements. */ + private _nodeBySeq = new Map() + /** The last processed seq. -1 forces a full rebuild on first access. */ + private _lastProcessedSeq = -1 + + constructor(private log: readonly SessionEvent[]) {} + + /** + * Reset to unprocessed state. Call after the log has been replaced + * wholesale (e.g. after Session seed). Not needed for normal appends — + * those are picked up incrementally. + */ + invalidate(): void { + this._lastProcessedSeq = -1 + this._nodes = [] + this._nodeBySeq.clear() + } + + /** The surface nodes in linked-list order (head to tail). */ + get nodes(): readonly SurfaceNode[] { + if (this._lastProcessedSeq < this.log.length - 1) this._processDelta() + return this._nodes + } + + /** Whether any event in the log carries `surfaceOp` markers. */ + get hasSurface(): boolean { + if (this._nodes.length > 0) return true + // Never processed anything — scan the whole log. + if (this._lastProcessedSeq === -1) return this.log.some(e => e.surfaceOp !== undefined) + // Processed up to _lastProcessedSeq without finding surface nodes; check + // only new events. + for (let i = this._lastProcessedSeq + 1; i < this.log.length; i++) { + if (this.log[i]?.surfaceOp !== undefined) return true + } + return false + } + + /** + * Process events from `_lastProcessedSeq + 1` through the end of the log, + * folding new surface markers into the existing linked list. + */ + private _processDelta(): void { + for (let i = this._lastProcessedSeq + 1; i < this.log.length; i++) { + const event = this.log[i] + if (event === undefined || event.surfaceOp === undefined) continue + + if (event.surfaceOp === 'append') { + const tail = this._nodes.length > 0 ? this._nodes[this._nodes.length - 1] : undefined + const node: SurfaceNode = { seq: event.seq, prev: tail?.seq ?? null, next: null } + if (tail) tail.next = event.seq + this._nodes.push(node) + this._nodeBySeq.set(event.seq, node) + } else { + this._replace(this._nodes, this._nodeBySeq, event.seq, event.surfaceOp) + } + } + this._lastProcessedSeq = this.log.length - 1 + } + + /** Apply a replace operation to the in-progress surface. */ + private _replace( + nodes: SurfaceNode[], + nodeBySeq: Map, + newSeq: number, + op: Extract, + ): void { + const startIdx = nodes.findIndex(n => n.seq === op.start) + if (startIdx === -1) { + throw new Error(`surface replace: start seq ${op.start} not found in surface`) + } + const endIdx = nodes.findIndex(n => n.seq === op.end) + if (endIdx === -1) { + throw new Error(`surface replace: end seq ${op.end} not found in surface`) + } + if (startIdx > endIdx) { + throw new Error(`surface replace: start seq ${op.start} (index ${startIdx}) is after end seq ${op.end} (index ${endIdx})`) + } + + // Remove shadowed nodes from `[startIdx, endIdx]` inclusive. + const count = endIdx - startIdx + 1 + const removed = nodes.splice(startIdx, count) + for (const r of removed) nodeBySeq.delete(r.seq) + + // Insert the new node where the removed range was. + const prevNode = startIdx > 0 ? nodes[startIdx - 1] : undefined + const nextNode = startIdx < nodes.length ? nodes[startIdx] : undefined + + const newNode: SurfaceNode = { + seq: newSeq, + prev: prevNode?.seq ?? null, + next: nextNode?.seq ?? null, + } + if (prevNode) prevNode.next = newSeq + if (nextNode) nextNode.prev = newSeq + nodes.splice(startIdx, 0, newNode) + nodeBySeq.set(newSeq, newNode) + } +} diff --git a/packages/session/src/types.ts b/packages/session/src/types.ts index d8ffaf9dc8..9f2cd87643 100644 --- a/packages/session/src/types.ts +++ b/packages/session/src/types.ts @@ -160,6 +160,34 @@ export interface SessionEventMap { export type SessionEventType = keyof SessionEventMap +/** + * How a session event entered the surface linked list. Absent for non-surface + * events (boundaries, chunks, usage, errors). + * + * - `'append'`: added to the tail — normal path for user/assistant/tool/context + * messages. + * - `{ op: 'replace', start, end }`: replaces surface nodes from `start` + * (inclusive) through `end` (inclusive) with this node. Both must exist as + * surface nodes in the current surface. `start === end` replaces a single + * node. The node's {@link SessionEvent.sourceEventSeqs} must include every + * shadowed surface node. Used by compaction and possible other manipulations. + */ +export type SurfaceOp = + | 'append' + | { op: 'replace'; start: number; end: number } + +/** + * Optional surface metadata passed to {@link Session.append}. + * `surfaceOp` controls how the event enters the surface linked list; + * `sourceEventSeqs` records the seq numbers of events that are provenance + * sources of this one (e.g. the `assistant/chunk` seqs behind an + * `assistant/message`, or the shadowed nodes behind a compaction replacement). + */ +export interface SurfaceAppendOpts { + surfaceOp?: SurfaceOp + sourceEventSeqs?: number[] +} + /** * One immutable entry in the session log. * @@ -174,5 +202,13 @@ export type SessionEvent = { /** Unix epoch milliseconds. */ time: number data: SessionEventMap[K] + /** + * Seq numbers of events that are provenance sources of this event + * (e.g. the `assistant/chunk` seqs that built an `assistant/message`, + * or the surface nodes shadowed by a compaction marker). + */ + sourceEventSeqs?: number[] + /** How this event entered the surface; absent for non-surface events. */ + surfaceOp?: SurfaceOp } }[T] diff --git a/packages/session/tests/repair.spec.ts b/packages/session/tests/repair.spec.ts index 893015b218..cf6fb2b51c 100644 --- a/packages/session/tests/repair.spec.ts +++ b/packages/session/tests/repair.spec.ts @@ -122,4 +122,36 @@ describe('interruptedTurnClosers', () => { const result = closers[0]! expect(result.type === 'tool/result' && result.data.callId).toBe('call-b') }) + + it('synthesized tool/result carries surfaceOp and sourceEventSeqs when tool/call was logged', () => { + const events: SessionEvent[] = [ + userTurnStart(1, 0), + { type: 'step/start', seq: 1, time: 1, data: { turn: 1, step: 1 } }, + { type: 'assistant/message', seq: 2, time: 2, data: { turn: 1, step: 1, content: [ + { type: 'tool-call', id: CallId('call-1'), name: 'bash', arguments: '{}' }, + ] } }, + { type: 'tool/call', seq: 3, time: 3, data: { turn: 1, step: 1, callId: CallId('call-1'), name: 'bash', arguments: '{}' } }, + ] + const closers = interruptedTurnClosers(events) + expect(closers.map(e => e.type)).toEqual(['tool/result', 'step/end', 'turn/end']) + const result = closers[0]! + expect(result.surfaceOp).toBe('append') + expect(result.sourceEventSeqs).toEqual([3]) + }) + + it('handles tool/call without a matching assistant/message entry gracefully', () => { + // A tool/call event exists in the log but no assistant/message registered + // the callId in pendingCalls (e.g., a plugin appended it directly, or the + // assistant/message from a prior step didn't have this call). The repair + // should still close the turn — it just won't synthesize a result for this + // call (there's nothing to answer). + const events: SessionEvent[] = [ + userTurnStart(1, 0), + { type: 'step/start', seq: 1, time: 1, data: { turn: 1, step: 1 } }, + { type: 'tool/call', seq: 2, time: 2, data: { turn: 1, step: 1, callId: CallId('orphan'), name: 'bash', arguments: '{}' } }, + ] + const closers = interruptedTurnClosers(events) + // No pending calls → no synthetic tool/result, just step/end + turn/end. + expect(closers.map(e => e.type)).toEqual(['step/end', 'turn/end']) + }) }) diff --git a/packages/session/tests/surface.spec.ts b/packages/session/tests/surface.spec.ts new file mode 100644 index 0000000000..d8f503c5e5 --- /dev/null +++ b/packages/session/tests/surface.spec.ts @@ -0,0 +1,324 @@ +import { describe, expect, it } from 'vitest' +import type { SessionEvent } from '@deepseek-ai/dsh-session' +import { Session, SessionId } from '@deepseek-ai/dsh-session' +import { CallId } from '@deepseek-ai/dsh-llm' + +/** Build a minimal session with turn boundaries and a single user message. */ +function surfaceSession(): Session { + const s = new Session(SessionId('ss')) + s.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + s.append('user/message', { content: [{ type: 'text', text: 'hello' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'hi' }] }, { surfaceOp: 'append' }) + s.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + return s +} + +describe('SurfaceManager', () => { + it('rebuilds a linked list from surfaceOp: append markers', () => { + const s = surfaceSession() + const nodes = s.surface.nodes + // Only the user/message and assistant/message carry surfaceOp: 'append'. + // The turn boundaries do not have surface markers. + expect(nodes.length).toBe(2) + expect(nodes[0]!.seq).toBe(1) // user/message (turn/start is seq 0) + expect(nodes[0]!.prev).toBeNull() + expect(nodes[0]!.next).toBe(2) // assistant/message (seq 2) + expect(nodes[1]!.seq).toBe(2) + expect(nodes[1]!.prev).toBe(1) + expect(nodes[1]!.next).toBeNull() + }) + + it('hasSurface returns false when no events have surfaceOp', () => { + const s = new Session(SessionId('nosurface')) + s.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + s.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }) + s.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + expect(s.surface.hasSurface).toBe(false) + }) + + it('hasSurface returns true when any event has surfaceOp', () => { + const s = surfaceSession() + expect(s.surface.hasSurface).toBe(true) + }) + + it('hasSurface detects surface markers that arrive after initial processing', () => { + // Start with no surface markers. Access nodes first to set _lastProcessedSeq + // (via delta processing), keeping _nodes empty. Then append a mix of non-surface + // and surface events, and verify hasSurface detects via the delta-only check. + const s = new Session(SessionId('late')) + s.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + s.append('user/message', { content: [{ type: 'text', text: 'hi' }], source: { kind: 'user' } }) + // Access nodes to trigger processing: sets _lastProcessedSeq = 1, _nodes = []. + expect(s.surface.nodes.length).toBe(0) + // Append non-surface events first (exercises the loop-continue branch), then + // a surface event (exercises the return-true branch). + s.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + s.append('turn/start', { turn: 2, trigger: { kind: 'continuation' } }) + s.append('assistant/message', { turn: 2, step: 1, content: [] }, { surfaceOp: 'append' }) + // hasSurface checks only new seqs [2, 3, 4]; skips 2 and 3 (non-surface), + // finds surfaceOp on seq 4 and returns true. + expect(s.surface.hasSurface).toBe(true) + }) + + it('invalidate resets to full rebuild', () => { + const s = surfaceSession() + expect(s.surface.nodes.length).toBe(2) + // After invalidate, the surface should rebuild from scratch on next access. + ;(s.surface).invalidate() + expect(s.surface.nodes.length).toBe(2) // same result, but rebuilt + }) + + it('empty surface yields empty nodes', () => { + const s = new Session(SessionId('empty')) + // Only turn boundaries, no surface nodes. + s.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + s.append('step/start', { turn: 1, step: 1 }) + s.append('step/end', { turn: 1, step: 1 }) + s.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + expect(s.surface.nodes.length).toBe(0) + expect(s.surface.hasSurface).toBe(false) + // deriveMessages returns empty array + expect(s.deriveMessages()).toEqual([]) + }) + + it('picks up new events incrementally (delta processing)', () => { + const s = surfaceSession() + expect(s.surface.nodes.length).toBe(2) + // Append another surface node + s.append('tool/result', { turn: 1, step: 1, callId: CallId('c1'), content: [{ type: 'text', text: 'ok' }], isError: false }, { surfaceOp: 'append' }) + expect(s.surface.nodes.length).toBe(3) + expect(s.surface.nodes[2]!.seq).toBe(4) // seq 4: after turn/end at seq 3 + expect(s.surface.nodes[2]!.prev).toBe(2) + expect(s.surface.nodes[1]!.next).toBe(4) + }) + + it('replays identically from a seeded log with surface markers', () => { + const original = surfaceSession() + original.append('tool/result', { turn: 1, step: 1, callId: CallId('c1'), content: [{ type: 'text', text: 'ok' }], isError: false }, { surfaceOp: 'append' }) + const replayed = new Session(SessionId('replay'), [...original.events]) + // Surface rebuilds from the seeded log's markers. + expect(replayed.surface.nodes.map(n => n.seq)).toEqual([1, 2, 4]) + expect(replayed.deriveMessages()).toEqual(original.deriveMessages()) + }) + + it('rebuild with replace operation splices out shadowed nodes', () => { + const s = surfaceSession() + // seq: 0=turn/start, 1=user, 2=assistant, 3=turn/end + // Surface nodes: seq 1 (user), seq 2 (assistant). + // Replace both with a compaction marker. Both 1 and 2 are valid surface seqs. + s.append('assistant/message', + { turn: 2, step: 1, content: [{ type: 'text', text: 'summary' }] }, + { surfaceOp: { op: 'replace', start: 1, end: 2 }, sourceEventSeqs: [1, 2] }, + ) + // Now the surface should have just the compaction node. + expect(s.surface.nodes.length).toBe(1) + expect(s.surface.nodes[0]!.seq).toBe(4) // seq of the compaction marker + expect(s.surface.nodes[0]!.prev).toBeNull() + expect(s.surface.nodes[0]!.next).toBeNull() + }) + + it('replace with both ends at real nodes splices only the range', () => { + const s = new Session(SessionId('range')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('user/message', { content: [{ type: 'text', text: 'b' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 1 + s.append('user/message', { content: [{ type: 'text', text: 'c' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 2 + // Replace seq 0 through 1 inclusive: shadow a and b, keep c. + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'summary' }] }, + { surfaceOp: { op: 'replace', start: 0, end: 1 }, sourceEventSeqs: [0, 1] }, + ) // seq 3 + expect(s.surface.nodes.map(n => n.seq)).toEqual([3, 2]) + // Links: 3 ↔ 2 + expect(s.surface.nodes[0]!.prev).toBeNull() + expect(s.surface.nodes[0]!.next).toBe(2) + expect(s.surface.nodes[1]!.prev).toBe(3) + expect(s.surface.nodes[1]!.next).toBeNull() + }) + + it('single-node replacement (start === end)', () => { + const s = new Session(SessionId('single')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('user/message', { content: [{ type: 'text', text: 'b' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 1 + // Replace only seq 1 (single node). + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'x' }] }, + { surfaceOp: { op: 'replace', start: 1, end: 1 }, sourceEventSeqs: [1] }, + ) // seq 2 + expect(s.surface.nodes.map(n => n.seq)).toEqual([0, 2]) + expect(s.surface.nodes[0]!.next).toBe(2) + expect(s.surface.nodes[1]!.prev).toBe(0) + }) + + it('throws when replace start is not found', () => { + const s = new Session(SessionId('bad-start')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'y' }] }, + { surfaceOp: { op: 'replace', start: 5, end: 0 }, sourceEventSeqs: [5, 0] }, + ) + expect(() => s.surface.nodes).toThrow(/surface replace: start seq 5 not found/) + }) + + it('throws when replace end is not found', () => { + const s = new Session(SessionId('bad-end')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'y' }] }, + { surfaceOp: { op: 'replace', start: 0, end: 99 }, sourceEventSeqs: [0] }, + ) + expect(() => s.surface.nodes).toThrow(/surface replace: end seq 99 not found/) + }) + + it('throws when start is after end', () => { + const s = new Session(SessionId('reversed')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('user/message', { content: [{ type: 'text', text: 'b' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 1 + // start=1, end=0 would be reversed order. + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'y' }] }, + { surfaceOp: { op: 'replace', start: 1, end: 0 }, sourceEventSeqs: [1, 0] }, + ) + expect(() => s.surface.nodes).toThrow(/start seq 1.*after end seq 0/) + }) + + it('sourceEventSeqs is snapshot so caller mutation does not affect logged event', () => { + const s = new Session(SessionId('immutable')) + const sources = [10, 20] + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'h' }] }, { surfaceOp: 'append', sourceEventSeqs: sources }) + // Mutate caller's array after append. + sources.push(30) + sources[0] = 99 + const logged = s.events[0]! + expect(logged.sourceEventSeqs).toEqual([10, 20]) + }) + + it('replace starting at non-head position links to previous node correctly', () => { + const s = new Session(SessionId('mid-replace')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 0 + s.append('user/message', { content: [{ type: 'text', text: 'b' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 1 + s.append('user/message', { content: [{ type: 'text', text: 'c' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) // seq 2 + // Replace the middle node (seq 1) only, keeping seq 0 and seq 2. + s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'x' }] }, + { surfaceOp: { op: 'replace', start: 1, end: 1 }, sourceEventSeqs: [1] }, + ) // seq 3 + expect(s.surface.nodes.map(n => n.seq)).toEqual([0, 3, 2]) + // Links: 0 → 3 → 2 + expect(s.surface.nodes[0]!.prev).toBeNull() + expect(s.surface.nodes[0]!.next).toBe(3) + expect(s.surface.nodes[1]!.prev).toBe(0) + expect(s.surface.nodes[1]!.next).toBe(2) + expect(s.surface.nodes[2]!.prev).toBe(3) + expect(s.surface.nodes[2]!.next).toBeNull() + }) + + it('surfaceOp replace object is snapshot so caller mutation is isolated', () => { + const s = new Session(SessionId('immutable-op')) + s.append('user/message', { content: [{ type: 'text', text: 'a' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + const op = { op: 'replace' as const, start: 0, end: 0 } + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 's' }] }, { surfaceOp: op, sourceEventSeqs: [0] }) + // Mutate caller's object after append. + op.start = 99 + const logged = s.events[1]! + expect(logged.surfaceOp).toEqual({ op: 'replace', start: 0, end: 0 }) + }) +}) + +describe('deriveMessages with surface', () => { + it('uses the surface path when surface markers are present', () => { + const s = surfaceSession() + const messages = s.deriveMessages() + expect(messages).toHaveLength(2) + expect(messages[0]!.role).toBe('user') + expect(messages[0]!.content[0]).toMatchObject({ type: 'text', text: 'hello' }) + expect(messages[1]!.role).toBe('assistant') + expect(messages[1]!.content[0]).toMatchObject({ type: 'text', text: 'hi' }) + }) + + it('falls back to linear scan when no surface markers exist', () => { + const s = new Session(SessionId('legacy')) + s.append('user/message', { content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } }) + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'a' }] }) + const messages = s.deriveMessages() + expect(messages).toHaveLength(2) + expect(messages[0]!.role).toBe('user') + expect(messages[1]!.role).toBe('assistant') + }) + + it('surface path skips non-surface events (chunks, boundaries)', () => { + const s = new Session(SessionId('filter')) + s.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) + s.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'h' } }) + s.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 1, text: 'i' } }) + s.append('user/message', { content: [{ type: 'text', text: 'hello' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'hi' }] }, { surfaceOp: 'append' }) + s.append('turn/end', { turn: 1, reason: { kind: 'completed' } }) + // Chunks and boundaries are NOT in the surface, so only 2 messages. + expect(s.deriveMessages()).toHaveLength(2) + }) + + it('deriveMessages via surface respects replace (shadowed nodes are excluded)', () => { + const s = new Session(SessionId('compacted')) + s.append('user/message', { content: [{ type: 'text', text: 'original' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + s.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'compacted' }] }, { surfaceOp: { op: 'replace', start: 0, end: 0 }, sourceEventSeqs: [0] }) + // Only the compaction node is visible. + const messages = s.deriveMessages() + expect(messages).toHaveLength(1) + expect(messages[0]!.content[0]).toMatchObject({ type: 'text', text: 'compacted' }) + }) + + it('context/message and steering/message appear on surface', () => { + const s = new Session(SessionId('ctx')) + s.append('context/message', { content: [{ type: 'text', text: 'file changed' }], source: { kind: 'plugin', plugin: 'watcher' } }, { surfaceOp: 'append' }) + s.append('steering/message', { turn: 1, content: [{ type: 'text', text: 'focus' }], source: { kind: 'user' } }, { surfaceOp: 'append' }) + const messages = s.deriveMessages() + expect(messages).toHaveLength(2) + expect(messages[0]!.content[0]).toMatchObject({ type: 'text', text: '' }) + expect(messages[1]!.content[0]).toMatchObject({ type: 'text', text: '' }) + }) +}) + +describe('Session.append surface opts', () => { + it('records sourceEventSeqs and surfaceOp on the event', () => { + const s = new Session(SessionId('opts')) + const event = s.append('assistant/message', + { turn: 1, step: 1, content: [{ type: 'text', text: 'h' }] }, + { surfaceOp: 'append', sourceEventSeqs: [3, 5, 7] }, + ) + expect(event.sourceEventSeqs).toEqual([3, 5, 7]) + expect(event.surfaceOp).toBe('append') + // The logged event matches the returned event. + expect(s.events[0]!.sourceEventSeqs).toEqual([3, 5, 7]) + expect(s.events[0]!.surfaceOp).toBe('append') + }) + + it('deriveMessages skips surface nodes whose event type is not message-producing', () => { + // A surface node with a type not handled by _deriveOneMessage (e.g., 'usage' + // placed on surface) should be skipped — the null-check in the surface + // derivation path is exercised. + const seed: SessionEvent[] = [ + { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } }, + { type: 'step/start', seq: 1, time: 2, data: { turn: 1, step: 1 } }, + { type: 'usage', seq: 2, time: 3, data: { turn: 1, step: 1, usage: { inputTokens: 0, outputTokens: 0 } }, surfaceOp: 'append' as const }, + { type: 'step/end', seq: 3, time: 4, data: { turn: 1, step: 1 } }, + { type: 'turn/end', seq: 4, time: 5, data: { turn: 1, reason: { kind: 'completed' } } }, + ] + const s = new Session(SessionId('nomessage'), seed) + // The usage event is on the surface but _deriveOneMessage returns null for it. + expect(s.deriveMessages()).toHaveLength(0) + }) + + it('append without surface opts produces an event without surface fields', () => { + const s = new Session(SessionId('noopts')) + s.append('user/message', { content: [{ type: 'text', text: 'x' }], source: { kind: 'user' } }) + expect(s.events[0]!.sourceEventSeqs).toBeUndefined() + expect(s.events[0]!.surfaceOp).toBeUndefined() + }) + + it('surfaceOp primitives are not cloned (they are immutable)', () => { + const s = new Session(SessionId('prim')) + const event = s.append('assistant/message', { turn: 1, step: 1, content: [] }, { surfaceOp: 'append' }) + // The string 'append' is a primitive — identity-preserving is fine. + expect(event.surfaceOp).toBe('append') + }) +})