import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises'; import { execFile } from 'node:child_process'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { promisify } from 'node:util'; import { INTERACTION_TAG_SESSION_ID, IAgentLifecycleService, IAgentContextMemoryService, IAgentConversationUndoParticipantRegistry, IAgentLoopService, IAgentScopeContext, IAgentTaskService, IEventBus, IFlagService, ISessionIndex, ISessionMetadata, ISessionLifecycleService, ISessionManager, IWorkspaceInstanceManager, LifecycleScope, interactions, makeAgentScopeContext, type AgentContext, type ContextMessage, type AgentConversationUndoParticipant, type Event2, TOWER_FLAG_ID, _setTowerFeatureAssembledForTests, type ISessionScopeHandle, type Scope, } from '@moonshot-ai/agent-core-v2'; import { TowerStore } from '@moonshot-ai/agent-core-v2/features/tower/protocol/index'; import { AgentTranscript, TranscriptStore, type AgentTranscriptSnapshot, type AppendOp, type FrameUpsertOp, type InteractionUpsertOp, type TranscriptFrame, type TranscriptOperation, type TranscriptTask, type TranscriptTurn, } from '@moonshot-ai/transcript'; import { afterEach, describe, expect, it, vi } from 'vitest'; import { bindSessionTranscript } from '../../src/services/transcript/coreBinding'; import { toWireQuestion } from '../../src/protocol/question-wire'; import type { AgentActivitySnapshot } from '@moonshot-ai/agent-core-v2/agent/loop/loop'; import type { LegacyActivityApproval } from '../../src/services/legacyStatus/legacyStatus'; import { AgentTranscriptProjector, type ProjectorBusEvent, } from '../../src/services/transcript/coreEventMap'; import { healTurnOps, TranscriptService, snapshotToOps, TRANSCRIPT_OPS_JOURNAL_CAPACITY, } from '../../src/services/transcript/transcriptService'; _setTowerFeatureAssembledForTests(true); const execFileAsync = promisify(execFile); function ev(payload: Record): ProjectorBusEvent { return payload as unknown as ProjectorBusEvent; } const TEST_SESSION_ID = 'session-test'; function turnOps(turnId: string, items: ReturnType): TranscriptTurn { const turn = items.find( (item): item is TranscriptTurn => item.kind === 'turn' && item.turnId === turnId, ); if (turn === undefined) throw new Error(`turn ${turnId} not found`); return turn; } function coldTranscriptService(home: string): TranscriptService { return new TranscriptService({ homeDir: home, core: { accessor: { get: (token: unknown) => { if (token === ISessionManager) return { get: () => undefined, list: () => [] }; if (token === IWorkspaceInstanceManager) { return { list: () => [], onDidChange: () => ({ dispose: () => undefined }) }; } if (token === ISessionIndex) return { get: async () => ({ workspaceId: 'ws' }) }; return undefined; }, }, } as unknown as Scope, }); } describe('AgentTranscriptProjector', () => { it('keeps cron steers grouped while undo removes only the appropriate suffix', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cron-undo-')); const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); const append = (role: string, text: string, kind?: string) => ({ type: 'context.append_message', message: { role, content: [{ type: 'text', text }], toolCalls: [], origin: kind === undefined ? undefined : { kind } }, }); const steer = (text: string) => ({ type: 'turn.steer', input: [{ type: 'text', text }], origin: { kind: 'cron_job' } }); const records: Record[] = [ append('user', 'first prompt', 'user'), append('assistant', 'first answer'), steer('retained cron'), append('user', 'retained cron', 'cron_job'), append('assistant', 'cron answer'), append('user', 'second prompt', 'user'), steer('removed cron'), append('user', 'removed cron', 'cron_job'), append('assistant', 'second answer'), ]; try { await mkdir(wireDir, { recursive: true }); const service = coldTranscriptService(home); const read = async () => { await writeFile(join(wireDir, 'wire.jsonl'), records.map((record) => JSON.stringify(record)).join('\n') + '\n'); return (await service.readColdSnapshot('s1', 'main'))!.items.filter((item) => item.kind === 'turn'); }; const before = await read(); expect(before.map((turn) => turn.prompt)).toEqual(['first prompt', 'second prompt']); records.push({ type: 'context.undo', count: 1 }); const after = await read(); expect(after.map((turn) => turn.prompt)).toEqual(['first prompt']); expect(after[0]!.steps.flatMap((step) => step.frames).filter((frame) => frame.kind === 'text').map((frame) => frame.text)).toEqual(['first answer', 'retained cron', 'cron answer']); records.push(append('user', 'removed cron', 'cron_job'), append('assistant', 'new cron answer')); expect((await read()).map((turn) => turn.prompt)).toEqual(['first prompt', 'removed cron']); } finally { await rm(home, { recursive: true, force: true }); } }); it('projects a full turn: headers, delta appends, flush, tool frames', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { const mapped = projector.map(event); ops.push(...mapped); tx.apply(mapped); }; feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1, stepId: 'u1' })); feed(ev({ type: 'assistant.delta', turnId: 1, delta: 'Hello' })); feed(ev({ type: 'assistant.delta', turnId: 1, delta: ' world' })); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'call_1', name: 'Bash', args: '{"command":"ls"}', display: { kind: 'command', command: 'ls' }, }), ); feed(ev({ type: 'tool.result', turnId: 1, toolCallId: 'call_1', output: 'file.txt' })); feed(ev({ type: 'turn.step.completed', turnId: 1, step: 1, stepId: 'u1' })); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'completed' })); const appends = ops.filter((op): op is AppendOp => op.op === 'append'); expect(appends.map((op) => [op.offset, op.text])).toEqual([ [0, 'Hello'], [5, ' world'], ]); const upserts = ops.filter((op): op is FrameUpsertOp => op.op === 'frame.upsert'); const flushUpsert = upserts.find( (op) => op.frame.kind === 'text' && op.frame.text === 'Hello world', ); expect(flushUpsert).toBeDefined(); const turn = turnOps('t1', tx.getItems()); expect(turn.state).toBe('completed'); expect(turn.origin).toEqual({ kind: 'user', payload: { kind: 'user' } }); expect(turn.endedAt).toBeTypeOf('string'); expect(turn.steps).toHaveLength(1); const step = turn.steps[0]!; expect(step.state).toBe('completed'); const text = step.frames.find((frame) => frame.kind === 'text'); expect(text).toMatchObject({ role: 'assistant', text: 'Hello world' }); const tool = step.frames.find((frame) => frame.kind === 'tool'); expect(tool).toMatchObject({ frameId: 't1.1.call_1', toolCallId: 'call_1', name: 'Bash', state: 'done', input: { command: 'ls' }, output: 'file.txt', display: { kind: 'command', command: 'ls' }, }); }); it('projects the live prompt from turn.started and keeps it through turn.ended', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => { tx.apply(projector.map(event)); }; feed(ev({ type: 'turn.started', turnId: 0, promptId: 'prompt-1', origin: { kind: 'user' }, prompt: 'fix the bug', })); feed(ev({ type: 'assistant.delta', turnId: 0, delta: 'on it' })); feed(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); const turn = turnOps('t0', tx.getItems()); expect(turn.triggerPromptId).toBe('prompt-1'); expect(turn.prompt).toBe('fix the bug'); expect(turn.state).toBe('completed'); }); it('projects the live prompt for subagent system triggers and keeps it through turn.ended', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => { tx.apply(projector.map(event)); }; feed( ev({ type: 'turn.started', turnId: 0, origin: { kind: 'system_trigger', name: 'subagent' }, prompt: 'scan the repo', promptAttachments: [{ kind: 'image', fileId: 'file_1' }], }), ); feed(ev({ type: 'assistant.delta', turnId: 0, delta: 'scanning' })); feed(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); const turn = turnOps('t0', tx.getItems()); expect(turn.prompt).toBe('scan the repo'); expect(turn.attachmentIds).toEqual(['t0.att1']); expect(turn.state).toBe('completed'); }); it('projects turn.started promptAttachments into attachment entities and turn.attachmentIds', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { const mapped = projector.map(event); ops.push(...mapped); tx.apply(mapped); }; feed( ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'what is this?', promptAttachments: [{ kind: 'image', fileId: 'file_1', name: 'photo.png' }], }), ); feed(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(ops.filter((op) => op.op === 'attachment.upsert')).toEqual([ { op: 'attachment.upsert', attachment: { attachmentId: 't0.att1', mediaType: 'image/*', name: 'photo.png', source: { kind: 'session_media', fileId: 'file_1' }, }, }, ]); const turn = turnOps('t0', tx.getItems()); expect(turn.prompt).toBe('what is this?'); expect(turn.attachmentIds).toEqual(['t0.att1']); expect(tx.getAttachment('t0.att1')).toEqual({ attachmentId: 't0.att1', mediaType: 'image/*', name: 'photo.png', source: { kind: 'session_media', fileId: 'file_1' }, }); }); it('projects turn.started file promptAttachments into path-sourced attachment entities', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { const mapped = projector.map(event); ops.push(...mapped); tx.apply(mapped); }; feed( ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'summarize this', promptAttachments: [ { kind: 'file', name: 'report.pdf', mediaType: 'application/pdf', size: 1234, path: '/data/report.pdf', }, ], }), ); feed(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(ops.filter((op) => op.op === 'attachment.upsert')).toEqual([ { op: 'attachment.upsert', attachment: { attachmentId: 't0.att1', mediaType: 'application/pdf', name: 'report.pdf', size: 1234, }, }, ]); const turn = turnOps('t0', tx.getItems()); expect(turn.prompt).toBe('summarize this'); expect(turn.attachmentIds).toEqual(['t0.att1']); expect(tx.getAttachment('t0.att1')).toEqual({ attachmentId: 't0.att1', mediaType: 'application/pdf', name: 'report.pdf', size: 1234, }); }); it('places late-attach deltas into the engine-reported active step', () => { const tx = new AgentTranscript('main'); const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { stepOrdinal: (turnId) => (turnId === 't0' ? 2 : undefined), }); const ops = projector.map(ev({ type: 'assistant.delta', turnId: 0, delta: 'late' })); tx.apply(ops); const turn = turnOps('t0', tx.getItems()); expect(turn.steps.map((s) => s.stepId)).toEqual(['t0.2']); expect(turn.steps[0]?.frames[0]).toMatchObject({ kind: 'text', text: 'late' }); }); it('adopts a backfilled stream frame on mid-turn attach instead of clobbering it', () => { const tx = new AgentTranscript('main'); tx.apply([ { op: 'turn.upsert', turn: { kind: 'turn', turnId: 't0', ordinal: 0, state: 'running', origin: { kind: 'user' } }, }, { op: 'step.upsert', turnId: 't0', step: { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'running' }, }, { op: 'frame.upsert', turnId: 't0', stepId: 't0.1', frame: { kind: 'text', frameId: 't0.1.f1', role: 'assistant', text: 'Hello ' }, }, ]); const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { stepFrames: (turnId, stepId) => tx.getTurn(turnId)?.steps.find((s) => s.stepId === stepId)?.frames, }); const ops = projector.map(ev({ type: 'assistant.delta', turnId: 0, delta: 'world' })); tx.apply(ops); expect(ops.some((op) => op.op === 'frame.upsert')).toBe(false); const append = ops.find((op): op is AppendOp => op.op === 'append'); expect(append && [append.offset, append.text]).toEqual([6, 'world']); const turn = turnOps('t0', tx.getItems()); const text = turn.steps[0]?.frames.find((frame) => frame.kind === 'text'); expect(text).toMatchObject({ text: 'Hello world' }); const next = projector.map(ev({ type: 'thinking.delta', turnId: 0, delta: 'hmm' })); const created = next.find((op): op is FrameUpsertOp => op.op === 'frame.upsert'); expect(created?.frame.frameId).toBe('t0.1.f2'); }); it('adopts a backfilled tool frame when the result arrives after a mid-bind attach', () => { const tx = new AgentTranscript('main'); tx.apply([ { op: 'turn.upsert', turn: { kind: 'turn', turnId: 't0', ordinal: 0, state: 'running', origin: { kind: 'user' } }, }, { op: 'step.upsert', turnId: 't0', step: { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'running' }, }, { op: 'frame.upsert', turnId: 't0', stepId: 't0.1', frame: { kind: 'tool', frameId: 't0.1.call_1', toolCallId: 'call_1', name: 'Bash', state: 'running', input: { command: 'ls' }, }, }, ]); const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { toolFrame: (toolCallId) => { for (const item of tx.getItems()) { if (item.kind !== 'turn') continue; for (const step of item.steps) { for (const frame of step.frames) { if (frame.kind === 'tool' && frame.toolCallId === toolCallId) { return { turnId: item.turnId, stepId: step.stepId, frame }; } } } } return undefined; }, }); const ops = projector.map(ev({ type: 'tool.result', toolCallId: 'call_1', output: 'file.txt' })); expect(ops).toHaveLength(1); tx.apply(ops); const turn = turnOps('t0', tx.getItems()); const tool = turn.steps[0]?.frames.find((frame) => frame.kind === 'tool'); expect(tool).toMatchObject({ toolCallId: 'call_1', state: 'done', output: 'file.txt' }); }); it('adopts a seeded parent tool frame when subagent.spawned links the child', () => { const tx = new AgentTranscript('main'); tx.apply([ { op: 'turn.upsert', turn: { kind: 'turn', turnId: 't0', ordinal: 0, state: 'running', origin: { kind: 'user' } }, }, { op: 'step.upsert', turnId: 't0', step: { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'running' }, }, { op: 'frame.upsert', turnId: 't0', stepId: 't0.1', frame: { kind: 'tool', frameId: 't0.1.call_agent', toolCallId: 'call_agent', name: 'Agent', state: 'running', input: { prompt: 'scan' }, }, }, ]); const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { toolFrame: (toolCallId) => { for (const item of tx.getItems()) { if (item.kind !== 'turn') continue; for (const step of item.steps) { for (const frame of step.frames) { if (frame.kind === 'tool' && frame.toolCallId === toolCallId) { return { turnId: item.turnId, stepId: step.stepId, frame }; } } } } return undefined; }, }); const ops = projector.map( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'explore', parentToolCallId: 'call_agent', runInBackground: false, }), ); tx.apply(ops); const turn = turnOps('t0', tx.getItems()); const tool = turn.steps[0]?.frames.find((frame) => frame.kind === 'tool'); expect(tool?.kind === 'tool' && tool.agentRefs).toEqual([{ agentId: 'agent-1', role: 'child' }]); }); it('gives live markers their own namespace so they never collide with backfilled markers', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply([{ op: 'marker.upsert', item: { kind: 'marker', markerId: 'm1', marker: 'skill' } }]); const ops = projector.map(ev({ type: 'compaction.started', trigger: 'auto' })); tx.apply(ops); const markers = tx .getItems() .filter((item): item is Extract => item.kind === 'marker'); expect(markers.map((m) => [m.markerId, m.marker])).toEqual([ ['m1', 'skill'], ['live-m1', 'compaction'], ]); }); it('snapshotToOps anchors standalone items so backfill keeps history order against live turns', () => { const snapshot: AgentTranscriptSnapshot = { interactions: [], attachments: [], todos: [], prompts: [], items: [ { kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, prompt: 'one', steps: [], }, { kind: 'marker', markerId: 'm1', marker: 'skill' }, { kind: 'turn', turnId: 't1', ordinal: 1, state: 'completed', origin: { kind: 'user' }, prompt: 'two', steps: [], }, { kind: 'taskref', refId: 'r1', taskId: 'bash-1' }, ], tasks: [], meta: {}, }; const ops = snapshotToOps(snapshot); expect(ops.find((op) => op.op === 'marker.upsert')).toMatchObject({ beforeTurn: 1 }); expect(ops.find((op) => op.op === 'taskref.upsert')).toMatchObject({ beforeTurn: 2 }); const tx = new AgentTranscript('main'); tx.apply([ { op: 'turn.upsert', turn: { kind: 'turn', turnId: 't2', ordinal: 2, state: 'running', origin: { kind: 'user' } }, }, ]); tx.apply(ops); expect( tx.getItems().map((item) => { if (item.kind === 'turn') return item.turnId; if (item.kind === 'marker') return item.markerId; return item.refId; }), ).toEqual(['t0', 'm1', 't1', 'r1', 't2']); }); it('snapshotToOps flattens attachment entities so backfilled attachmentIds never dangle', () => { const snapshot: AgentTranscriptSnapshot = { interactions: [], attachments: [ { attachmentId: 'att_1', mediaType: 'image/*', name: 'shot.png', source: { kind: 'file', fileId: 'file_1' }, }, ], todos: [], prompts: [], items: [ { kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, prompt: 'what is this?', attachmentIds: ['att_1'], steps: [], }, ], tasks: [], meta: {}, }; const ops = snapshotToOps(snapshot); expect(ops.filter((op) => op.op === 'attachment.upsert')).toEqual([ { op: 'attachment.upsert', attachment: snapshot.attachments[0] }, ]); const tx = new AgentTranscript('main'); tx.apply(ops); expect(tx.getAttachment('att_1')).toEqual(snapshot.attachments[0]); expect(turnOps('t0', tx.getItems()).attachmentIds).toEqual(['att_1']); }); it('flushes open frames on turn.ended even without step completion', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed(ev({ type: 'thinking.delta', turnId: 1, delta: 'hmm' })); feed(ev({ type: 'assistant.delta', turnId: 1, delta: 'partial' })); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'cancelled' })); const turn = turnOps('t1', tx.getItems()); expect(turn.state).toBe('cancelled'); const step = turn.steps[0]!; expect(step.state).toBe('interrupted'); expect(step.frames).toContainEqual( expect.objectContaining({ kind: 'thinking', text: 'hmm' }), ); expect(step.frames).toContainEqual( expect.objectContaining({ kind: 'text', text: 'partial' }), ); }); it('marks a user-cancelled turn with an interruption marker, but not programmatic aborts', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'hi' })); feed( ev({ type: 'turn.ended', turnId: 0, reason: 'cancelled', interruptReason: 'user_cancelled' }), ); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' }, prompt: 'again' })); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'cancelled', interruptReason: 'aborted' })); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' }, prompt: 'legacy' })); feed(ev({ type: 'turn.ended', turnId: 2, reason: 'cancelled' })); const markers = tx .getItems() .filter((item): item is Extract => item.kind === 'marker'); expect(markers).toHaveLength(1); expect(markers[0]).toMatchObject({ marker: 'interruption', payload: { turnId: 0, reason: 'user_cancelled' }, }); }); it('carries usage / finishReason / the full timing breakdown on turn.step.completed', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'turn.step.completed', turnId: 1, step: 1, usage: { inputOther: 100, output: 20, inputCacheRead: 30, inputCacheCreation: 40 }, rawFinishReason: 'tool_calls', llmFirstTokenLatencyMs: 120, llmStreamDurationMs: 900, llmRequestBuildMs: 10, llmServerFirstTokenMs: 110, llmServerDecodeMs: 800, llmClientConsumeMs: 100, llmClientBlockedMs: 40, }), ); const step = turnOps('t1', tx.getItems()).steps[0]!; expect(step.state).toBe('completed'); expect(step.usage).toEqual({ inputOther: 100, output: 20, inputCacheRead: 30, inputCacheCreation: 40, }); expect(step.finishReason).toBe('tool_calls'); expect(step.timing).toEqual({ llmFirstTokenLatencyMs: 120, llmStreamDurationMs: 900, llmRequestBuildMs: 10, llmServerFirstTokenMs: 110, llmServerDecodeMs: 800, llmClientConsumeMs: 100, llmClientBlockedMs: 40, }); feed(ev({ type: 'turn.step.started', turnId: 1, step: 2 })); feed( ev({ type: 'turn.step.completed', turnId: 1, step: 2, finishReason: 'stop', rawFinishReason: 'raw_stop', providerFinishReason: 'provider_stop', }), ); expect(turnOps('t1', tx.getItems()).steps[1]!.finishReason).toBe('stop'); feed(ev({ type: 'turn.step.started', turnId: 1, step: 3 })); feed( ev({ type: 'turn.step.completed', turnId: 1, step: 3, providerFinishReason: 'length' }), ); expect(turnOps('t1', tx.getItems()).steps[2]!.finishReason).toBe('length'); }); it('carries endReason / endMessage on turn.step.interrupted', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'turn.step.interrupted', turnId: 1, step: 1, reason: 'aborted', message: 'user cancelled', }), ); const step = turnOps('t1', tx.getItems()).steps[0]!; expect(step.state).toBe('interrupted'); expect(step.endReason).toBe('aborted'); expect(step.endMessage).toBe('user cancelled'); }); it('sets retry on turn.step.retrying and clears it at the terminal upsert', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const step = (): TranscriptTurn['steps'][number] => turnOps('t1', tx.getItems()).steps[0]!; feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'turn.step.retrying', turnId: 1, step: 1, failedAttempt: 1, nextAttempt: 2, maxAttempts: 3, delayMs: 2000, errorName: 'ProviderRateLimitError', errorMessage: '429 too many requests', statusCode: 429, }), ); expect(step().state).toBe('running'); expect(step().retry).toEqual({ failedAttempt: 1, nextAttempt: 2, maxAttempts: 3, delayMs: 2000, errorName: 'ProviderRateLimitError', errorMessage: '429 too many requests', statusCode: 429, }); feed(ev({ type: 'turn.step.completed', turnId: 1, step: 1 })); expect(step().state).toBe('completed'); expect(step().retry).toBeUndefined(); }); it('fills durationMs / error / accumulated step usage on turn.ended', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'turn.step.completed', turnId: 1, step: 1, usage: { inputOther: 100, output: 10, inputCacheRead: 5, inputCacheCreation: 50 }, }), ); feed(ev({ type: 'turn.step.started', turnId: 1, step: 2 })); feed( ev({ type: 'turn.step.completed', turnId: 1, step: 2, usage: { inputOther: 200, output: 20, inputCacheRead: 0, inputCacheCreation: 25 }, }), ); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'completed', durationMs: 4200 })); const turn = turnOps('t1', tx.getItems()); expect(turn.durationMs).toBe(4200); expect(turn.usage).toEqual({ inputTokens: 375, cachedTokens: 5, outputTokens: 30 }); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' } })); feed( ev({ type: 'turn.ended', turnId: 2, reason: 'failed', durationMs: 50, error: { code: 'internal', message: 'kaboom', retryable: false }, }), ); const failed = turnOps('t2', tx.getItems()); expect(failed.state).toBe('failed'); expect(failed.error).toBe('kaboom'); expect(failed.usage).toBeUndefined(); }); it('takes the turn header endedAt from the turn.ended event time', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'completed', time: 1_700_000_000_000 })); expect(turnOps('t1', tx.getItems()).endedAt).toBe(new Date(1_700_000_000_000).toISOString()); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' } })); feed(ev({ type: 'turn.ended', turnId: 2, reason: 'completed' })); expect(turnOps('t2', tx.getItems()).endedAt).toBeTypeOf('string'); }); it('accumulates tool.call.delta into inputText, kept across tool.call.started', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const toolFrame = (toolCallId: string): TranscriptFrame | undefined => turnOps('t1', tx.getItems()) .steps.flatMap((step) => step.frames) .find((frame) => frame.kind === 'tool' && frame.toolCallId === toolCallId); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'tool.call.delta', turnId: 1, toolCallId: 'c1', name: 'Bash', argumentsPart: '{"comm', }), ); feed(ev({ type: 'tool.call.delta', turnId: 1, toolCallId: 'c1', argumentsPart: 'and":"ls"}' })); expect(toolFrame('c1')).toMatchObject({ kind: 'tool', frameId: 't1.1.c1', name: 'Bash', state: 'running', inputText: '{"command":"ls"}', }); feed(ev({ type: 'tool.call.delta', turnId: 1, toolCallId: 'c2', argumentsPart: '{}' })); expect(toolFrame('c2')).toMatchObject({ name: '', inputText: '{}' }); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'c1', name: 'Bash', args: { command: 'ls' }, }), ); expect(toolFrame('c1')).toMatchObject({ input: { command: 'ls' }, inputText: '{"command":"ls"}', }); feed(ev({ type: 'tool.call.delta', turnId: 1, toolCallId: 'c1', argumentsPart: '\n' })); expect(toolFrame('c1')).toMatchObject({ inputText: '{"command":"ls"}\n' }); feed(ev({ type: 'tool.result', turnId: 1, toolCallId: 'c1', output: 'file.txt' })); expect(toolFrame('c1')).toMatchObject({ state: 'done', inputText: '{"command":"ls"}\n' }); }); it('overwrites tool frame progress and drops progress for unknown calls', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); expect( projector.map( ev({ type: 'tool.progress', turnId: 1, toolCallId: 'ghost', update: { kind: 'stdout', text: 'x' }, }), ), ).toEqual([]); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'c1', name: 'Bash', args: {} }), ); feed( ev({ type: 'tool.progress', turnId: 1, toolCallId: 'c1', update: { kind: 'stdout', text: 'line1' }, }), ); const tool = (): TranscriptFrame | undefined => turnOps('t1', tx.getItems()).steps[0]!.frames.find((frame) => frame.kind === 'tool'); expect(tool()).toMatchObject({ progress: { kind: 'stdout', text: 'line1' } }); feed( ev({ type: 'tool.progress', turnId: 1, toolCallId: 'c1', update: { kind: 'progress', percent: 40 }, }), ); expect(tool()).toMatchObject({ progress: { kind: 'progress', percent: 40 } }); expect((tool() as { progress?: Record }).progress?.['text']).toBeUndefined(); }); it('marks tool.result errors and keeps the display payload', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'c1', name: 'Read', args: { path: '/x' }, display: { kind: 'file', path: '/x' }, }), ); feed(ev({ type: 'tool.result', turnId: 1, toolCallId: 'c1', output: 'ENOENT', isError: true })); const tool = turnOps('t1', tx.getItems()).steps[0]!.frames.find((f) => f.kind === 'tool'); expect(tool).toMatchObject({ state: 'error', output: 'ENOENT', error: 'ENOENT', display: { kind: 'file', path: '/x' }, }); }); it('projects process tasks as shell tasks with streaming output', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { const mapped = projector.map(event); ops.push(...mapped); tx.apply(mapped); }; const started = { taskId: 'bash-1', kind: 'process', description: 'ls -la', status: 'running', detached: false, startedAt: 1_700_000_000_000, endedAt: null, }; feed(ev({ type: 'task.started', info: started })); feed(ev({ type: 'shell.started', commandId: 'cmd-1', taskId: 'bash-1' })); feed(ev({ type: 'shell.output', commandId: 'cmd-1', update: { kind: 'stdout', text: 'a\n' } })); feed(ev({ type: 'shell.output', commandId: 'cmd-1', update: { kind: 'stderr', text: 'b\n' } })); feed( ev({ type: 'task.terminated', info: { ...started, status: 'completed', endedAt: 1_700_000_001_000 }, }), ); expect(ops.some((op) => op.op === 'taskref.upsert' && op.item.taskId === 'bash-1')).toBe(true); const appends = ops.filter((op): op is AppendOp => op.op === 'append'); expect(appends.map((op) => [op.offset, op.text])).toEqual([ [0, 'a\n'], [2, 'b\n'], ]); const task = tx.getTask('bash-1'); expect(task).toMatchObject({ kind: 'shell', state: 'completed', detached: false, description: 'ls -la', outputTail: 'a\nb\n', }); expect( projector.map( ev({ type: 'shell.output', commandId: 'cmd-1', update: { kind: 'progress', percent: 50 } }), ), ).toEqual([]); }); it('fills the shell task output from late stderr chunks before completing', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'shell.started', commandId: 'c1', taskId: 'task-1' }))); tx.apply( projector.map(ev({ type: 'shell.output', commandId: 'c1', update: { kind: 'stderr', text: 'boom' } })), ); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c1', isError: true }))); expect(tx.getTask('task-1')).toMatchObject({ state: 'failed', outputTail: 'boom' }); }); it('routes shell output/completion via the event taskId when shell.started was missed', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.map( ev({ type: 'shell.output', commandId: 'c1', taskId: 'task-1', update: { kind: 'stdout', text: 'hello' } }), ), ); expect(tx.getTask('task-1')).toMatchObject({ kind: 'shell', state: 'running', outputTail: 'hello' }); expect(tx.getItems()).toContainEqual(expect.objectContaining({ kind: 'taskref', taskId: 'task-1' })); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c1', taskId: 'task-1', isError: false }))); expect(tx.getTask('task-1')).toMatchObject({ state: 'completed', outputTail: 'hello' }); }); it('emits a taskref when only shell.completed arrives for a command', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c1', taskId: 'task-1', isError: true }))); expect(tx.getTask('task-1')).toMatchObject({ kind: 'shell', state: 'failed' }); expect(tx.getItems()).toContainEqual(expect.objectContaining({ kind: 'taskref', taskId: 'task-1' })); }); it('projects no-taskId shell failures under a synthetic per-command task id', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.map(ev({ type: 'shell.output', commandId: 'c1', update: { kind: 'stderr', text: 'boom' } })), ); expect(tx.getTask('shell-c1')).toMatchObject({ kind: 'shell', state: 'running', outputTail: 'boom' }); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c1', isError: true }))); expect(tx.getTask('shell-c1')).toMatchObject({ state: 'failed', outputTail: 'boom' }); expect(tx.getItems()).toContainEqual(expect.objectContaining({ kind: 'taskref', taskId: 'shell-c1' })); }); it('marks a foreground shell task terminal on shell.completed', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'shell.started', commandId: 'c1', taskId: 'task-1' }))); expect(tx.getTask('task-1')?.state).toBe('running'); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c1', isError: false }))); expect(tx.getTask('task-1')).toMatchObject({ kind: 'shell', state: 'completed' }); expect(tx.getTask('task-1')?.endedAt).toBeTypeOf('string'); tx.apply(projector.map(ev({ type: 'shell.started', commandId: 'c2', taskId: 'task-2' }))); tx.apply(projector.map(ev({ type: 'shell.completed', commandId: 'c2', isError: true }))); expect(tx.getTask('task-2')?.state).toBe('failed'); }); it('ignores task.notified (it re-surfaces as an origin:task turn)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); expect(projector.map(ev({ type: 'task.notified', taskId: 't' }))).toEqual([]); }); it('links spawned subagents to the spawning tool frame (member for swarm)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'call_swarm', name: 'AgentSwarm', args: {}, }), ); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-0', subagentName: 'worker', parentToolCallId: 'call_swarm', description: 'scan the repo', swarmIndex: 0, runInBackground: false, model: 'example-model', thinkingEffort: 'high', }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-0', resultSummary: 'done' })); const tool = turnOps('t1', tx.getItems()).steps[0]!.frames.find((f) => f.kind === 'tool'); expect(tool).toMatchObject({ agentRefs: [{ agentId: 'agent-0', role: 'member' }], }); const task = tx.getTask('agent-0'); expect(task).toMatchObject({ kind: 'subagent', state: 'completed', agentId: 'agent-0', description: 'scan the repo', detached: false, model: 'example-model', thinkingEffort: 'high', }); }); it('keys an Agent-tool subagent row by its registered task id and folds the lifecycle', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'explore', parentToolCallId: 'call-1', description: 'Inspect files', runInBackground: true, taskId: 'task-9', }), ); feed( ev({ type: 'task.started', info: { taskId: 'task-9', kind: 'agent', description: 'Inspect files', status: 'running', detached: true, agentId: 'agent-1', startedAt: 1_700_000_000_000, endedAt: null, }, }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); feed( ev({ type: 'task.terminated', info: { taskId: 'task-9', kind: 'agent', description: 'Inspect files', status: 'completed', detached: true, agentId: 'agent-1', startedAt: 1_700_000_000_000, endedAt: 1_700_000_001_000, }, }), ); expect(tx.getTask('task-9')).toMatchObject({ kind: 'subagent', state: 'completed', agentId: 'agent-1', description: 'Inspect files', detached: true, resultSummary: 'done', }); expect(tx.getTask('agent-1')).toBeUndefined(); }); it('drops the stale task mapping when a child respawns without a task id', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'explore', parentToolCallId: 'call-1', description: 'Inspect files', runInBackground: true, taskId: 'task-9', }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'worker', parentToolCallId: 'call-2', description: 'scan again', runInBackground: false, }), ); feed(ev({ type: 'subagent.started', subagentId: 'agent-1' })); expect(tx.getTask('task-9')).toMatchObject({ state: 'completed', resultSummary: 'done' }); expect(tx.getTask('agent-1')).toMatchObject({ kind: 'subagent', state: 'running' }); }); it('resets the agent-run generation stamps when a completed subagent is resumed for rework', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'worker', parentToolCallId: 'call-1', description: 'tower worker', runInBackground: true, }), ); feed( ev({ type: 'task.started', info: { taskId: 'task-9', kind: 'agent', description: 'tower worker', status: 'running', detached: true, agentId: 'agent-1', startedAt: 1_700_000_000_000, endedAt: null, }, }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); feed( ev({ type: 'task.terminated', info: { taskId: 'task-9', kind: 'agent', description: 'tower worker', status: 'completed', detached: true, agentId: 'agent-1', startedAt: 1_700_000_000_000, endedAt: 1_700_000_001_000, }, }), ); const firstRun = tx.getTask('agent-1'); expect(firstRun).toMatchObject({ state: 'completed', resultSummary: 'done' }); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'worker', parentToolCallId: 'call-2', description: 'tower worker', runInBackground: true, taskId: 'task-10', }), ); feed(ev({ type: 'subagent.started', subagentId: 'agent-1' })); const resumed = tx.getTask('agent-1'); expect(resumed).toMatchObject({ kind: 'subagent', state: 'running' }); expect(resumed?.endedAt).toBeUndefined(); expect(resumed?.resultSummary).toBeUndefined(); expect(resumed?.error).toBeUndefined(); expect((resumed?.startedAt ?? '') >= (firstRun?.endedAt ?? '')).toBe(true); expect(tx.getTask('task-9')).toMatchObject({ state: 'completed', resultSummary: 'done', endedAt: new Date(1_700_000_001_000).toISOString(), }); expect(tx.getTask('task-10')).toMatchObject({ kind: 'subagent', state: 'running', agentId: 'agent-1', }); }); it('clears the terminal stamps when a terminal subagent is respawned without a task id', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'worker', parentToolCallId: 'call-1', description: 'scan', runInBackground: false, }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); const finished = tx.getTask('agent-1'); expect(finished).toMatchObject({ state: 'completed', resultSummary: 'done' }); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'worker', parentToolCallId: 'call-2', description: 'scan again', runInBackground: false, }), ); const respawned = tx.getTask('agent-1'); expect(respawned).toMatchObject({ kind: 'subagent', state: 'running' }); expect(respawned?.endedAt).toBeUndefined(); expect(respawned?.resultSummary).toBeUndefined(); expect((respawned?.startedAt ?? '') >= (finished?.endedAt ?? '')).toBe(true); }); it('keeps the run stamps when subagent.started arrives for an already-running task', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'subagent.spawned', subagentId: 'agent-1', subagentName: 'explore', parentToolCallId: 'call-1', description: 'Inspect files', runInBackground: true, taskId: 'task-9', }), ); feed(ev({ type: 'subagent.started', subagentId: 'agent-1' })); const startedAt = tx.getTask('task-9')?.startedAt; feed(ev({ type: 'subagent.started', subagentId: 'agent-1' })); expect(tx.getTask('task-9')).toMatchObject({ state: 'running', startedAt }); }); it('recovers the agent → task association from a backfilled task.started', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'task.started', info: { taskId: 'task-9', kind: 'agent', description: 'Inspect files', status: 'running', detached: true, agentId: 'agent-1', startedAt: 1_700_000_000_000, endedAt: null, }, }), ); feed(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); expect(tx.getTask('task-9')).toMatchObject({ state: 'completed', resultSummary: 'done' }); expect(tx.getTask('agent-1')).toBeUndefined(); }); it('projects goal updates into meta.goal plus an inline marker', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const snapshot = { goalId: 'g1', objective: 'ship it', status: 'active', completionCriterion: 'tests green', turnsUsed: 3, tokensUsed: 1234, wallClockMs: 5000, budget: { tokenBudget: 50000 }, }; const ops = projector.map(ev({ type: 'goal.updated', snapshot, change: { kind: 'lifecycle' } })); tx.apply(ops); expect(tx.getMeta().goal).toEqual({ objective: 'ship it', status: 'active', completionCriterion: 'tests green', budgetUsed: 1234, budgetLimit: 50000, }); const marker = tx.getItems().find((item) => item.kind === 'marker'); expect(marker).toMatchObject({ marker: 'goal', payload: { snapshot } }); const clearedOps = projector.map(ev({ type: 'goal.updated', snapshot: null })); expect(clearedOps[0]).toEqual({ op: 'meta.merge', meta: { goal: null } }); tx.apply(clearedOps); expect(tx.getMeta().goal).toBeUndefined(); }); it('mirrors plan / swarm mode slices into meta.modes (only when provided)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'agent.status.updated', planMode: true }))); tx.apply(projector.map(ev({ type: 'agent.status.updated', swarmMode: true }))); expect(tx.getMeta().modes).toEqual({ plan: {}, swarm: {} }); tx.apply(projector.map(ev({ type: 'agent.status.updated', planMode: false }))); expect(tx.getMeta().modes).toEqual({ swarm: {} }); tx.apply(projector.map(ev({ type: 'agent.status.updated', swarmMode: false }))); expect(tx.getMeta().modes).toBeUndefined(); }); it('mirrors the tower mode slice into meta.modes', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'agent.status.updated', towerMode: true }))); expect(tx.getMeta().modes).toEqual({ tower: {} }); tx.apply(projector.map(ev({ type: 'agent.status.updated', towerMode: false }))); expect(tx.getMeta().modes).toBeUndefined(); }); it('mirrors status slices into meta.agent (shallow-merged across slices)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const usageOnly = projector.map(ev({ type: 'agent.status.updated', usage: {} })); expect(usageOnly).toEqual([{ op: 'meta.merge', meta: { agent: { usage: {} } } }]); feed(ev({ type: 'agent.status.updated', model: 'k2', thinkingEffort: 'high' })); feed( ev({ type: 'agent.status.updated', usage: { total: { inputOther: 1, output: 2, inputCacheRead: 3, inputCacheCreation: 4 }, }, }), ); feed( ev({ type: 'agent.status.updated', contextTokens: 1000, maxContextTokens: 200000, contextUsage: 0.5, }), ); feed(ev({ type: 'agent.status.updated', permission: 'yolo' })); expect(tx.getMeta().agent).toEqual({ model: 'k2', thinkingEffort: 'high', usage: { total: { inputOther: 1, output: 2, inputCacheRead: 3, inputCacheCreation: 4 } }, contextTokens: 1000, maxContextTokens: 200000, contextUsage: 0.5, permission: 'yolo', }); feed(ev({ type: 'agent.status.updated', model: 'k3' })); expect(tx.getMeta().agent).toMatchObject({ model: 'k3', thinkingEffort: 'high' }); }); it('maps domain events into meta.agent.phase', () => { let snapshot: AgentActivitySnapshot = {}; let approvals: readonly LegacyActivityApproval[] = []; const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { activitySnapshot: () => snapshot, pendingApprovals: () => approvals, }); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const runningTurn = (overrides: Record): AgentActivitySnapshot => ({ turn: { turnId: 1, phase: 'running', step: 1, ending: false, activeToolCalls: [], since: 1000, ...overrides, }, }); snapshot = runningTurn({}); feed(ev({ type: 'turn.started', agentId: 'main', turnId: 1, origin: { kind: 'user' } })); expect(tx.getMeta().agent?.phase).toEqual({ kind: 'running', turnId: 1, step: 1, stepId: '', since: 1000, }); snapshot = runningTurn({ phase: 'retrying', retry: { failedAttempt: 1, nextAttempt: 2, maxAttempts: 10, delayMs: 500 }, }); feed( ev({ type: 'turn.step.retrying', agentId: 'main', turnId: 1, step: 1, failedAttempt: 1, nextAttempt: 2, maxAttempts: 10, delayMs: 500, errorName: 'status', errorMessage: 'boom', }), ); expect(tx.getMeta().agent?.phase).toMatchObject({ kind: 'retrying', failedAttempt: 1, nextAttempt: 2, maxAttempts: 10, }); approvals = [{ approvalId: 'ap1', toolCallId: 'c1', since: 1500 }]; feed( ev({ type: 'permission.approval.requested', agentId: 'main', turnId: 1, toolCallId: 'c1', id: 'ap1', }), ); expect(tx.getMeta().agent?.phase).toEqual({ kind: 'awaiting_approval', turnId: 1, step: 1, approval: { approvalId: 'ap1', toolCallId: 'c1' }, since: 1500, }); feed(ev({ type: 'turn.ended', agentId: 'main', turnId: 1, reason: 'completed', durationMs: 100 })); expect(tx.getMeta().agent?.phase).toMatchObject({ kind: 'ended', turnId: 1, reason: 'completed', durationMs: 100, }); snapshot = {}; feed(ev({ type: 'permission.approval.resolved', agentId: 'main', turnId: 1, toolCallId: 'c1', id: 'ap1', decision: 'approved' })); expect(tx.getMeta().agent?.phase).toMatchObject({ kind: 'ended', turnId: 1 }); }); it('projects plan.revision as a marker and refines the active plan badge', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID, { resolvePlanRevisionKey: (key) => `sessions/w/s/agents/main/${key}`, }); const tx = new AgentTranscript('main'); const revision = { type: 'plan.revision', id: 'plan-1', version: 1, key: 'plan/plan-1/v1.md', sha256: 'deadbeef', bytes: 128, }; tx.apply(projector.map(ev(revision))); expect(tx.getMeta().modes).toBeUndefined(); tx.apply(projector.map(ev({ type: 'agent.status.updated', planMode: true }))); expect(tx.getMeta().modes).toEqual({ plan: {} }); tx.apply( projector.map(ev({ ...revision, version: 2, key: 'plan/plan-1/v2.md' })), ); expect(tx.getMeta().modes).toEqual({ plan: { reviewPath: 'sessions/w/s/agents/main/plan/plan-1/v2.md', version: 2 }, }); const markers = tx .getItems() .filter((item) => item.kind === 'marker' && item.marker === 'plan.revision'); expect(markers.map((item) => item.kind === 'marker' && item.markerId)).toEqual([ 'live-m1', 'live-m2', ]); expect(markers[1]).toMatchObject({ payload: { id: 'plan-1', version: 2, path: 'sessions/w/s/agents/main/plan/plan-1/v2.md', sha256: 'deadbeef', bytes: 128, }, }); tx.apply(projector.map(ev({ type: 'agent.status.updated', planMode: false }))); expect(tx.getMeta().modes).toBeUndefined(); expect( tx.getItems().filter((item) => item.kind === 'marker' && item.marker === 'plan.revision'), ).toHaveLength(2); }); it('projects skill / plugin-command / cron / compaction / hook / undo markers', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'skill.activated', activationId: 'a1', skillName: 'gen-docs', trigger: 'user-slash' })); feed( ev({ type: 'plugin_command.activated', activationId: 'a2', pluginId: 'p', commandName: 'c', trigger: 'user-slash', }), ); feed(ev({ type: 'cron.fired', origin: { kind: 'cron_job', jobId: 'j1' }, prompt: 'ping' })); feed(ev({ type: 'compaction.started', trigger: 'auto' })); feed(ev({ type: 'compaction.completed', result: { kept: 3 } })); feed(ev({ type: 'hook.result', hookEvent: 'SessionStart', content: 'hook says hi' })); feed( ev({ type: 'hook.result', turnId: 3, hookEvent: 'UserPromptSubmit', content: 'blocked by hook', blocked: true, }), ); feed(ev({ type: 'context.spliced', start: 1, deleteCount: 2, messages: [] })); const markers = tx .getItems() .filter((item): item is Extract => item.kind === 'marker'); expect(markers.map((m) => m.marker)).toEqual([ 'skill', 'skill', 'cron.fired', 'compaction', 'compaction', 'hook', 'hook', 'undo', ]); expect(markers[1]!.payload).toMatchObject({ variant: 'plugin_command' }); expect(markers[3]!.payload).toMatchObject({ phase: 'started' }); expect(markers[4]!.payload).toMatchObject({ phase: 'completed' }); expect(markers[5]!.payload).toEqual({ hookEvent: 'SessionStart', content: 'hook says hi' }); expect(markers[6]!.payload).toEqual({ turnId: 3, hookEvent: 'UserPromptSubmit', content: 'blocked by hook', blocked: true, }); expect(markers[7]!.payload).toMatchObject({ start: 1, deleteCount: 2 }); }); it('does not infer removed turns from an undo count', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); expect(projector.map(ev({ type: 'context.undone', agentId: 'main', turns: 1, fromTurnId: 0 }))).toEqual([]); }); it('projects error / warning events as notice markers outside any step', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.map(ev({ type: 'error', code: 'mcp.failed', message: 'boom', retryable: false })), ); tx.apply(projector.map(ev({ type: 'warning', message: 'AGENTS.md oversized' }))); const markers = tx .getItems() .filter((item): item is Extract => item.kind === 'marker'); expect(markers).toHaveLength(2); expect(markers[0]).toMatchObject({ marker: 'notice', payload: { level: 'error', message: 'boom', event: { code: 'mcp.failed' } }, }); expect(markers[1]).toMatchObject({ marker: 'notice', payload: { level: 'warning', message: 'AGENTS.md oversized' }, }); }); it('emits interactions as global entities only (no inline frame), back-links on resolve', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 2, step: 1 })); feed( ev({ type: 'tool.call.started', turnId: 2, toolCallId: 'call_9', name: 'Bash', args: {}, }), ); const request = { toolCallId: 'call_9', toolName: 'Bash', action: 'run', display: { kind: 'command', command: 'rm -rf /tmp/x' }, }; tx.apply( projector.mapInteractionRequested({ id: 'apr-1', kind: 'approval', payload: request, createdAt: 1000, }), ); expect(turnOps('t2', tx.getItems()).steps[0]!.frames.map((f) => f.kind)).toEqual(['tool']); expect(tx.getInteraction('apr-1')).toMatchObject({ interactionId: 'apr-1', interactionKind: 'approval', toolCallId: 'call_9', state: 'pending', request, }); expect(tx.listPendingInteractions()).toEqual(['apr-1']); tx.apply(projector.mapInteractionResolved('apr-1', { decision: 'approved', scope: 'session' })); const tool = turnOps('t2', tx.getItems()).steps[0]!.frames.find((f) => f.kind === 'tool'); expect(tool).toMatchObject({ approvalId: 'apr-1' }); expect(turnOps('t2', tx.getItems()).steps[0]!.frames.map((f) => f.kind)).toEqual(['tool']); expect(tx.getInteraction('apr-1')).toMatchObject({ state: 'approved', response: { decision: 'approved', scope: 'session' }, }); expect(tx.listPendingInteractions()).toEqual([]); }); it('surfaces a mid-turn task notification as a user input frame linked to the task', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const notified = (): ProjectorBusEvent => ev({ type: 'task.notified', notificationType: 'task.completed', title: 'Background process completed', body: 'pnpm test — 42 passed', severity: 'info', sourceKind: 'background_task', sourceId: 'task_1', }); tx.apply(projector.map(notified())); expect(tx.getItems()).toHaveLength(0); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); tx.apply(projector.map(notified())); const frames = turnOps('t1', tx.getItems()).steps[0]!.frames; const frame = frames.find((f) => f.kind === 'text' && f.role === 'user'); expect(frame).toMatchObject({ kind: 'text', role: 'user', taskId: 'task_1' }); expect(frame?.kind === 'text' && frame.text).toContain('Background process completed'); }); it('attaches a between-steps task notification to the following step', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const notified = (sourceId: string): ProjectorBusEvent => ev({ type: 'task.notified', notificationType: 'task.completed', title: 'Background agent completed', body: 'inspect done.', severity: 'info', sourceKind: 'background_task', sourceId, }); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed(ev({ type: 'turn.step.completed', turnId: 1, step: 1 })); feed(notified('task_1')); feed(notified('task_2')); expect(turnOps('t1', tx.getItems()).steps[0]!.frames).toHaveLength(0); feed(ev({ type: 'turn.step.started', turnId: 1, step: 2 })); const steps = turnOps('t1', tx.getItems()).steps; expect(steps).toHaveLength(2); expect(steps[1]!.frames.map((f) => f.kind === 'text' && f.taskId)).toEqual(['task_1', 'task_2']); }); it('drops a task notification that is the turn prompt itself', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'task', taskId: 'task_1' } })); feed( ev({ type: 'task.notified', notificationType: 'task.completed', title: 'Background agent completed', body: 'inspect done.', severity: 'info', sourceKind: 'background_task', sourceId: 'task_1', }), ); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); expect(turnOps('t1', tx.getItems()).steps[0]!.frames).toHaveLength(0); }); it('keeps a different task’s notification in a task-origin turn', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'task', taskId: 'task_1' } })); feed( ev({ type: 'task.notified', notificationType: 'task.completed', title: 'Background agent completed', body: 'second task done.', severity: 'info', sourceKind: 'background_task', sourceId: 'task_2', }), ); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); const frames = turnOps('t1', tx.getItems()).steps[0]!.frames; expect(frames.map((f) => f.kind === 'text' && f.taskId)).toEqual(['task_2']); }); it('drops a buffered task notification when the turn ends before the next step', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed(ev({ type: 'turn.step.completed', turnId: 1, step: 1 })); feed( ev({ type: 'task.notified', notificationType: 'task.completed', title: 'Background agent completed', body: 'inspect done.', severity: 'info', sourceKind: 'background_task', sourceId: 'task_1', }), ); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'completed' })); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 2, step: 1 })); expect(turnOps('t2', tx.getItems()).steps[0]!.frames).toHaveLength(0); }); it('replaces the global todo document on a confirmed TodoList write', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); feed(ev({ type: 'turn.step.started', turnId: 1, step: 1 })); feed(ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'call_read', name: 'TodoList', args: {} })); feed(ev({ type: 'tool.result', toolCallId: 'call_read', output: '2 todos' })); expect(tx.getTodo('todo')).toBeUndefined(); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'call_write', name: 'TodoList', args: { todos: [{ title: 'write tests', status: 'in_progress' }, { title: 'ship', status: 'pending' }] }, }), ); const writeFrame = turnOps('t1', tx.getItems()).steps[0]!.frames.find( (f) => f.kind === 'tool' && f.toolCallId === 'call_write', ); expect(writeFrame?.kind === 'tool' && writeFrame.todoId).toBe('todo'); feed(ev({ type: 'tool.result', toolCallId: 'call_write', output: 'updated' })); expect(tx.getTodo('todo')?.items).toEqual([ { title: 'write tests', status: 'in_progress' }, { title: 'ship', status: 'pending' }, ]); feed( ev({ type: 'tool.call.started', turnId: 1, toolCallId: 'call_fail', name: 'TodoList', args: { todos: [] }, }), ); feed(ev({ type: 'tool.result', toolCallId: 'call_fail', output: 'boom', isError: true })); expect(tx.getTodo('todo')?.items).toHaveLength(2); }); it('emits an unanchored entity when the payload has no toolCallId', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.mapInteractionRequested({ id: 'q1', kind: 'question', payload: { questions: [{ question: 'Pick', options: [] }] }, createdAt: 1000, }), ); expect(tx.getItems()).toHaveLength(0); const entity = tx.getInteraction('q1'); expect(entity).toMatchObject({ interactionKind: 'question', state: 'pending' }); expect(entity?.toolCallId).toBeUndefined(); expect(tx.listPendingInteractions()).toEqual(['q1']); tx.apply(projector.mapInteractionResolved('q1', null)); expect(tx.getInteraction('q1')).toMatchObject({ state: 'dismissed' }); expect(tx.listPendingInteractions()).toEqual([]); }); it('projects question requests onto the wire shape with stable question/option ids', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.mapInteractionRequested({ id: 'q-wire', kind: 'question', payload: { toolCallId: 'call_q', turnId: 3, questions: [ { question: 'Pick one', header: 'h', body: 'b', multiSelect: false, otherLabel: 'Other', otherDescription: 'free text', options: [{ label: 'A', description: 'first' }, { label: 'B' }], }, ], }, createdAt: 7000, }), ); const entity = tx.getInteraction('q-wire'); expect(entity?.toolCallId).toBe('call_q'); expect(entity?.request).toEqual({ question_id: 'q-wire', session_id: TEST_SESSION_ID, questions: [ { id: 'q_0', question: 'Pick one', header: 'h', body: 'b', multi_select: false, allow_other: true, other_label: 'Other', other_description: 'free text', options: [ { id: 'opt_0_0', label: 'A', description: 'first' }, { id: 'opt_0_1', label: 'B' }, ], }, ], created_at: new Date(7000).toISOString(), turn_id: 3, tool_call_id: 'call_q', }); tx.apply(projector.mapInteractionResolved('q-wire', { q_0: 'A' })); expect(tx.getInteraction('q-wire')).toMatchObject({ state: 'answered' }); }); it('keeps a malformed question payload raw instead of failing the projection', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply( projector.mapInteractionRequested({ id: 'q-raw', kind: 'question', payload: { toolCallId: 'call_x' }, createdAt: 1000, }), ); const entity = tx.getInteraction('q-raw'); expect(entity?.toolCallId).toBe('call_x'); expect(entity?.request).toEqual({ toolCallId: 'call_x' }); }); it('preserves queue metadata through prompt lifecycle updates', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const metadata = [{ display_text: 'Save button', kimi_code_composer: { version: 1, doc: { type: 'doc' } } }]; feed(ev({ type: 'prompt.submitted', promptId: 'p1', userMessageId: 'm1', status: 'queued', content: [{ type: 'text', text: 'wire' }], clientMetadata: metadata, createdAt: '2026-01-01T00:00:00.000Z' })); feed(ev({ type: 'prompt.queued', promptId: 'p1', content: [{ type: 'text', text: 'wire' }], queueLength: 1, clientMetadata: metadata })); feed(ev({ type: 'prompt.started', promptId: 'p1' })); feed(ev({ type: 'prompt.completed', promptId: 'p1', finishedAt: '2026-01-01T00:00:02.000Z', reason: 'completed' })); expect(tx.getPrompt('p1')?.clientMetadata).toEqual(metadata); }); it('projects prompt submitted/completed/aborted/steered as global queue entities', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'prompt.submitted', promptId: 'p1', userMessageId: 'm1', status: 'running', content: [{ type: 'text', text: 'first' }], createdAt: '2026-01-01T00:00:00.000Z', }), ); feed( ev({ type: 'prompt.submitted', promptId: 'p2', userMessageId: 'm2', status: 'queued', content: [{ type: 'text', text: 'second' }], createdAt: '2026-01-01T00:00:01.000Z', }), ); expect(tx.getPrompt('p1')).toMatchObject({ status: 'running', userMessageId: 'm1' }); expect(tx.getPrompt('p2')).toMatchObject({ status: 'queued' }); feed(ev({ type: 'prompt.started', promptId: 'p2' })); expect(tx.getPrompt('p2')).toMatchObject({ status: 'running', userMessageId: 'm2', content: [{ type: 'text', text: 'second' }], createdAt: '2026-01-01T00:00:01.000Z', }); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [ { type: 'text', text: 'first' }, { type: 'text', text: 'second' }, ], steeredAt: '2026-01-01T00:00:02.000Z', }), ); expect(tx.getPrompt('p1')).toMatchObject({ status: 'running', steeredAt: '2026-01-01T00:00:02.000Z', content: [ { type: 'text', text: 'first' }, { type: 'text', text: 'second' }, ], }); expect(tx.getPrompt('p2')).toMatchObject({ status: 'completed', userMessageId: 'm2', steeredAt: '2026-01-01T00:00:02.000Z', finishedAt: '2026-01-01T00:00:02.000Z', }); feed( ev({ type: 'prompt.completed', promptId: 'p1', finishedAt: '2026-01-01T00:00:10.000Z', reason: 'completed', }), ); expect(tx.getPrompt('p1')).toMatchObject({ status: 'completed', finishedAt: '2026-01-01T00:00:10.000Z', content: [ { type: 'text', text: 'first' }, { type: 'text', text: 'second' }, ], }); feed(ev({ type: 'prompt.aborted', promptId: 'p3', abortedAt: '2026-01-01T00:00:03.000Z' })); expect(tx.getPrompt('p3')).toEqual({ promptId: 'p3', status: 'aborted', createdAt: '2026-01-01T00:00:03.000Z', finishedAt: '2026-01-01T00:00:03.000Z', }); feed( ev({ type: 'prompt.completed', promptId: 'p4', finishedAt: '2026-01-01T00:00:04.000Z', reason: 'failed', }), ); expect(tx.getPrompt('p4')).toEqual({ promptId: 'p4', status: 'failed', createdAt: '2026-01-01T00:00:04.000Z', finishedAt: '2026-01-01T00:00:04.000Z', }); }); it('projects prompt.steered media content to the wire shape (no daemon ref or path leak)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [ { type: 'text', text: 'look at this' }, { type: 'image_url', imageUrl: { url: 'kimi-file://f_img1?path=%2Fabs%2Fsession%2Fmedia%2Ff_img1.png' }, }, ], steeredAt: '2026-01-01T00:00:02.000Z', }), ); const prompt = tx.getPrompt('p1'); expect(prompt?.content).toEqual([ { type: 'text', text: 'look at this' }, { type: 'image', source: { kind: 'session_media', file_id: 'f_img1' } }, ]); }); it('projects turn.steer as a user frame at the next step start, pairing promptIds from prompt.steered', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 3, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 3, step: 1 })); feed(ev({ type: 'turn.step.completed', turnId: 3, step: 1 })); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [ { type: 'text', text: 'steered in' }, { type: 'video_url', videoUrl: { url: 'kimi-file://f_vid2', name: 'queued.mp4' } }, ], steeredAt: '2026-01-01T00:00:02.000Z', }), ); feed( ev({ type: 'turn.steer', input: [ { type: 'text', text: 'steered in' }, { type: 'video_url', videoUrl: { url: 'kimi-file://f_vid2', name: 'queued.mp4' } }, ], origin: { kind: 'user' }, }), ); expect(turnOps('t3', tx.getItems()).steps).toHaveLength(1); feed(ev({ type: 'turn.step.started', turnId: 3, step: 2 })); const turn = turnOps('t3', tx.getItems()); expect(turn.steps).toHaveLength(2); const frame = turn.steps[1]?.frames[0]; expect(frame).toMatchObject({ kind: 'text', role: 'user', text: 'steered in', promptIds: ['p2'], origin: { kind: 'user' }, }); expect(frame?.kind === 'text' ? frame.attachmentIds : undefined).toHaveLength(1); const attachmentId = frame?.kind === 'text' ? frame.attachmentIds?.[0] : undefined; expect(attachmentId === undefined ? undefined : tx.getAttachment(attachmentId)).toMatchObject({ mediaType: 'video/*', name: 'queued.mp4', source: { kind: 'session_media', fileId: 'f_vid2' }, }); }); it('projects turn.steer into the running step immediately, with daemon media as attachments', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { const mapped = projector.map(event); ops.push(...mapped); tx.apply(mapped); }; feed(ev({ type: 'turn.started', turnId: 4, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 4, step: 1 })); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2', 'p3'], content: [ { type: 'text', text: 'look at this' }, { type: 'image_url', imageUrl: { url: 'kimi-file://f_img9?path=%2Fabs%2Fsession%2Fmedia%2Ff_img9.png', name: 'architecture.png', }, }, ], steeredAt: '2026-01-01T00:00:02.000Z', }), ); feed( ev({ type: 'turn.steer', input: [ { type: 'text', text: 'private instructions' }, { type: 'text', text: 'private instructions' }, { type: 'text', text: 'look at this' }, { type: 'image_url', imageUrl: { url: 'kimi-file://f_img9?path=%2Fabs%2Fsession%2Fmedia%2Ff_img9.png', name: 'architecture.png', }, }, ], origin: { kind: 'user', skillActivations: [ { activationId: 'a1', skillName: 'deploy', skillPath: '/private/deploy/SKILL.md' }, { activationId: 'a2', skillName: 'review', skillArgs: 'strict', skillPath: '/private/review/SKILL.md' }, ], attachments: [{ name: 'secret.txt', mediaType: 'text/plain', size: 12, path: '/private/secret.txt', }], }, }), ); const attachmentOp = ops.find((op) => op.op === 'attachment.upsert'); expect(attachmentOp).toMatchObject({ attachment: { mediaType: 'image/*', name: 'architecture.png', source: { kind: 'session_media', fileId: 'f_img9' }, }, }); const frame = turnOps('t4', tx.getItems()).steps[0]?.frames[0]; expect(frame).toMatchObject({ kind: 'text', role: 'user', text: 'look at this', promptIds: ['p2', 'p3'], origin: { kind: 'user', skillActivations: [ { skillName: 'deploy' }, { skillName: 'review', skillArgs: 'strict' }, ], }, }); expect(JSON.stringify(frame)).not.toContain('/private/'); const attachmentIds = frame?.kind === 'text' ? frame.attachmentIds : undefined; expect(attachmentIds).toHaveLength(2); expect(attachmentIds?.[0]).toBe( attachmentOp?.op === 'attachment.upsert' ? attachmentOp.attachment.attachmentId : undefined, ); expect(tx.getAttachment(attachmentIds![1]!)).toEqual({ attachmentId: attachmentIds![1], mediaType: 'text/plain', name: 'secret.txt', size: 12, }); expect(JSON.stringify([...tx.getAttachments().values()])).not.toContain('/private/'); }); it('records a user slash skill activation steered into a running turn without taking a queued prompt id', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const origin = { kind: 'skill_activation' as const, activationId: 'act-1', trigger: 'user-slash' as const, skillName: 'example-skill', skillArgs: 'args' }; feed(ev({ type: 'turn.started', turnId: 5, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 5, step: 1 })); feed(ev({ type: 'prompt.steered', activePromptId: 'active', promptIds: ['queued'], content: [{ type: 'text', text: 'queued input' }], steeredAt: '2026-01-01T00:00:02.000Z' })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'User activated the skill' }], origin })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'queued input' }], origin: { kind: 'user' } })); const frames = turnOps('t5', tx.getItems()).steps[0]!.frames; expect(frames[0]).toMatchObject({ role: 'user', text: 'User activated the skill', origin: { kind: 'skill_activation', trigger: 'user-slash', skillName: 'example-skill' } }); expect((frames[0] as { promptIds?: readonly string[] }).promptIds).toBeUndefined(); expect(frames[1]).toMatchObject({ text: 'queued input', promptIds: ['queued'] }); }); it('projects origin file attachments on steered frames for user prompts and slash skills', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const file = { name: 'notes.pdf', mediaType: 'application/pdf', size: 42, path: '/tmp/notes.pdf' }; feed(ev({ type: 'turn.started', turnId: 7, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 7, step: 1 })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'User activated the skill' }], origin: { kind: 'skill_activation', activationId: 'act-3', trigger: 'user-slash', skillName: 'example-skill', attachments: [file] } })); feed(ev({ type: 'prompt.steered', activePromptId: 'active', promptIds: ['queued'], content: [{ type: 'text', text: 'see file' }], steeredAt: '2026-01-01T00:00:02.000Z' })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'see file' }], origin: { kind: 'user', attachments: [file] } })); const frames = turnOps('t7', tx.getItems()).steps[0]!.frames as { attachmentIds?: readonly string[] }[]; expect(frames).toHaveLength(2); for (const frame of frames) { expect(frame.attachmentIds).toHaveLength(1); expect(tx.getAttachment(frame.attachmentIds![0]!)).toMatchObject({ name: 'notes.pdf', mediaType: 'application/pdf', size: 42 }); } }); it('still ignores a model-triggered skill activation steer', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 6, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 6, step: 1 })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'model tool activation' }], origin: { kind: 'skill_activation', activationId: 'act-2', trigger: 'model-tool', skillName: 'example-skill' } })); expect(turnOps('t6', tx.getItems()).steps[0]!.frames).toHaveLength(0); }); it('projects a live user turn payload without server-local paths', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const clientMetadata = [{ display_text: 'Visible prompt' }]; tx.apply(projector.map(ev({ type: 'turn.started', turnId: 9, prompt: 'visible prompt', origin: { kind: 'user', clientMetadata, skillActivations: [{ activationId: 'a1', skillName: 'deploy', skillArgs: 'now', skillPath: '/private/deploy/SKILL.md' }], attachments: [{ name: 'notes.pdf', mediaType: 'application/pdf', size: 42, path: '/private/notes.pdf' }], }, }))); const turn = turnOps('t9', tx.getItems()); expect(turn.origin).toEqual({ kind: 'user', payload: { kind: 'user', clientMetadata, skillActivations: [{ skillName: 'deploy', skillArgs: 'now' }] } }); expect(JSON.stringify(turn)).not.toContain('/private/'); }); it('keeps client metadata on a steered slash skill frame', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); const clientMetadata = [{ display_text: 'Save button' }]; feed(ev({ type: 'turn.started', turnId: 8, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 8, step: 1 })); feed(ev({ type: 'turn.steer', input: [{ type: 'text', text: 'User skill context' }], origin: { kind: 'skill_activation', activationId: 'act-4', trigger: 'user-slash', skillName: 'example-skill', clientMetadata } })); expect(turnOps('t8', tx.getItems()).steps[0]!.frames[0]).toMatchObject({ role: 'user', origin: { kind: 'skill_activation', skillName: 'example-skill', clientMetadata }, }); }); it('ignores turn.steer for non-user origins and for turns that are not running', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 5, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 5, step: 1 })); feed( ev({ type: 'turn.steer', input: [{ type: 'text', text: 'backgrounded output' }], origin: { kind: 'injection', variant: 'shell_command_backgrounded' }, }), ); expect(turnOps('t5', tx.getItems()).steps[0]?.frames).toHaveLength(0); feed(ev({ type: 'turn.step.completed', turnId: 5, step: 1 })); feed(ev({ type: 'turn.ended', turnId: 5, reason: 'completed' })); feed( ev({ type: 'turn.steer', input: [{ type: 'text', text: 'too late' }], origin: { kind: 'user' }, }), ); expect( turnOps('t5', tx.getItems()).steps.flatMap((step) => step.frames), ).toHaveLength(0); }); it('flushes a pending steer into the last step when the turn ends before the next step', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 6, origin: { kind: 'user' }, prompt: 'active' })); feed(ev({ type: 'turn.step.started', turnId: 6, step: 1 })); feed(ev({ type: 'turn.step.completed', turnId: 6, step: 1 })); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [{ type: 'text', text: 'last word' }], steeredAt: '2026-01-01T00:00:02.000Z', }), ); feed( ev({ type: 'turn.steer', input: [{ type: 'text', text: 'last word' }], origin: { kind: 'user' }, }), ); feed(ev({ type: 'turn.ended', turnId: 6, reason: 'cancelled', interruptReason: 'user_cancelled' })); const turn = turnOps('t6', tx.getItems()); const lastStep = turn.steps.at(-1); expect(lastStep?.frames.at(-1)).toMatchObject({ kind: 'text', role: 'user', text: 'last word', promptIds: ['p2'], origin: { kind: 'user' }, }); }); it('flushes a pending steer into a user-only step when the turn ends before its first step', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'turn.started', turnId: 7, origin: { kind: 'user' }, prompt: 'active' })); feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [{ type: 'text', text: 'last word' }], steeredAt: '2026-01-01T00:00:02.000Z', }), ); feed( ev({ type: 'turn.steer', input: [ { type: 'text', text: 'private instructions' }, { type: 'text', text: 'last word' }, ], origin: { kind: 'user', skillActivations: [{ activationId: 'a1', skillName: 'review', skillArgs: 'strict' }], }, }), ); feed(ev({ type: 'turn.ended', turnId: 7, reason: 'cancelled', interruptReason: 'user_cancelled' })); const turn = turnOps('t7', tx.getItems()); expect(turn.steps).toHaveLength(1); expect(turn.steps[0]).toMatchObject({ state: 'interrupted' }); expect(turn.steps[0]?.frames[0]).toMatchObject({ kind: 'text', role: 'user', text: 'last word', promptIds: ['p2'], origin: { kind: 'user', skillActivations: [{ skillName: 'review', skillArgs: 'strict' }], }, }); }); it('buffers turn.steer seen before the projector ever saw turn.started (mid-turn attach)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const ops: TranscriptOperation[] = []; const feed = (event: ProjectorBusEvent): void => { ops.push(...projector.map(event)); }; feed( ev({ type: 'prompt.steered', activePromptId: 'p1', promptIds: ['p2'], content: [{ type: 'text', text: 'steered mid-attach' }], steeredAt: '2026-01-01T00:00:02.000Z', }), ); feed( ev({ type: 'turn.steer', input: [{ type: 'text', text: 'steered mid-attach' }], origin: { kind: 'user' }, }), ); expect(ops).toHaveLength(2); expect(ops.every((op) => op.op === 'prompt.upsert')).toBe(true); feed(ev({ type: 'turn.step.started', turnId: 3, step: 2 })); const frameOp = ops.find((op) => op.op === 'frame.upsert'); expect(frameOp).toMatchObject({ turnId: 't3', stepId: 't3.2', frame: { kind: 'text', role: 'user', text: 'steered mid-attach', promptIds: ['p2'], origin: { kind: 'user' }, }, }); }); it('readColdSnapshot answers empty for path-hostile agent ids without touching disk', async () => { const service = new TranscriptService({ homeDir: '/nonexistent-home', core: { accessor: { get: (token: unknown) => { if (token === ISessionManager) { return { get: () => undefined, list: () => [] }; } if (token === IWorkspaceInstanceManager) { return { list: () => [], onDidChange: () => ({ dispose: () => undefined }) }; } if (token === ISessionIndex) return { get: async () => ({ workspaceId: 'ws' }) }; return undefined; }, }, } as unknown as Scope, }); for (const hostile of ['../../main', '..', 'a/b', 'a\\b']) { const snapshot = await service.readColdSnapshot('s1', hostile); expect(snapshot?.items).toEqual([]); } }); it('readColdSnapshot folds task/todo/goal/plan/interaction records into the cold snapshot', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-facts-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' }, }, time: 1000, }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'running' }], toolCalls: [{ type: 'function', id: 'call_1', name: 'Bash', arguments: '{"command":"ls"}' }], }, time: 2000, }, { type: 'tools.update_store', key: 'todo', value: [{ title: 'write tests', status: 'in_progress' }], time: 3000, }, { type: 'goal.create', goalId: 'g1', objective: 'fix the bug', time: 4000 }, { type: 'plan_mode.enter', id: 'plan-1', time: 5000 }, { type: 'task.started', info: { taskId: 'task_1', kind: 'process', description: 'pnpm test', status: 'running', startedAt: 6000, endedAt: null, }, time: 6000, }, { type: 'task.terminated', info: { taskId: 'task_1', kind: 'process', description: 'pnpm test', status: 'completed', startedAt: 6000, endedAt: 9000, }, outputTail: '42 passed', time: 9000, }, { type: 'interaction.request', id: 'apr-1', kind: 'approval', toolCallId: 'call_1', request: { toolName: 'Bash' }, time: 7000, }, { type: 'interaction.resolved', id: 'apr-1', response: { decision: 'approved' }, time: 8000, }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const service = new TranscriptService({ homeDir: home, core: { accessor: { get: (token: unknown) => { if (token === ISessionManager) return { get: () => undefined, list: () => [] }; if (token === IWorkspaceInstanceManager) { return { list: () => [], onDidChange: () => ({ dispose: () => undefined }) }; } if (token === ISessionIndex) return { get: async () => ({ workspaceId: 'ws' }) }; return undefined; }, }, } as unknown as Scope, }); const snapshot = await service.readColdSnapshot('s1', 'main'); expect(snapshot).toBeDefined(); expect(snapshot!.tasks).toEqual([ { taskId: 'task_1', kind: 'shell', state: 'completed', detached: true, description: 'pnpm test', agentId: undefined, outputTail: '42 passed', startedAt: new Date(6000).toISOString(), endedAt: new Date(9000).toISOString(), }, ]); expect(snapshot!.todos).toEqual([ { todoId: 'todo', items: [{ title: 'write tests', status: 'in_progress' }], updatedAt: new Date(3000).toISOString(), }, ]); expect(snapshot!.meta.goal).toMatchObject({ objective: 'fix the bug', status: 'active' }); expect(snapshot!.meta.modes).toEqual({ plan: {} }); expect(snapshot!.interactions).toEqual([ { interactionId: 'apr-1', interactionKind: 'approval', toolCallId: 'call_1', state: 'approved', request: { toolName: 'Bash' }, response: { decision: 'approved' }, }, ]); const standalone = snapshot!.items.filter((item) => item.kind !== 'turn'); expect(standalone).toEqual([ expect.objectContaining({ kind: 'marker', marker: 'goal', markerId: 'm1' }), expect.objectContaining({ kind: 'marker', marker: 'plan.enter', markerId: 'm2' }), expect.objectContaining({ kind: 'taskref', refId: 'ref-task_1', taskId: 'task_1' }), ]); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot projects question requests onto the wire shape without rewriting the log', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-question-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records = [ { type: 'interaction.request', id: 'q-cold', kind: 'question', toolCallId: 'call_q', request: { toolCallId: 'call_q', questions: [{ question: 'Pick', options: [{ label: 'A' }, { label: 'B' }] }], }, time: 7000, }, { type: 'interaction.request', id: 'q-inner', kind: 'question', request: { toolCallId: 'call_inner', questions: [{ question: 'Inner', options: [{ label: 'X' }] }], }, time: 8500, }, { type: 'interaction.request', id: 'q-bad', kind: 'question', request: { toolName: 'nope' }, time: 8000, }, { type: 'interaction.request', id: 'apr-1', kind: 'approval', toolCallId: 'call_1', request: { toolName: 'Bash' }, time: 9000, }, ]; const wireFile = join(wireDir, 'wire.jsonl'); const content = `${records.map((r) => JSON.stringify(r)).join('\n')}\n`; await writeFile(wireFile, content); const service = coldTranscriptService(home); const snapshot = await service.readColdSnapshot('s1', 'main'); const byId = new Map(snapshot!.interactions.map((i) => [i.interactionId, i])); expect(byId.get('q-cold')).toMatchObject({ interactionKind: 'question', toolCallId: 'call_q', state: 'cancelled', }); expect(byId.get('q-cold')?.request).toEqual({ question_id: 'q-cold', session_id: 's1', questions: [ { id: 'q_0', question: 'Pick', options: [ { id: 'opt_0_0', label: 'A' }, { id: 'opt_0_1', label: 'B' }, ], allow_other: true, }, ], created_at: new Date(7000).toISOString(), tool_call_id: 'call_q', }); expect(byId.get('q-bad')?.request).toEqual({ toolName: 'nope' }); expect(byId.get('q-inner')).toMatchObject({ toolCallId: 'call_inner' }); expect(byId.get('q-inner')?.request).toMatchObject({ question_id: 'q-inner', tool_call_id: 'call_inner', }); expect(byId.get('apr-1')?.request).toEqual({ toolName: 'Bash' }); await expect(readFile(wireFile, 'utf-8')).resolves.toBe(content); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot derives meta.activity from the final turn state when no live session exists', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-activity-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const write = async (records: unknown[]): Promise => writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const user = { type: 'context.append_message', message: { id: 'prompt-1', role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }; const assistant = { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, time: 2000 }; const boundary = { type: 'turn.prompt', input: [{ type: 'text', text: 'hi' }], origin: { kind: 'user' }, promptId: 'prompt-1', time: 500 }; await write([boundary, user, assistant, { type: 'turn.ended', turnId: 0, reason: 'completed', time: 3000 }]); const ended = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); expect(ended!.meta.activity).toBe('idle'); expect(ended!.items[0]).toMatchObject({ triggerPromptId: 'prompt-1' }); await write([boundary, user, assistant]); const dangling = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); expect(dangling!.meta.activity).toBe('idle'); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot opens a task-origin turn only when the wire has the turn.prompt boundary', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-taskturn-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const notification = '\nTitle: Background agent completed\nSeverity: info\ninspect done.\n'; const taskOrigin = { kind: 'task', taskId: 'task_9', status: 'completed', notificationId: 'n1' }; const opening = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, time: 2000 }, ]; const boundary = { type: 'turn.prompt', input: [{ type: 'text', text: notification }], origin: taskOrigin, time: 3000 }; const delivered = { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: notification }], toolCalls: [], origin: taskOrigin }, time: 4000 }; const reply = { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'reporting back' }], toolCalls: [] }, time: 5000 }; const write = async (records: unknown[]): Promise => writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); await write([...opening, boundary, delivered, reply]); const withBoundary = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const originKinds = withBoundary!.items .filter((item) => item.kind === 'turn') .map((item) => (item.kind === 'turn' ? item.origin.kind : '')); expect(originKinds).toEqual(['user', 'task']); await write([...opening, delivered, reply]); const withoutBoundary = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const turns = withoutBoundary!.items.filter((item) => item.kind === 'turn'); expect(turns).toHaveLength(1); const turn = turns[0]; if (turn?.kind !== 'turn') throw new Error('expected turn'); expect( turn.steps.flatMap((step) => step.frames).some((f) => f.kind === 'text' && f.role === 'user'), ).toBe(true); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot folds a steered user message into its turn instead of opening a new one', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'active' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, time: 2000 }, { type: 'turn.steer', input: [{ type: 'text', text: 'steered in' }], origin: { kind: 'user' }, time: 3000 }, { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'steered in' }], toolCalls: [], origin: { kind: 'user' } }, time: 3001 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, time: 4000 }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const snapshot = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const turns = snapshot!.items.filter((item) => item.kind === 'turn'); expect(turns).toHaveLength(1); const turn = turns[0]; if (turn?.kind !== 'turn') throw new Error('expected turn'); expect(turn.steps).toHaveLength(2); expect(turn.steps[1]?.frames[0]).toMatchObject({ kind: 'text', role: 'user', text: 'steered in', origin: { kind: 'user' }, }); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot preserves safe bundled skill provenance before the first step', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-bundled-steer-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const origin = { kind: 'user', skillActivations: [ { activationId: 'a1', skillName: 'deploy', skillPath: '/private/deploy/SKILL.md' }, { activationId: 'a2', skillName: 'review', skillArgs: 'strict', skillPath: '/private/review/SKILL.md' }, ], attachments: [{ name: 'secret.txt', mediaType: 'text/plain', size: 12, path: '/private/secret.txt', }], }; const content = [ { type: 'text', text: 'private instructions' }, { type: 'text', text: 'private instructions' }, { type: 'text', text: 'steered in' }, ]; const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'active' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, { type: 'turn.steer', input: content, origin, time: 3000 }, { type: 'context.append_message', message: { role: 'user', content, toolCalls: [], origin }, time: 3001 }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const snapshot = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const turn = snapshot?.items.find((item) => item.kind === 'turn'); if (turn?.kind !== 'turn') throw new Error('expected turn'); const frame = turn.steps.flatMap((step) => step.frames).find( (candidate) => candidate.kind === 'text' && candidate.role === 'user', ); expect(frame).toMatchObject({ kind: 'text', role: 'user', text: 'steered in', origin: { kind: 'user', skillActivations: [ { skillName: 'deploy' }, { skillName: 'review', skillArgs: 'strict' }, ], }, }); expect(JSON.stringify(frame)).not.toContain('/private/'); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot keeps skill-activation steers as skill markers instead of user frames', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-skillsteer-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const skillOrigin = { kind: 'skill_activation', activationId: 'a1', skillName: 'write-tui', trigger: 'model-tool', skillSource: 'project', }; const nestedOrigin = { ...skillOrigin, activationId: 'a2', skillName: 'design', trigger: 'nested-skill' }; const skillText = 'Skill tool loaded instructions for this request. Follow them.'; const nestedText = 'Nested skill instructions.'; const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'active' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, time: 2000 }, { type: 'turn.steer', input: [{ type: 'text', text: skillText }], origin: skillOrigin, time: 3000 }, { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: skillText }], toolCalls: [], origin: skillOrigin }, time: 3001 }, { type: 'turn.steer', input: [{ type: 'text', text: nestedText }], origin: nestedOrigin, time: 3002 }, { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: nestedText }], toolCalls: [], origin: nestedOrigin }, time: 3003 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, time: 4000 }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const snapshot = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const markers = snapshot!.items.filter((item) => item.kind === 'marker'); expect(markers).toHaveLength(2); expect(markers.every((item) => item.kind === 'marker' && item.marker === 'skill')).toBe(true); const turns = snapshot!.items.filter((item) => item.kind === 'turn'); expect(turns).toHaveLength(1); const turn = turns[0]; if (turn?.kind !== 'turn') throw new Error('expected turn'); expect(turn.prompt).toBe('active'); const userFrames = turn.steps .flatMap((step) => step.frames) .filter((frame) => frame.kind === 'text' && frame.role === 'user'); expect(userFrames).toHaveLength(0); } finally { await rm(home, { recursive: true, force: true }); } }); it('readColdSnapshot drops undone task-turn boundaries so a redelivered notification folds', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-undoboundary-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const notification = '\nTitle: Background agent completed\nSeverity: info\ninspect done.\n'; const taskOrigin = { kind: 'task', taskId: 'task_9', status: 'completed', notificationId: 'n1' }; const opening = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, time: 2000 }, ]; const boundary = { type: 'turn.prompt', input: [{ type: 'text', text: notification }], origin: taskOrigin, time: 3000 }; const delivered = { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: notification }], toolCalls: [], origin: taskOrigin }, time: 4000 }; const reply = { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'reporting back' }], toolCalls: [] }, time: 5000 }; const again = { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'again' }], toolCalls: [], origin: { kind: 'user' } }, time: 7000 }; const answer2 = { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'answer2' }], toolCalls: [] }, time: 8000 }; const redelivered = { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: notification }], toolCalls: [], origin: taskOrigin }, time: 9000 }; const write = async (records: unknown[]): Promise => writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); await write([...opening, boundary, delivered, reply, again, answer2, redelivered]); const withBoundary = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const originKinds = withBoundary!.items .filter((item) => item.kind === 'turn') .map((item) => (item.kind === 'turn' ? item.origin.kind : '')); expect(originKinds).toEqual(['user', 'task', 'user', 'task']); await write([...opening, boundary, delivered, reply, { type: 'context.undo', count: 1, time: 6000 }, again, answer2, redelivered]); const withUndo = await coldTranscriptService(home).readColdSnapshot('s1', 'main'); const undoTurns = withUndo!.items.filter((item) => item.kind === 'turn'); expect(undoTurns).toHaveLength(1); const undoTurn = undoTurns[0]; if (undoTurn?.kind !== 'turn') throw new Error('expected turn'); expect(undoTurn.origin.kind).toBe('user'); expect( undoTurn.steps .flatMap((step) => step.frames) .some((f) => f.kind === 'text' && f.text.includes('inspect done')), ).toBe(false); } finally { await rm(home, { recursive: true, force: true }); } }); it('gates the cold tower mode badge behind the tower experiment flag', async () => { const home = await mkdtemp(join(tmpdir(), 'transcript-cold-tower-')); try { const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' }, }, time: 1000, }, { type: 'tower_mode.enter', time: 2000 }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); const serviceWith = ( flagOn: boolean, opts: { cwd?: string; liveSessionIds?: string[] } = {}, ) => new TranscriptService({ homeDir: home, core: { accessor: { get: (token: unknown) => { if (token === ISessionManager) { return { get: (id: string) => (opts.liveSessionIds?.includes(id) ? {} : undefined), list: () => [], }; } if (token === IWorkspaceInstanceManager) { return { list: () => [], onDidChange: () => ({ dispose: () => undefined }) }; } if (token === ISessionIndex) { return { get: async () => ({ workspaceId: 'ws', cwd: opts.cwd }) }; } if (token === IFlagService) { return { enabled: (id: string) => flagOn && id === TOWER_FLAG_ID }; } return undefined; }, }, } as unknown as Scope, }); const withFlag = await serviceWith(true).readColdSnapshot('s1', 'main'); expect(withFlag!.meta.modes).toEqual({ tower: {} }); const withoutFlag = await serviceWith(false).readColdSnapshot('s1', 'main'); expect(withoutFlag!.meta.modes).toBeUndefined(); _setTowerFeatureAssembledForTests(false); try { const notAssembled = await serviceWith(true).readColdSnapshot('s1', 'main'); expect(notAssembled!.meta.modes).toBeUndefined(); } finally { _setTowerFeatureAssembledForTests(true); } const repo = await mkdtemp(join(tmpdir(), 'tower-cold-owner-')); try { await execFileAsync('git', ['init', '-b', 'main'], { cwd: repo }); await execFileAsync('git', ['config', 'user.email', 'tower-test@example.com'], { cwd: repo }); await execFileAsync('git', ['config', 'user.name', 'Tower Test'], { cwd: repo }); await writeFile(join(repo, 'README.md'), '# fixture\n'); await execFileAsync('git', ['add', 'README.md'], { cwd: repo }); await execFileAsync('git', ['commit', '-m', 'initial'], { cwd: repo }); await new TowerStore(repo).init('session-b'); const adoptedByLive = await serviceWith(true, { cwd: repo, liveSessionIds: ['session-b'], }).readColdSnapshot('s1', 'main'); expect(adoptedByLive!.meta.modes).toBeUndefined(); const adoptedByDead = await serviceWith(true, { cwd: repo }).readColdSnapshot('s1', 'main'); expect(adoptedByDead!.meta.modes).toEqual({ tower: {} }); } finally { await rm(repo, { recursive: true, force: true }); } const childDir = join(home, 'sessions', 'ws', 's1', 'agents', 'worker-1'); await mkdir(childDir, { recursive: true }); await writeFile( join(childDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`, ); const child = await serviceWith(true).readColdSnapshot('s1', 'worker-1'); expect(child!.meta.modes).toBeUndefined(); } finally { await rm(home, { recursive: true, force: true }); } }); it('folds blocked turn endings into failed (engine wire contract)', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); tx.apply(projector.map(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' } }))); tx.apply(projector.map(ev({ type: 'turn.ended', turnId: 0, reason: 'blocked' }))); expect(turnOps('t0', tx.getItems()).state).toBe('failed'); }); it('tracks the prompt queue from submitted/queued through terminal', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'prompt.submitted', promptId: 'p1', userMessageId: 'p1', status: 'running', content: [{ type: 'text', text: 'now' }], createdAt: '2026-08-20T00:00:00.000Z' })); expect(tx.getPrompt('p1')).toMatchObject({ status: 'running' }); feed(ev({ type: 'prompt.queued', promptId: 'p2', content: [{ type: 'text', text: 'later' }], queueLength: 1 })); expect(tx.getPrompt('p2')).toMatchObject({ status: 'queued' }); feed(ev({ type: 'prompt.completed', promptId: 'p2', finishedAt: '2026-08-20T00:00:01.000Z', reason: 'completed' })); expect(tx.getPrompt('p2')).toMatchObject({ status: 'completed' }); feed(ev({ type: 'prompt.aborted', promptId: 'p1', abortedAt: '2026-08-20T00:00:02.000Z' })); expect(tx.getPrompt('p1')).toMatchObject({ status: 'aborted' }); }); it('mirrors turn liveness into meta.activity', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); expect(tx.getMeta().activity).toBeUndefined(); feed(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); expect(tx.getMeta().activity).toBe('turn'); feed(ev({ type: 'turn.ended', turnId: 1, reason: 'completed' })); expect(tx.getMeta().activity).toBe('idle'); feed(ev({ type: 'turn.started', turnId: 2, origin: { kind: 'user' } })); expect(tx.getMeta().activity).toBe('turn'); feed(ev({ type: 'turn.ended', turnId: 2, reason: 'failed' })); expect(tx.getMeta().activity).toBe('idle'); }); it('maps cron / task origins onto the turn header', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed( ev({ type: 'turn.started', turnId: 1, origin: { kind: 'cron_job', jobId: 'job-9', cron: '* * * * *' }, }), ); feed( ev({ type: 'turn.started', turnId: 2, origin: { kind: 'task', taskId: 'bash-1', status: 'completed', notificationId: 'n1' }, }), ); expect(turnOps('t1', tx.getItems()).origin).toEqual({ kind: 'cron', taskId: 'job-9', payload: { kind: 'cron_job', jobId: 'job-9', cron: '* * * * *' }, }); expect(turnOps('t2', tx.getItems()).origin).toEqual({ kind: 'task', taskId: 'bash-1', payload: { kind: 'task', taskId: 'bash-1', status: 'completed', notificationId: 'n1' }, }); }); it('treats subagent.started/failed/suspended within the running→failed vocabulary', () => { const projector = new AgentTranscriptProjector('main', TEST_SESSION_ID); const tx = new AgentTranscript('main'); const feed = (event: ProjectorBusEvent): void => void tx.apply(projector.map(event)); feed(ev({ type: 'subagent.started', subagentId: 'agent-1' })); expect(tx.getTask('agent-1')).toMatchObject({ kind: 'subagent', state: 'running' }); feed(ev({ type: 'subagent.suspended', subagentId: 'agent-1', reason: 'approval' })); expect(tx.getTask('agent-1')).toMatchObject({ state: 'running', stateReason: 'approval' }); feed(ev({ type: 'subagent.failed', subagentId: 'agent-1', error: 'boom' })); expect(tx.getTask('agent-1')).toMatchObject({ state: 'failed', error: 'boom' }); feed( ev({ type: 'subagent.completed', subagentId: 'agent-2', resultSummary: 'found 3 files', usage: { inputOther: 10, output: 5, inputCacheRead: 2, inputCacheCreation: 1 }, }), ); expect(tx.getTask('agent-2')).toMatchObject({ state: 'completed', resultSummary: 'found 3 files', usage: { inputOther: 10, output: 5, inputCacheRead: 2, inputCacheCreation: 1 }, }); }); }); describe('AgentTranscript transcript task vocabulary', () => { it('documents the task states used by the projector', () => { const states: Array = [ 'running', 'completed', 'failed', 'timed_out', 'killed', 'lost', ]; expect(states).toHaveLength(6); }); }); describe('bindSessionTranscript', () => { class FakeBus { private readonly handlers = new Set<(event: Event2) => void>(); subscribe(cb: (event: Event2) => void): { dispose: () => void } { this.handlers.add(cb); return { dispose: () => this.handlers.delete(cb) }; } emit(event: Event2): void { for (const cb of this.handlers) cb(event); } } interface FakeAgentHandle { readonly id: string; readonly context: AgentContext; readonly bus: FakeBus; contextMessages: ContextMessage[]; readonly undoParticipants: Map; readonly accessor: { get: (token: unknown) => unknown }; } class FakeAgents { private readonly handles = new Map(); private readonly createHandlers = new Set<(context: AgentContext) => void>(); private readonly closeHandlers = new Set<(context: AgentContext) => void>(); list(): AgentContext[] { return [...this.handles.values()].map((handle) => handle.context); } get(agentId: string): AgentContext | undefined { return this.handles.get(agentId)?.context; } handleOf(agentId: string): FakeAgentHandle | undefined { return this.handles.get(agentId); } byId(id: string): FakeAgentHandle | undefined { return this.handles.get(id); } onDidCreate(cb: (context: AgentContext) => void): { dispose: () => void } { this.createHandlers.add(cb); return { dispose: () => this.createHandlers.delete(cb) }; } onDidClose(cb: (context: AgentContext) => void): { dispose: () => void } { this.closeHandlers.add(cb); return { dispose: () => this.closeHandlers.delete(cb) }; } add(id: string, opts?: { loopStatus?: { state?: 'idle' | 'running'; activeTurnId?: number }; tasks?: readonly unknown[]; activePromptId?: string }): FakeAgentHandle { const bus = this.handles.get(id)?.bus ?? new FakeBus(); const undoParticipants = new Map(); const scope = makeAgentScopeContext({ agentId: id, agentScope: `agents/${id}`, generation: 1, }); let activity: AgentActivitySnapshot = {}; bus.subscribe((event) => { if (event.type === 'turn.started') { activity = { turn: { turnId: (event as { turnId?: number }).turnId ?? 0, phase: 'running', step: 1, ending: false, activeToolCalls: [], since: 0, }, }; } else if (event.type === 'turn.ended') { activity = {}; } }); const handle: FakeAgentHandle = { id, context: scope.agentContext, bus, contextMessages: [], undoParticipants, accessor: { get: (token: unknown) => { if (token === IAgentScopeContext) return scope; if (token === IEventBus) return bus; if (token === IAgentContextMemoryService) return { get: () => handle.contextMessages }; if (token === IAgentConversationUndoParticipantRegistry) { return { register: (participant: AgentConversationUndoParticipant) => { undoParticipants.set(participant.id, participant); return { dispose: () => { undoParticipants.delete(participant.id); } }; }, }; } if (token === IAgentLoopService) { const active = opts?.activePromptId === undefined ? undefined : { id: opts.activePromptId, userMessageId: opts.activePromptId, createdAt: '2026-01-01T00:00:00.000Z', state: 'running' as const, message: { role: 'user' as const, content: [{ type: 'text' as const, text: 'hi' }], toolCalls: [], origin: { kind: 'user' as const }, }, launched: Promise.resolve(undefined), completion: new Promise(() => {}), }; return { snapshot: () => ({ state: opts?.loopStatus?.state ?? (activity.turn === undefined ? 'idle' : 'running'), activeTurnId: opts?.loopStatus?.activeTurnId ?? activity.turn?.turnId, activePromptId: opts?.activePromptId, queue: [], notificationCount: 0, paused: false, hasPendingRequests: activity.turn !== undefined, turn: activity.turn, activeTraceId: undefined, }), promptHandle: (id: string) => (active?.id === id ? active : undefined), }; } if (token === IAgentTaskService) { return { list: () => opts?.tasks ?? [] }; } return undefined; }, }, }; this.handles.set(id, handle); for (const cb of this.createHandlers) cb(handle.context); return handle; } remove(id: string): void { const removed = this.handles.get(id); this.handles.delete(id); if (removed !== undefined) { for (const cb of this.closeHandlers) cb(removed.context); } } } function fakeSession(manager: FakeAgents): ISessionScopeHandle { return { id: 's1', accessor: { get: (token: unknown) => { if (token === IAgentLifecycleService) return manager; if (token === ISessionMetadata) return { read: async () => ({ agents: {} }) }; return undefined; }, }, } as unknown as ISessionScopeHandle; } afterEach(() => { interactions.purgeSession('s1'); }); it('registers pre-bind pendings without frames and replays an early resolve at seed time', () => { const agents = new FakeAgents(); interactions.enqueue({ id: 'apr-1', kind: 'approval', payload: { toolCallId: 'call_1' }, tags: { agentId: 'main', sessionId: 's1', turnId: 0 }, }); const store = new TranscriptStore('s1'); const ops: TranscriptOperation[] = []; const binding = bindSessionTranscript(store, fakeSession(agents), undefined, (event) => ops.push(...event.ops), ); expect(ops).toHaveLength(0); interactions.respond('apr-1', { decision: 'approved' }); expect(ops).toHaveLength(0); binding.seedPendingInteractions(); const states = ops .filter((op): op is InteractionUpsertOp => op.op === 'interaction.upsert') .map((op) => op.interaction.state); expect(states).toEqual(['pending', 'approved']); binding.dispose(); }); it('keeps the materialized transcript and roster entry when an agent is disposed', () => { const agents = new FakeAgents(); const store = new TranscriptStore('s1'); const binding = bindSessionTranscript( store, fakeSession(agents), ); const sub = agents.add('sub-1'); agents.add('main'); sub.bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'scan' })); sub.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(store.getAgent('sub-1')?.getItems()).toHaveLength(1); agents.remove('sub-1'); expect(store.getAgent('sub-1')?.getItems()).toHaveLength(1); const descriptor = store.agents().find((a) => a.agentId === 'sub-1'); expect(descriptor).toBeDefined(); expect(typeof descriptor?.disposedAt).toBe('string'); expect(store.agents().find((a) => a.agentId === 'main')?.disposedAt).toBeUndefined(); binding.dispose(); }); it('seeds pre-attach Agent task mappings so a late-bound projector folds the lifecycle', () => { const agents = new FakeAgents(); agents.add('main', { tasks: [ { taskId: 'task-9', kind: 'agent', agentId: 'agent-1', status: 'running', description: 'Inspect', detached: false, startedAt: 1_700_000_000_000, }, ], }); const store = new TranscriptStore('s1'); const binding = bindSessionTranscript( store, fakeSession(agents), ); expect(store.getAgent('main')?.getTask('task-9')).toMatchObject({ kind: 'subagent', state: 'running', detached: false, description: 'Inspect', agentId: 'agent-1', }); agents.byId('main')!.bus.emit(ev({ type: 'subagent.completed', subagentId: 'agent-1', resultSummary: 'done' })); expect(store.getAgent('main')?.getTask('task-9')).toMatchObject({ state: 'completed', resultSummary: 'done', detached: false, }); expect(store.getAgent('main')?.getTask('agent-1')).toBeUndefined(); binding.dispose(); }); const SHOT_PNG_UPLOAD = { type: 'file', file_id: 'file_1', media_type: 'image/png', name: 'shot.png', }; async function seedWireHome(attachment?: Record): Promise { const home = await mkdtemp(join(tmpdir(), 'transcript-overlay-')); const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records: Record[] = [ { type: 'context.append_message', message: { role: 'user', content: attachment === undefined ? [{ type: 'text', text: 'hi' }] : [{ type: 'text', text: 'what is this?' }, attachment], toolCalls: [], origin: { kind: 'user' }, }, time: new Date().toISOString(), }, ]; if (attachment !== undefined) { records.push({ type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'a screenshot' }], toolCalls: [], }, time: new Date().toISOString(), }); } await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); return home; } function fakeCoreWithAgents(agents: FakeAgents): Scope { const sessionLifecycle = { onDidCloseSession: () => ({ dispose: () => undefined }), onDidArchiveSession: () => ({ dispose: () => undefined }), get: (sid: string) => (sid === 's1' ? fakeSession(agents) : undefined), }; const handler = { id: 'ws', kind: 'program', accessor: { get: (t: unknown) => (t === ISessionLifecycleService ? sessionLifecycle : undefined), }, dispose: () => undefined, }; return { accessor: { get: (token: unknown) => { if (token === ISessionManager) { return { get: sessionLifecycle.get, list: () => [sessionLifecycle.get('s1')] }; } if (token === IWorkspaceInstanceManager) { return { list: () => [{ program: { accessor: handler.accessor } }], onDidChange: () => ({ dispose: () => undefined }), }; } if (token === ISessionIndex) return { get: async () => ({ workspaceId: 'ws' }) }; return undefined; }, }, } as unknown as Scope; } it('stops projecting for an agent once it is disposed', () => { const agents = new FakeAgents(); const store = new TranscriptStore('s1'); const binding = bindSessionTranscript( store, fakeSession(agents), ); const sub = agents.add('sub-1'); sub.bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'scan' })); expect(store.getAgent('sub-1')?.getItems()).toHaveLength(1); agents.remove('sub-1'); sub.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(store.getAgent('sub-1')?.getItems()[0]).toMatchObject({ kind: 'turn', state: 'running' }); binding.dispose(); }); it('heals a kind-mismatched frame instead of skipping it on length', () => { const snapshotTurn: TranscriptTurn = { kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, steps: [ { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'completed', frames: [ { kind: 'thinking', frameId: 't0.1.f1', text: 'hmm' }, { kind: 'text', frameId: 't0.1.f2', role: 'assistant', text: 'Hello world' }, ], }, ], }; const liveTurn: TranscriptTurn = { kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, steps: [ { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'completed', frames: [{ kind: 'text', frameId: 't0.1.f1', role: 'assistant', text: 'world' }], }, ], }; const frames = healTurnOps(snapshotTurn, liveTurn) .filter((op): op is FrameUpsertOp => op.op === 'frame.upsert') .map((op) => op.frame); expect(frames).toContainEqual(expect.objectContaining({ kind: 'thinking', frameId: 't0.1.f1', text: 'hmm' })); expect(frames).toContainEqual(expect.objectContaining({ kind: 'text', frameId: 't0.1.f2', text: 'Hello world' })); }); it('heals missing tool frames and missed results, keeps richer live ones', () => { const makeTurn = (frames: TranscriptTurn['steps'][number]['frames']): TranscriptTurn => ({ kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, steps: [ { kind: 'step', stepId: 't0.1', turnId: 't0', ordinal: 1, state: 'completed', frames }, ], }); const snapshotTurn = makeTurn([ { kind: 'tool', frameId: 't0.1.call_1', toolCallId: 'call_1', name: 'Bash', state: 'done', input: { command: 'ls' }, output: 'a.txt' }, { kind: 'tool', frameId: 't0.1.call_2', toolCallId: 'call_2', name: 'Read', state: 'done', input: {}, output: 'x' }, { kind: 'tool', frameId: 't0.1.call_3', toolCallId: 'call_3', name: 'Bash', state: 'done', input: {}, output: 'y' }, ]); const liveTurn = makeTurn([ { kind: 'tool', frameId: 't0.1.call_1', toolCallId: 'call_1', name: 'Bash', state: 'running', input: { command: 'ls' }, display: { kind: 'command', command: 'ls' } }, { kind: 'tool', frameId: 't0.1.call_2', toolCallId: 'call_2', name: 'Read', state: 'done', input: {}, output: 'live-out' }, ]); const frames = healTurnOps(snapshotTurn, liveTurn) .filter((op): op is FrameUpsertOp => op.op === 'frame.upsert') .map((op) => op.frame); expect(frames).toHaveLength(2); expect(frames).toContainEqual( expect.objectContaining({ frameId: 't0.1.call_1', state: 'done', output: 'a.txt', display: { kind: 'command', command: 'ls' }, }), ); expect(frames).toContainEqual(expect.objectContaining({ frameId: 't0.1.call_3', output: 'y' })); }); it('heal keeps the live attachment ids over the snapshot cold ids', () => { const makeTurn = (attachmentIds: string[] | undefined): TranscriptTurn => ({ kind: 'turn', turnId: 't0', ordinal: 0, state: 'completed', origin: { kind: 'user' }, attachmentIds, steps: [], }); const header = healTurnOps(makeTurn(['att_1']), makeTurn(['t0.att1'])).find( (op) => op.op === 'turn.upsert', ); expect(header).toMatchObject({ turn: { attachmentIds: ['t0.att1'] } }); const fallback = healTurnOps(makeTurn(['att_1']), makeTurn(undefined)).find( (op) => op.op === 'turn.upsert', ); expect(fallback).toMatchObject({ turn: { attachmentIds: ['att_1'] } }); }); it('heal keeps the live trigger prompt id over a cold turn without one', () => { const makeTurn = (triggerPromptId: string | undefined): TranscriptTurn => ({ kind: 'turn', turnId: 't0', triggerPromptId, ordinal: 0, state: 'completed', origin: { kind: 'user' }, steps: [], }); const header = healTurnOps(makeTurn(undefined), makeTurn('prompt-1')).find( (op) => op.op === 'turn.upsert', ); expect(header).toMatchObject({ turn: { triggerPromptId: 'prompt-1' } }); }); it('terminal turn.upsert inherits the backfilled header when the projector missed turn.started', () => { const agents = new FakeAgents(); const store = new TranscriptStore('s1'); const ops: TranscriptOperation[] = []; const binding = bindSessionTranscript( store, fakeSession(agents), undefined, (event) => ops.push(...event.ops), ); const main = agents.add('main'); store.ensureAgent('main').apply([ { op: 'attachment.upsert', attachment: { attachmentId: 'att_1', mediaType: 'image/*', name: 'shot.png', source: { kind: 'file', fileId: 'f_1' }, }, }, { op: 'turn.upsert', turn: { kind: 'turn', turnId: 't0', ordinal: 0, state: 'running', origin: { kind: 'user' }, prompt: 'hi', attachmentIds: ['att_1'], startedAt: '2026-08-04T00:00:00.000Z', }, }, ]); main.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); const terminal = ops.filter((op) => op.op === 'turn.upsert'); expect(terminal).toHaveLength(1); expect(terminal[0]).toMatchObject({ turn: { turnId: 't0', state: 'completed', origin: { kind: 'user' }, prompt: 'hi', attachmentIds: ['att_1'], startedAt: '2026-08-04T00:00:00.000Z', }, }); expect(store.getAgent('main')?.getTurn('t0')).toMatchObject({ state: 'completed', prompt: 'hi', attachmentIds: ['att_1'], }); binding.dispose(); }); it('seeds pending interactions per agent, not before that agent is backfilled', () => { const agents = new FakeAgents(); interactions.enqueue({ id: 'q-main', kind: 'question', payload: { toolCallId: 'call_main' }, tags: { agentId: 'main', sessionId: 's1', turnId: 0 } }); interactions.enqueue({ id: 'q-sub', kind: 'question', payload: { toolCallId: 'call_sub' }, tags: { agentId: 'sub-1', sessionId: 's1', turnId: 0 } }); const store = new TranscriptStore('s1'); const byAgent = new Map(); const binding = bindSessionTranscript(store, fakeSession(agents), undefined, (event) => { byAgent.set(event.agentId, [...(byAgent.get(event.agentId) ?? []), ...event.ops]); }); binding.seedPendingInteractions('main'); expect([...byAgent.keys()]).toEqual(['main']); binding.seedPendingInteractions('sub-1'); expect([...byAgent.keys()].toSorted()).toEqual(['main', 'sub-1']); binding.dispose(); }); it('projects live question entities with the same wire shape as the legacy question event', () => { const agents = new FakeAgents(); const asked = interactions.enqueue({ id: 'q-parity', kind: 'question', payload: { questions: [ { question: 'Pick one', options: [{ label: 'A', description: 'first' }, { label: 'B' }], }, ], }, tags: { agentId: 'main', sessionId: 's1', turnId: 0 }, }); const store = new TranscriptStore('s1'); const binding = bindSessionTranscript(store, fakeSession(agents)); binding.seedPendingInteractions('main'); const entity = store.getAgent('main')?.getInteraction('q-parity'); expect(entity?.state).toBe('pending'); expect(entity?.request).toEqual(toWireQuestion(asked, 's1')); binding.dispose(); }); it('defers pendings created before their owning agent is seeded', () => { const agents = new FakeAgents(); const store = new TranscriptStore('s1'); const byAgent = new Map(); const binding = bindSessionTranscript(store, fakeSession(agents), undefined, (event) => { byAgent.set(event.agentId, [...(byAgent.get(event.agentId) ?? []), ...event.ops]); }); interactions.enqueue({ id: 'q-sub', kind: 'question', payload: { toolCallId: 'call_sub' }, tags: { agentId: 'sub-1', sessionId: 's1', turnId: 0 } }); expect(byAgent.size).toBe(0); binding.seedPendingInteractions('main'); expect(byAgent.size).toBe(0); binding.seedPendingInteractions('sub-1'); expect([...byAgent.keys()]).toEqual(['sub-1']); binding.dispose(); }); it('announces pendings from live-created agents immediately (their projector is complete)', () => { const agents = new FakeAgents(); const store = new TranscriptStore('s1'); const byAgent = new Map(); const binding = bindSessionTranscript(store, fakeSession(agents), undefined, (event) => { byAgent.set(event.agentId, [...(byAgent.get(event.agentId) ?? []), ...event.ops]); }); agents.add('sub-1'); interactions.enqueue({ id: 'q1', kind: 'question', payload: { toolCallId: 'call_q1' }, tags: { agentId: 'sub-1', sessionId: 's1', turnId: 0 } }); expect([...byAgent.keys()]).toEqual(['sub-1']); binding.dispose(); }); async function seedWireHomeWithTool(): Promise { const home = await mkdtemp(join(tmpdir(), 'transcript-backfill-live-')); const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); await mkdir(wireDir, { recursive: true }); const records = [ { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' }, }, time: new Date().toISOString(), }, { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'Hello ' }], toolCalls: [{ type: 'function', id: 'call_1', name: 'Bash', arguments: '{"command":"ls"}' }], }, time: new Date().toISOString(), }, { type: 'context.append_message', message: { role: 'tool', content: [{ type: 'text', text: 'a.txt' }], toolCallId: 'call_1', toolCalls: [], }, time: new Date().toISOString(), }, ]; await writeFile(join(wireDir, 'wire.jsonl'), `${records.map((r) => JSON.stringify(r)).join('\n')}\n`); return home; } async function waitFor(condition: () => boolean, timeoutMs = 2000): Promise { const deadline = Date.now() + timeoutMs; while (!condition()) { if (Date.now() > deadline) throw new Error('waitFor timed out'); await new Promise((resolve) => setTimeout(resolve, 20)); } } it('subscribes the bus for an agent whose projector was seeded before its handle existed', () => { const agents = new FakeAgents(); interactions.enqueue({ id: 'q-sub', kind: 'question', payload: { toolCallId: 'call_sub' }, tags: { agentId: 'sub-1', sessionId: 's1', turnId: 0 } }); const store = new TranscriptStore('s1'); const byAgent = new Map(); const binding = bindSessionTranscript(store, fakeSession(agents), undefined, (event) => { byAgent.set(event.agentId, [...(byAgent.get(event.agentId) ?? []), ...event.ops]); }); binding.seedPendingInteractions('sub-1'); expect(byAgent.get('sub-1')?.map((op) => op.op)).toEqual(['interaction.upsert']); const sub = agents.add('sub-1'); sub.bus.emit(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' } })); expect(byAgent.get('sub-1')!.length).toBeGreaterThan(1); binding.dispose(); }); it('seeds active prompt identity before a late-bound turn ends', () => { const agents = new FakeAgents(); const main = agents.add('main', { loopStatus: { state: 'running', activeTurnId: 0 }, activePromptId: 'prompt-1', }); const store = new TranscriptStore('s1'); const binding = bindSessionTranscript(store, fakeSession(agents)); main.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(store.getAgent('main')?.getTurn('t0')).toMatchObject({ state: 'completed', triggerPromptId: 'prompt-1', }); binding.dispose(); }); it('overlays the in-flight turn as running after a backfill', async () => { const home = await seedWireHome(); try { const agents = new FakeAgents(); agents.add('main', { loopStatus: { state: 'running', activeTurnId: 0 }, activePromptId: 'prompt-1', }); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents), }); const store = service.forSessionLive('s1'); await service.whenReady('s1'); expect(store?.getAgent('main')?.getTurn('t0')).toMatchObject({ state: 'running', triggerPromptId: 'prompt-1', prompt: 'hi', }); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('merges the backfill live-first: live frame fields and longer text survive', async () => { const home = await seedWireHomeWithTool(); try { const agents = new FakeAgents(); agents.add('main', { loopStatus: { state: 'running', activeTurnId: 0 } }); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents), }); const store = service.forSessionLive('s1'); const bus = agents.byId('main')!.bus; bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'hi' })); bus.emit(ev({ type: 'turn.step.started', turnId: 0, step: 1 })); bus.emit(ev({ type: 'assistant.delta', turnId: 0, delta: 'Hello world' })); bus.emit( ev({ type: 'tool.call.started', turnId: 0, toolCallId: 'call_1', name: 'Bash', args: { command: 'ls' }, display: { kind: 'command', command: 'ls' }, }), ); await service.whenReady('s1'); const turn = store?.getAgent('main')?.getTurn('t0'); expect(turn?.state).toBe('running'); const text = turn?.steps[0]?.frames.find((f) => f.kind === 'text'); expect(text).toMatchObject({ text: 'Hello world' }); const tool = turn?.steps[0]?.frames.find((f) => f.kind === 'tool'); expect(tool).toMatchObject({ state: 'done', output: 'a.txt', display: { kind: 'command', command: 'ls' }, }); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('re-asserts running when the backfill rebuilds the live turn completed', async () => { const home = await seedWireHome(); try { const agents = new FakeAgents(); agents.add('main', { loopStatus: { state: 'running', activeTurnId: 0 } }); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents), }); const store = service.forSessionLive('s1'); agents .byId('main')! .bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'live hi' })); await service.whenReady('s1'); expect(store?.getAgent('main')?.getTurn('t0')).toMatchObject({ state: 'running', prompt: 'live hi', }); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('backfill + overlay keep the live attachment ids and drop the cold counterpart entity', async () => { const home = await seedWireHome(SHOT_PNG_UPLOAD); try { const agents = new FakeAgents(); agents.add('main', { loopStatus: { state: 'running', activeTurnId: 0 } }); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents), }); const store = service.forSessionLive('s1'); agents.byId('main')!.bus.emit( ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'live prompt', promptAttachments: [{ kind: 'image', fileId: 'file_1' }], }), ); await service.whenReady('s1'); const agent = store?.getAgent('main'); expect(agent?.getTurn('t0')).toMatchObject({ state: 'running', prompt: 'live prompt', attachmentIds: ['t0.att1'], }); expect(agent?.getAttachment('t0.att1')).toBeDefined(); expect(agent?.getAttachment('att_1')).toBeUndefined(); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it('post-turn heal keeps the live attachment ids and never upserts the cold counterparts', async () => { const home = await seedWireHome(SHOT_PNG_UPLOAD); try { const agents = new FakeAgents(); agents.add('main', { loopStatus: { state: 'idle' } }); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents), }); const store = service.forSessionLive('s1'); const batches: TranscriptOperation[][] = []; service.onSessionOps('s1', (event) => { if (event.agentId === 'main') batches.push([...event.ops]); }); const bus = agents.byId('main')!.bus; bus.emit( ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' }, prompt: 'live prompt', promptAttachments: [{ kind: 'image', fileId: 'file_1' }], }), ); await service.whenReady('s1'); bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); await waitFor(() => batches.some((batch) => batch.some((op) => op.op === 'step.upsert' && op.turnId === 't0')), ); const healBatch = batches.find((batch) => batch.some((op) => op.op === 'step.upsert' && op.turnId === 't0'), )!; expect(healBatch.find((op) => op.op === 'turn.upsert')).toMatchObject({ turn: { attachmentIds: ['t0.att1'] }, }); const attachmentUpserts = batches.flatMap((batch) => batch.filter((op) => op.op === 'attachment.upsert'), ); expect(attachmentUpserts).toEqual([ { op: 'attachment.upsert', attachment: expect.objectContaining({ attachmentId: 't0.att1' }) }, ]); const agent = store?.getAgent('main'); expect(agent?.getTurn('t0')).toMatchObject({ state: 'completed', attachmentIds: ['t0.att1'], }); expect(agent?.getAttachment('t0.att1')).toBeDefined(); expect(agent?.getAttachment('att_1')).toBeUndefined(); service.dropSession('s1'); } finally { await rm(home, { recursive: true, force: true }); } }); it.each([1, 2])('removes undone content even when the history read fails %i times', async (failures) => { const home = await seedWireHome(); const agents = new FakeAgents(); const main = agents.add('main'); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents) }); try { const store = service.forSessionLive('s1')!; await service.whenReady('s1'); expect(store.getAgent('main')!.getItems().length).toBeGreaterThan(0); main.contextMessages = [{ role: 'user', content: [{ type: 'text', text: 'hi' }], toolCalls: [], origin: { kind: 'user' } }]; main.bus.emit(ev({ type: 'turn.started', turnId: 1, origin: { kind: 'user' }, prompt: 'undone prompt' })); const read = vi.spyOn(service, 'readColdSnapshot'); for (let i = 0; i < failures; i++) read.mockRejectedValueOnce(new Error('temporary read failure')); await main.undoParticipants.get('transcript')!.reconcileAfterUndo(); expect(store.getAgent('main')!.getItems().filter((item) => item.kind === 'turn').map((turn) => turn.prompt)).toEqual(['hi']); expect(service.getOpsSince('s1', 'main', 0)?.batches.at(-1)?.ops[0]?.op).toBe('reset'); read.mockRestore(); } finally { service.dropSession('s1'); await rm(home, { recursive: true, force: true }); } }); it('does not restore undone messages from an earlier history read', async () => { const home = await seedWireHome(); const agents = new FakeAgents(); const main = agents.add('main'); const service = new TranscriptService({ homeDir: home, core: fakeCoreWithAgents(agents) }); try { const store = service.forSessionLive('s1')!; await service.whenReady('s1'); const stale = await service.readColdSnapshot('s1', 'main'); let release!: (snapshot: AgentTranscriptSnapshot | undefined) => void; const pending = new Promise((resolve) => { release = resolve; }); const read = vi.spyOn(service, 'readColdSnapshot').mockImplementationOnce(() => pending); main.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); await waitFor(() => read.mock.calls.length === 1); await writeFile(join(home, 'sessions', 'ws', 's1', 'agents', 'main', 'wire.jsonl'), ''); await main.undoParticipants.get('transcript')!.reconcileAfterUndo(); expect(store.getAgent('main')?.getItems()).toEqual([]); release(stale); await new Promise((resolve) => setTimeout(resolve, 0)); expect(store.getAgent('main')?.getItems()).toEqual([]); read.mockRestore(); } finally { service.dropSession('s1'); await rm(home, { recursive: true, force: true }); } }); describe('op journal', () => { it('assigns consecutive per-agent seqs and serves catch-up from the journal', async () => { const agents = new FakeAgents(); const main = agents.add('main'); const service = new TranscriptService({ homeDir: '/nonexistent-home', core: fakeCoreWithAgents(agents), }); service.forSessionLive('s1'); await service.whenReady('s1'); const base = service.getSeqWatermark('s1', 'main'); const seen: number[] = []; service.onSessionOps('s1', (_event, seq) => seen.push(seq)); main.bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' } })); main.bus.emit(ev({ type: 'turn.ended', turnId: 0, reason: 'completed' })); expect(seen).toEqual([base + 1, base + 2]); expect(service.getSeqWatermark('s1', 'main')).toBe(base + 2); const catchup = service.getOpsSince('s1', 'main', base); expect(catchup?.complete).toBe(true); expect(catchup?.latestSeq).toBe(base + 2); expect(catchup?.batches.map((batch) => batch.seq)).toEqual([base + 1, base + 2]); expect(service.getOpsSince('s1', 'main', base + 2)).toMatchObject({ batches: [], latestSeq: base + 2, complete: true, }); expect(service.getOpsSince('s1', 'main', base + 3)?.complete).toBe(false); const sub = agents.add('sub-1'); sub.bus.emit(ev({ type: 'turn.started', turnId: 0, origin: { kind: 'user' } })); expect(service.getSeqWatermark('s1', 'sub-1')).toBe(1); expect(service.getOpsSince('s1', 'sub-1', 0)?.batches.map((batch) => batch.seq)).toEqual([1]); expect(service.getSeqWatermark('s1', 'nope')).toBe(0); expect(service.getOpsSince('nope-session', 'main', 0)).toBeUndefined(); service.dropSession('s1'); }); it('marks catch-up incomplete once the bounded journal evicts old batches', async () => { const agents = new FakeAgents(); const main = agents.add('main'); const service = new TranscriptService({ homeDir: '/nonexistent-home', core: fakeCoreWithAgents(agents), }); service.forSessionLive('s1'); await service.whenReady('s1'); const base = service.getSeqWatermark('s1', 'main'); for (let turnId = 1; turnId <= TRANSCRIPT_OPS_JOURNAL_CAPACITY + 1; turnId++) { main.bus.emit(ev({ type: 'turn.started', turnId, origin: { kind: 'user' } })); } const watermark = service.getSeqWatermark('s1', 'main'); expect(watermark).toBe(base + TRANSCRIPT_OPS_JOURNAL_CAPACITY + 1); const evicted = service.getOpsSince('s1', 'main', base); expect(evicted?.complete).toBe(false); expect(evicted?.latestSeq).toBe(watermark); expect(evicted?.batches).toHaveLength(TRANSCRIPT_OPS_JOURNAL_CAPACITY); const recent = service.getOpsSince('s1', 'main', watermark - 10); expect(recent?.complete).toBe(true); expect(recent?.batches.map((batch) => batch.seq)).toEqual( Array.from({ length: 10 }, (_, i) => watermark - 9 + i), ); service.dropSession('s1'); }); }); });