返回 DeepSeek-Reasonix
turn_loop.go
根目录 / internal / control / turn_loop.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "time"
8
9 "reasonix/internal/agent"
10 "reasonix/internal/event"
11 "reasonix/internal/extension"
12 "reasonix/internal/provider"
13 "reasonix/internal/session"
14 )
15
16 type queuedTurnKind int
17
18 const (
19 queuedUser queuedTurnKind = iota
20 queuedGoal
21 )
22
23 type queuedTurn struct {
24 kind queuedTurnKind
25 body func(ctx context.Context) error
26 onStart func()
27 goalRound *goalRoundReservation
28 // admissionCtx is only used for the synchronous durability boundary. The
29 // model turn has its own controller-owned context and survives RPC return.
30 admissionCtx context.Context
31 }
32
33 // turnLoop is the session-scoped execution authority. Controller.mu guards it.
34 type turnLoop struct {
35 phase session.RuntimePhase
36 cancel context.CancelFunc
37 done chan struct{}
38 recoveryFanout bool // watchdog terminal publication still owns the stores
39 turnID string
40 token uint64
41 lastToken uint64
42 pending []queuedTurn
43 wake bool
44 generation uint64
45 runtime *session.Runtime
46 cancelRequested bool
47 finishingBound turnFinishingBoundary
48 }
49
50 type controllerExecution struct {
51 c *Controller
52 }
53
54 func (e controllerExecution) Snapshot() session.RuntimeSnapshot {
55 if e.c == nil {
56 return session.RuntimeSnapshot{Phase: session.RuntimeIdle}
57 }
58 e.c.mu.Lock()
59 defer e.c.mu.Unlock()
60 if op := e.c.maintenance; op != nil {
61 return session.RuntimeSnapshot{Phase: maintenanceRuntimePhase(op.activity), Activity: session.MaintenanceActivity}
62 }
63 return session.RuntimeSnapshot{Phase: e.c.turns.phase, Activity: e.c.turns.activityNameLocked()}
64 }
65
66 func (e controllerExecution) Cancel() bool {
67 if e.c == nil {
68 return false
69 }
70 if _, present, cancelled := e.c.signalMaintenanceCancel(); present {
71 return cancelled
72 }
73 return e.c.signalTurnCancel()
74 }
75
76 func (t *turnLoop) activityNameLocked() string {
77 switch t.phase {
78 case session.RuntimeCancelling:
79 return "cancelling"
80 case session.RuntimeRecoveryRequired:
81 return "recovery_required"
82 case session.RuntimeRunning, session.RuntimeFinalizing:
83 if t.cancelRequested {
84 return "cancelling"
85 }
86 return "turn"
87 default:
88 return ""
89 }
90 }
91
92 func (c *Controller) bodyActiveLocked() bool {
93 switch c.turns.phase {
94 case session.RuntimeRunning, session.RuntimeCancelling:
95 return true
96 default:
97 return false
98 }
99 }
100
101 func (c *Controller) finalizingLocked() bool {
102 return c.turns.phase == session.RuntimeFinalizing
103 }
104
105 func (c *Controller) cancelRequestedLocked() bool {
106 if c.closed {
107 return false
108 }
109 return c.turns.cancelRequested || c.turns.phase == session.RuntimeCancelling
110 }
111
112 func (c *Controller) recoveryRequiredLocked() bool {
113 return c.turns.phase == session.RuntimeRecoveryRequired
114 }
115
116 func (c *Controller) bindExecutionControl() {
117 _, runtime, exclusive := c.v3Binding()
118 if !exclusive || runtime == nil {
119 return
120 }
121 snap := runtime.StateSnapshot()
122 gen := runtime.BindExecution(controllerExecution{c: c})
123 c.executionGeneration.Store(gen)
124 c.mu.Lock()
125 c.turns.generation = gen
126 c.turns.runtime = runtime
127 if snap.Phase == session.RuntimeRecoveryRequired {
128 c.turns.phase = session.RuntimeRecoveryRequired
129 }
130 c.mu.Unlock()
131 }
132
133 // ExecutionGeneration returns the session-runtime execution generation owned
134 // by this controller. Zero means the controller is a prepared replacement that
135 // has not been published as the execution owner.
136 func (c *Controller) ExecutionGeneration() uint64 {
137 if c == nil {
138 return 0
139 }
140 return c.executionGeneration.Load()
141 }
142
143 // ActivateSessionExecution publishes this controller as the exact execution
144 // owner. expectedGeneration is zero for a previously unbound runtime and the
145 // outgoing controller generation for a fail-atomic replacement. The method is
146 // deliberately callback- and I/O-free so hosts may invoke it in their final
147 // pointer-swap critical section.
148 func (c *Controller) ActivateSessionExecution(expectedGeneration uint64) error {
149 if c == nil {
150 return session.ErrSessionNotRunning
151 }
152 c.turnEvents.commitMu.Lock()
153 defer c.turnEvents.commitMu.Unlock()
154 c.mu.Lock()
155 defer c.mu.Unlock()
156 runtime := c.turns.runtime
157 if runtime == nil {
158 return nil
159 }
160 if generation := c.turns.generation; generation != 0 && runtime.OwnsExecution(generation) {
161 return nil
162 }
163 var generation uint64
164 if expectedGeneration == 0 {
165 generation = runtime.BindExecution(controllerExecution{c: c})
166 } else if pending := c.turnEvents.pendingExecutionCommit; pending != nil {
167 c.turnEvents.pendingExecutionCommit = nil
168 var err error
169 generation, _, err = runtime.ReplaceExecutionAndCommit(expectedGeneration, controllerExecution{c: c}, *pending)
170 if err != nil {
171 return err
172 }
173 } else {
174 generation = runtime.ReplaceExecution(expectedGeneration, controllerExecution{c: c})
175 }
176 if generation == 0 {
177 return session.ErrRuntimeBusy
178 }
179 c.turns.generation = generation
180 c.executionGeneration.Store(generation)
181 return nil
182 }
183
184 // ActivateControllerReplacement transfers session execution ownership when a
185 // host commits a controller pointer swap. Controllers without a shared Runtime
186 // need no additional activation.
187 func ActivateControllerReplacement(old, next *Controller) error {
188 if next == nil {
189 return session.ErrSessionNotRunning
190 }
191 if old != nil && old.background.scope != nil && old.background.scope == next.background.scope {
192 if old.workspaceRoot != next.workspaceRoot {
193 return fmt.Errorf("background ownership cannot cross workspaces")
194 }
195 before, beforeOK := old.SessionRef()
196 after, afterOK := next.SessionRef()
197 if beforeOK && afterOK && before.SessionID != after.SessionID {
198 return fmt.Errorf("background ownership cannot cross sessions")
199 }
200 if ModelReplacementBlocked(old) {
201 return fmt.Errorf("session changed while preparing replacement")
202 }
203 }
204 if old != nil && old.background.scope != nil && old.background.scope == next.background.scope && old.ToolApprovalMode() != next.ToolApprovalMode() && len(old.jobs.Running()) > 0 {
205 return fmt.Errorf("permission changes require stopping background jobs before replacement")
206 }
207 _, nextRuntime, nextExclusive := next.v3Binding()
208 if !nextExclusive || nextRuntime == nil {
209 retireBackgroundCallbacks(old, next)
210 next.PublishBackgroundScope()
211 return nil
212 }
213 expected := uint64(0)
214 if old != nil {
215 _, oldRuntime, oldExclusive := old.v3Binding()
216 if oldExclusive && oldRuntime == nextRuntime {
217 expected = old.ExecutionGeneration()
218 }
219 }
220 if err := next.ActivateSessionExecution(expected); err != nil {
221 return err
222 }
223 retireBackgroundCallbacks(old, next)
224 next.PublishBackgroundScope()
225 return nil
226 }
227
228 func retireBackgroundCallbacks(old, next *Controller) {
229 if old != nil && old != next && old.background.scope != nil && old.background.scope == next.background.scope {
230 old.mu.Lock()
231 old.background.retired = true
232 old.mu.Unlock()
233 }
234 }
235
236 // ActivateSessionAPIReplacement is the host-facing form used at a final
237 // controller pointer swap. Non-Controller implementations have no exclusive
238 // Runtime ownership to transfer and are left unchanged.
239 func ActivateSessionAPIReplacement(old, next SessionAPI) error {
240 concreteNext, ok := next.(*Controller)
241 if !ok || concreteNext == nil {
242 return nil
243 }
244 concreteOld, _ := old.(*Controller)
245 if err := ActivateControllerReplacement(concreteOld, concreteNext); err != nil {
246 return err
247 }
248 if concreteOld != nil && concreteOld.attachmentScope() != "" && concreteOld.attachmentScope() == concreteNext.attachmentScope() && concreteOld.workspaceRoot == concreteNext.workspaceRoot {
249 concreteOld.attachmentService().Drafts().CopyScopeTo(concreteOld.attachmentScope(), concreteNext.attachmentService().Drafts())
250 }
251 return nil
252 }
253
254 func (c *Controller) unbindExecutionControl(runtime *session.Runtime) {
255 if runtime == nil {
256 return
257 }
258 c.mu.Lock()
259 gen := c.turns.generation
260 c.mu.Unlock()
261 runtime.UnbindExecution(gen)
262 }
263
264 func (c *Controller) noteExecutionLocked(phase session.RuntimePhase, activity string) {
265 if c.turns.runtime == nil || c.turns.generation == 0 {
266 return
267 }
268 c.turns.runtime.NoteExecution(c.turns.generation, phase, activity)
269 }
270
271 func (c *Controller) currentTurnToken() (token uint64, turnID string, active bool) {
272 c.mu.Lock()
273 defer c.mu.Unlock()
274 switch c.turns.phase {
275 case session.RuntimeIdle, session.RuntimeClosed:
276 return c.turns.lastToken, "", false
277 default:
278 return c.turns.token, c.turns.turnID, true
279 }
280 }
281
282 func (c *Controller) discardLateTurnEvent(e event.Event) bool {
283 if e.TurnID == "" || !lateBusinessEvent(e.Kind) {
284 return false
285 }
286 _, turnID, active := c.currentTurnToken()
287 if active {
288 return e.TurnID != turnID
289 }
290 return true
291 }
292
293 func (c *Controller) startTurnLocked(parent context.Context, next queuedTurn) (ctx context.Context, cancel context.CancelFunc, admitted bool) {
294 if parent == nil {
295 parent = context.Background()
296 }
297 if c.turns.runtime != nil && !c.turns.runtime.BeginExecution(c.turns.generation, "turn") {
298 return nil, nil, false
299 }
300 ctx, cancel = context.WithCancel(extension.ContextWithRuntimeOwner(c.withAuthentication(parent), c.runtimeOwner))
301 c.turns.cancel = cancel
302 c.turns.done = make(chan struct{})
303 c.turns.finishingBound.beginIdle()
304 c.turns.phase = session.RuntimeRunning
305 c.turns.cancelRequested = false
306 c.turns.token++
307 ctx = context.WithValue(ctx, executionTokenKey{}, c.turns.token)
308 c.turns.turnID = ""
309 return ctx, cancel, true
310 }
311
312 func (c *Controller) popNextPendingLocked() (queuedTurn, bool) {
313 if len(c.turns.pending) == 0 {
314 c.turns.wake = false
315 return queuedTurn{}, false
316 }
317 next := c.turns.pending[0]
318 c.turns.pending = c.turns.pending[1:]
319 c.turns.wake = len(c.turns.pending) > 0
320 return next, true
321 }
322
323 func (c *Controller) queueTurnLocked(item queuedTurn) {
324 // Harness wakeRequested: input that cannot join the current activity is
325 // claimed once when the body converges. Close clears this queue so a
326 // disposed session never starts a latched turn.
327 c.turns.pending = append(c.turns.pending, item)
328 c.turns.wake = true
329 }
330
331 func (c *Controller) signalTurnCancel() bool {
332 _, _, cancelled := c.signalTurnCancelIdentity()
333 return cancelled
334 }
335
336 func (c *Controller) signalTurnCancelIdentity() (uint64, string, bool) {
337 c.mu.Lock()
338 token, turnID := c.turns.token, c.turns.turnID
339 cancel := c.turns.cancel
340 first := cancel != nil && c.turns.phase == session.RuntimeRunning
341 if cancel != nil && (c.turns.phase == session.RuntimeRunning || c.turns.phase == session.RuntimeCancelling) {
342 c.turns.phase = session.RuntimeCancelling
343 c.turns.cancelRequested = true
344 }
345 done := c.turns.done
346 c.mu.Unlock()
347 if cancel == nil {
348 return token, turnID, false
349 }
350 cancel()
351 c.recordLifecycle("cancel_signalled", "turn", turnID, 0, "")
352 if first {
353 c.startCancellationWatchdog(done)
354 }
355 return token, turnID, true
356 }
357
358 func (c *Controller) enterRecoveryLocked(reason string) {
359 c.turns.phase = session.RuntimeRecoveryRequired
360 c.noteExecutionLocked(session.RuntimeRecoveryRequired, reason)
361 if c.turns.runtime != nil {
362 c.turns.runtime.RequireRecovery(reason)
363 }
364 }
365
366 func (c *Controller) spawnGuardedTurn(ctx context.Context, cancel context.CancelFunc, item queuedTurn) {
367 ctx, completion := withGuardedTurnCompletion(ctx)
368 admissionCtx := item.admissionCtx
369 if admissionCtx == nil {
370 admissionCtx = context.Background()
371 }
372 body := c.prepareTurnAdmissionWithGoalRound(admissionCtx, item.body, item.goalRound)
373 if ledger := c.turnEventLedger(); ledger != nil {
374 c.mu.Lock()
375 c.turns.turnID = ledger.ActiveTurnID()
376 c.mu.Unlock()
377 }
378 c.liveness.reset(time.Now())
379 c.autosaveWG.Go(func() {
380 c.autosaveWhileRunning(ctx)
381 })
382 go func() {
383 defer cancel()
384 defer func() {
385 c.finishGoalRoundActivity(item.goalRound)
386 c.kickGoalDriver()
387 }()
388 defer func() {
389 if r := recover(); r != nil {
390 err := fmt.Errorf("internal error: %v", r)
391 item.goalRound.setResult(err, false)
392 c.finishGuardedTurn(err, completion)
393 }
394 }()
395 err := body(ctx)
396 if item.goalRound != nil {
397 item.goalRound.setResult(err, errors.Is(ctx.Err(), context.Canceled) && c.CancelRequested())
398 }
399 c.finishGuardedTurn(explainError(err), completion)
400 }()
401 }
402
403 func (c *Controller) cancellationGrace() time.Duration {
404 if c != nil && c.testCancelGrace > 0 {
405 return c.testCancelGrace
406 }
407 return 15 * time.Second
408 }
409
410 func (c *Controller) finishGuardedTurn(err error, completion *guardedTurnCompletion) {
411 c.authentication.recordFailure(err, c.ModelRef())
412 c.memory.clearAutoRemember()
413 c.mu.Lock()
414 cancelRequested := c.turns.cancelRequested
415 if c.turns.phase == session.RuntimeRecoveryRequired {
416 c.turns.cancel = nil
417 closing := c.closed
418 c.mu.Unlock()
419 // The cancellation watchdog already committed the recovery terminal.
420 // A closing controller must not emit another terminal after its ledger
421 // and session binding have been finalized.
422 if !closing {
423 c.emitTurnDoneEvent(err, cancelRequested, completion)
424 }
425 c.mu.Lock()
426 // Keep the owned turn live through terminal fanout. Close must not
427 // release stores while that fanout can still publish durable events.
428 if c.turns.done != nil {
429 close(c.turns.done)
430 c.turns.done = nil
431 }
432 closing = c.closed && !c.turns.recoveryFanout
433 c.turns.finishingBound.endIdle()
434 c.mu.Unlock()
435 if closing {
436 c.finalizeControllerClose()
437 }
438 c.refreshRuntimeState(event.Event{})
439 return
440 }
441 if c.turns.done != nil {
442 close(c.turns.done)
443 c.turns.done = nil
444 }
445 c.turns.phase = session.RuntimeFinalizing
446 c.turns.finishingBound.begin(true)
447 c.turns.cancel = nil
448 c.noteExecutionLocked(session.RuntimeFinalizing, "turn")
449 c.mu.Unlock()
450
451 c.refreshRuntimeState(event.Event{})
452 defer func() {
453 c.mu.Lock()
454 c.turns.finishingBound.end()
455 c.turns.cancelRequested = false
456 if c.turns.phase == session.RuntimeRecoveryRequired {
457 closing := c.closed
458 c.turns.finishingBound.endIdle()
459 c.mu.Unlock()
460 if closing {
461 c.finalizeControllerClose()
462 }
463 c.refreshRuntimeState(event.Event{})
464 return
465 }
466 if ledger := c.turnEventLedger(); ledger != nil && ledger.CurrentStatus() == event.TurnRecoveryRequired {
467 c.enterRecoveryLocked("terminal")
468 c.turns.finishingBound.endIdle()
469 c.mu.Unlock()
470 c.refreshRuntimeState(event.Event{})
471 return
472 }
473 if c.closed {
474 c.turns.lastToken = c.turns.token
475 c.turns.phase = session.RuntimeClosed
476 c.turns.turnID = ""
477 c.noteExecutionLocked(session.RuntimeIdle, "")
478 c.turns.finishingBound.endIdle()
479 c.mu.Unlock()
480 c.finalizeControllerClose()
481 c.refreshRuntimeState(event.Event{})
482 return
483 }
484 if authErr := c.authentication.admissionError(); authErr != nil {
485 for _, pending := range c.turns.pending {
486 if pending.goalRound != nil {
487 pending.goalRound.setResult(authErr, false)
488 }
489 }
490 c.turns.pending = nil
491 c.turns.wake = false
492 c.turns.lastToken = c.turns.token
493 c.turns.phase = session.RuntimeIdle
494 c.turns.turnID = ""
495 c.noteExecutionLocked(session.RuntimeIdle, "")
496 c.turns.finishingBound.endIdle()
497 c.mu.Unlock()
498 c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Code: "authentication_not_ready", Text: authErr.Error()})
499 c.refreshRuntimeState(event.Event{})
500 return
501 }
502 next, ok := c.popNextPendingLocked()
503 if !ok {
504 c.turns.lastToken = c.turns.token
505 c.turns.phase = session.RuntimeIdle
506 c.turns.turnID = ""
507 c.noteExecutionLocked(session.RuntimeIdle, "")
508 c.turns.finishingBound.endIdle()
509 c.mu.Unlock()
510 c.maybeDispatchInbox()
511 c.refreshRuntimeState(event.Event{})
512 return
513 }
514 ctx, cancel, admitted := c.startTurnLocked(context.Background(), next)
515 if !admitted {
516 c.enterRecoveryLocked("execution_owner_lost")
517 c.turns.finishingBound.endIdle()
518 c.mu.Unlock()
519 c.refreshRuntimeState(event.Event{})
520 return
521 }
522 c.mu.Unlock()
523 if next.onStart != nil {
524 next.onStart()
525 }
526 c.spawnGuardedTurn(ctx, cancel, next)
527 c.refreshRuntimeState(event.Event{})
528 }()
529 c.emitTurnDoneEvent(err, cancelRequested, completion)
530 }
531
532 func (c *Controller) emitTurnDoneEvent(err error, cancelRequested bool, completion *guardedTurnCompletion) {
533 c.inbox.mu.Lock()
534 activeInboxID := ""
535 for id := range c.inbox.activeItemIDs {
536 activeInboxID = id
537 break
538 }
539 c.inbox.mu.Unlock()
540 done := event.Event{
541 Kind: event.TurnDone,
542 Err: err,
543 Cancelled: cancelRequested,
544 Outcome: turnOutcome(err),
545 CheckpointTurn: c.validatedCheckpointTurn(completion),
546 Receipt: c.executor.CompletionReceipt(),
547 ItemID: activeInboxID,
548 }
549 if done.CheckpointTurn != nil {
550 changes := completion.checkpoint.store.FreezeTurnChanges(*done.CheckpointTurn)
551 if done.Receipt == nil && (len(changes.Files) > 0 || len(changes.Reasons) > 0) {
552 done.Receipt = &event.CompletionReceipt{AssessmentKind: "facts", Verdict: "unknown"}
553 }
554 if done.Receipt != nil {
555 receipt := *done.Receipt
556 receipt.Diff = changes.Summary()
557 receipt.Interrupted = cancelRequested
558 done.Receipt = &receipt
559 }
560 }
561 done.Receipt = bindCompletionLogSources(done.Receipt, c.History())
562 done = c.applyTurnDoneProtocol(done, cancelRequested)
563 c.applyToolRecoveryTurnStatus(&done, completion)
564 var readErr *agent.IncompleteReadError
565 if errors.As(err, &readErr) {
566 done.ReadPause = readErr.Pause
567 }
568 done.Diagnostic = provider.DiagnoseFailure(err)
569 done.Detail = provider.FailureDiagnosticDetail(done.Diagnostic)
570 if !cancelRequested {
571 done.ProtocolRecovery = c.executor.PendingProtocolRecovery()
572 }
573 var readinessErr *agent.FinalReadinessError
574 if errors.As(err, &readinessErr) {
575 done.Readiness = &event.FinalReadiness{Attempts: readinessErr.Attempts, Missing: append([]string(nil), readinessErr.Missing...)}
576 }
577 c.onInboxTurnDone()
578 c.sink.Emit(done)
579 }
580
581 func (c *Controller) startCancellationWatchdog(done chan struct{}) {
582 if c == nil || done == nil {
583 return
584 }
585 go func() {
586 timer := time.NewTimer(c.cancellationGrace())
587 defer timer.Stop()
588 select {
589 case <-timer.C:
590 case <-done:
591 return
592 }
593
594 c.mu.Lock()
595 stillRunning := c.turns.done == done && (c.turns.phase == session.RuntimeRunning || c.turns.phase == session.RuntimeCancelling)
596 turnID := c.turns.turnID
597 if stillRunning {
598 c.turns.recoveryFanout = true
599 c.enterRecoveryLocked("cancellation_grace_expired")
600 }
601 c.mu.Unlock()
602 if !stillRunning {
603 return
604 }
605 if turnID == "" {
606 if ledger := c.turnEventLedger(); ledger != nil {
607 turnID = ledger.ActiveTurnID()
608 }
609 }
610 recovery := &event.RecoveryStatus{
611 State: "recovery_required",
612 Phase: "cancellation_grace_expired",
613 Reason: "cancellation_grace_expired",
614 RequiresUserDecision: true,
615 }
616 _ = c.emitTurnEventChecked(event.Event{
617 Kind: event.TurnDone,
618 TurnID: turnID,
619 Status: event.TurnRecoveryRequired,
620 Cancelled: true,
621 Outcome: "unknown",
622 Recovery: recovery,
623 })
624 c.mu.Lock()
625 c.turns.recoveryFanout = false
626 closing := c.closed && c.turns.done == nil && !c.finalizingLocked()
627 c.mu.Unlock()
628 if closing {
629 c.finalizeControllerClose()
630 }
631 c.refreshRuntimeState(event.Event{})
632 }()
633 }
634
634 lines GO