返回 DeepSeek-Reasonix
turn_events.go
根目录 / internal / control / turn_events.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "sync"
9 "sync/atomic"
10 "time"
11
12 "reasonix/internal/agent"
13 "reasonix/internal/event"
14 "reasonix/internal/evidence"
15 "reasonix/internal/session"
16 "reasonix/internal/sessioninbox"
17 "reasonix/internal/transcript"
18 "reasonix/internal/turnevent"
19 )
20
21 // turnEventSink persists lifecycle envelopes before frontend publication.
22 // Provider-facing transcript messages remain a separate artifact.
23 type turnEventSink struct {
24 event.AuditForwarder
25 innerMu sync.RWMutex
26 inner event.Sink
27 stream event.Sink
28 c *Controller
29 publish atomic.Int32
30 }
31
32 type turnEventDurableSink struct{ owner *turnEventSink }
33
34 // turnEventState has an independent lock so ledger I/O never holds c.mu.
35 type turnEventState struct {
36 mu sync.RWMutex
37 ledger *turnevent.Ledger
38 err error
39 v3 *session.Session
40 v3Path string
41 // v3Runtime pins the session instance the cached store belongs to. A
42 // reclaim closes the old runtime and a later takeover re-opens the same
43 // identity, so the path key alone would keep serving the closed store.
44 v3Runtime *session.Runtime
45 v3Release func(context.Context) error
46 v3Err error
47 projection *transcript.Projection
48 projectionErr error
49 commitMu sync.Mutex
50 persistMu sync.Mutex
51 projectionPath string
52 pendingCheckpoint *transcript.Checkpoint
53 projectionPersistedThrough uint64
54 projectionWriteErr error
55 volatileTodos []event.Todo
56 volatileTodoWritten bool
57 // pendingExecutionCommit is prepared by an unpublished hot-rebuild
58 // candidate and consumed atomically with Runtime execution activation.
59 // commitMu owns it and its queue reservation.
60 pendingExecutionCommit *session.PreparedBatch
61 pendingTermination *TerminationPlan
62 turnMessageIDs map[string]bool
63 finalizedTurn string
64 terminationBoundary *terminationBoundary
65 openStream openStreamOutput
66 }
67
68 // projectVolatileTodo keeps the same event-derived projection for controllers
69 // that have not acquired a session path yet. It is a cache of successful
70 // lifecycle events, never a second writable todo state machine.
71 func (c *Controller) projectVolatileTodo(e event.Event) {
72 if c == nil {
73 return
74 }
75 c.turnEvents.mu.Lock()
76 defer c.turnEvents.mu.Unlock()
77 switch {
78 case e.Kind == event.TurnStarted:
79 c.turnEvents.volatileTodos = []event.Todo{}
80 c.turnEvents.volatileTodoWritten = false
81 case e.Kind == event.ToolResult && e.Tool.TodoWritten:
82 c.turnEvents.volatileTodos = append([]event.Todo(nil), e.Tool.Todos...)
83 c.turnEvents.volatileTodoWritten = true
84 }
85 }
86
87 func (c *Controller) volatileTodoState() ([]event.Todo, bool) {
88 if c == nil {
89 return []event.Todo{}, false
90 }
91 c.turnEvents.mu.RLock()
92 defer c.turnEvents.mu.RUnlock()
93 return append([]event.Todo(nil), c.turnEvents.volatileTodos...), c.turnEvents.volatileTodoWritten
94 }
95
96 func newTurnEventSink(inner event.Sink, c *Controller) *turnEventSink {
97 s := &turnEventSink{inner: inner, c: c}
98 s.stream = event.Coalesce(&turnEventDurableSink{owner: s}, event.DefaultStreamDeltaWindow)
99 s.AuditForwarder = event.AuditForwarder{Inner: s.stream}
100 return s
101 }
102
103 func (s *turnEventSink) InboxChanged(snap sessioninbox.InboxSnapshot) {
104 if s != nil {
105 notifyInboxChanged(s.innerSnapshot(), snap)
106 }
107 }
108
109 var _ event.OptionalSinkCapabilities = (*turnEventSink)(nil)
110 var _ event.CheckedSink = (*turnEventSink)(nil)
111 var _ event.OptionalSinkCapabilities = (*turnEventDurableSink)(nil)
112 var _ event.CheckedSink = (*turnEventDurableSink)(nil)
113
114 func (s *turnEventSink) Emit(e event.Event) {
115 if s == nil {
116 return
117 }
118 s.observe(e)
119 if turnEventSynchronousBarrier(e.Kind) {
120 if err := event.EmitChecked(s.stream, e); err != nil {
121 s.fail(err)
122 }
123 return
124 }
125 s.stream.Emit(e)
126 }
127
128 // observe feeds every raw event to the ledger's routing and to the liveness
129 // tracker before ordering, so silence is measured from real emission time.
130 func (s *turnEventSink) observe(e event.Event) {
131 if s.c == nil {
132 return
133 }
134 if ledger := s.c.turnEventLedger(); ledger != nil {
135 ledger.ObserveRawEvent(e)
136 }
137 s.c.turnEvents.openStream.observe(e)
138 s.c.liveness.observe(e, time.Now())
139 }
140
141 func turnEventSynchronousBarrier(kind event.Kind) bool {
142 switch kind {
143 case event.ToolDispatch, event.ToolStarted, event.ToolResult, event.AskRequest, event.ApprovalRequest,
144 event.MCPInteractionRequest, event.PromptAnswered, event.TurnStatusChanged,
145 event.TurnStarted, event.TurnDone, event.SessionOperation:
146 return true
147 default:
148 return false
149 }
150 }
151
152 func (s *turnEventSink) EmitChecked(e event.Event) error {
153 if s == nil {
154 return nil
155 }
156 s.observe(e)
157 var err error
158 if s.publish.Load() > 0 && e.Kind == event.PromptAnswered {
159 // A frontend may answer during prompt publication, so the coalescer cannot
160 // wait on itself. Only that already-ordered PromptAnswered barrier may use
161 // this re-entrant path; other checked events preserve coalescer ordering.
162 err = (&turnEventDurableSink{owner: s}).EmitChecked(e)
163 } else {
164 err = event.EmitChecked(s.stream, e)
165 }
166 if err != nil {
167 s.fail(err)
168 }
169 return err
170 }
171
172 func (s *turnEventSink) fail(err error) {
173 if s != nil && s.c != nil && err != nil {
174 s.c.failTurnEventLedger(err)
175 }
176 }
177
178 func (s *turnEventSink) innerSnapshot() event.Sink {
179 if s == nil {
180 return nil
181 }
182 s.innerMu.RLock()
183 defer s.innerMu.RUnlock()
184 return s.inner
185 }
186
187 func (s *turnEventSink) setInner(inner event.Sink) {
188 if s == nil {
189 return
190 }
191 s.innerMu.Lock()
192 s.inner = inner
193 s.innerMu.Unlock()
194 }
195
196 func (s *turnEventSink) publishInner(e event.Event) {
197 inner := s.innerSnapshot()
198 if inner == nil {
199 return
200 }
201 s.publish.Add(1)
202 defer s.publish.Add(-1)
203 inner.Emit(e)
204 }
205
206 // emitChecked persists before publish and returns durability failures to the
207 // admission boundary. It also suppresses the executor's duplicate TurnStarted
208 // because the controller has already committed that transition before the
209 // provider goroutine is launched.
210 func (s *turnEventSink) persistAndPublish(e event.Event) error {
211 if s == nil || s.c == nil {
212 return nil
213 }
214 if e.Kind == event.SessionOperation {
215 if err := s.c.persistMaintenanceOperation(e); err != nil {
216 return err
217 }
218 return s.publishOutsideTurn(s.c.turnEventLedger(), e)
219 }
220 if e.RecoveryCheckpoint {
221 return s.c.CheckpointSession(context.Background(), agent.CheckpointBeforeTopTool)
222 }
223 if err := s.c.stampToolRecoveryEvent(e); err != nil {
224 return err
225 }
226 ledger := s.c.turnEventLedger()
227 if ledger == nil {
228 s.c.projectVolatileTodo(e)
229 s.c.refreshRuntimeState(e)
230 s.publishInner(e)
231 return nil
232 }
233 if staleTurnStatus(e, ledger) {
234 return nil
235 }
236 // Outside-turn notices are not lifecycle records and must pass through after
237 // bootstrap or a terminal event.
238 if ledger.ActiveTurnID() == "" {
239 return s.publishOutsideTurn(ledger, e)
240 }
241 if e.Kind == event.TurnStarted && ledger.CurrentStatus() == event.TurnInProgress {
242 return nil
243 }
244 status := publicationTurnStatus(e, ledger)
245 if e.WriteIntent {
246 return nil
247 }
248 // No frontend callback runs while commitMu is held. Prompt publication
249 // can synchronously reenter this sink to append PromptAnswered.
250 stamped, envelope, ok, err := s.commitEnvelope(ledger, e, status)
251 if err != nil {
252 return err
253 }
254 if !ok {
255 return nil
256 }
257 if err := s.c.flushSubmissionStart(s.c.submissionAdmissionContext(), e.Kind); err != nil {
258 return err
259 }
260 projectionSaved := true
261 if e.Kind == event.TurnDone {
262 if store := s.c.sessionEventStore(); store != nil {
263 if _, err := store.Flush(context.Background()); err != nil {
264 return err
265 }
266 }
267 s.c.captureTranscriptCheckpoint(ledger, envelope.TranscriptDigest)
268 if err := s.c.persistTranscriptCheckpoint(ledger); err != nil {
269 projectionSaved = false
270 slog.Warn("controller: persist transcript display checkpoint", "err", err)
271 }
272 }
273 if _, runtime, exclusive := s.c.v3Binding(); exclusive && runtime != nil {
274 if err := runtime.PublishTranscriptFrame(envelope); err != nil {
275 return err
276 }
277 }
278 s.c.recordTurnLifecycle(stamped)
279 s.c.refreshRuntimeState(stamped)
280 s.publishInner(stamped)
281 if e.Kind == event.TurnDone && !ledger.ProjectionAckRequired() && projectionSaved {
282 if err := ledger.AcknowledgeProjection(stamped.TurnID); err != nil {
283 return err
284 }
285 }
286 return nil
287 }
288
289 func lateBusinessEvent(kind event.Kind) bool {
290 switch kind {
291 case event.ToolDispatch, event.ToolStarted, event.ToolProgress, event.ToolResult,
292 event.AskRequest, event.ApprovalRequest, event.MCPInteractionRequest,
293 event.PromptAnswered, event.TurnStarted, event.TurnStatusChanged, event.TurnDone:
294 return true
295 default:
296 return false
297 }
298 }
299
300 func (s *turnEventSink) commitEnvelope(ledger *turnevent.Ledger, e event.Event, status event.TurnStatus) (event.Event, turnevent.Envelope, bool, error) {
301 if e.Kind == event.TurnDone {
302 s.c.snapshotMu.Lock()
303 defer s.c.snapshotMu.Unlock()
304 }
305 s.c.turnEvents.commitMu.Lock()
306 defer s.c.turnEvents.commitMu.Unlock()
307 if s.c.discardLateTurnEvent(e) {
308 slog.Info("controller: discarded late turn event", "kind", e.Kind, "turnId", e.TurnID)
309 return e, turnevent.Envelope{}, false, nil
310 }
311 ctx := context.Background()
312 if e.Kind == event.TurnDone {
313 var cancel context.CancelFunc
314 ctx, cancel = context.WithTimeout(ctx, terminationFlushTimeout)
315 defer cancel()
316 }
317 if e.Kind == event.Notice && e.Code == event.NoticeCodeMCPToolsList && e.MessageID == "" {
318 if store := s.c.sessionEventStore(); store != nil {
319 e.MessageID = fmt.Sprintf("notice:%s:%d", store.ID(), store.EventSequence()+1)
320 }
321 }
322 if err := s.c.appendSessionEventLocked(ctx, e); err != nil {
323 return e, turnevent.Envelope{}, false, err
324 }
325 if e.Kind == event.TurnDone {
326 e.ReadCompletion = s.c.updateTurnLedgerTranscript(ledger)
327 }
328 stamped, envelope, ok, err := ledger.AppendEnvelope(e, status)
329 if err != nil || !ok || stamped.Sequence == 0 {
330 return stamped, envelope, ok, err
331 }
332 s.c.turnEvents.mu.RLock()
333 projection := s.c.turnEvents.projection
334 s.c.turnEvents.mu.RUnlock()
335 _, _, exclusive := s.c.v3Binding()
336 if projection != nil && !exclusive {
337 if projectionErr := projection.Apply(envelope); projectionErr != nil {
338 s.c.turnEvents.mu.Lock()
339 s.c.turnEvents.projectionErr = projectionErr
340 s.c.turnEvents.mu.Unlock()
341 }
342 }
343 return stamped, envelope, true, nil
344 }
345
346 func (s *turnEventDurableSink) Emit(e event.Event) {
347 _ = s.EmitChecked(e)
348 }
349
350 func (s *turnEventDurableSink) EmitChecked(e event.Event) error {
351 if s == nil || s.owner == nil {
352 return nil
353 }
354 err := s.owner.persistAndPublish(e)
355 if err == nil {
356 return nil
357 }
358 if classifyCommitError(err) == commitLifecycle {
359 slog.Info("controller: lifecycle event commit", "err", err, "kind", e.Kind)
360 return nil
361 }
362 if e.Kind == event.TurnDone && classifyCommitError(err) != commitOwnership {
363 s.owner.c.disarmGoalLifecycle("persistence-error")
364 s.owner.c.mu.Lock()
365 s.owner.c.enterRecoveryLocked("terminal_commit_failed")
366 s.owner.c.mu.Unlock()
367 if _, runtime, exclusive := s.owner.c.v3Binding(); exclusive && runtime != nil {
368 runtime.Transcript().PersistenceFailed()
369 }
370 }
371 // Async stream callers cannot observe checked errors. Fail the Turn here so
372 // a poisoned WAL immediately cancels provider, prompt, and process work.
373 slog.Error("controller: append turn event ledger", "err", err, "kind", e.Kind)
374 s.owner.fail(err)
375 if e.Kind == event.TurnDone {
376 // The durable terminal failed, so publish a sequence-free control-plane
377 // failure only to release UI state. It is never treated as ledger truth.
378 e.Err = errors.Join(e.Err, err)
379 e.Status = event.TurnFailed
380 if inner := s.owner.innerSnapshot(); inner != nil {
381 inner.Emit(e)
382 }
383 }
384 return err
385 }
386
387 func (s *turnEventDurableSink) inner() event.Sink {
388 if s == nil || s.owner == nil {
389 return nil
390 }
391 return s.owner.innerSnapshot()
392 }
393
394 func (s *turnEventDurableSink) RecordDelegationAudit(a evidence.DelegationAudit) {
395 event.RecordDelegationAudit(s.inner(), a)
396 }
397 func (s *turnEventDurableSink) RecordReadinessAudit(a evidence.ReadinessAudit) {
398 event.RecordReadinessAudit(s.inner(), a)
399 }
400 func (s *turnEventDurableSink) RecordAnchorSafetyAudit(a event.AnchorSafetyAudit) {
401 event.RecordAnchorSafetyAudit(s.inner(), a)
402 }
403 func (s *turnEventDurableSink) RecordTurnCompletion() { event.RecordTurnCompletion(s.inner()) }
404 func (s *turnEventDurableSink) RecordContractShadow(a event.ContractShadowAudit) {
405 event.RecordContractShadow(s.inner(), a)
406 }
407 func (s *turnEventDurableSink) RecordCompletionReport(a event.CompletionReportAudit) {
408 event.RecordCompletionReport(s.inner(), a)
409 }
410 func (s *turnEventDurableSink) RecordMemoryRecall(a event.MemoryRecallAudit) {
411 event.RecordMemoryRecall(s.inner(), a)
412 }
413 func (s *turnEventDurableSink) RecordDelegationAdmission(a event.DelegationAdmissionAudit) {
414 event.RecordDelegationAdmission(s.inner(), a)
415 }
416 func (s *turnEventDurableSink) RecordOutcomeProgress(a evidence.OutcomeSample) {
417 event.RecordOutcomeProgress(s.inner(), a)
418 }
419 func (s *turnEventDurableSink) RecordProtocolRecovery(a event.ProtocolRecoveryAudit) {
420 event.RecordProtocolRecovery(s.inner(), a)
421 }
422 func (s *turnEventDurableSink) RecordWorkspaceMutation(a event.WorkspaceMutation) {
423 event.RecordWorkspaceMutation(s.inner(), a)
424 }
425 func (s *turnEventDurableSink) RecordRunBudget(a event.RunBudgetSample) {
426 event.RecordRunBudget(s.inner(), a)
427 }
428 func (s *turnEventDurableSink) RecordSubagentLifecycle(a event.SubagentLifecycleInfo) {
429 event.RecordSubagentLifecycle(s.inner(), a)
430 }
431
432 func terminalTurnStatus(e event.Event) event.TurnStatus {
433 if e.Recovery != nil && e.Recovery.State == "recovery_required" {
434 return event.TurnRecoveryRequired
435 }
436 if e.Cancelled || errors.Is(e.Err, context.Canceled) {
437 return event.TurnInterrupted
438 }
439 if e.Err != nil {
440 return event.TurnFailed
441 }
442 return event.TurnCompleted
443 }
444
445 func (c *Controller) turnEventLedger() *turnevent.Ledger {
446 if c == nil {
447 return nil
448 }
449 c.turnEvents.mu.RLock()
450 defer c.turnEvents.mu.RUnlock()
451 return c.turnEvents.ledger
452 }
453
454 func (c *Controller) turnEventLedgerError() error {
455 if c == nil {
456 return nil
457 }
458 c.turnEvents.mu.RLock()
459 defer c.turnEvents.mu.RUnlock()
460 return c.turnEvents.err
461 }
462
463 func (c *Controller) applyTurnDoneProtocol(done event.Event, cancelRequested bool) event.Event {
464 if cancelRequested {
465 // Interruption is a terminal state, not a send failure; partial text is
466 // already display-only by this point.
467 done.Err = nil
468 }
469 return done
470 }
471
472 func (c *Controller) turnEventRuntimeStatus() (string, event.TurnStatus, uint64, uint64) {
473 ledger := c.turnEventLedger()
474 if ledger == nil {
475 return "", "", 0, 0
476 }
477 latest, replayAfter := ledger.ProjectionCursor()
478 return ledger.ActiveTurnID(), ledger.CurrentStatus(), latest, replayAfter
479 }
480
481 func (c *Controller) rebindTurnEvents(sessionPath string) {
482 defer c.refreshRuntimeState(event.Event{})
483 if c == nil {
484 return
485 }
486 desiredV3Path := sessionDirectory(sessionPath)
487 ledgerID := agent.BranchID(sessionPath)
488 var desiredRuntime *session.Runtime
489 if _, runtime, _ := c.v3Binding(); runtime != nil {
490 ref := runtime.Ref()
491 desiredV3Path = "session:" + ref.HostID + "/" + ref.SessionID
492 ledgerID = ref.SessionID
493 desiredRuntime = runtime
494 }
495 c.turnEvents.mu.RLock()
496 currentV3, currentV3Path := c.turnEvents.v3, c.turnEvents.v3Path
497 currentV3Runtime := c.turnEvents.v3Runtime
498 c.turnEvents.mu.RUnlock()
499 v3, releaseV3, v3Err := currentV3, (func(context.Context) error)(nil), error(nil)
500 // The runtime pin matters for exclusive sessions: a reclaim closes the
501 // old instance and the takeover re-opens the same identity, so the path
502 // alone cannot tell a live store from the closed one it replaced.
503 if currentV3 == nil || currentV3Path != desiredV3Path || currentV3Runtime != desiredRuntime {
504 v3, releaseV3, v3Err = c.openSessionEventStore(sessionPath)
505 }
506 ledger, err := c.openTurnLedger(sessionPath, ledgerID, v3Err)
507 if err != nil {
508 // Normalize platform-specific open errors behind the same storage
509 // sentinel used by append failures. Keep the original error in the
510 // chain so unsupported-schema callers can still inspect its type.
511 err = fmt.Errorf("%w: %w", turnevent.ErrTurnLedgerUnavailable, err)
512 slog.Warn("controller: open v3 session event store", "err", err, "session", agent.BranchID(sessionPath))
513 c.turnEvents.mu.Lock()
514 previousV3, previousRelease := c.turnEvents.v3, c.turnEvents.v3Release
515 previousLedger := c.turnEvents.ledger
516 c.turnEvents.ledger = nil
517 c.turnEvents.err = err
518 c.turnEvents.v3 = nil
519 c.turnEvents.v3Path = ""
520 c.turnEvents.v3Runtime = nil
521 c.turnEvents.v3Release = nil
522 c.turnEvents.v3Err = err
523 c.turnEvents.mu.Unlock()
524 if previousLedger != nil {
525 if closeErr := previousLedger.Close(); closeErr != nil {
526 slog.Warn("controller: close ledger after failed rebind", "err", closeErr)
527 }
528 }
529 // Fail admission closed without losing the compatibility writer's
530 // cleanup owner. Service-backed runtimes remain host-owned.
531 if previousV3 != nil && !c.sessionEngineEnabled() {
532 var closeErr error
533 if previousRelease != nil {
534 closeErr = previousRelease(context.Background())
535 } else {
536 closeErr = previousV3.Close(context.Background())
537 }
538 if closeErr != nil {
539 slog.Warn("controller: close session after failed rebind", "err", closeErr)
540 }
541 }
542 return
543 }
544 c.turnEvents.mu.Lock()
545 c.turnEvents.volatileTodos = []event.Todo{}
546 c.turnEvents.volatileTodoWritten = false
547 previous := c.turnEvents.ledger
548 previousV3 := c.turnEvents.v3
549 previousV3Release := c.turnEvents.v3Release
550 c.turnEvents.ledger = ledger
551 c.turnEvents.err = nil
552 c.turnEvents.v3 = v3
553 c.turnEvents.v3Path = desiredV3Path
554 c.turnEvents.v3Runtime = desiredRuntime
555 if releaseV3 != nil {
556 c.turnEvents.v3Release = releaseV3
557 }
558 c.turnEvents.v3Err = nil
559 if c.turnEvents.projection != nil {
560 c.turnEvents.projection.CloseFollowers()
561 }
562 c.turnEvents.projection = nil
563 c.turnEvents.projectionErr = nil
564 c.turnEvents.projectionPath = sessionPath
565 c.turnEvents.pendingCheckpoint = nil
566 c.turnEvents.projectionPersistedThrough = 0
567 c.turnEvents.projectionWriteErr = nil
568 c.turnEvents.mu.Unlock()
569 if !c.sessionEngineEnabled() {
570 c.bindAttachmentService()
571 }
572 var projection *transcript.Projection
573 var projectionErr error
574 if !c.sessionEngineEnabled() {
575 projection, projectionErr = c.restoreTranscriptProjection(sessionPath, ledger)
576 }
577 c.turnEvents.mu.Lock()
578 c.turnEvents.projection, c.turnEvents.projectionErr = projection, projectionErr
579 c.turnEvents.mu.Unlock()
580 if previous != nil && previous != ledger {
581 if closeErr := previous.Close(); closeErr != nil {
582 slog.Warn("controller: close previous turn event ledger", "err", closeErr)
583 }
584 }
585 // An exclusive v3 handle belongs to SessionRuntime. Runtime publication
586 // closes the exact previous instance through SessionService after the new
587 // binding is visible; this compatibility cleanup must never close it early.
588 if previousV3 != nil && previousV3 != v3 && !c.sessionEngineEnabled() {
589 var closeErr error
590 if previousV3Release != nil {
591 closeErr = previousV3Release(context.Background())
592 } else {
593 closeErr = previousV3.Close(context.Background())
594 }
595 if closeErr != nil {
596 slog.Warn("controller: flush and close previous v3 session", "err", closeErr)
597 }
598 }
599 }
600
601 func (c *Controller) openTurnLedger(path, id string, storeErr error) (*turnevent.Ledger, error) {
602 if storeErr == nil && c.NativeLegacySession() && path != "" {
603 return turnevent.Open(path, id)
604 }
605 return turnevent.NewMemory(id), storeErr
606 }
607
608 func classifyCommitError(err error) commitFailureKind {
609 if err == nil {
610 return commitOK
611 }
612 if errors.Is(err, errTerminationDurability) {
613 return commitUnexpected
614 }
615 if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
616 return commitLifecycle
617 }
618 if errors.Is(err, session.ErrStaleActivity) || errors.Is(err, session.ErrOperationConflict) {
619 return commitLifecycle
620 }
621 if errors.Is(err, session.ErrSessionNotRunning) || errors.Is(err, session.ErrStaleGeneration) ||
622 errors.Is(err, session.ErrStaleExecution) || errors.Is(err, session.ErrReadOnly) || errors.Is(err, session.ErrRuntimeRetiring) {
623 return commitOwnership
624 }
625 return commitUnexpected
626 }
627
628 type commitFailureKind int
629
630 const (
631 commitOK commitFailureKind = iota
632 commitLifecycle
633 commitOwnership
634 commitUnexpected
635 )
636
637 func (c *Controller) failTurnEventLedger(err error) {
638 defer c.refreshRuntimeState(event.Event{})
639 if c == nil || err == nil {
640 return
641 }
642 switch classifyCommitError(err) {
643 case commitLifecycle:
644 slog.Info("controller: lifecycle commit result", "err", err)
645 return
646 case commitOwnership:
647 c.turnEvents.mu.Lock()
648 if c.turnEvents.err == nil {
649 c.turnEvents.err = err
650 }
651 c.turnEvents.mu.Unlock()
652 c.signalTurnCancel()
653 c.promptOwner.CancelAll()
654 c.approval.clearAll()
655 return
656 }
657 c.turnEvents.mu.Lock()
658 if c.turnEvents.err == nil {
659 c.turnEvents.err = err
660 }
661 c.turnEvents.mu.Unlock()
662 c.signalTurnCancel()
663 c.promptOwner.CancelAll()
664 c.approval.clearAll()
665 }
666
667 // staleTurnStatus reports a status stamped for a turn that has since reached
668 // its terminal event; cancelling is sticky, so it must not reach the next turn.
669 func staleTurnStatus(e event.Event, ledger *turnevent.Ledger) bool {
670 return e.Kind == event.TurnStatusChanged && e.TurnID != "" && e.TurnID != ledger.ActiveTurnID()
671 }
672
673 // emitTurnStatus stamps the transition with the turn that requested it so the
674 // ledger can drop it if that turn already reached its terminal event.
675 func (c *Controller) emitTurnStatus(status event.TurnStatus, turnID string) {
676 if c == nil || status == "" {
677 return
678 }
679 c.sink.Emit(event.Event{Kind: event.TurnStatusChanged, Status: status, TurnID: turnID})
680 }
681
682 // emitTurnEventChecked reaches the lifecycle sink below the inbox observer so
683 // admission can fail closed on disk errors instead of starting an unledgered
684 // provider request. Lifecycle events do not participate in inbox notice logic.
685 func (c *Controller) emitTurnEventChecked(e event.Event) error {
686 if c == nil {
687 return nil
688 }
689 if e.ItemID != "" && e.TurnID == "" {
690 if identity, ok := c.promptOwner.Identity(e.ItemID); ok {
691 e.TurnID = identity.TurnID
692 e.PromptKind = string(identity.Kind)
693 }
694 } else if e.ItemID != "" && e.PromptKind == "" {
695 if identity, ok := c.promptOwner.Identity(e.ItemID); ok {
696 e.PromptKind = string(identity.Kind)
697 }
698 }
699 return event.EmitChecked(c.sink, e)
700 }
701
702 // SetTurnEventRoutingMetadata attaches desktop routing identity to lifecycle
703 // envelopes only. It never changes provider-visible prompts or tool schemas.
704 func (c *Controller) SetTurnEventRoutingMetadata(runtimeEpoch, submissionID string) {
705 c.promptEpochMu.Lock()
706 c.promptRuntimeEpoch = runtimeEpoch
707 c.promptEpochMu.Unlock()
708 if ledger := c.turnEventLedger(); ledger != nil {
709 ledger.RequireProjectionAck(true)
710 ledger.SetRoutingMetadata(runtimeEpoch, submissionID)
711 }
712 c.BindTranscriptRuntimeEpoch(runtimeEpoch)
713 }
714
715 // TurnEventsAfter returns the durable lifecycle suffix used by reconnecting
716 // frontends to close sequence gaps.
717 func (c *Controller) TurnEventsAfter(after uint64) ([]turnevent.Envelope, error) {
718 ledger := c.turnEventLedger()
719 if ledger == nil {
720 return []turnevent.Envelope{}, nil
721 }
722 return ledger.EventsAfter(after)
723 }
724
725 func (c *Controller) TurnEventReplay(after uint64) (turnevent.ReplayView, error) {
726 ledger := c.turnEventLedger()
727 if ledger == nil {
728 return turnevent.ReplayView{Events: []turnevent.Envelope{}}, nil
729 }
730 return ledger.Replay(after)
731 }
732
733 func (c *Controller) AcknowledgeTurnProjection(turnID string) error {
734 ledger := c.turnEventLedger()
735 if ledger == nil {
736 return nil
737 }
738 if err := c.persistTranscriptCheckpoint(ledger); err != nil {
739 return err
740 }
741 return ledger.AcknowledgeProjection(turnID)
742 }
743
744 func (c *Controller) ObserveTurnProjectionRetry() {
745 if ledger := c.turnEventLedger(); ledger != nil {
746 ledger.ObserveProjectionRetry()
747 }
748 }
749
750 func (c *Controller) PendingTurnProjections() []turnevent.PendingProjection {
751 ledger := c.turnEventLedger()
752 if ledger == nil {
753 return []turnevent.PendingProjection{}
754 }
755 return ledger.PendingProjections()
756 }
757
758 func (c *Controller) TurnEventMetrics() turnevent.MetricsSnapshot {
759 ledger := c.turnEventLedger()
760 if ledger == nil {
761 return turnevent.MetricsSnapshot{}
762 }
763 return ledger.MetricsSnapshot()
764 }
765
766 func (c *Controller) DrainTurnEventMetrics() turnevent.MetricsSnapshot {
767 ledger := c.turnEventLedger()
768 if ledger == nil {
769 return turnevent.MetricsSnapshot{}
770 }
771 return ledger.DrainMetrics()
772 }
773
774 func publicationTurnStatus(e event.Event, ledger *turnevent.Ledger) event.TurnStatus {
775 status := e.Status
776 if status == "" {
777 status = ledger.CurrentStatus()
778 }
779 switch e.Kind {
780 case event.TurnStarted:
781 status = event.TurnInProgress
782 case event.AskRequest, event.ApprovalRequest, event.MCPInteractionRequest:
783 status = event.TurnWaitingUser
784 case event.TurnDone:
785 status = terminalTurnStatus(e)
786 case event.TurnStatusChanged:
787 // The emitter supplied the exact transition in e.Status.
788 }
789 return status
790 }
791
791 lines GO