diff --git a/.changeset/steer-frame-prompt-ids.md b/.changeset/steer-frame-prompt-ids.md new file mode 100644 index 00000000000..18156ce52fb --- /dev/null +++ b/.changeset/steer-frame-prompt-ids.md @@ -0,0 +1,6 @@ +--- +"@moonshot-ai/kap-server": patch +"@moonshot-ai/transcript": patch +--- + +Stamp cold-rebuilt steer frames with the promptIds paired from prompt.steered records, and flush pending steer/notification frames in arrival order so post-turn heal can no longer drop prompt pairing or misorder frames. diff --git a/packages/kap-server/src/services/transcript/coreEventMap.ts b/packages/kap-server/src/services/transcript/coreEventMap.ts index ebe2c334f1b..40b9d62b34d 100644 --- a/packages/kap-server/src/services/transcript/coreEventMap.ts +++ b/packages/kap-server/src/services/transcript/coreEventMap.ts @@ -193,12 +193,15 @@ export interface ToolFrameRecord { export class AgentTranscriptProjector { private currentTurn: TurnHeader | undefined; private currentStep: StepHeader | undefined; - private pendingTaskNotifications: { text: string; taskId: string | undefined }[] = []; - private pendingSteers: { - input: readonly ContentPart[]; - promptIds: readonly string[] | undefined; - origin: TranscriptUserOrigin; - }[] = []; + private pendingUserFrames: ( + | { kind: 'notification'; text: string; taskId: string | undefined } + | { + kind: 'steer'; + input: readonly ContentPart[]; + promptIds: readonly string[] | undefined; + origin: TranscriptUserOrigin; + } + )[] = []; private unpairedSteerPromptIds: string[][] = []; private readonly stepOrdinals = new Map(); private frameOrdinal = 0; @@ -403,8 +406,7 @@ export class AgentTranscriptProjector { startedAt: nowIso(), }; this.currentStep = undefined; - this.pendingTaskNotifications = []; - this.pendingSteers = []; + this.pendingUserFrames = []; this.openText = undefined; this.openThinking = undefined; ops.push({ op: 'turn.upsert', turn: this.currentTurn }); @@ -428,7 +430,8 @@ export class AgentTranscriptProjector { this.currentStep = step; ops.push({ op: 'step.upsert', turnId: step.turnId, step }); } - if (this.currentStep === undefined && this.pendingSteers.length > 0) { + const pendingSteers = this.pendingUserFrames.filter((pending) => pending.kind === 'steer'); + if (this.currentStep === undefined && pendingSteers.length > 0) { const ordinal = (this.stepOrdinals.get(turnId) ?? this.lookups?.stepOrdinal?.(turnId) ?? 0) + 1; const step: StepHeader = { kind: 'step', @@ -443,7 +446,7 @@ export class AgentTranscriptProjector { ops.push({ op: 'step.upsert', turnId, step }); } if (this.currentStep !== undefined) { - for (const pending of this.pendingSteers) { + for (const pending of pendingSteers) { this.steerUserFrame( ops, turnId, @@ -454,7 +457,7 @@ export class AgentTranscriptProjector { ); } } - this.pendingSteers = []; + this.pendingUserFrames = []; const prev = this.currentTurn?.turnId === turnId ? this.currentTurn : this.lookups?.turn?.(turnId); const state = mapTurnEndState(event.reason); @@ -476,7 +479,6 @@ export class AgentTranscriptProjector { ops.push({ op: 'turn.upsert', turn: this.currentTurn }); ops.push({ op: 'meta.merge', meta: { activity: 'idle' } }); this.currentStep = undefined; - this.pendingTaskNotifications = []; if (event.reason === 'cancelled' && event.interruptReason === 'user_cancelled') { ops.push( this.markerOp('interruption', { turnId: event.turnId, reason: event.interruptReason }), @@ -523,25 +525,25 @@ export class AgentTranscriptProjector { this.openText = undefined; this.openThinking = undefined; const ops: TranscriptOperation[] = [{ op: 'step.upsert', turnId, step: this.currentStep }]; - for (const pending of this.pendingTaskNotifications) { - ops.push({ - op: 'frame.upsert', - turnId, - stepId, - frame: { - kind: 'text', - frameId: `${stepId}.f${++this.frameOrdinal}`, - role: 'user', - text: pending.text, - taskId: pending.taskId, - }, - }); - } - this.pendingTaskNotifications = []; - for (const pending of this.pendingSteers) { + for (const pending of this.pendingUserFrames) { + if (pending.kind === 'notification') { + ops.push({ + op: 'frame.upsert', + turnId, + stepId, + frame: { + kind: 'text', + frameId: `${stepId}.f${++this.frameOrdinal}`, + role: 'user', + text: pending.text, + taskId: pending.taskId, + }, + }); + continue; + } this.steerUserFrame(ops, turnId, stepId, pending.input, pending.promptIds, pending.origin); } - this.pendingSteers = []; + this.pendingUserFrames = []; return ops; } @@ -892,7 +894,7 @@ export class AgentTranscriptProjector { return [{ op: 'frame.upsert', turnId: turn.turnId, stepId: step.stepId, frame }]; } if (turn.origin?.kind === 'task' && (turn.origin.taskId === undefined || turn.origin.taskId === event.sourceId)) return []; - this.pendingTaskNotifications.push({ text, taskId: event.sourceId }); + this.pendingUserFrames.push({ kind: 'notification', text, taskId: event.sourceId }); return []; } @@ -1458,7 +1460,8 @@ export class AgentTranscriptProjector { ); return ops; } - this.pendingSteers.push({ + this.pendingUserFrames.push({ + kind: 'steer', input, promptIds: this.unpairedSteerPromptIds.shift(), origin: frameOrigin, diff --git a/packages/kap-server/src/services/transcript/transcriptService.ts b/packages/kap-server/src/services/transcript/transcriptService.ts index 887f1c1c90f..9a05cfd17a5 100644 --- a/packages/kap-server/src/services/transcript/transcriptService.ts +++ b/packages/kap-server/src/services/transcript/transcriptService.ts @@ -469,28 +469,76 @@ export class TranscriptService { } const messages = [...reduceContextTranscript(records).entries]; const taskOriginTurnTaskIds = new Set(); - const steeredContents = new Map>(); - const anchorStack: { taskIdsSnapshot: Set }[] = []; + let steeredContents = new Map>(); + const steeredPromptIds: (readonly string[] | undefined)[] = []; + const pendingSteerPromptIds: (readonly string[])[] = []; + const anchorStack: { + taskIdsSnapshot: Set; + steeredSnapshot: Map>; + pendingSteerCount: number; + steeredIdsCount: number; + }[] = []; let anchorFloor = 0; + let floorSteerState: + | { + steeredSnapshot: Map>; + pendingSteerCount: number; + steeredIdsCount: number; + } + | undefined; let sawTurnPrompt = false; for (const record of records) { if (record.type === 'context.undo') { const count = typeof record['count'] === 'number' ? (record['count'] as number) : 0; + let poppedAny = false; for (let i = 0; i < count && anchorStack.length > anchorFloor; i++) { const popped = anchorStack.pop()!; + poppedAny = true; taskOriginTurnTaskIds.clear(); for (const id of popped.taskIdsSnapshot) taskOriginTurnTaskIds.add(id); } + if (poppedAny) { + const top = + anchorStack.length > anchorFloor ? anchorStack[anchorStack.length - 1] : undefined; + pendingSteerPromptIds.length = top?.pendingSteerCount ?? floorSteerState?.pendingSteerCount ?? 0; + steeredPromptIds.length = top?.steeredIdsCount ?? floorSteerState?.steeredIdsCount ?? 0; + steeredContents = new Map( + [...(top?.steeredSnapshot ?? floorSteerState?.steeredSnapshot ?? [])].map( + ([key, byKind]) => [key, new Map(byKind)], + ), + ); + } continue; } if (record.type === 'context.clear') { anchorFloor = anchorStack.length; + floorSteerState = { + steeredSnapshot: new Map( + [...steeredContents].map(([key, byKind]) => [key, new Map(byKind)]), + ), + pendingSteerCount: pendingSteerPromptIds.length, + steeredIdsCount: steeredPromptIds.length, + }; continue; } if (record.type === 'context.append_message') { const message = (record as { message?: ContextMessage }).message; if (message !== undefined && isUndoAnchor(message)) { - anchorStack.push({ taskIdsSnapshot: new Set(taskOriginTurnTaskIds) }); + anchorStack.push({ + taskIdsSnapshot: new Set(taskOriginTurnTaskIds), + steeredSnapshot: new Map( + [...steeredContents].map(([key, byKind]) => [key, new Map(byKind)]), + ), + pendingSteerCount: pendingSteerPromptIds.length, + steeredIdsCount: steeredPromptIds.length, + }); + } + continue; + } + if (record.type === 'prompt.steered') { + const promptIds = record['promptIds']; + if (Array.isArray(promptIds) && promptIds.every((id) => typeof id === 'string')) { + pendingSteerPromptIds.push(promptIds as readonly string[]); } continue; } @@ -503,6 +551,11 @@ export class TranscriptService { const byKind = steeredContents.get(key) ?? new Map(); byKind.set(kind, (byKind.get(kind) ?? 0) + 1); steeredContents.set(key, byKind); + steeredPromptIds.push( + kind === 'user' && pendingSteerPromptIds.length > 0 + ? pendingSteerPromptIds.shift() + : undefined, + ); } continue; } @@ -519,7 +572,9 @@ export class TranscriptService { } const base = groupMessagesIntoSnapshot( messages, - sawTurnPrompt || steeredContents.size > 0 ? { taskOriginTurnTaskIds, steeredContents } : undefined, + sawTurnPrompt || steeredContents.size > 0 + ? { taskOriginTurnTaskIds, steeredContents, steeredPromptIds } + : undefined, ); const folded = foldWireRecordFacts(projectQuestionInteractionRecords(records, sessionId), base, { resolvePlanRevisionKey: (key) => diff --git a/packages/kap-server/test/services/transcript.test.ts b/packages/kap-server/test/services/transcript.test.ts index a31cd4a76c2..c190f9063b4 100644 --- a/packages/kap-server/test/services/transcript.test.ts +++ b/packages/kap-server/test/services/transcript.test.ts @@ -2105,6 +2105,41 @@ describe('AgentTranscriptProjector', () => { }); }); + it('flushes queued steers and task notifications at step start in arrival order', () => { + 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: 4, origin: { kind: 'user' }, prompt: 'active' })); + feed(ev({ type: 'turn.step.started', turnId: 4, step: 1 })); + feed(ev({ type: 'turn.step.completed', turnId: 4, step: 1 })); + feed( + ev({ + type: 'turn.steer', + input: [{ type: 'text', text: 'steered in' }], + origin: { kind: 'user' }, + }), + ); + 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: 4, step: 2 })); + const frames = turnOps('t4', tx.getItems()).steps[1]!.frames; + expect(frames.map((f) => f.kind === 'text' && 'text' in f && f.text)).toEqual([ + 'steered in', + 'Background agent completed\ninspect done.', + ]); + }); + 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'); @@ -2680,6 +2715,233 @@ describe('AgentTranscriptProjector', () => { } }); + it('readColdSnapshot stamps the steered frame with promptIds paired from prompt.steered', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-ids-')); + 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: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_steered'], content: [{ type: 'text', text: 'steered in' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2500 }, + { 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 turn = snapshot!.items.find((item) => item.kind === 'turn'); + if (turn?.kind !== 'turn') throw new Error('expected turn'); + expect(turn.steps[1]?.frames[0]).toMatchObject({ + kind: 'text', + role: 'user', + text: 'steered in', + promptIds: ['msg_steered'], + }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot pairs promptIds when the steer carries bundled skill blocks', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-skill-')); + try { + const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); + await mkdir(wireDir, { recursive: true }); + const activation = { activationId: 'a1', skillName: 'deploy' }; + 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: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_steered'], content: [{ type: 'text', text: 'deploy now' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2500 }, + { type: 'turn.steer', input: [{ type: 'text', text: '/deploy' }, { type: 'text', text: 'deploy now' }], origin: { kind: 'user', skillActivations: [activation] }, time: 3000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: '/deploy' }, { type: 'text', text: 'deploy now' }], toolCalls: [], origin: { kind: 'user', skillActivations: [activation] } }, time: 3001 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'done' }], 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 turn = snapshot!.items.find((item) => item.kind === 'turn'); + if (turn?.kind !== 'turn') throw new Error('expected turn'); + expect(turn.steps[1]?.frames[0]).toMatchObject({ + kind: 'text', + role: 'user', + text: 'deploy now', + promptIds: ['msg_steered'], + }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot rolls back steer id queues across context.undo', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-undo-')); + 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: 'one' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, time: 2000 }, + { type: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_a'], content: [{ type: 'text', text: 'same' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2500 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'same' }], origin: { kind: 'user' }, time: 3000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, time: 3001 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, time: 4000 }, + { type: 'context.undo', count: 1, time: 5000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'two' }], toolCalls: [], origin: { kind: 'user' } }, time: 6000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working2' }], toolCalls: [] }, time: 6500 }, + { type: 'prompt.steered', activePromptId: 'msg_active2', promptIds: ['msg_b'], content: [{ type: 'text', text: 'same' }], steeredAt: '2026-09-01T10:01:00.000Z', time: 7000 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'same' }], origin: { kind: 'user' }, time: 7500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, time: 7501 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted2' }], toolCalls: [] }, time: 8000 }, + ]; + 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(2); + const turn = turns[1]; + if (turn?.kind !== 'turn') throw new Error('expected turn'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames).toHaveLength(1); + expect(steerFrames[0]).toMatchObject({ text: 'same', promptIds: ['msg_b'] }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot keeps user steer ids when a non-user turn.steer lands in between', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-nonuser-')); + try { + const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); + await mkdir(wireDir, { recursive: true }); + const skillOrigin = { kind: 'skill_activation', trigger: 'user-slash', skillName: 'deploy' }; + 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: 1500 }, + { type: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_u'], content: [{ type: 'text', text: 'user steer' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2000 }, + { type: 'turn.steer', input: [{ type: 'text', text: '/deploy' }], origin: skillOrigin, time: 2500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: '/deploy' }], toolCalls: [], origin: skillOrigin }, time: 2501 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'user steer' }], origin: { kind: 'user' }, time: 3000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'user steer' }], 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 turn = snapshot!.items.find((item) => item.kind === 'turn'); + if (turn?.kind !== 'turn') throw new Error('expected turn'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames).toHaveLength(2); + expect(steerFrames[0]).toMatchObject({ text: '/deploy' }); + expect(steerFrames[0]!.kind === 'text' && steerFrames[0]!.role === 'user' && steerFrames[0]!.promptIds).toBeUndefined(); + expect(steerFrames[1]).toMatchObject({ text: 'user steer', promptIds: ['msg_u'] }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot keeps the surviving turn\'s steer ids when undo drops the newest prompt', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-survive-')); + 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: 'one' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, time: 1500 }, + { type: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_s'], content: [{ type: 'text', text: 'same' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2000 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'same' }], origin: { kind: 'user' }, time: 2500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, time: 3001 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, time: 3500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'two' }], toolCalls: [], origin: { kind: 'user' } }, time: 4000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working2' }], toolCalls: [] }, time: 4500 }, + { type: 'context.undo', count: 1, time: 5000 }, + ]; + 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'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames).toHaveLength(1); + expect(steerFrames[0]).toMatchObject({ text: 'same', promptIds: ['msg_s'] }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot keeps user steer ids when a marker-folded steer lands in between', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-marker-')); + try { + const wireDir = join(home, 'sessions', 'ws', 's1', 'agents', 'main'); + await mkdir(wireDir, { recursive: true }); + const skillOrigin = { kind: 'skill_activation', trigger: 'model-tool', skillName: 'deploy' }; + 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: 1500 }, + { type: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_u'], content: [{ type: 'text', text: 'user steer' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2000 }, + { type: 'turn.steer', input: [{ type: 'text', text: '/deploy' }], origin: skillOrigin, time: 2500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: '/deploy' }], toolCalls: [], origin: skillOrigin }, time: 2501 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'user steer' }], origin: { kind: 'user' }, time: 3000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'user steer' }], 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 turn = snapshot!.items.find((item) => item.kind === 'turn'); + if (turn?.kind !== 'turn') throw new Error('expected turn'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames).toHaveLength(1); + expect(steerFrames[0]).toMatchObject({ text: 'user steer', promptIds: ['msg_u'] }); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + it('readColdSnapshot preserves pre-clear steer ids when undo drops the first post-clear prompt', async () => { + const home = await mkdtemp(join(tmpdir(), 'transcript-cold-steer-clear-')); + 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: 'one' }], toolCalls: [], origin: { kind: 'user' } }, time: 1000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, time: 1500 }, + { type: 'prompt.steered', activePromptId: 'msg_active', promptIds: ['msg_s'], content: [{ type: 'text', text: 'same' }], steeredAt: '2026-09-01T10:00:00.000Z', time: 2000 }, + { type: 'turn.steer', input: [{ type: 'text', text: 'same' }], origin: { kind: 'user' }, time: 2500 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, time: 3001 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, time: 3500 }, + { type: 'context.clear', time: 4000 }, + { type: 'context.append_message', message: { role: 'user', content: [{ type: 'text', text: 'two' }], toolCalls: [], origin: { kind: 'user' } }, time: 5000 }, + { type: 'context.append_message', message: { role: 'assistant', content: [{ type: 'text', text: 'working2' }], toolCalls: [] }, time: 5500 }, + { type: 'context.undo', count: 1, time: 6000 }, + ]; + 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'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames).toHaveLength(1); + expect(steerFrames[0]).toMatchObject({ text: 'same', promptIds: ['msg_s'] }); + } 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 { diff --git a/packages/transcript/src/history/groupTurns.ts b/packages/transcript/src/history/groupTurns.ts index 57d7a34fa7a..a6921c69047 100644 --- a/packages/transcript/src/history/groupTurns.ts +++ b/packages/transcript/src/history/groupTurns.ts @@ -71,6 +71,7 @@ export function groupMessagesIntoSnapshot( options?: { readonly taskOriginTurnTaskIds?: ReadonlySet; readonly steeredContents?: ReadonlyMap>; + readonly steeredPromptIds?: readonly (readonly string[] | undefined)[]; }, ): AgentTranscriptSnapshot { const items: TranscriptItem[] = []; @@ -78,6 +79,8 @@ export function groupMessagesIntoSnapshot( const steeredContents = new Map( [...(options?.steeredContents ?? [])].map(([key, byKind]) => [key, new Map(byKind)]), ); + const steeredPromptIds = options?.steeredPromptIds ?? []; + let steeredPromptIdIndex = 0; let turn: TurnDraft | undefined; let pendingNotificationFrames: { text: string; @@ -225,6 +228,16 @@ export function groupMessagesIntoSnapshot( if (!isTaskOrigin) prevNonTaskRole = message.role; if (message.role === 'user') { + const contentKey = JSON.stringify(message.content ?? []); + const steerKind = originKind ?? 'user'; + const steeredByKind = steeredContents.get(contentKey); + const steeredRemaining = steeredByKind?.get(steerKind) ?? 0; + const matchedSteer = steeredByKind !== undefined && steeredRemaining > 0; + const steeredPromptId = matchedSteer ? steeredPromptIds[steeredPromptIdIndex] : undefined; + if (matchedSteer) { + steeredByKind.set(steerKind, steeredRemaining - 1); + steeredPromptIdIndex += 1; + } if (originKind !== undefined && HIDDEN_USER_ORIGINS.has(originKind)) { if (opensOwnTurn(message)) { const opening = @@ -240,12 +253,7 @@ export function groupMessagesIntoSnapshot( pushMarker(markerKey, { text: textOf(message), origin: message.origin }); continue; } - const contentKey = JSON.stringify(message.content ?? []); - const steerKind = originKind ?? 'user'; - const steeredByKind = steeredContents.get(contentKey); - const steeredRemaining = steeredByKind?.get(steerKind) ?? 0; - if (steeredByKind !== undefined && steeredRemaining > 0) { - steeredByKind.set(steerKind, steeredRemaining - 1); + if (matchedSteer) { const bundled = bundledSkillActivations(message); const parts = message.content ?? []; bundled.forEach((activation, index) => { @@ -260,6 +268,7 @@ export function groupMessagesIntoSnapshot( text: opening.text, taskId: undefined, attachmentIds: opening.attachmentIds, + promptIds: steeredPromptId, origin: projectTranscriptUserOrigin(message.origin), steered: true, }); diff --git a/packages/transcript/test/layers.test.ts b/packages/transcript/test/layers.test.ts index 79aa9bd7910..ae7c54474e8 100644 --- a/packages/transcript/test/layers.test.ts +++ b/packages/transcript/test/layers.test.ts @@ -548,6 +548,51 @@ describe('groupMessagesIntoSnapshot (cold path)', () => { }); }); + it('stamps the folded steer frame with the prompt ids from the steer queue', () => { + const snapshot = groupMessagesIntoSnapshot( + [ + { role: 'user', content: [{ type: 'text', text: 'active' }], toolCalls: [], origin: { kind: 'user' } }, + { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, + { role: 'user', content: [{ type: 'text', text: 'steered in' }], toolCalls: [], origin: { kind: 'user' } }, + { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, + ], + { steeredContents: new Map([[JSON.stringify([{ type: 'text', text: 'steered in' }]), new Map([['user', 1]])]]), steeredPromptIds: [['prompt_a']] }, + ); + + const turn = snapshot.items[0]; + if (turn?.kind !== 'turn') throw new Error('expected turn'); + expect(turn.steps[1]?.frames[0]).toMatchObject({ + kind: 'text', + role: 'user', + text: 'steered in', + promptIds: ['prompt_a'], + }); + }); + + it('pairs repeated steers of identical content with their prompt ids in order', () => { + const snapshot = groupMessagesIntoSnapshot( + [ + { role: 'user', content: [{ type: 'text', text: 'active' }], toolCalls: [], origin: { kind: 'user' } }, + { role: 'assistant', content: [{ type: 'text', text: 'working' }], toolCalls: [] }, + { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, + { role: 'assistant', content: [{ type: 'text', text: 'noted' }], toolCalls: [] }, + { role: 'user', content: [{ type: 'text', text: 'same' }], toolCalls: [], origin: { kind: 'user' } }, + { role: 'assistant', content: [{ type: 'text', text: 'noted again' }], toolCalls: [] }, + ], + { steeredContents: new Map([[JSON.stringify([{ type: 'text', text: 'same' }]), new Map([['user', 2]])]]), steeredPromptIds: [['prompt_a'], ['prompt_b']] }, + ); + + const turn = snapshot.items[0]; + if (turn?.kind !== 'turn') throw new Error('expected turn'); + const steerFrames = turn.steps.flatMap((step) => + step.frames.filter((frame) => frame.kind === 'text' && frame.role === 'user'), + ); + expect(steerFrames.map((frame) => frame.kind === 'text' && frame.role === 'user' ? frame.promptIds : undefined)).toEqual([ + ['prompt_a'], + ['prompt_b'], + ]); + }); + it('keeps a trailing steered message visible by appending it to the last step', () => { const snapshot = groupMessagesIntoSnapshot( [