返回 DeepSeek-Reasonix
runtime_state.go
根目录 / internal / control / runtime_state.go
1 package control
2
3 import (
4 "crypto/rand"
5 "encoding/hex"
6 "log/slog"
7 "reflect"
8 "sync"
9 "sync/atomic"
10
11 "reasonix/internal/agent"
12 "reasonix/internal/event"
13 "reasonix/internal/jobs"
14 "reasonix/internal/session"
15 "reasonix/internal/turnevent"
16 )
17
18 // RuntimeStateReader is optional so older embedders of SessionAPI keep working.
19 type RuntimeStateReader interface {
20 RuntimeStateSnapshot() event.RuntimeStateSnapshot
21 }
22
23 type controllerRuntimeState struct {
24 mu sync.Mutex // serializes sampling, commit and publication order; never held by observers
25 published atomic.Pointer[event.RuntimeStateSnapshot]
26 snapshot event.RuntimeStateSnapshot
27 ledger *turnevent.Ledger
28 path string
29 activity string
30 sink event.Sink
31 pending *event.RuntimeStateSnapshot
32 draining bool
33 jobUnsubscribe func()
34 }
35
36 func newRuntimeStateEpoch() string {
37 var bytes [16]byte
38 if _, err := rand.Read(bytes[:]); err != nil {
39 panic(err)
40 }
41 return hex.EncodeToString(bytes[:])
42 }
43
44 // RuntimeStateSnapshot first commits a fresh projection from the current owners
45 // and then returns that immutable boundary. This closes the small callback lag
46 // after a background job starts or exits without making readers combine fields
47 // from separate snapshots.
48 func (c *Controller) RuntimeStateSnapshot() event.RuntimeStateSnapshot {
49 if c == nil {
50 return event.RuntimeStateSnapshot{Todos: []event.Todo{}, Interactions: []event.PendingInteraction{}}
51 }
52 c.refreshRuntimeState(event.Event{})
53 c.runtimeState.mu.Lock()
54 defer c.runtimeState.mu.Unlock()
55 return cloneRuntimeState(c.runtimeState.snapshot)
56 }
57
58 func cloneRuntimeState(in event.RuntimeStateSnapshot) event.RuntimeStateSnapshot {
59 out := in
60 out.Todos = append([]event.Todo{}, in.Todos...)
61 out.Interactions = append([]event.PendingInteraction{}, in.Interactions...)
62 if in.Recovery != nil {
63 recovery := *in.Recovery
64 out.Recovery = &recovery
65 }
66 if in.Maintenance != nil {
67 maintenance := *in.Maintenance
68 out.Maintenance = &maintenance
69 }
70 if in.Goal != nil {
71 goal := *in.Goal
72 if in.Goal.MaxGoalRounds != nil {
73 limit := *in.Goal.MaxGoalRounds
74 goal.MaxGoalRounds = &limit
75 }
76 if in.Goal.BlockedReason != nil {
77 reason := *in.Goal.BlockedReason
78 goal.BlockedReason = &reason
79 }
80 out.Goal = &goal
81 }
82 return out
83 }
84
85 func (c *Controller) initializeRuntimeState() {
86 c.runtimeState.mu.Lock()
87 c.runtimeState.sink = c.sink
88 c.runtimeState.mu.Unlock()
89 c.refreshRuntimeState(event.Event{})
90 if c.jobs != nil {
91 // A manager may be shared across a controller rebuild. Subscribe to all
92 // session transitions and filter against the current committed binding.
93 _, stop := c.jobs.SubscribeRuntime("", func(state jobs.RuntimeState) {
94 if c.receivesBackgroundRuntimeEvents() {
95 c.refreshRuntimeState(event.Event{})
96 }
97 })
98 c.runtimeState.mu.Lock()
99 c.runtimeState.jobUnsubscribe = stop
100 c.runtimeState.mu.Unlock()
101 c.refreshRuntimeState(event.Event{})
102 }
103 }
104
105 // refreshRuntimeState is a commit boundary, not a read-side workaround. Every
106 // lifecycle and job boundary calls it after releasing its owning locks. A
107 // single sampler re-reads current owners instead of replaying stale booleans.
108 func (c *Controller) refreshRuntimeState(e event.Event) {
109 c.refreshRuntimeStateAttempt(e, 0)
110 }
111
112 // refreshRuntimeStateAttempt retries an unstable multi-owner sample at most
113 // three times. Lifecycle callbacks will sample again after later transitions;
114 // a perpetually changing runtime must not grow the stack or monopolize a caller.
115 func (c *Controller) refreshRuntimeStateAttempt(e event.Event, attempt int) {
116 if c == nil {
117 return
118 }
119 r := &c.runtimeState
120 r.mu.Lock()
121 if r.sink == nil {
122 r.mu.Unlock()
123 return
124 } // construction has not finished
125 c.mu.Lock()
126 running, finishing, closed, cancelling, path := c.bodyActiveLocked(), c.finalizingLocked(), c.closed, c.cancelRequestedLocked(), c.sessionPath
127 maintenance := c.maintenanceSnapshotLocked()
128 c.mu.Unlock()
129 _, v3Runtime, exclusiveSession := c.v3Binding()
130 var v3RuntimeSnapshot session.RuntimeSnapshot
131 if exclusiveSession && v3Runtime != nil {
132 v3RuntimeSnapshot = v3Runtime.StateSnapshot()
133 }
134 ledger := c.turnEventLedger()
135 initialized := r.snapshot.SchemaVersion == 1
136 base, activity := r.snapshot, r.activity
137 if r.snapshot.ProjectionEpoch == "" || r.path != path || r.ledger != ledger {
138 base = event.RuntimeStateSnapshot{ProjectionEpoch: newRuntimeStateEpoch(), RuntimeEpoch: newRuntimeStateEpoch()}
139 activity = ""
140 }
141 next := base
142 next.SchemaVersion = 1
143 goalView, goalErr := c.goalLifecycleView()
144 next.Goal = goalView
145 next.GoalError = ""
146 if goalErr != nil {
147 next.GoalError = goalErr.Error()
148 }
149 if ref, ok := c.SessionRef(); ok {
150 next.HostID = ref.HostID
151 next.SessionID = ref.SessionID
152 next.SessionCodec = session.Codec
153 next.RuntimeEpoch = v3RuntimeSnapshot.Epoch
154 next.ActivityRevision = v3RuntimeSnapshot.ActivityRevision
155 } else {
156 next.HostID = ""
157 next.SessionID = ""
158 next.SessionCodec = ""
159 }
160 if ledger != nil {
161 next.TurnID, next.TurnStatus, next.TurnEventSeq = ledger.RuntimeIdentity()
162 }
163 v3Snapshot, hasV3Snapshot := c.sessionStateSnapshot()
164 applyRuntimeSessionState(&next, v3Snapshot, hasV3Snapshot)
165 next.HeadID = agent.BranchID(path)
166 if exclusiveSession {
167 next.HeadID = ""
168 }
169 if hasV3Snapshot {
170 // The typed v3 projection above is authoritative.
171 } else if ledger != nil {
172 next.Todos, next.TodoWritten = ledger.TodoState()
173 } else {
174 next.Todos, next.TodoWritten = c.volatileTodoState()
175 }
176 // Keep the empty wire shape stable. Ledger projections intentionally use a
177 // nil backing slice internally, while every public snapshot promises [].
178 // Normalizing before the semantic comparison prevents a read from creating
179 // a new revision solely because nil and an empty slice differ to reflect.
180 if next.Todos == nil {
181 next.Todos = []event.Todo{}
182 }
183 setRuntimePhase(&next, exclusiveSession, v3Runtime, v3RuntimeSnapshot, running, finishing, closed, cancelling)
184 applyMaintenanceRuntimeState(&next, maintenance, running, finishing, closed, cancelling)
185 identities, promptRevision := c.promptOwner.IdentitiesRevision()
186 next.PendingPrompt = len(identities) > 0
187 next.Interactions = make([]event.PendingInteraction, len(identities))
188 for i, identity := range identities {
189 next.Interactions[i] = event.PendingInteraction{RequestID: identity.PromptID, ToolCallID: identity.ToolCallID, Kind: string(identity.Kind), HeadID: next.HeadID, TurnID: identity.TurnID, RuntimeEpoch: identity.RuntimeEpoch}
190 }
191 // Compatibility consumers may still use Cancellable as a button-state
192 // hint. Derive it only from the authoritative phase/request projection so a
193 // worker crossing into the finishing window cannot make an accepted cancel
194 // briefly look unavailable.
195 next.Cancellable = next.Phase == "executing" || next.Phase == "cancelling" || next.PendingPrompt
196 if maintenance != nil && (maintenance.Activity == "finalizing" || maintenance.Activity == "recovery_required") {
197 next.Cancellable = false
198 }
199 next.BackgroundJobs = 0
200 if c.jobs != nil {
201 next.BackgroundJobs = len(c.jobs.RunningForSession(c.parentSessionID()))
202 }
203 // Sampling owners is off their locks. Do not commit a mixture if the
204 // admission/close/binding boundary advanced while another owner was read.
205 stable := c.runtimeBoundaryStable(running, finishing, closed, cancelling, path, maintenance)
206 currentGoal, currentGoalErr := c.goalLifecycleView()
207 stable = stable && reflect.DeepEqual(goalView, currentGoal)
208 stable = stable && ((goalErr == nil && currentGoalErr == nil) || (goalErr != nil && currentGoalErr != nil && goalErr.Error() == currentGoalErr.Error()))
209 if exclusiveSession && v3Runtime != nil {
210 _, currentRuntime, currentExclusive := c.v3Binding()
211 stable = stable && currentExclusive && currentRuntime == v3Runtime && currentRuntime.StateSnapshot().ActivityRevision == v3RuntimeSnapshot.ActivityRevision
212 }
213 if !stable || ledger != c.turnEventLedger() || promptRevision != c.promptOwner.Revision() {
214 r.mu.Unlock()
215 if attempt < 2 {
216 c.refreshRuntimeStateAttempt(event.Event{}, attempt+1)
217 }
218 return
219 }
220 if closed && !running && next.BackgroundJobs == 0 && r.jobUnsubscribe != nil {
221 stop := r.jobUnsubscribe
222 r.jobUnsubscribe = nil
223 defer stop()
224 }
225 activity = runtimeActivity(next, e, activity)
226 next.Activity = activity
227 setRuntimeRecovery(&next, v3Snapshot, hasV3Snapshot, ledger, activity)
228 // Token deltas do not need runtime notifications. Keep the last published
229 // watermark until a semantic state changes, avoiding a second token stream.
230 compare := next
231 compare.TurnEventSeq = r.snapshot.TurnEventSeq
232 if reflect.DeepEqual(compare, r.snapshot) {
233 r.mu.Unlock()
234 return
235 }
236 next.Revision++
237 r.commitSnapshot(next)
238 r.path, r.ledger, r.activity = path, ledger, activity
239 defer slog.Debug("runtime state committed", "source", "controller", "epoch", next.RuntimeEpoch[:8], "revision", next.Revision, "phase", next.Phase)
240 if !initialized {
241 r.mu.Unlock()
242 return
243 }
244 pending := cloneRuntimeState(next)
245 r.pending = &pending
246 if r.draining {
247 r.mu.Unlock()
248 return
249 }
250 r.draining = true
251 r.mu.Unlock()
252 go c.publishRuntimeState()
253 }
254
255 func (c *Controller) publishRuntimeState() {
256 r := &c.runtimeState
257 for {
258 r.mu.Lock()
259 if r.pending == nil {
260 r.draining = false
261 r.mu.Unlock()
262 return
263 }
264 snapshot, sink := cloneRuntimeState(*r.pending), r.sink
265 r.pending = nil
266 r.mu.Unlock()
267 event.PublishRuntimeState(sink, snapshot)
268 }
269 }
270
271 func runtimeActivity(state event.RuntimeStateSnapshot, e event.Event, activity string) string {
272 if state.Maintenance != nil {
273 switch state.Maintenance.Activity {
274 case "cancelling":
275 return "stopping_compaction"
276 case "finalizing":
277 return "saving_compaction"
278 case "recovery_required":
279 return "maintenance_recovery_required"
280 default:
281 return "compacting"
282 }
283 }
284 if state.PendingPrompt {
285 return "waiting_input"
286 }
287 if state.Phase == "cancelling" {
288 return "cancelling"
289 }
290 if state.Phase == "recovery_required" {
291 return "recovery_required"
292 }
293 if state.Phase == "executing" {
294 if e.TurnID == "" || e.TurnID == state.TurnID {
295 switch e.Kind {
296 case event.Text, event.Message:
297 activity = "streaming"
298 case event.TurnStarted, event.Reasoning, event.ToolDispatch, event.ToolProgress, event.ToolResult, event.CompactionStarted, event.Retrying:
299 activity = "thinking"
300 }
301 }
302 if activity == "" {
303 activity = "thinking"
304 }
305 } else {
306 activity = ""
307 }
308 return activity
309 }
310
311 func setRuntimePhase(next *event.RuntimeStateSnapshot, exclusiveSession bool, v3Runtime *session.Runtime, v3RuntimeSnapshot session.RuntimeSnapshot, running, finishing, closed, cancelling bool) {
312 next.Phase = "idle"
313 if closed {
314 next.Phase = "closed"
315 return
316 }
317 if exclusiveSession && v3Runtime != nil {
318 switch v3RuntimeSnapshot.Phase {
319 case session.RuntimeRunning:
320 next.Phase = "executing"
321 case session.RuntimeCancelling:
322 next.Phase = "cancelling"
323 case session.RuntimeFinalizing:
324 next.Phase = "finishing"
325 case session.RuntimeRecoveryRequired:
326 next.Phase = "recovery_required"
327 case session.RuntimeClosed:
328 next.Phase = "closed"
329 }
330 } else {
331 switch {
332 case next.TurnStatus == event.TurnRecoveryRequired:
333 next.Phase = "recovery_required"
334 case cancelling:
335 next.Phase = "cancelling"
336 case running:
337 next.Phase = "executing"
338 case finishing:
339 next.Phase = "finishing"
340 case closed:
341 next.Phase = "closed"
342 }
343 }
344 }
345
346 func applyRuntimeSessionState(next *event.RuntimeStateSnapshot, v3Snapshot session.Snapshot, hasV3Snapshot bool) {
347 if hasV3Snapshot {
348 snapshot := v3Snapshot
349 next.CommittedSeq = snapshot.EventSequence
350 next.DurableSeq = snapshot.DurableSequence
351 next.Persistence = string(snapshot.PersistenceStatus)
352 next.PersistenceErr = snapshot.PersistenceError
353 next.Todos = append([]event.Todo(nil), snapshot.Projection.Todos...)
354 next.TodoWritten = snapshot.Projection.TodoWritten
355 if snapshot.Projection.TurnID != "" {
356 next.TurnID = snapshot.Projection.TurnID
357 next.TurnStatus = snapshot.Projection.TurnStatus
358 }
359 if snapshot.Projection.Recovery != nil && snapshot.Projection.Recovery.State == "recovery_required" {
360 next.TurnStatus = event.TurnRecoveryRequired
361 }
362 } else {
363 next.CommittedSeq = next.TurnEventSeq
364 next.DurableSeq = next.TurnEventSeq
365 next.Persistence = "unavailable"
366 }
367 }
368
369 func setRuntimeRecovery(next *event.RuntimeStateSnapshot, v3Snapshot session.Snapshot, hasV3Snapshot bool, ledger *turnevent.Ledger, activity string) {
370 next.Recovery = nil
371 if next.Phase == "recovery_required" {
372 if hasV3Snapshot && v3Snapshot.Projection.Recovery != nil {
373 recovery := *v3Snapshot.Projection.Recovery
374 next.Recovery = &recovery
375 } else if ledger != nil {
376 next.Recovery = ledger.RecoveryStatus()
377 }
378 if next.Recovery == nil {
379 next.Recovery = &event.RecoveryStatus{State: "recovery_required", Phase: activity, Reason: "runtime state requires recovery"}
380 } else if next.Recovery.Phase == "" {
381 next.Recovery.Phase = activity
382 }
383 }
384 }
385
386 func (c *Controller) runtimeBoundaryStable(running, finishing, closed, cancelling bool, path string, maintenance *event.MaintenanceState) bool {
387 c.mu.Lock()
388 defer c.mu.Unlock()
389 return running == c.bodyActiveLocked() && finishing == c.finalizingLocked() && closed == c.closed && cancelling == c.cancelRequestedLocked() && path == c.sessionPath && reflect.DeepEqual(maintenance, c.maintenanceSnapshotLocked())
390 }
391
392 func applyMaintenanceRuntimeState(next *event.RuntimeStateSnapshot, maintenance *event.MaintenanceState, running, finishing, closed, cancelling bool) {
393 next.Maintenance = maintenance
394 if maintenance != nil && !closed {
395 switch maintenance.Activity {
396 case "cancelling":
397 next.Phase = "cancelling"
398 case "finalizing":
399 next.Phase = "finishing"
400 case "recovery_required":
401 next.Phase = "recovery_required"
402 default:
403 next.Phase = "executing"
404 }
405 }
406 // Close is immediately authoritative for the public controller view even
407 // while the session runtime remains in its private finalizing barrier. The
408 // latter keeps commit authority alive until TurnDone is durable; exposing it
409 // here would make a closed controller look runnable again.
410 next.Running = (running || finishing || maintenance != nil) && !closed
411 next.CancelRequested = (cancelling || maintenance != nil && maintenance.Activity == "cancelling") && !closed
412 }
413
413 lines GO