| 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 |