Download packages/kap-server/test/services/transcript.test.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 163 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/kap-server/test/services/transcript.test.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/kap-server/test/services/transcript.test.ts
-
curl -L -o transcript.test.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/kap-server/test/services/transcript.test.ts
163 kB
| 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<string, unknown>): ProjectorBusEvent { | |
| return payload as unknown as ProjectorBusEvent; | |
| } | |
| const TEST_SESSION_ID = 'session-test'; | |
| function turnOps(turnId: string, items: ReturnType<AgentTranscript['getItems']>): 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<string, unknown>[] = [ | |
| 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<typeof item, { kind: 'marker' }> => 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<typeof item, { kind: 'marker' }> => 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<string, unknown> }).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<string, unknown>): 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<typeof item, { kind: 'marker' }> => 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<typeof item, { kind: 'marker' }> => 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: '<skill-loaded name="deploy">private instructions</skill-loaded>' }, | |
| { type: 'text', text: '<skill-loaded name="review">private instructions</skill-loaded>' }, | |
| { 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: '<skill-loaded name="review">private instructions</skill-loaded>' }, | |
| { 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<void> => | |
| 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 = | |
| '<notification id="task:task_9:completed" category="task" type="task.completed" source_kind="background_task" source_id="task_9">\nTitle: Background agent completed\nSeverity: info\ninspect done.\n</notification>'; | |
| 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<void> => | |
| 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: '<skill-loaded name="deploy">private instructions</skill-loaded>' }, | |
| { type: 'text', text: '<skill-loaded name="review">private instructions</skill-loaded>' }, | |
| { 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 = | |
| '<notification id="task:task_9:completed" category="task" type="task.completed" source_kind="background_task" source_id="task_9">\nTitle: Background agent completed\nSeverity: info\ninspect done.\n</notification>'; | |
| 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<void> => | |
| 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<TranscriptTask['state']> = [ | |
| 'running', | |
| 'completed', | |
| 'failed', | |
| 'timed_out', | |
| 'killed', | |
| 'lost', | |
| ]; | |
| expect(states).toHaveLength(6); | |
| }); | |
| }); | |
| describe('bindSessionTranscript', () => { | |
| class FakeBus { | |
| private readonly handlers = new Set<(event: Event2<any>) => void>(); | |
| subscribe(cb: (event: Event2<any>) => void): { dispose: () => void } { | |
| this.handlers.add(cb); | |
| return { dispose: () => this.handlers.delete(cb) }; | |
| } | |
| emit(event: Event2<any>): 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<string, AgentConversationUndoParticipant>; | |
| readonly accessor: { get: (token: unknown) => unknown }; | |
| } | |
| class FakeAgents { | |
| private readonly handles = new Map<string, FakeAgentHandle>(); | |
| 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<string, AgentConversationUndoParticipant>(); | |
| 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<string, unknown>): Promise<string> { | |
| 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<string, unknown>[] = [ | |
| { | |
| 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<string, TranscriptOperation[]>(); | |
| 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<string, TranscriptOperation[]>(); | |
| 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<string, TranscriptOperation[]>(); | |
| 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<string> { | |
| 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<void> { | |
| 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<string, TranscriptOperation[]>(); | |
| 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<AgentTranscriptSnapshot | undefined>((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'); | |
| }); | |
| }); | |
| }); | |