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