返回 DeepSeek-Reasonix
maintenance.go
根目录 / internal / control / maintenance.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "strings"
8 "sync"
9 "time"
10
11 "reasonix/internal/agent"
12 "reasonix/internal/event"
13 "reasonix/internal/extension"
14 "reasonix/internal/session"
15 )
16
17 var (
18 ErrMaintenanceBusy = fmt.Errorf("%w: session maintenance is already running", ErrTurnRunning)
19 ErrMaintenanceRecovery = fmt.Errorf("%w: session maintenance requires recovery", ErrRecoveryRequired)
20 )
21
22 type controllerMaintenance struct {
23 maintenanceIdentity
24 maintenanceResult
25 publishMu sync.Mutex
26 activity string
27 revision uint64
28 cancel context.CancelFunc
29 done chan struct{}
30 safeToRelease bool
31 cancelSignalled bool
32 terminal bool
33 }
34
35 // Identity and the pre-operation checkpoint are immutable after admission.
36 type maintenanceIdentity struct {
37 id string
38 kind string
39 runtimeEpoch string
40 sessionPath string
41 projectionVersion uint64
42 }
43
44 // Result fields are frozen together at each serialized publication boundary.
45 type maintenanceResult struct {
46 status string
47 errorCode string
48 detail string
49 inputTokens int
50 applied bool
51 resultTokens int
52 messages int
53 summary string
54 archive string
55 }
56
57 func (c *Controller) beginMaintenance(parent context.Context, kind string) (*controllerMaintenance, context.Context, error) {
58 if parent == nil {
59 parent = context.Background()
60 }
61 if err := c.ensureWriteAuthorityReady(); err != nil {
62 return nil, nil, err
63 }
64 runtimeEpoch := c.RuntimeStateSnapshot().RuntimeEpoch
65 c.mu.Lock()
66 verb := strings.ReplaceAll(kind, "_", " ")
67 if strings.HasPrefix(kind, "summarize_") {
68 verb = "summarize"
69 }
70 switch {
71 case c.closed:
72 c.mu.Unlock()
73 return nil, nil, errors.New("controller is closed")
74 case c.bodyActiveLocked() || c.finalizingLocked():
75 c.mu.Unlock()
76 return nil, nil, fmt.Errorf("cannot %s while a turn is running", verb)
77 case c.rotating:
78 c.mu.Unlock()
79 return nil, nil, errRotationInProgress
80 case c.maintenance != nil:
81 err := ErrMaintenanceBusy
82 if c.maintenance.activity == "recovery_required" {
83 err = ErrMaintenanceRecovery
84 }
85 c.mu.Unlock()
86 return nil, nil, err
87 case c.recoveryRequiredLocked():
88 c.mu.Unlock()
89 return nil, nil, ErrMaintenanceRecovery
90 }
91 if c.rejectDrainingGenerationLocked() {
92 c.mu.Unlock()
93 return nil, nil, ErrRuntimeDraining
94 }
95 if c.turns.runtime != nil && !c.turns.runtime.BeginExecution(c.turns.generation, session.MaintenanceActivity) {
96 c.mu.Unlock()
97 return nil, nil, ErrMaintenanceBusy
98 }
99 work, cancel := context.WithCancel(extension.ContextWithRuntimeOwner(c.withAuthentication(parent), c.runtimeOwner))
100 before := c.ContextMaintenanceSnapshot()
101 op := &controllerMaintenance{
102 maintenanceIdentity: maintenanceIdentity{
103 id: "maintenance-" + newRuntimeStateEpoch(), kind: kind,
104 runtimeEpoch: runtimeEpoch, sessionPath: c.sessionPath,
105 projectionVersion: before.ProjectionVersion,
106 },
107 maintenanceResult: maintenanceResult{status: "running", inputTokens: before.ProjectedTokens},
108 activity: "running", cancel: cancel, done: make(chan struct{}),
109 }
110 c.maintenance = op
111 c.mu.Unlock()
112 if err := c.emitMaintenanceOperation(op, "running", "", "", before, false); err != nil {
113 cancel()
114 _ = c.retainMaintenanceRecovery(op, "operation_persist_failed", err, before, false)
115 return nil, nil, fmt.Errorf("persist maintenance start: %w", err)
116 }
117 c.refreshRuntimeState(event.Event{})
118 return op, work, nil
119 }
120
121 // startCompactAsync is used by /compact so registration happens before the
122 // command handler returns. Desktop/HTTP/CLI callers keep the synchronous
123 // Compact API and wait on the same lifecycle.
124 func (c *Controller) startCompactAsync(instructions string) error {
125 if err := c.authentication.admissionError(); err != nil {
126 return err
127 }
128 if c.executor == nil {
129 return nil
130 }
131 op, ctx, err := c.beginMaintenance(context.Background(), "compact")
132 if err != nil {
133 return err
134 }
135 go func() {
136 _ = c.executeMaintenance(op, ctx, func(work context.Context) error {
137 return c.executor.CompactNow(work, instructions)
138 })
139 }()
140 return nil
141 }
142
143 func (c *Controller) executeMaintenance(op *controllerMaintenance, ctx context.Context, work func(context.Context) error) error {
144 err := work(ctx)
145 c.authentication.recordFailure(err, c.ModelRef())
146 after := c.ContextMaintenanceSnapshot()
147 applied := after.ProjectionVersion > op.projectionVersion
148
149 // Once the worker returns, expose the non-cancellable persistence boundary.
150 if emitErr := c.emitMaintenanceOperation(op, "finalizing", "", "", after, applied); emitErr != nil {
151 return c.retainMaintenanceRecovery(op, "operation_persist_failed", emitErr, after, applied)
152 }
153 c.refreshRuntimeState(event.Event{})
154
155 status, code, detail := "completed", "", ""
156 if err == nil && applied {
157 // The final projection commit won the boundary. A Stop that arrived
158 // afterwards is idempotent and must not relabel completed work.
159 } else if errors.Is(err, context.Canceled) || errors.Is(ctx.Err(), context.Canceled) {
160 status, code = "cancelled", "cancelled"
161 if applied {
162 status = "partially_completed"
163 }
164 } else if err != nil {
165 status, code, detail = "failed", "summary_failed", err.Error()
166 } else if !applied {
167 status, code = "noop", "no_history"
168 }
169
170 // Manual projection changes and their operation record share one finishing
171 // boundary. A save failure retains maintenance ownership and blocks inbox
172 // dispatch so the UI cannot report a safe idle session.
173 if applied {
174 if saveErr := c.SnapshotRewrite(); saveErr != nil {
175 return c.retainMaintenanceRecovery(op, "save_failed", saveErr, after, applied)
176 }
177 }
178
179 if emitErr := c.emitMaintenanceOperation(op, status, code, detail, after, applied); emitErr != nil {
180 return c.retainMaintenanceRecovery(op, "operation_persist_failed", emitErr, after, applied)
181 }
182 c.finishMaintenanceOperation(op)
183 return err
184 }
185
186 func (c *Controller) retainMaintenanceRecovery(op *controllerMaintenance, code string, cause error, snapshot agent.ContextMaintenanceSnapshot, applied bool) error {
187 if cause == nil {
188 cause = ErrMaintenanceRecovery
189 }
190 // Publication still owns the persistence resources even though the worker
191 // has exited. Close may release them only after this final attempt settles.
192 _ = c.emitMaintenanceOperation(op, "recovery_required", code, cause.Error(), snapshot, applied)
193 c.mu.Lock()
194 closed := c.closed
195 if c.maintenance == op {
196 op.safeToRelease = true
197 if closed {
198 c.maintenance = nil
199 close(op.done)
200 }
201 }
202 c.mu.Unlock()
203 c.refreshRuntimeState(event.Event{})
204 if closed {
205 c.finalizeControllerClose()
206 }
207 return cause
208 }
209
210 func (c *Controller) emitMaintenanceOperation(op *controllerMaintenance, status, code, detail string, snapshot agent.ContextMaintenanceSnapshot, applied bool) error {
211 if op == nil {
212 return nil
213 }
214 op.publishMu.Lock()
215 defer op.publishMu.Unlock()
216 summary := ""
217 if applied && c.executor != nil {
218 summary = c.executor.LastCompactionSummary()
219 }
220 c.mu.Lock()
221 if op.terminal {
222 c.mu.Unlock()
223 return nil
224 }
225 if c.maintenance != op || !maintenanceTransitionAllowed(op, status, code) {
226 c.mu.Unlock()
227 return nil
228 }
229 activity := op.activity
230 switch status {
231 case "running", "cancelling", "finalizing", "recovery_required":
232 activity = status
233 }
234 info := &event.SessionOperationInfo{
235 OperationID: op.id, OperationRevision: op.revision + 1, RuntimeEpoch: op.runtimeEpoch,
236 Kind: op.kind, Activity: activity, Status: status,
237 ErrorCode: code, Detail: detail, Applied: op.applied || applied,
238 InputTokens: op.inputTokens, ResultTokens: snapshot.ProjectedTokens,
239 Messages: op.messages, Summary: op.summary, Archive: op.archive,
240 }
241 if summary != "" {
242 info.Summary = summary
243 }
244 if snapshot.LastReceipt != nil && info.Applied {
245 info.Messages, info.Archive = snapshot.LastReceipt.CoveredCount, snapshot.LastReceipt.Archive
246 }
247 terminal := maintenanceTerminalStatus(status)
248 if !terminal {
249 c.applyMaintenanceInfoLocked(op, info)
250 }
251 c.mu.Unlock()
252 err := event.EmitChecked(c.sink, event.Event{Kind: event.SessionOperation, SessionOperation: info})
253 if err == nil && terminal {
254 c.mu.Lock()
255 c.applyMaintenanceInfoLocked(op, info)
256 op.terminal = true
257 c.mu.Unlock()
258 }
259 return err
260 }
261
262 func maintenanceRuntimePhase(activity string) session.RuntimePhase {
263 switch activity {
264 case "cancelling":
265 return session.RuntimeCancelling
266 case "finalizing":
267 return session.RuntimeFinalizing
268 case "recovery_required":
269 return session.RuntimeRecoveryRequired
270 default:
271 return session.RuntimeRunning
272 }
273 }
274
275 func (c *Controller) applyMaintenanceInfoLocked(op *controllerMaintenance, info *event.SessionOperationInfo) {
276 op.activity, op.status, op.revision = info.Activity, info.Status, info.OperationRevision
277 op.errorCode, op.detail, op.applied = info.ErrorCode, info.Detail, info.Applied
278 op.resultTokens, op.messages, op.summary, op.archive = info.ResultTokens, info.Messages, info.Summary, info.Archive
279 c.noteExecutionLocked(maintenanceRuntimePhase(op.activity), session.MaintenanceActivity)
280 }
281
282 func maintenanceTransitionAllowed(op *controllerMaintenance, status, code string) bool {
283 if op.terminal {
284 return false
285 }
286 if status == "running" {
287 return op.revision == 0
288 }
289 if status == "cancelling" {
290 return op.activity == "running"
291 }
292 if status == "recovery_required" {
293 return code != "cancel_timeout" || op.cancelSignalled && (op.activity == "running" || op.activity == "cancelling")
294 }
295 // Only the settled worker may leave timeout recovery; persistence recovery
296 // is retained until its owner is explicitly recovered or closed.
297 if op.activity == "recovery_required" && op.errorCode != "cancel_timeout" {
298 return false
299 }
300 return status == "finalizing" || maintenanceTerminalStatus(status)
301 }
302
303 func (c *Controller) signalMaintenanceCancel() (id string, present, cancelled bool) {
304 c.mu.Lock()
305 op := c.maintenance
306 if op == nil {
307 c.mu.Unlock()
308 return "", false, false
309 }
310 id = op.id
311 if op.activity == "finalizing" || op.activity == "recovery_required" || op.terminal {
312 c.mu.Unlock()
313 return id, true, false
314 }
315 if op.cancelSignalled {
316 c.mu.Unlock()
317 return id, true, false
318 }
319 op.cancelSignalled = true
320 cancel := op.cancel
321 applied := op.applied
322 c.mu.Unlock()
323 // Signal first. Event delivery and every queue/disk operation happen after
324 // the provider/adapter has observed cancellation.
325 cancel()
326 c.recordLifecycle("cancel_signalled", "maintenance", id, 0, "")
327 c.startMaintenanceCancellationWatchdog(op)
328 // Cancellation receipts must not wait for the event lane or a runtime
329 // sampler. Publication revalidates this exact operation and its phase.
330 go func() {
331 _ = c.emitMaintenanceOperation(op, "cancelling", "", "", c.ContextMaintenanceSnapshot(), applied)
332 c.refreshRuntimeState(event.Event{})
333 }()
334 return id, true, true
335 }
336
337 func (c *Controller) startMaintenanceCancellationWatchdog(op *controllerMaintenance) {
338 if c == nil || op == nil {
339 return
340 }
341 go func() {
342 timer := time.NewTimer(c.cancellationGrace())
343 defer timer.Stop()
344 select {
345 case <-op.done:
346 return
347 case <-timer.C:
348 }
349 c.mu.Lock()
350 if c.maintenance != op || !op.cancelSignalled || op.activity != "running" && op.activity != "cancelling" {
351 c.mu.Unlock()
352 return
353 }
354 applied := op.applied
355 c.mu.Unlock()
356 _ = c.emitMaintenanceOperation(op, "recovery_required", "cancel_timeout", "the compaction worker did not stop within the cancellation grace period", c.ContextMaintenanceSnapshot(), applied)
357 c.refreshRuntimeState(event.Event{})
358 }()
359 }
360
361 func (c *Controller) maintenanceSnapshotLocked() *event.MaintenanceState {
362 if c.maintenance == nil {
363 return nil
364 }
365 op := c.maintenance
366 return &event.MaintenanceState{
367 OperationID: op.id, OperationRevision: op.revision, RuntimeEpoch: op.runtimeEpoch,
368 Kind: op.kind, Activity: op.activity, Status: op.status,
369 ErrorCode: op.errorCode, Detail: op.detail, Applied: op.applied,
370 InputTokens: op.inputTokens, ResultTokens: op.resultTokens, Messages: op.messages,
371 }
372 }
373
374 func maintenanceTerminalStatus(status string) bool {
375 switch status {
376 case "completed", "noop", "cancelled", "partially_completed", "failed", "interrupted":
377 return true
378 default:
379 return false
380 }
381 }
382
383 // finishMaintenanceOperation releases the maintenance owner and reserves the
384 // oldest volatile follow-up under the same lock. This closes the idle gap in
385 // which a new submission could otherwise jump ahead of accepted guidance.
386 func (c *Controller) finishMaintenanceOperation(op *controllerMaintenance) {
387 var next queuedTurn
388 var runCtx context.Context
389 var cancel context.CancelFunc
390 var started bool
391 c.mu.Lock()
392 if c.maintenance != op {
393 c.mu.Unlock()
394 return
395 }
396 if runtime := c.turns.runtime; runtime != nil && !runtime.FinishMaintenanceExecution(c.turns.generation) {
397 c.mu.Unlock()
398 _ = c.retainMaintenanceRecovery(op, "execution_owner_lost", ErrMaintenanceRecovery, c.ContextMaintenanceSnapshot(), op.applied)
399 return
400 }
401 c.maintenance = nil
402 close(op.done)
403 closed := c.closed
404 if !closed && !c.recoveryRequiredLocked() && c.authentication.admissionError() == nil {
405 if candidate, ok := c.popNextPendingLocked(); ok {
406 next = candidate
407 runCtx, cancel, started = c.startTurnLocked(context.Background(), next)
408 if !started {
409 c.turns.pending = append([]queuedTurn{next}, c.turns.pending...)
410 c.turns.wake = true
411 c.enterRecoveryLocked("execution_owner_lost")
412 }
413 }
414 }
415 if !started {
416 c.noteExecutionLocked(session.RuntimeIdle, "")
417 }
418 c.mu.Unlock()
419 op.cancel()
420 c.refreshRuntimeState(event.Event{})
421 if closed {
422 c.finalizeControllerClose()
423 return
424 }
425 if started {
426 if next.onStart != nil {
427 next.onStart()
428 }
429 c.spawnGuardedTurn(runCtx, cancel, next)
430 return
431 }
432 c.mu.Lock()
433 recovery := c.recoveryRequiredLocked()
434 c.mu.Unlock()
435 if !recovery && c.authentication.admissionError() == nil {
436 c.maybeDispatchInbox()
437 c.kickGoalDriver()
438 }
439 }
440
441 // ActiveMaintenanceOperationID lets result-capable hosts correlate a
442 // management-command receipt with the authoritative operation event. Empty is
443 // valid when the maintenance task reached a terminal state before the caller
444 // received its receipt.
445 func (c *Controller) ActiveMaintenanceOperationID() string {
446 if c == nil {
447 return ""
448 }
449 c.mu.Lock()
450 defer c.mu.Unlock()
451 if c.maintenance == nil {
452 return ""
453 }
454 return c.maintenance.id
455 }
456
456 lines GO