| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "fmt" |
| 7 | "time" |
| 8 | |
| 9 | "reasonix/internal/agent" |
| 10 | "reasonix/internal/event" |
| 11 | "reasonix/internal/extension" |
| 12 | "reasonix/internal/provider" |
| 13 | "reasonix/internal/session" |
| 14 | ) |
| 15 | |
| 16 | type queuedTurnKind int |
| 17 | |
| 18 | const ( |
| 19 | queuedUser queuedTurnKind = iota |
| 20 | queuedGoal |
| 21 | ) |
| 22 | |
| 23 | type queuedTurn struct { |
| 24 | kind queuedTurnKind |
| 25 | body func(ctx context.Context) error |
| 26 | onStart func() |
| 27 | goalRound *goalRoundReservation |
| 28 | // admissionCtx is only used for the synchronous durability boundary. The |
| 29 | // model turn has its own controller-owned context and survives RPC return. |
| 30 | admissionCtx context.Context |
| 31 | } |
| 32 | |
| 33 | // turnLoop is the session-scoped execution authority. Controller.mu guards it. |
| 34 | type turnLoop struct { |
| 35 | phase session.RuntimePhase |
| 36 | cancel context.CancelFunc |
| 37 | done chan struct{} |
| 38 | recoveryFanout bool // watchdog terminal publication still owns the stores |
| 39 | turnID string |
| 40 | token uint64 |
| 41 | lastToken uint64 |
| 42 | pending []queuedTurn |
| 43 | wake bool |
| 44 | generation uint64 |
| 45 | runtime *session.Runtime |
| 46 | cancelRequested bool |
| 47 | finishingBound turnFinishingBoundary |
| 48 | } |
| 49 | |
| 50 | type controllerExecution struct { |
| 51 | c *Controller |
| 52 | } |
| 53 | |
| 54 | func (e controllerExecution) Snapshot() session.RuntimeSnapshot { |
| 55 | if e.c == nil { |
| 56 | return session.RuntimeSnapshot{Phase: session.RuntimeIdle} |
| 57 | } |
| 58 | e.c.mu.Lock() |
| 59 | defer e.c.mu.Unlock() |
| 60 | if op := e.c.maintenance; op != nil { |
| 61 | return session.RuntimeSnapshot{Phase: maintenanceRuntimePhase(op.activity), Activity: session.MaintenanceActivity} |
| 62 | } |
| 63 | return session.RuntimeSnapshot{Phase: e.c.turns.phase, Activity: e.c.turns.activityNameLocked()} |
| 64 | } |
| 65 | |
| 66 | func (e controllerExecution) Cancel() bool { |
| 67 | if e.c == nil { |
| 68 | return false |
| 69 | } |
| 70 | if _, present, cancelled := e.c.signalMaintenanceCancel(); present { |
| 71 | return cancelled |
| 72 | } |
| 73 | return e.c.signalTurnCancel() |
| 74 | } |
| 75 | |
| 76 | func (t *turnLoop) activityNameLocked() string { |
| 77 | switch t.phase { |
| 78 | case session.RuntimeCancelling: |
| 79 | return "cancelling" |
| 80 | case session.RuntimeRecoveryRequired: |
| 81 | return "recovery_required" |
| 82 | case session.RuntimeRunning, session.RuntimeFinalizing: |
| 83 | if t.cancelRequested { |
| 84 | return "cancelling" |
| 85 | } |
| 86 | return "turn" |
| 87 | default: |
| 88 | return "" |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | func (c *Controller) bodyActiveLocked() bool { |
| 93 | switch c.turns.phase { |
| 94 | case session.RuntimeRunning, session.RuntimeCancelling: |
| 95 | return true |
| 96 | default: |
| 97 | return false |
| 98 | } |
| 99 | } |
| 100 | |
| 101 | func (c *Controller) finalizingLocked() bool { |
| 102 | return c.turns.phase == session.RuntimeFinalizing |
| 103 | } |
| 104 | |
| 105 | func (c *Controller) cancelRequestedLocked() bool { |
| 106 | if c.closed { |
| 107 | return false |
| 108 | } |
| 109 | return c.turns.cancelRequested || c.turns.phase == session.RuntimeCancelling |
| 110 | } |
| 111 | |
| 112 | func (c *Controller) recoveryRequiredLocked() bool { |
| 113 | return c.turns.phase == session.RuntimeRecoveryRequired |
| 114 | } |
| 115 | |
| 116 | func (c *Controller) bindExecutionControl() { |
| 117 | _, runtime, exclusive := c.v3Binding() |
| 118 | if !exclusive || runtime == nil { |
| 119 | return |
| 120 | } |
| 121 | snap := runtime.StateSnapshot() |
| 122 | gen := runtime.BindExecution(controllerExecution{c: c}) |
| 123 | c.executionGeneration.Store(gen) |
| 124 | c.mu.Lock() |
| 125 | c.turns.generation = gen |
| 126 | c.turns.runtime = runtime |
| 127 | if snap.Phase == session.RuntimeRecoveryRequired { |
| 128 | c.turns.phase = session.RuntimeRecoveryRequired |
| 129 | } |
| 130 | c.mu.Unlock() |
| 131 | } |
| 132 | |
| 133 | // ExecutionGeneration returns the session-runtime execution generation owned |
| 134 | // by this controller. Zero means the controller is a prepared replacement that |
| 135 | // has not been published as the execution owner. |
| 136 | func (c *Controller) ExecutionGeneration() uint64 { |
| 137 | if c == nil { |
| 138 | return 0 |
| 139 | } |
| 140 | return c.executionGeneration.Load() |
| 141 | } |
| 142 | |
| 143 | // ActivateSessionExecution publishes this controller as the exact execution |
| 144 | // owner. expectedGeneration is zero for a previously unbound runtime and the |
| 145 | // outgoing controller generation for a fail-atomic replacement. The method is |
| 146 | // deliberately callback- and I/O-free so hosts may invoke it in their final |
| 147 | // pointer-swap critical section. |
| 148 | func (c *Controller) ActivateSessionExecution(expectedGeneration uint64) error { |
| 149 | if c == nil { |
| 150 | return session.ErrSessionNotRunning |
| 151 | } |
| 152 | c.turnEvents.commitMu.Lock() |
| 153 | defer c.turnEvents.commitMu.Unlock() |
| 154 | c.mu.Lock() |
| 155 | defer c.mu.Unlock() |
| 156 | runtime := c.turns.runtime |
| 157 | if runtime == nil { |
| 158 | return nil |
| 159 | } |
| 160 | if generation := c.turns.generation; generation != 0 && runtime.OwnsExecution(generation) { |
| 161 | return nil |
| 162 | } |
| 163 | var generation uint64 |
| 164 | if expectedGeneration == 0 { |
| 165 | generation = runtime.BindExecution(controllerExecution{c: c}) |
| 166 | } else if pending := c.turnEvents.pendingExecutionCommit; pending != nil { |
| 167 | c.turnEvents.pendingExecutionCommit = nil |
| 168 | var err error |
| 169 | generation, _, err = runtime.ReplaceExecutionAndCommit(expectedGeneration, controllerExecution{c: c}, *pending) |
| 170 | if err != nil { |
| 171 | return err |
| 172 | } |
| 173 | } else { |
| 174 | generation = runtime.ReplaceExecution(expectedGeneration, controllerExecution{c: c}) |
| 175 | } |
| 176 | if generation == 0 { |
| 177 | return session.ErrRuntimeBusy |
| 178 | } |
| 179 | c.turns.generation = generation |
| 180 | c.executionGeneration.Store(generation) |
| 181 | return nil |
| 182 | } |
| 183 | |
| 184 | // ActivateControllerReplacement transfers session execution ownership when a |
| 185 | // host commits a controller pointer swap. Controllers without a shared Runtime |
| 186 | // need no additional activation. |
| 187 | func ActivateControllerReplacement(old, next *Controller) error { |
| 188 | if next == nil { |
| 189 | return session.ErrSessionNotRunning |
| 190 | } |
| 191 | if old != nil && old.background.scope != nil && old.background.scope == next.background.scope { |
| 192 | if old.workspaceRoot != next.workspaceRoot { |
| 193 | return fmt.Errorf("background ownership cannot cross workspaces") |
| 194 | } |
| 195 | before, beforeOK := old.SessionRef() |
| 196 | after, afterOK := next.SessionRef() |
| 197 | if beforeOK && afterOK && before.SessionID != after.SessionID { |
| 198 | return fmt.Errorf("background ownership cannot cross sessions") |
| 199 | } |
| 200 | if ModelReplacementBlocked(old) { |
| 201 | return fmt.Errorf("session changed while preparing replacement") |
| 202 | } |
| 203 | } |
| 204 | if old != nil && old.background.scope != nil && old.background.scope == next.background.scope && old.ToolApprovalMode() != next.ToolApprovalMode() && len(old.jobs.Running()) > 0 { |
| 205 | return fmt.Errorf("permission changes require stopping background jobs before replacement") |
| 206 | } |
| 207 | _, nextRuntime, nextExclusive := next.v3Binding() |
| 208 | if !nextExclusive || nextRuntime == nil { |
| 209 | retireBackgroundCallbacks(old, next) |
| 210 | next.PublishBackgroundScope() |
| 211 | return nil |
| 212 | } |
| 213 | expected := uint64(0) |
| 214 | if old != nil { |
| 215 | _, oldRuntime, oldExclusive := old.v3Binding() |
| 216 | if oldExclusive && oldRuntime == nextRuntime { |
| 217 | expected = old.ExecutionGeneration() |
| 218 | } |
| 219 | } |
| 220 | if err := next.ActivateSessionExecution(expected); err != nil { |
| 221 | return err |
| 222 | } |
| 223 | retireBackgroundCallbacks(old, next) |
| 224 | next.PublishBackgroundScope() |
| 225 | return nil |
| 226 | } |
| 227 | |
| 228 | func retireBackgroundCallbacks(old, next *Controller) { |
| 229 | if old != nil && old != next && old.background.scope != nil && old.background.scope == next.background.scope { |
| 230 | old.mu.Lock() |
| 231 | old.background.retired = true |
| 232 | old.mu.Unlock() |
| 233 | } |
| 234 | } |
| 235 | |
| 236 | // ActivateSessionAPIReplacement is the host-facing form used at a final |
| 237 | // controller pointer swap. Non-Controller implementations have no exclusive |
| 238 | // Runtime ownership to transfer and are left unchanged. |
| 239 | func ActivateSessionAPIReplacement(old, next SessionAPI) error { |
| 240 | concreteNext, ok := next.(*Controller) |
| 241 | if !ok || concreteNext == nil { |
| 242 | return nil |
| 243 | } |
| 244 | concreteOld, _ := old.(*Controller) |
| 245 | if err := ActivateControllerReplacement(concreteOld, concreteNext); err != nil { |
| 246 | return err |
| 247 | } |
| 248 | if concreteOld != nil && concreteOld.attachmentScope() != "" && concreteOld.attachmentScope() == concreteNext.attachmentScope() && concreteOld.workspaceRoot == concreteNext.workspaceRoot { |
| 249 | concreteOld.attachmentService().Drafts().CopyScopeTo(concreteOld.attachmentScope(), concreteNext.attachmentService().Drafts()) |
| 250 | } |
| 251 | return nil |
| 252 | } |
| 253 | |
| 254 | func (c *Controller) unbindExecutionControl(runtime *session.Runtime) { |
| 255 | if runtime == nil { |
| 256 | return |
| 257 | } |
| 258 | c.mu.Lock() |
| 259 | gen := c.turns.generation |
| 260 | c.mu.Unlock() |
| 261 | runtime.UnbindExecution(gen) |
| 262 | } |
| 263 | |
| 264 | func (c *Controller) noteExecutionLocked(phase session.RuntimePhase, activity string) { |
| 265 | if c.turns.runtime == nil || c.turns.generation == 0 { |
| 266 | return |
| 267 | } |
| 268 | c.turns.runtime.NoteExecution(c.turns.generation, phase, activity) |
| 269 | } |
| 270 | |
| 271 | func (c *Controller) currentTurnToken() (token uint64, turnID string, active bool) { |
| 272 | c.mu.Lock() |
| 273 | defer c.mu.Unlock() |
| 274 | switch c.turns.phase { |
| 275 | case session.RuntimeIdle, session.RuntimeClosed: |
| 276 | return c.turns.lastToken, "", false |
| 277 | default: |
| 278 | return c.turns.token, c.turns.turnID, true |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | func (c *Controller) discardLateTurnEvent(e event.Event) bool { |
| 283 | if e.TurnID == "" || !lateBusinessEvent(e.Kind) { |
| 284 | return false |
| 285 | } |
| 286 | _, turnID, active := c.currentTurnToken() |
| 287 | if active { |
| 288 | return e.TurnID != turnID |
| 289 | } |
| 290 | return true |
| 291 | } |
| 292 | |
| 293 | func (c *Controller) startTurnLocked(parent context.Context, next queuedTurn) (ctx context.Context, cancel context.CancelFunc, admitted bool) { |
| 294 | if parent == nil { |
| 295 | parent = context.Background() |
| 296 | } |
| 297 | if c.turns.runtime != nil && !c.turns.runtime.BeginExecution(c.turns.generation, "turn") { |
| 298 | return nil, nil, false |
| 299 | } |
| 300 | ctx, cancel = context.WithCancel(extension.ContextWithRuntimeOwner(c.withAuthentication(parent), c.runtimeOwner)) |
| 301 | c.turns.cancel = cancel |
| 302 | c.turns.done = make(chan struct{}) |
| 303 | c.turns.finishingBound.beginIdle() |
| 304 | c.turns.phase = session.RuntimeRunning |
| 305 | c.turns.cancelRequested = false |
| 306 | c.turns.token++ |
| 307 | ctx = context.WithValue(ctx, executionTokenKey{}, c.turns.token) |
| 308 | c.turns.turnID = "" |
| 309 | return ctx, cancel, true |
| 310 | } |
| 311 | |
| 312 | func (c *Controller) popNextPendingLocked() (queuedTurn, bool) { |
| 313 | if len(c.turns.pending) == 0 { |
| 314 | c.turns.wake = false |
| 315 | return queuedTurn{}, false |
| 316 | } |
| 317 | next := c.turns.pending[0] |
| 318 | c.turns.pending = c.turns.pending[1:] |
| 319 | c.turns.wake = len(c.turns.pending) > 0 |
| 320 | return next, true |
| 321 | } |
| 322 | |
| 323 | func (c *Controller) queueTurnLocked(item queuedTurn) { |
| 324 | // Harness wakeRequested: input that cannot join the current activity is |
| 325 | // claimed once when the body converges. Close clears this queue so a |
| 326 | // disposed session never starts a latched turn. |
| 327 | c.turns.pending = append(c.turns.pending, item) |
| 328 | c.turns.wake = true |
| 329 | } |
| 330 | |
| 331 | func (c *Controller) signalTurnCancel() bool { |
| 332 | _, _, cancelled := c.signalTurnCancelIdentity() |
| 333 | return cancelled |
| 334 | } |
| 335 | |
| 336 | func (c *Controller) signalTurnCancelIdentity() (uint64, string, bool) { |
| 337 | c.mu.Lock() |
| 338 | token, turnID := c.turns.token, c.turns.turnID |
| 339 | cancel := c.turns.cancel |
| 340 | first := cancel != nil && c.turns.phase == session.RuntimeRunning |
| 341 | if cancel != nil && (c.turns.phase == session.RuntimeRunning || c.turns.phase == session.RuntimeCancelling) { |
| 342 | c.turns.phase = session.RuntimeCancelling |
| 343 | c.turns.cancelRequested = true |
| 344 | } |
| 345 | done := c.turns.done |
| 346 | c.mu.Unlock() |
| 347 | if cancel == nil { |
| 348 | return token, turnID, false |
| 349 | } |
| 350 | cancel() |
| 351 | c.recordLifecycle("cancel_signalled", "turn", turnID, 0, "") |
| 352 | if first { |
| 353 | c.startCancellationWatchdog(done) |
| 354 | } |
| 355 | return token, turnID, true |
| 356 | } |
| 357 | |
| 358 | func (c *Controller) enterRecoveryLocked(reason string) { |
| 359 | c.turns.phase = session.RuntimeRecoveryRequired |
| 360 | c.noteExecutionLocked(session.RuntimeRecoveryRequired, reason) |
| 361 | if c.turns.runtime != nil { |
| 362 | c.turns.runtime.RequireRecovery(reason) |
| 363 | } |
| 364 | } |
| 365 | |
| 366 | func (c *Controller) spawnGuardedTurn(ctx context.Context, cancel context.CancelFunc, item queuedTurn) { |
| 367 | ctx, completion := withGuardedTurnCompletion(ctx) |
| 368 | admissionCtx := item.admissionCtx |
| 369 | if admissionCtx == nil { |
| 370 | admissionCtx = context.Background() |
| 371 | } |
| 372 | body := c.prepareTurnAdmissionWithGoalRound(admissionCtx, item.body, item.goalRound) |
| 373 | if ledger := c.turnEventLedger(); ledger != nil { |
| 374 | c.mu.Lock() |
| 375 | c.turns.turnID = ledger.ActiveTurnID() |
| 376 | c.mu.Unlock() |
| 377 | } |
| 378 | c.liveness.reset(time.Now()) |
| 379 | c.autosaveWG.Go(func() { |
| 380 | c.autosaveWhileRunning(ctx) |
| 381 | }) |
| 382 | go func() { |
| 383 | defer cancel() |
| 384 | defer func() { |
| 385 | c.finishGoalRoundActivity(item.goalRound) |
| 386 | c.kickGoalDriver() |
| 387 | }() |
| 388 | defer func() { |
| 389 | if r := recover(); r != nil { |
| 390 | err := fmt.Errorf("internal error: %v", r) |
| 391 | item.goalRound.setResult(err, false) |
| 392 | c.finishGuardedTurn(err, completion) |
| 393 | } |
| 394 | }() |
| 395 | err := body(ctx) |
| 396 | if item.goalRound != nil { |
| 397 | item.goalRound.setResult(err, errors.Is(ctx.Err(), context.Canceled) && c.CancelRequested()) |
| 398 | } |
| 399 | c.finishGuardedTurn(explainError(err), completion) |
| 400 | }() |
| 401 | } |
| 402 | |
| 403 | func (c *Controller) cancellationGrace() time.Duration { |
| 404 | if c != nil && c.testCancelGrace > 0 { |
| 405 | return c.testCancelGrace |
| 406 | } |
| 407 | return 15 * time.Second |
| 408 | } |
| 409 | |
| 410 | func (c *Controller) finishGuardedTurn(err error, completion *guardedTurnCompletion) { |
| 411 | c.authentication.recordFailure(err, c.ModelRef()) |
| 412 | c.memory.clearAutoRemember() |
| 413 | c.mu.Lock() |
| 414 | cancelRequested := c.turns.cancelRequested |
| 415 | if c.turns.phase == session.RuntimeRecoveryRequired { |
| 416 | c.turns.cancel = nil |
| 417 | closing := c.closed |
| 418 | c.mu.Unlock() |
| 419 | // The cancellation watchdog already committed the recovery terminal. |
| 420 | // A closing controller must not emit another terminal after its ledger |
| 421 | // and session binding have been finalized. |
| 422 | if !closing { |
| 423 | c.emitTurnDoneEvent(err, cancelRequested, completion) |
| 424 | } |
| 425 | c.mu.Lock() |
| 426 | // Keep the owned turn live through terminal fanout. Close must not |
| 427 | // release stores while that fanout can still publish durable events. |
| 428 | if c.turns.done != nil { |
| 429 | close(c.turns.done) |
| 430 | c.turns.done = nil |
| 431 | } |
| 432 | closing = c.closed && !c.turns.recoveryFanout |
| 433 | c.turns.finishingBound.endIdle() |
| 434 | c.mu.Unlock() |
| 435 | if closing { |
| 436 | c.finalizeControllerClose() |
| 437 | } |
| 438 | c.refreshRuntimeState(event.Event{}) |
| 439 | return |
| 440 | } |
| 441 | if c.turns.done != nil { |
| 442 | close(c.turns.done) |
| 443 | c.turns.done = nil |
| 444 | } |
| 445 | c.turns.phase = session.RuntimeFinalizing |
| 446 | c.turns.finishingBound.begin(true) |
| 447 | c.turns.cancel = nil |
| 448 | c.noteExecutionLocked(session.RuntimeFinalizing, "turn") |
| 449 | c.mu.Unlock() |
| 450 | |
| 451 | c.refreshRuntimeState(event.Event{}) |
| 452 | defer func() { |
| 453 | c.mu.Lock() |
| 454 | c.turns.finishingBound.end() |
| 455 | c.turns.cancelRequested = false |
| 456 | if c.turns.phase == session.RuntimeRecoveryRequired { |
| 457 | closing := c.closed |
| 458 | c.turns.finishingBound.endIdle() |
| 459 | c.mu.Unlock() |
| 460 | if closing { |
| 461 | c.finalizeControllerClose() |
| 462 | } |
| 463 | c.refreshRuntimeState(event.Event{}) |
| 464 | return |
| 465 | } |
| 466 | if ledger := c.turnEventLedger(); ledger != nil && ledger.CurrentStatus() == event.TurnRecoveryRequired { |
| 467 | c.enterRecoveryLocked("terminal") |
| 468 | c.turns.finishingBound.endIdle() |
| 469 | c.mu.Unlock() |
| 470 | c.refreshRuntimeState(event.Event{}) |
| 471 | return |
| 472 | } |
| 473 | if c.closed { |
| 474 | c.turns.lastToken = c.turns.token |
| 475 | c.turns.phase = session.RuntimeClosed |
| 476 | c.turns.turnID = "" |
| 477 | c.noteExecutionLocked(session.RuntimeIdle, "") |
| 478 | c.turns.finishingBound.endIdle() |
| 479 | c.mu.Unlock() |
| 480 | c.finalizeControllerClose() |
| 481 | c.refreshRuntimeState(event.Event{}) |
| 482 | return |
| 483 | } |
| 484 | if authErr := c.authentication.admissionError(); authErr != nil { |
| 485 | for _, pending := range c.turns.pending { |
| 486 | if pending.goalRound != nil { |
| 487 | pending.goalRound.setResult(authErr, false) |
| 488 | } |
| 489 | } |
| 490 | c.turns.pending = nil |
| 491 | c.turns.wake = false |
| 492 | c.turns.lastToken = c.turns.token |
| 493 | c.turns.phase = session.RuntimeIdle |
| 494 | c.turns.turnID = "" |
| 495 | c.noteExecutionLocked(session.RuntimeIdle, "") |
| 496 | c.turns.finishingBound.endIdle() |
| 497 | c.mu.Unlock() |
| 498 | c.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Code: "authentication_not_ready", Text: authErr.Error()}) |
| 499 | c.refreshRuntimeState(event.Event{}) |
| 500 | return |
| 501 | } |
| 502 | next, ok := c.popNextPendingLocked() |
| 503 | if !ok { |
| 504 | c.turns.lastToken = c.turns.token |
| 505 | c.turns.phase = session.RuntimeIdle |
| 506 | c.turns.turnID = "" |
| 507 | c.noteExecutionLocked(session.RuntimeIdle, "") |
| 508 | c.turns.finishingBound.endIdle() |
| 509 | c.mu.Unlock() |
| 510 | c.maybeDispatchInbox() |
| 511 | c.refreshRuntimeState(event.Event{}) |
| 512 | return |
| 513 | } |
| 514 | ctx, cancel, admitted := c.startTurnLocked(context.Background(), next) |
| 515 | if !admitted { |
| 516 | c.enterRecoveryLocked("execution_owner_lost") |
| 517 | c.turns.finishingBound.endIdle() |
| 518 | c.mu.Unlock() |
| 519 | c.refreshRuntimeState(event.Event{}) |
| 520 | return |
| 521 | } |
| 522 | c.mu.Unlock() |
| 523 | if next.onStart != nil { |
| 524 | next.onStart() |
| 525 | } |
| 526 | c.spawnGuardedTurn(ctx, cancel, next) |
| 527 | c.refreshRuntimeState(event.Event{}) |
| 528 | }() |
| 529 | c.emitTurnDoneEvent(err, cancelRequested, completion) |
| 530 | } |
| 531 | |
| 532 | func (c *Controller) emitTurnDoneEvent(err error, cancelRequested bool, completion *guardedTurnCompletion) { |
| 533 | c.inbox.mu.Lock() |
| 534 | activeInboxID := "" |
| 535 | for id := range c.inbox.activeItemIDs { |
| 536 | activeInboxID = id |
| 537 | break |
| 538 | } |
| 539 | c.inbox.mu.Unlock() |
| 540 | done := event.Event{ |
| 541 | Kind: event.TurnDone, |
| 542 | Err: err, |
| 543 | Cancelled: cancelRequested, |
| 544 | Outcome: turnOutcome(err), |
| 545 | CheckpointTurn: c.validatedCheckpointTurn(completion), |
| 546 | Receipt: c.executor.CompletionReceipt(), |
| 547 | ItemID: activeInboxID, |
| 548 | } |
| 549 | if done.CheckpointTurn != nil { |
| 550 | changes := completion.checkpoint.store.FreezeTurnChanges(*done.CheckpointTurn) |
| 551 | if done.Receipt == nil && (len(changes.Files) > 0 || len(changes.Reasons) > 0) { |
| 552 | done.Receipt = &event.CompletionReceipt{AssessmentKind: "facts", Verdict: "unknown"} |
| 553 | } |
| 554 | if done.Receipt != nil { |
| 555 | receipt := *done.Receipt |
| 556 | receipt.Diff = changes.Summary() |
| 557 | receipt.Interrupted = cancelRequested |
| 558 | done.Receipt = &receipt |
| 559 | } |
| 560 | } |
| 561 | done.Receipt = bindCompletionLogSources(done.Receipt, c.History()) |
| 562 | done = c.applyTurnDoneProtocol(done, cancelRequested) |
| 563 | c.applyToolRecoveryTurnStatus(&done, completion) |
| 564 | var readErr *agent.IncompleteReadError |
| 565 | if errors.As(err, &readErr) { |
| 566 | done.ReadPause = readErr.Pause |
| 567 | } |
| 568 | done.Diagnostic = provider.DiagnoseFailure(err) |
| 569 | done.Detail = provider.FailureDiagnosticDetail(done.Diagnostic) |
| 570 | if !cancelRequested { |
| 571 | done.ProtocolRecovery = c.executor.PendingProtocolRecovery() |
| 572 | } |
| 573 | var readinessErr *agent.FinalReadinessError |
| 574 | if errors.As(err, &readinessErr) { |
| 575 | done.Readiness = &event.FinalReadiness{Attempts: readinessErr.Attempts, Missing: append([]string(nil), readinessErr.Missing...)} |
| 576 | } |
| 577 | c.onInboxTurnDone() |
| 578 | c.sink.Emit(done) |
| 579 | } |
| 580 | |
| 581 | func (c *Controller) startCancellationWatchdog(done chan struct{}) { |
| 582 | if c == nil || done == nil { |
| 583 | return |
| 584 | } |
| 585 | go func() { |
| 586 | timer := time.NewTimer(c.cancellationGrace()) |
| 587 | defer timer.Stop() |
| 588 | select { |
| 589 | case <-timer.C: |
| 590 | case <-done: |
| 591 | return |
| 592 | } |
| 593 | |
| 594 | c.mu.Lock() |
| 595 | stillRunning := c.turns.done == done && (c.turns.phase == session.RuntimeRunning || c.turns.phase == session.RuntimeCancelling) |
| 596 | turnID := c.turns.turnID |
| 597 | if stillRunning { |
| 598 | c.turns.recoveryFanout = true |
| 599 | c.enterRecoveryLocked("cancellation_grace_expired") |
| 600 | } |
| 601 | c.mu.Unlock() |
| 602 | if !stillRunning { |
| 603 | return |
| 604 | } |
| 605 | if turnID == "" { |
| 606 | if ledger := c.turnEventLedger(); ledger != nil { |
| 607 | turnID = ledger.ActiveTurnID() |
| 608 | } |
| 609 | } |
| 610 | recovery := &event.RecoveryStatus{ |
| 611 | State: "recovery_required", |
| 612 | Phase: "cancellation_grace_expired", |
| 613 | Reason: "cancellation_grace_expired", |
| 614 | RequiresUserDecision: true, |
| 615 | } |
| 616 | _ = c.emitTurnEventChecked(event.Event{ |
| 617 | Kind: event.TurnDone, |
| 618 | TurnID: turnID, |
| 619 | Status: event.TurnRecoveryRequired, |
| 620 | Cancelled: true, |
| 621 | Outcome: "unknown", |
| 622 | Recovery: recovery, |
| 623 | }) |
| 624 | c.mu.Lock() |
| 625 | c.turns.recoveryFanout = false |
| 626 | closing := c.closed && c.turns.done == nil && !c.finalizingLocked() |
| 627 | c.mu.Unlock() |
| 628 | if closing { |
| 629 | c.finalizeControllerClose() |
| 630 | } |
| 631 | c.refreshRuntimeState(event.Event{}) |
| 632 | }() |
| 633 | } |
| 634 |