返回 DeepSeek-Reasonix
run_loop.go
根目录 / internal / agent / run_loop.go
1 package agent
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "slices"
9 "strings"
10 "time"
11
12 "reasonix/internal/event"
13 "reasonix/internal/evidence"
14 "reasonix/internal/provider"
15 "reasonix/internal/runtimepolicy"
16 )
17
18 // streamedTurn is one provider completion collected by stream. Keeping the
19 // result together makes the missing-reasoning recovery path explicit: the
20 // first, malformed completion is never committed before a safe replacement is
21 // available, and a failed recovery can still fall back to the complete first
22 // response without re-running any tool.
23 type streamedTurn struct {
24 messageID string
25 displayReasoning string
26 settledAttemptID string
27 settledAttempt int
28 text string
29 reasoning string
30 signature string
31 reasoningID string
32 reasoningStatus string
33 reasoningComplete bool
34 reasoningState provider.ReasoningState
35 thinkingBlocks []provider.ThinkingBlock
36 calls []provider.ToolCall
37 responsesItems []json.RawMessage
38 serverSearch []provider.ServerSearchCall
39 usage *provider.Usage
40 interrupted bool
41 partialToolStarted bool
42 partialCalls []provider.ToolCall
43 maxArgChars int // peak streaming tool-arg size for failed-attempt estimates
44 err error
45 }
46
47 func (s streamedTurn) assistantMessage() provider.Message {
48 return provider.Message{
49 ID: s.messageID,
50 Role: provider.RoleAssistant, Content: s.text, ReasoningContent: s.reasoning,
51 ReasoningState: s.reasoningState, ThinkingBlocks: s.thinkingBlocks,
52 ReasoningSignature: s.signature, ReasoningID: s.reasoningID, ReasoningStatus: s.reasoningStatus,
53 ToolCalls: s.calls, ResponsesItems: s.responsesItems, ServerSearch: s.serverSearch,
54 }
55 }
56
57 // beginRunTurn handles evidence scope, delivery classification, background-job
58 // evidence re-lease, and the initial user-turn persistence. Callers still own
59 // all Run-level defers (workspace lease, evidence commit, delivery checkpoint,
60 // steer queue, active-turn timestamp).
61 func (a *Agent) beginRunTurn(ctx context.Context, input string, pinned pinnedRevisionPlan) (rawInput string, state *turnRuntime, err error) {
62 rawInput = RawUserInput(ctx, input)
63 providerInput := input
64 // A fresh user turn starts from zeroed per-turn host state; the new turn's
65 // values are computed below. Cross-turn state (checkpoint, scope, failure
66 // budgets) lives in taskRuntime and is reconciled there.
67 a.stragglers.drain(ctx, parallelStragglerGrace)
68 a.turn = turnRuntime{}
69 a.reads.runGen++
70 a.reads.tasks = newReadTasks(a.sess.path, a.reads.runGen)
71 a.reads.deliveries = make(map[string]readDelivery)
72 a.reads.visible = nil
73 scope, scoped := DeliveryExecutionScopeFromContext(ctx)
74 if a.task.ledger != nil {
75 switch {
76 case scoped && a.task.scopeID == scope.ID:
77 a.task.ledger.ResetBackgroundLeases()
78 default:
79 a.resetTurnEvidence()
80 }
81 }
82 if scoped {
83 a.task.scopeID = scope.ID
84 } else {
85 a.task.scopeID = ""
86 }
87 a.turn.deliveryScopeActive = scoped
88 if scoped && a.task.checkpoint.ScopeID != scope.ID {
89 a.task.checkpoint = evidence.DeliveryCheckpoint{ScopeID: scope.ID}
90 }
91 a.leasePendingBackgroundEvidence(ctx)
92 // Use the owning task text for explicit action constraints and recovery.
93 // Child framing must not be interpreted as an instruction from the user.
94 a.turn.turnInput = a.classifierTaskText
95 if scoped && strings.TrimSpace(scope.TaskText) != "" {
96 a.turn.turnInput = scope.TaskText
97 } else if strings.TrimSpace(a.turn.turnInput) == "" {
98 a.turn.turnInput = rawInput
99 }
100 a.turn.recoveryTaskSummary = boundedRecoveryTaskSummary(a.turn.turnInput)
101 if constraints, ok := runtimepolicy.FromContext(ctx); ok {
102 a.turn.constraints = constraints
103 } else {
104 a.turn.constraints = runtimepolicy.ParseConstraints(runtimepolicy.StripQuotedConstraints(a.turn.turnInput))
105 if a.planMode.Load() {
106 a.turn.constraints.PlanModeReadOnly = true
107 a.turn.constraints.ForbidMutation = true
108 }
109 }
110 if inherited, ok := runtimepolicy.InheritedFromContext(ctx); ok && !a.readOnlyExecution {
111 a.turn.constraints = mergeInheritedConstraints(a.turn.constraints, inherited.Constraints)
112 if inherited.PlanReadOnly {
113 a.turn.constraints.PlanModeReadOnly = true
114 a.turn.constraints.ForbidMutation = true
115 }
116 } else if a.inheritedExec != nil && !a.readOnlyExecution {
117 a.turn.constraints = mergeInheritedConstraints(a.turn.constraints, a.inheritedExec.Constraints)
118 if a.inheritedExec.PlanReadOnly {
119 a.turn.constraints.PlanModeReadOnly = true
120 a.turn.constraints.ForbidMutation = true
121 }
122 }
123 a.turn.engine = runtimepolicy.NewEngine(a.turn.constraints)
124 // A cancelled/error turn leaves a provider-excluded recovery record at the
125 // transcript tail. Fold its bounded facts into this new user turn exactly
126 // once; the user's raw text remains the source above.
127 a.ensureUnreplayableHistoryRecovery()
128 providerInput = withInterruptedRecovery(providerInput, a.verifyInterruptedWrites(ctx, a.pendingInterruptedRecovery()))
129 a.task.prepareScope(scoped, scope.ID)
130 a.svc.sink.Emit(event.Event{Kind: event.TurnStarted})
131 a.emitTurnPhase(event.TurnPhaseWorking)
132 input = a.prepareProviderTurn(ctx, providerInput)
133 userCreatedAt := time.Now().UnixMilli()
134 a.activeTurnCreatedAt.Store(userCreatedAt)
135 rawContent := rawInput
136 if rawContent == "" {
137 rawContent = a.turn.turnInput
138 }
139 userMessage := provider.Message{
140 ID: turnUserMessageID(ctx, a.sess.conversation),
141 Role: provider.RoleUser, Origin: inputMessageOrigin(ctx), Content: input, RawContent: rawContent,
142 Images: userImages(ctx), ImageInputs: userImageInputs(ctx), VisionSummary: VisionSummaryFromContext(ctx), CreatedAt: userCreatedAt,
143 }
144 if err := userMessage.ValidateImageFields(); err != nil {
145 return rawInput, nil, err
146 }
147 if err := a.appendPinnedRevisionAndUser(ctx, pinned, userMessage); err != nil {
148 return rawInput, nil, err
149 }
150 a.admitUserMessage(ctx, userMessage)
151
152 // The loop fields join the classification computed above rather than
153 // opening a second object: one turn, one turnRuntime. The zero values the
154 // old literal spelled out are already there from the reset at the top.
155 state = &a.turn
156 state.input = input
157 state.budget = runBudget{started: time.Now()}
158 return rawInput, state, nil
159 }
160
161 // runToolLoop owns the main tool-round budget and dispatches each streamed
162 // assistant turn into final-response or tool-round handling.
163 func (a *Agent) runToolLoop(ctx context.Context, state *turnRuntime) (runErr error) {
164 releaseMCPListObserver := a.activateMCPListObserver()
165 defer releaseMCPListObserver()
166 ctx = a.withAgentContext(ctx)
167 truncatedRounds := 0
168 for step := 0; state.runMaxSteps <= 0 || step < state.runMaxSteps || state.graceRound; step++ {
169 // Consume a queued steer and persist it to the session so it
170 // survives tab switches and history replay. The model sees it as
171 // guidance (with a prefix), not a new task. One cache miss per
172 // steer is unavoidable — the model must see the new instruction.
173 if text, itemID, ok := a.consumeSteer(); ok {
174 steerMessage := provider.Message{
175 ID: NewMessageID(),
176 Role: provider.RoleUser, Origin: provider.MessageOriginUser,
177 Content: a.withTurnPreferences(midTurnSteerMessage(text)), RawContent: text,
178 }
179 if err := a.appendCommittedMessages(ctx, "mid-turn-steer", steerMessage); err != nil {
180 return err
181 }
182 a.svc.sink.Emit(event.Event{Kind: event.Steer, MessageID: steerMessage.ID, Text: text, ItemID: itemID})
183 } else if itemID != "" {
184 // Loader failed after dequeue: durable entry stays for inspection
185 // (unapplied path marks uncertain + pause via the notice sink).
186 a.RecordUnappliedSteer("(body load failed)", itemID)
187 }
188 schemas := a.providerToolSchemas()
189 prefixShape := a.capturePrefixShape(schemas)
190 prevPrefixShape := a.sess.lastPrefixShape
191 if !a.sess.haveLastPrefixShape {
192 prevPrefixShape = prefixShape
193 }
194 // Drain reasons queued since the previous capture (compaction,
195 // snip/prune, rewind, guardian merge) so CompareShape can attribute
196 // any prefix change to the operation that actually caused it, instead
197 // of a generic rewrite signal that also fires on local-only metadata
198 // edits.
199 contentReasons := a.sess.conversation.DrainContentRewriteReasons()
200
201 // Prefix shape is captured once before sampling and frozen for the
202 // whole attempt lifecycle — stream retries must not rewrite session
203 // history mid-round, so the shape stays stable across body replays.
204 streamed := a.streamWithSamplingRecovery(ctx, step+1)
205 text, reasoning, calls, usage := streamed.text, streamed.reasoning, streamed.calls, streamed.usage
206 partialCalls, err := streamed.partialCalls, streamed.err
207 cacheDiagnostics := CompareShape(prevPrefixShape, prefixShape, usage, contentReasons)
208 a.attachSessionContextDiagnostics(&cacheDiagnostics)
209 if err != nil {
210 quote := a.emitTurnUsage(usage, &cacheDiagnostics)
211 a.observeRunBudget(state, usage, quote)
212 if msg, ok := finishReasonMessage(usage); ok {
213 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: msg})
214 }
215 // Exhausted stream retries (or a non-retryable error): persist one
216 // bounded LocalOnly recovery record for the next real user message.
217 // Intermediate failed attempts never wrote session state.
218 a.recordInterruptedDisplay(text, reasoning, partialCalls, true, err, state.workDurationMs(), streamed.messageID)
219 // A broken provider stream can otherwise look like a silent hang
220 // followed only by the generic interrupted-turn notice (#9560).
221 if code, msg := streamInterruptNotice(err); msg != "" {
222 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Code: code, Text: msg})
223 }
224 return err
225 }
226 a.sess.lastPrefixShape = prefixShape
227 a.sess.haveLastPrefixShape = true
228 quote := a.emitTurnUsage(usage, &cacheDiagnostics)
229 a.observeRunBudget(state, usage, quote)
230 if msg, ok := finishReasonMessage(usage); ok {
231 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: msg})
232 }
233
234 // Commit clean terminal attempts, preserving provider reasoning contracts.
235 calls = a.withPreviewFileDiffs(ctx, calls)
236 if err := assignRecoveryCallIDs(calls); err != nil {
237 return err
238 }
239 assistant := streamed.assistantMessage()
240 assistant.ToolCalls = calls
241 assistant.WorkDurationMs = state.workDurationMs()
242 if err := a.appendCommittedMessages(ctx, "assistant-attempt", assistant); err != nil {
243 if errors.Is(err, context.Canceled) {
244 a.recordInterruptedDisplay(text, reasoning, partialCalls, true, err, state.workDurationMs(), streamed.messageID)
245 }
246 return err
247 }
248 a.publishCommittedSample(streamed)
249
250 if len(calls) == 0 {
251 cont, ferr := a.handleFinalResponse(ctx, state, text, reasoning, usage)
252 if !cont {
253 return ferr
254 }
255 continue
256 }
257
258 if usage != nil && usage.FinishReason == "length" {
259 truncatedRounds++
260 if err := a.recordTruncatedToolResults(withMessageIdentity(ctx, streamed.messageID), calls); err != nil {
261 return err
262 }
263 if truncatedRounds > maxToolArgumentRepairs {
264 return fmt.Errorf("tool arguments remained truncated after three recovery rounds")
265 }
266 continue
267 }
268 truncatedRounds = 0
269
270 // Invariant: executeBatch only ever receives tool calls from a
271 // committed sampling attempt (clean terminal + response intercept).
272 cont, terr := a.handleToolRound(withMessageIdentity(ctx, streamed.messageID), state, step, text, reasoning, calls, usage)
273 if !cont {
274 return terr
275 }
276 }
277 // Only reached when a positive maxSteps guard is configured. The work so far
278 // is already in the session, so the user can just send another message to pick
279 // up where it left off.
280 return a.gracePause(state)
281 }
282
283 func (a *Agent) emitProtocolRetry(attempt int, hasFallback bool) {
284 maxAttempts := 1
285 if hasFallback {
286 maxAttempts = 2
287 }
288 a.svc.sink.Emit(event.Event{
289 Kind: event.Retrying, RetryAttempt: attempt, RetryMax: maxAttempts,
290 RetryScope: event.RetryScopeProtocol,
291 })
292 }
293
294 func (a *Agent) emitStreamAttempt(id string, action event.StreamAttemptAction, attempt int, reason string, err error) {
295 if reason == "" && err != nil {
296 reason = provider.StreamInterruptReason(err)
297 }
298 a.svc.sink.Emit(event.Event{
299 Kind: event.StreamAttempt,
300 MessageID: id,
301 AttemptID: id,
302 StreamAttempt: event.StreamAttemptInfo{
303 ID: id, Action: action, Attempt: attempt, Max: maxSamplingAttempts, Reason: reason,
304 },
305 })
306 }
307
308 func newStreamAttemptID(_ int) string {
309 // A successful attempt retains this local identity when its message is
310 // committed. Failed attempts have distinct identities and cannot alias it.
311 return NewMessageID()
312 }
313
314 // handleFinalResponse processes a no-tool assistant turn: recovery pause,
315 // readiness boundary, empty-final retry, executor handoff nudge, steer drain,
316 // and final compaction. cont=true continues the tool loop; cont=false returns
317 // err from Run (err may be nil for a clean final answer).
318 func (a *Agent) handleFinalResponse(ctx context.Context, state *turnRuntime, text, reasoning string, usage *provider.Usage) (cont bool, err error) {
319 if state.graceRound {
320 // Explicit max_steps and spend budgets are user-selected boundaries.
321 // Preserve the summary, then return a resumable pause so Goal does not
322 // immediately open another Run and silently bypass the chosen limit.
323 a.contextManager().ObserveUsage(usage)
324 return false, a.gracePause(state)
325 }
326 if !hasVisibleFinalAnswer(text) {
327 // A reasoning-only clean stop ends the turn, except where it would leave
328 // tool results with no visible synthesis: that case, and callers that
329 // require visible output, get the bounded synthetic retry. A truly empty
330 // response is classified before this function and retried unchanged.
331 if a.requireVisibleFinal {
332 state.terminal.emptyFinalBlocks++
333 if state.terminal.emptyFinalBlocks >= maxEmptyFinalBlocks {
334 return false, fmt.Errorf("model finished without a visible final answer %d times", state.terminal.emptyFinalBlocks)
335 }
336 return a.retryEmptyFinal(ctx, reasoning, usage)
337 }
338 if state.usedAnyTool && state.terminal.emptyFinalBlocks == 0 && silentSinceLastToolRound(a.sess.conversation.Messages) {
339 state.terminal.emptyFinalBlocks++
340 return a.retryEmptyFinal(ctx, reasoning, usage)
341 }
342 }
343 a.emitTurnShadows(a.turn.turnInput)
344 if !a.closeSteerIntakeIfIdle() {
345 return true, nil
346 }
347 // A final-answer turn skips compaction, so a large context
348 // carries into the next turn un-folded and can overflow the model window.
349 // No-op below the trigger, so normal turns keep their warm cache.
350 a.contextManager().ObserveUsage(usage)
351 a.closeTurnPhase()
352 return false, nil // model gave a final answer
353 }
354
355 func (a *Agent) retryEmptyFinal(ctx context.Context, reasoning string, usage *provider.Usage) (cont bool, err error) {
356 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelInfo, Code: event.NoticeCodeEmptyFinal, Text: emptyFinalNotice(), Detail: emptyFinalNoticeDetail(a.svc.prov.Name(), usage, len(reasoning))})
357 if err := a.appendCommittedMessages(ctx, "empty-final-retry", HostGeneratedUserMessage(a.withTurnPreferences(emptyFinalRetryMessage()))); err != nil {
358 return false, err
359 }
360 a.contextManager().ObserveUsage(usage)
361 return true, nil
362 }
363
364 // silentSinceLastToolRound reports whether the transcript holds a tool result
365 // with no visible assistant text after it.
366 func silentSinceLastToolRound(messages []provider.Message) bool {
367 for _, message := range slices.Backward(messages) {
368 if message.LocalOnly {
369 continue
370 }
371 switch message.Role {
372 case provider.RoleTool:
373 return true
374 case provider.RoleAssistant:
375 if hasVisibleFinalAnswer(message.Content) {
376 return false
377 }
378 }
379 }
380 return false
381 }
382
383 // handleToolRound executes a tool batch, persists tool messages, handles
384 // cancellation, todo stall tracking, recovery finalization pause, and the
385 // max-steps grace round. cont=true continues the tool loop; cont=false returns
386 // err from Run.
387 func (a *Agent) handleToolRound(ctx context.Context, state *turnRuntime, step int, text, reasoning string, calls []provider.ToolCall, usage *provider.Usage) (cont bool, err error) {
388 state.terminal.emptyFinalBlocks = 0
389 state.usedAnyTool = true
390
391 boundaryFinalizer := a.allowsBoundaryTurnFinalizer(ctx, state, calls)
392 if boundaryErr, stop := a.stopUnexecutedBoundaryCalls(ctx, state, calls, usage); stop {
393 return false, boundaryErr
394 }
395
396 // The phase pair around the batch is what makes the accounting mean its
397 // names: it bills this round's wait to the provider and the batch to tools.
398 a.emitTurnPhase(event.TurnPhaseChecking)
399 batch := a.executeBatch(ctx, state, calls)
400 a.emitTurnPhase(event.TurnPhaseWorking)
401 if batch.err != nil {
402 // Any completed results are already stored; a failed durability barrier
403 // prevents starting the next tool.
404 return false, batch.err
405 }
406 if a.successfulTurnFinalizer(ctx, calls, batch) {
407 // submit_plan is the planner's data-bearing final answer. Its paired tool
408 // result is stored, so another acknowledgement adds no host value and can
409 // turn a valid bounded plan into a max-steps pause.
410 a.contextManager().ObserveUsage(usage)
411 a.closeTurnPhase()
412 return false, nil
413 }
414 if boundaryFinalizer {
415 // The one allowed boundary finalizer ran but was rejected or blocked.
416 // Preserve the one-grace-round contract instead of opening an unbounded
417 // loop of malformed terminal submissions.
418 a.contextManager().ObserveUsage(usage)
419 return false, a.gracePause(state)
420 }
421 // The prompt only grows from here; compact before the next turn so it
422 // stays within the model's window.
423 a.contextManager().ObserveUsage(usage)
424
425 // Spend is checked before rounds: it is the axis a runaway is actually
426 // reported in, so on the turns both would catch it should be the one named.
427 if axis, detail := a.task.budget.exceeded(a.taskBudgetLimit(ctx)); axis != "" {
428 if err := a.armFinalizationRound(ctx, state, landCause{kind: "task_budget", axis: axis, detail: detail}); err != nil {
429 return false, err
430 }
431 return true, nil
432 }
433 if state.runMaxSteps > 0 && step+1 >= state.runMaxSteps {
434 if err := a.armFinalizationRound(ctx, state, landCause{kind: "max_steps", detail: fmt.Sprintf(
435 "budget (%s=%d) exhausted: one grace round to finalize", state.runMaxStepsKey, state.runMaxSteps)}); err != nil {
436 return false, err
437 }
438 }
439 return true, nil
440 }
441
442 func (a *Agent) pairUnexecutedGraceCalls(ctx context.Context, calls []provider.ToolCall, msg string) error {
443 messages := make([]provider.Message, 0, len(calls))
444 for _, call := range calls {
445 messages = append(messages, provider.Message{Role: provider.RoleTool, Content: msg, ToolCallID: call.ID, Name: call.Name})
446 }
447 return a.appendCommittedMessages(ctx, "unexecuted-grace-tools", messages...)
448 }
449
450 func (a *Agent) publishCommittedSample(streamed streamedTurn) {
451 // Publish settlement only after the complete message is accepted by
452 // the business log. Recovery must never observe an end without a result.
453 if streamed.text != "" || streamed.displayReasoning != "" {
454 a.svc.sink.Emit(event.Event{Kind: event.Message, MessageID: streamed.messageID, AttemptID: streamed.messageID,
455 Text: DisplayAssistantText(streamed.text), Reasoning: streamed.displayReasoning})
456 }
457 if streamed.settledAttemptID != "" {
458 a.emitStreamAttempt(streamed.settledAttemptID, event.StreamAttemptCommit, streamed.settledAttempt, "", nil)
459 }
460 }
461
461 lines GO