返回 DeepSeek-Reasonix
synchronous_turn.go
根目录 / internal / control / synchronous_turn.go
1 package control
2
3 import (
4 "context"
5 "errors"
6
7 "reasonix/internal/event"
8 "reasonix/internal/session"
9 )
10
11 // runSynchronousTurn owns the blocking transport lifecycle. Durable steer
12 // acknowledgement is intentionally shared with the asynchronous completion
13 // path, while follow-up dispatch remains owned by the synchronous frontend so
14 // its response sink stays bound for every queued turn.
15 func (c *Controller) runSynchronousTurn(
16 ctx context.Context,
17 onAdmitted func() error,
18 run func(context.Context) error,
19 ) error {
20 releaseAdmission := c.trySubmissionAdmissionLock()
21 if releaseAdmission == nil {
22 return ErrTurnRunning
23 }
24 defer releaseAdmission()
25 if err := c.authentication.admissionError(); err != nil {
26 return err
27 }
28 if err := c.ensureWriteAuthorityReady(); err != nil {
29 return err
30 }
31 if ledger := c.turnEventLedger(); ledger != nil && ledger.CurrentStatus() == event.TurnRecoveryRequired {
32 return ErrRecoveryRequired
33 }
34 parent := ctx
35 c.mu.Lock()
36 // Finishing is part of the gate: TurnDone is still fanning out. Closed
37 // seals a torn-down controller. Blocking callers get an error rather than
38 // parking because they already own and enforce the request boundary.
39 if c.maintenance != nil {
40 err := ErrMaintenanceBusy
41 if c.maintenance.activity == "recovery_required" {
42 err = ErrMaintenanceRecovery
43 }
44 c.mu.Unlock()
45 return err
46 }
47 if c.bodyActiveLocked() || c.finalizingLocked() || c.rotating || c.closed || c.recoveryRequiredLocked() {
48 c.mu.Unlock()
49 return ErrTurnRunning
50 }
51 if c.rejectDrainingGenerationLocked() {
52 c.mu.Unlock()
53 c.emitDrainingNotice()
54 return ErrRuntimeDraining
55 }
56 ctx, cancel, admitted := c.startTurnLocked(ctx, queuedTurn{})
57 if !admitted {
58 c.mu.Unlock()
59 c.emitDrainingNotice()
60 return ErrRuntimeDraining
61 }
62 c.mu.Unlock()
63 if parent != nil {
64 stop := context.AfterFunc(parent, func() { c.signalTurnCancel() })
65 defer stop()
66 }
67 c.refreshRuntimeState(event.Event{})
68 finish := func() { c.finishSynchronousTurn(cancel) }
69 if onAdmitted != nil {
70 if err := onAdmitted(); err != nil {
71 finish()
72 return err
73 }
74 }
75 defer event.RecordTurnCompletion(c.sink)
76 defer func() {
77 finish()
78 c.onInboxTurnDone()
79 }()
80 // Blocking transports use the same host-owned turn boundary as interactive
81 // submissions. The agent's TurnStarted notification is deliberately a
82 // duplicate display event; it cannot create runtime ownership or a durable
83 // turn by itself.
84 run = c.prepareTurnAdmission(run)
85 releaseAdmission()
86 runErr := run(ctx)
87 c.authentication.recordFailure(runErr, c.ModelRef())
88 // Keep the execution binding through the synchronous terminal commit just
89 // like the asynchronous loop. Close may make the public controller view
90 // closed here, but it cannot release the ledger/session underneath TurnDone.
91 c.mu.Lock()
92 if c.turns.done != nil {
93 close(c.turns.done)
94 c.turns.done = nil
95 }
96 c.turns.cancel = nil
97 if c.turns.phase != session.RuntimeRecoveryRequired {
98 c.turns.phase = session.RuntimeFinalizing
99 c.turns.finishingBound.begin(true)
100 c.noteExecutionLocked(session.RuntimeFinalizing, "turn")
101 }
102 c.mu.Unlock()
103 c.refreshRuntimeState(event.Event{})
104 if ledger := c.turnEventLedger(); ledger != nil && ledger.ActiveTurnID() != "" && !ledger.CurrentStatus().Terminal() {
105 cancelled := errors.Is(ctx.Err(), context.Canceled)
106 done := event.Event{Kind: event.TurnDone, Err: runErr, Cancelled: cancelled, Outcome: turnOutcome(runErr)}
107 if cancelled {
108 done.Status = event.TurnInterrupted
109 } else if runErr != nil {
110 done.Status = event.TurnFailed
111 } else {
112 done.Status = event.TurnCompleted
113 }
114 if terminalErr := c.emitTurnEventChecked(done); terminalErr != nil {
115 runErr = errors.Join(runErr, terminalErr)
116 }
117 }
118 return runErr
119 }
120
121 func (c *Controller) finishSynchronousTurn(cancel context.CancelFunc) {
122 c.mu.Lock()
123 if c.turns.done != nil {
124 close(c.turns.done)
125 c.turns.done = nil
126 }
127 c.turns.cancel = nil
128 c.turns.cancelRequested = false
129 closing := c.closed
130 recovery := c.turns.phase == session.RuntimeRecoveryRequired
131 c.turns.finishingBound.end()
132 c.turns.finishingBound.endIdle()
133 if !recovery {
134 c.turns.lastToken = c.turns.token
135 if closing {
136 c.turns.phase = session.RuntimeClosed
137 } else {
138 c.turns.phase = session.RuntimeIdle
139 }
140 c.turns.turnID = ""
141 c.noteExecutionLocked(session.RuntimeIdle, "")
142 }
143 c.mu.Unlock()
144 if closing {
145 c.finalizeControllerClose()
146 }
147 c.refreshRuntimeState(event.Event{})
148 c.kickGoalDriver()
149 cancel()
150 }
151
151 lines GO