| 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'); |
| }); |
| }); |
| }); |
|
|