返回 DeepSeek-Reasonix
execution.go
根目录 / internal / session / execution.go
1 package session
2
3 import (
4 "context"
5 )
6
7 // ExecutionControl is the session-scoped turn-loop bound to a Runtime.
8 // The loop is the execution authority; Runtime only forwards Cancel and
9 // caches the last NoteExecution from the bound generation.
10 type ExecutionControl interface {
11 Snapshot() RuntimeSnapshot
12 Cancel() bool
13 }
14
15 const MaintenanceActivity = "maintenance"
16
17 // FinishMaintenanceExecution releases a settled maintenance worker into the
18 // finalizing handoff barrier. Only its exact owner may clear a cancellation
19 // timeout; durable business recovery and failed persistence remain fenced.
20 func (r *Runtime) FinishMaintenanceExecution(generation uint64) bool {
21 if r == nil || generation == 0 {
22 return false
23 }
24 r.mu.Lock()
25 defer r.mu.Unlock()
26 owner := r.execution.Load()
27 if owner == nil || owner.generation != generation || r.activity != MaintenanceActivity || r.phase == RuntimeClosed {
28 return false
29 }
30 state := r.session.StateSnapshot()
31 if state.PersistenceStatus == PersistenceFailed || state.PersistenceStatus == PersistenceUncertain ||
32 state.Projection.Recovery != nil && state.Projection.Recovery.State == "recovery_required" {
33 return false
34 }
35 r.phase = RuntimeFinalizing
36 r.canceling.Store(false)
37 r.revision.Add(1)
38 return true
39 }
40
41 type executionBinding struct {
42 generation uint64
43 control ExecutionControl
44 }
45
46 // BindExecution installs the first generation-scoped turn-loop. Replacing an
47 // existing owner is always explicit through ReplaceExecution so constructing a
48 // candidate controller cannot steal Stop or mutation authority.
49 func (r *Runtime) BindExecution(control ExecutionControl) uint64 {
50 if r == nil || control == nil {
51 return 0
52 }
53 r.mu.Lock()
54 defer r.mu.Unlock()
55 if r.phase == RuntimeClosed || r.phase == RuntimeRecoveryRequired || r.execution.Load() != nil {
56 return 0
57 }
58 gen := r.bindGen.Add(1)
59 r.execution.Store(&executionBinding{generation: gen, control: control})
60 r.revision.Add(1)
61 return gen
62 }
63
64 // ReplaceExecution transfers an idle runtime from the exact expected
65 // generation to a replacement turn-loop. The lock acquisition is the
66 // linearization point: generations are allocated only after ownership and
67 // phase validation, so they can never be published out of order.
68 func (r *Runtime) ReplaceExecution(expectedGeneration uint64, control ExecutionControl) uint64 {
69 if r == nil || expectedGeneration == 0 || control == nil {
70 return 0
71 }
72 r.mu.Lock()
73 defer r.mu.Unlock()
74 cur := r.execution.Load()
75 if cur == nil || cur.generation != expectedGeneration || r.phase != RuntimeIdle {
76 return 0
77 }
78 gen := r.bindGen.Add(1)
79 r.execution.Store(&executionBinding{generation: gen, control: control})
80 r.revision.Add(1)
81 return gen
82 }
83
84 // ReplaceExecutionAndCommit atomically accepts a candidate's prevalidated
85 // session mutation and transfers an idle Runtime to that candidate. A failed
86 // commit leaves the outgoing execution owner untouched.
87 func (r *Runtime) ReplaceExecutionAndCommit(expectedGeneration uint64, control ExecutionControl, prepared PreparedBatch) (uint64, Commit, error) {
88 if r == nil || expectedGeneration == 0 || control == nil {
89 prepared.Release()
90 return 0, Commit{}, ErrStaleExecution
91 }
92 r.mu.Lock()
93 defer r.mu.Unlock()
94 cur := r.execution.Load()
95 if cur == nil || cur.generation != expectedGeneration || r.phase != RuntimeIdle {
96 prepared.Release()
97 return 0, Commit{}, ErrStaleExecution
98 }
99 commit, err := r.session.CommitPrepared(prepared)
100 if err != nil {
101 return 0, Commit{}, err
102 }
103 gen := r.bindGen.Add(1)
104 r.execution.Store(&executionBinding{generation: gen, control: control})
105 r.revision.Add(1)
106 return gen, commit, nil
107 }
108
109 // OwnsExecution reports whether generation is the currently bound turn-loop.
110 // It is a snapshot only; mutations that require fencing must validate again
111 // while holding r.mu.
112 func (r *Runtime) OwnsExecution(generation uint64) bool {
113 if r == nil || generation == 0 {
114 return false
115 }
116 cur := r.execution.Load()
117 return cur != nil && cur.generation == generation
118 }
119
120 // BeginExecution atomically admits a turn for the exact execution owner. A
121 // queued turn may move directly from finalizing to running; no observer sees a
122 // false idle boundary between turns.
123 func (r *Runtime) BeginExecution(generation uint64, activity string) bool {
124 if r == nil || generation == 0 {
125 return false
126 }
127 r.mu.Lock()
128 defer r.mu.Unlock()
129 cur := r.execution.Load()
130 if cur == nil || cur.generation != generation {
131 return false
132 }
133 if r.phase != RuntimeIdle && r.phase != RuntimeFinalizing {
134 return false
135 }
136 r.phase = RuntimeRunning
137 r.activity = activity
138 r.canceling.Store(false)
139 r.revision.Add(1)
140 return true
141 }
142
143 // CommitPreparedForExecution accepts a prepared batch only while generation
144 // is the exact execution owner. Preparation remains outside the Runtime lock;
145 // the short ownership check and in-memory Session acceptance form one commit
146 // boundary so a controller cutover cannot race a stale writer into the log.
147 func (r *Runtime) CommitPreparedForExecution(generation uint64, prepared PreparedBatch) (Commit, error) {
148 if r == nil || generation == 0 {
149 prepared.Release()
150 return Commit{}, ErrStaleExecution
151 }
152 r.mu.Lock()
153 defer r.mu.Unlock()
154 cur := r.execution.Load()
155 if cur == nil || cur.generation != generation {
156 prepared.Release()
157 return Commit{}, ErrStaleExecution
158 }
159 return r.session.CommitPrepared(prepared)
160 }
161
162 // CommitPreparedForTurn makes the owner/turn check and terminal acceptance one
163 // critical section. A prepared completion cannot close a successor's turn.
164 func (r *Runtime) CommitPreparedForTurn(generation uint64, turnID string, prepared PreparedBatch) (Commit, error) {
165 if r == nil || generation == 0 {
166 prepared.Release()
167 return Commit{}, ErrStaleExecution
168 }
169 r.mu.Lock()
170 defer r.mu.Unlock()
171 cur := r.execution.Load()
172 if cur == nil || cur.generation != generation || r.session.StateSnapshot().Projection.TurnID != turnID {
173 prepared.Release()
174 return Commit{}, ErrStaleExecution
175 }
176 return r.session.CommitPrepared(prepared)
177 }
178
179 // UnbindExecution releases only the exact generation. A superseded controller
180 // cannot clear the replacement's control binding.
181 func (r *Runtime) UnbindExecution(generation uint64) {
182 if r == nil || generation == 0 {
183 return
184 }
185 r.mu.Lock()
186 defer r.mu.Unlock()
187 cur := r.execution.Load()
188 if cur == nil || cur.generation != generation || (r.phase.busy() && r.phase != RuntimeRecoveryRequired) {
189 return
190 }
191 r.execution.Store(nil)
192 r.revision.Add(1)
193 }
194
195 // NoteExecution records a phase transition from the exact bound generation.
196 // Validation and mutation share r.mu so a replacement cannot land between
197 // them and let an older controller alter the new owner's phase.
198 func (r *Runtime) NoteExecution(generation uint64, phase RuntimePhase, activity string) {
199 if r == nil || generation == 0 {
200 return
201 }
202 r.mu.Lock()
203 cur := r.execution.Load()
204 if cur == nil || cur.generation != generation {
205 r.mu.Unlock()
206 return
207 }
208 if r.phase == RuntimeClosed {
209 r.mu.Unlock()
210 return
211 }
212 if r.phase == RuntimeRecoveryRequired && phase != RuntimeRecoveryRequired && phase != RuntimeClosed {
213 r.mu.Unlock()
214 return
215 }
216 if r.phase == phase && r.activity == activity {
217 r.mu.Unlock()
218 return
219 }
220 r.phase = phase
221 r.activity = activity
222 r.revision.Add(1)
223 if phase == RuntimeCancelling || phase == RuntimeRecoveryRequired {
224 r.canceling.Store(true)
225 } else {
226 r.canceling.Store(false)
227 }
228 owner := r.owner
229 r.mu.Unlock()
230 if phase == RuntimeIdle && owner != nil {
231 _ = owner.closeIfUnbound(context.Background(), r)
232 }
233 }
234
235 func (r *Runtime) loadExecution() *executionBinding {
236 if r == nil {
237 return nil
238 }
239 return r.execution.Load()
240 }
241
242 // Busy reports whether a phase still owns live execution or finalization
243 // work. Hosts use this shared definition so finalizing sessions cannot vanish
244 // from running lists before their terminal commit completes.
245 func (p RuntimePhase) Busy() bool {
246 switch p {
247 case RuntimeRunning, RuntimeCancelling, RuntimeFinalizing, RuntimeRecoveryRequired:
248 return true
249 default:
250 return false
251 }
252 }
253
254 func (p RuntimePhase) busy() bool { return p.Busy() }
255
256 func (r *Runtime) executionBusy() bool {
257 if r == nil {
258 return false
259 }
260 r.mu.Lock()
261 phase := r.phase
262 bound := r.execution.Load() != nil
263 r.mu.Unlock()
264 return phase.busy() && bound
265 }
266
266 lines GO