返回源码地图

packages/core/session/src/repair.ts

main snapshot · da00f7f5358f · 正文引用章节 05;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Pending tool-result recovery shared by failed live steps, interrupted logs,
3 * and fork seeds. Tail repair preserves closed steps and supplies only missing
4 * tool results and lifecycle boundaries, with cause-specific retry guidance.
5 * @module @deepseek-ai/dsh-session/repair
6 */
7
8import { brandString } from '@deepseek-ai/dsh-brand'
9import type { MessageId, ToolCallId, ToolResultMessage } from '@deepseek-ai/dsh-llm'
10import { deepFreeze } from '@deepseek-ai/dsh-util-values'
11import { SessionSeq } from './types.ts'
12import type { SessionEvent, SessionSeq as SessionSeqType } from './types.ts'
13
14/** Recovery code for an assistant tool request that never reached a recorded call start. */
15export const TOOL_NOT_STARTED = 'TOOL_NOT_STARTED'
16
17/** Recovery code for a recorded tool call whose completed outcome was not durably recorded. */
18export const TOOL_OUTCOME_UNKNOWN = 'TOOL_OUTCOME_UNKNOWN'
19
20/**
21 * Why an open tail turn is closed with synthetic events: `interrupted` is
22 * crash recovery over a persisted log; `forked` is a fork seed cut inside the
23 * source's open turn. The cause selects the synthetic `turn/end` reason, the
24 * model-visible wording of synthetic error tool results, and the
25 * deterministic synthetic message-id prefix. The error codes
26 * ({@link TOOL_NOT_STARTED} / {@link TOOL_OUTCOME_UNKNOWN}) are shared: both
27 * causes state the same fact about the call's recorded lifecycle.
28 */
29export type OpenTurnCloseCause = { readonly kind: 'interrupted' } | { readonly kind: 'forked' }
30
31/** Model-visible wording of the synthetic error tool results, keyed by cause. */
32const CLOSER_TEXT = {
33 interrupted: {
34 started: 'The tool call was interrupted after it was recorded, but no result was durably recorded. Its outcome is unknown. Decide whether to retry from the tool semantics: retry only if the operation is read-only or idempotent; if it may have side effects, first verify external state or ask the user. Do not retry blindly.',
35 notStarted: 'The tool call was interrupted before the Harness recorded it as started. Retry it if it is still needed.',
36 },
37 forked: {
38 started: 'The history inherited by this branch records this tool call starting but does not include its result. The parent session may have completed it after the fork point. Decide whether to retry from the tool semantics: retry only if the operation is read-only or idempotent; if it may have side effects, first verify external state or ask the user. Do not retry blindly.',
39 notStarted: 'The history inherited by this branch has no record of this tool call starting. The parent session may have executed it after the fork point. Decide whether to retry from the tool semantics: retry only if the operation is read-only or idempotent; if it may have side effects, first verify external state or ask the user. Do not retry blindly.',
40 },
41} as const
42
43/**
44 * Return deterministic synthetic events that close an open tail turn. Unmatched
45 * calls receive error results first, followed by an open `step/end` and a
46 * cause-specific `turn/end`; sequences continue the log and timestamps reuse the
47 * last real event. A balanced or empty log returns no events.
48 *
49 * @param events - the loaded durable log to scan (a valid committed prefix, possibly with a crash tail).
50 * @param cause - the owning close operation: selects result wording and the turn-ending reason.
51 * @returns the synthetic closer events to append after `events`, in order; empty when the log is already balanced.
52 */
53export function openTurnClosers(events: readonly SessionEvent[], cause: OpenTurnCloseCause): SessionEvent[] {
54 let openTurn: number | null = null
55 let openStep: number | null = null
56 const recovery = new ToolCallRecovery(cause)
57 for (const event of events) {
58 recovery.observe(event)
59 switch (event.type) {
60 case 'turn/start':
61 openTurn = event.data.turn
62 openStep = null
63 break
64 case 'turn/end':
65 openTurn = null
66 openStep = null
67 break
68 case 'step/start':
69 openStep = event.data.step
70 break
71 case 'step/end':
72 openStep = null
73 break
74 // Other event types do not move the turn/step boundary cursor.
75 default:
76 break
77 }
78 }
79
80 // Balanced log (no open tail turn): nothing to close. An open turn implies
81 // `events` is non-empty (its turn/start was logged), so `last` exists.
82 const last = events.at(-1)
83 if (openTurn === null || last === undefined) return []
84
85 // The last real event supplies the seq base and the timestamp for the
86 // synthetic closers (reusing the last timestamp keeps them deterministic and
87 // never invents a "future" time).
88 const closers: SessionEvent[] = recovery.results()
89 let seq = last.seq + closers.length + 1
90 const time = last.time
91
92 // Close an open step before its turn.
93 if (openStep !== null) {
94 closers.push({ type: 'step/end', seq: SessionSeq(seq++), time, data: { turn: openTurn, step: openStep } })
95 }
96 closers.push({ type: 'turn/end', seq: SessionSeq(seq++), time, data: { turn: openTurn, reason: { kind: cause.kind } } })
97 return closers
98}
99
100/**
101 * Track unanswered assistant tool requests from one Session's committed events.
102 * Observe from the start of the owned step or replay prefix, and recover before
103 * its step closes. This state retains pending identities, not event history.
104 */
105export class ToolCallRecovery {
106 private readonly pendingCalls = new Map<ToolCallId, { turn: number; step: number; callSeq?: SessionSeqType }>()
107 private last: Pick<SessionEvent, 'seq' | 'time'> | undefined
108
109 /** @param cause - defaults to interrupted live/crash recovery; fork-seed construction supplies its own cause. */
110 constructor(private readonly cause: OpenTurnCloseCause = { kind: 'interrupted' }) {}
111
112 /**
113 * Consume the next committed event; closed steps and turn boundaries discard pending requests.
114 * @param event - the next event from the same Session, in sequence order.
115 */
116 observe(event: SessionEvent): void {
117 this.last = { seq: event.seq, time: event.time }
118 switch (event.type) {
119 case 'turn/start':
120 case 'turn/end':
121 case 'step/end':
122 this.pendingCalls.clear()
123 break
124 case 'assistant/message':
125 for (const block of event.data.message.content) {
126 if (block.type === 'tool-call') {
127 this.pendingCalls.set(block.id, { turn: event.data.turn, step: event.data.step })
128 }
129 }
130 break
131 case 'tool/call': {
132 const entry = this.pendingCalls.get(event.data.callId)
133 if (entry) entry.callSeq = event.seq
134 break
135 }
136 case 'tool/result': {
137 const callId = event.data.message.source.callId
138 const entry = this.pendingCalls.get(callId)
139 if (event.surfaceOp === 'append' && entry !== undefined
140 && entry.turn === event.data.turn && entry.step === event.data.step) {
141 this.pendingCalls.delete(callId)
142 }
143 break
144 }
145 // SessionEvent is merge-extensible; unrelated events retain pending requests.
146 default:
147 break
148 }
149 }
150
151 /**
152 * Build conservative error results in assistant order without changing tracked state.
153 * Sequences follow the latest observed event and timestamps reuse its time.
154 * Callers commit the results and observe those commits before recovering again.
155 * @returns pending tool-result events, empty when no request remains unanswered.
156 */
157 results(): SessionEvent<'tool/result'>[] {
158 if (this.last === undefined) return []
159 let seq = this.last.seq + 1
160 const time = this.last.time
161 const results: SessionEvent<'tool/result'>[] = []
162
163 const text = CLOSER_TEXT[this.cause.kind]
164 // Close calls before their step: providers reject dangling assistant calls,
165 // and Map insertion order preserves their transcript order.
166 for (const [callId, { turn, step, callSeq }] of this.pendingCalls) {
167 const started = callSeq !== undefined
168 const message: ToolResultMessage = deepFreeze({
169 id: brandString<MessageId>(`${this.cause.kind}-tool-result-${callId}-${seq}`),
170 role: 'tool',
171 toolCallId: callId,
172 isError: true,
173 source: { kind: 'tool', callId },
174 content: [{
175 type: 'text',
176 text: started ? text.started : text.notStarted,
177 }],
178 })
179 results.push({
180 type: 'tool/result',
181 seq: SessionSeq(seq++),
182 time,
183 data: {
184 turn,
185 step,
186 message,
187 error: started
188 ? { name: 'ToolOutcomeUnknownError', code: TOOL_OUTCOME_UNKNOWN }
189 : { name: 'ToolNotStartedError', code: TOOL_NOT_STARTED },
190 },
191 surfaceOp: 'append',
192 ...started ? { sourceEventSeqs: [callSeq] } : {},
193 })
194 }
195
196 return results
197 }
198}
199
200/**
201 * Crash-recovery entry point: synthetic closers that balance a persisted log
202 * whose tail turn was interrupted. Used by crash-recovery callers; fork
203 * seeds receive their `forked`-cause closers through `buildForkSeed` in
204 * `./fork.ts`.
205 *
206 * @param events - the persisted log to scan, possibly ending inside an open turn.
207 * @returns the synthetic `interrupted` closer events to append after `events`; empty when the log is already balanced.
208 */
209export function interruptedTurnClosers(events: readonly SessionEvent[]): SessionEvent[] {
210 return openTurnClosers(events, { kind: 'interrupted' })
211}