返回 DeepSeek-Reasonix
inbox.go
根目录 / internal / control / inbox.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "maps"
9 "path/filepath"
10 "strings"
11 "sync"
12 "time"
13
14 "reasonix/internal/event"
15 "reasonix/internal/sessioninbox"
16 "reasonix/internal/sessiontemp"
17 "reasonix/internal/store"
18 )
19
20 // TurnAdmission is the exported classification of TrySubmitInboxItem /
21 // TrySteerInboxItem results.
22 type TurnAdmission string
23
24 const (
25 AdmissionStarted TurnAdmission = "started"
26 AdmissionSteerAccepted TurnAdmission = "steer_accepted"
27 AdmissionQueuedFollowup TurnAdmission = "queued_followup"
28 AdmissionRejectedBusy TurnAdmission = "rejected_busy"
29 AdmissionRejectedRotating TurnAdmission = "rejected_rotating"
30 AdmissionRejectedClosed TurnAdmission = "rejected_closed"
31 AdmissionRejectedCapacity TurnAdmission = "rejected_capacity"
32 )
33
34 // InboxRequest is the frontend-facing enqueue payload.
35 type InboxRequest struct {
36 ExpectedSessionPath string // optional exact-session fence; never persisted
37 Intent sessioninbox.InboxIntent
38 Display string
39 Raw string
40 Submit string
41 Format string
42 Source string
43 Idempotency string
44 Invocations []InvocationRequest
45 Extra map[string]string
46 // FreezeRefs lists workspace-relative paths to freeze at enqueue time.
47 FreezeRefs []string
48 Attachments []SubmissionAttachment
49 }
50
51 // Inbox port on SessionAPI.
52 type Inbox interface {
53 EnqueueInbox(req InboxRequest) (sessioninbox.InboxReceipt, error)
54 InboxSnapshot() sessioninbox.InboxSnapshot
55 ReadInboxItem(id string) (sessioninbox.InboxItemMeta, sessioninbox.PromptEnvelope, error)
56 UpdateInboxItem(id string, display, raw, submit string) (sessioninbox.InboxItemMeta, error)
57 AppendInboxItem(id, text, idempotency string, extra map[string]string) (sessioninbox.InboxItemMeta, error)
58 DeleteInboxItem(id string) error
59 CancelWithInboxItems(ids []string, source string) error
60 CancelWithInboxItemsResult(ids []string, source string) (InboxCancelResult, error)
61 MoveInboxItem(id string, toIndex int) error
62 SetInboxPaused(paused bool) error
63 RetryInboxItem(id string) error
64 RefreshInboxReferences(id string) error
65 TrySubmitInboxItem(id string) (sessioninbox.InboxReceipt, error)
66 RunInboxTurn(ctx context.Context, id string) error
67 TrySteerInboxItem(id string) (sessioninbox.InboxReceipt, error)
68 TryEnqueueAndSteer(req InboxRequest) (sessioninbox.InboxReceipt, error)
69 TryEnqueueFollowup(req InboxRequest) (sessioninbox.InboxReceipt, error)
70 }
71
72 // Compile-time port satisfaction.
73 var _ Inbox = (*Controller)(nil)
74
75 // inboxState is controller-owned inbox wiring (disk store + active items).
76 type inboxState struct {
77 prepareMu sync.Mutex
78 // admissionMu serializes competing admission state machines. Snapshot
79 // recovery and completion never hold it across Store I/O.
80 admissionMu sync.Mutex
81 // scanMu joins autonomous sidecar reads at shutdown without waiting for a
82 // dispatcher that may itself retire this controller during host admission.
83 scanMu sync.Mutex
84 mu sync.Mutex
85 store *sessioninbox.Store
86 // tempLease pins the process-local inbox used by an exclusive v3 Runtime.
87 // It is not recovery state and is deleted with the session temp generation.
88 tempLease *sessiontemp.Lease
89 closed bool // seals new sidecar opens when controller teardown starts
90 // activeItemIDs includes the running follow-up and every accepted steer.
91 // TurnDone durable-acks the set so multi-steer rounds leave no orphans.
92 activeItemIDs map[string]struct{}
93 // activeOwnership mirrors activeItemIDs for lock-free recovery checks while
94 // the Store owns its transaction lock. admittingOwnership covers the narrow
95 // durable-claim -> active-registration transition.
96 activeOwnership sync.Map
97 admittingOwnership sync.Map
98 dispatching bool
99 dispatchPending bool
100 // Retry bookkeeping is guarded by mu. Retries are bounded so a persistent
101 // disk or materialization failure cannot create a hot background loop.
102 dispatchRetryAttempts int
103 dispatchRetryScheduled bool
104 // beforePreparedAdmission is a deterministic test hook for the gap between
105 // durable preparation and Controller admission. Production leaves it nil.
106 beforePreparedAdmission func()
107 // beforeCompletionSnapshot exposes the slow snapshot boundary without
108 // changing production behavior.
109 beforeCompletionSnapshot func()
110 // beforeCompletionAck exposes the ownership-to-ack boundary to race tests.
111 beforeCompletionAck func()
112 // beforeSnapshotRead exposes the final Store snapshot boundary to lock tests.
113 beforeSnapshotRead func()
114 // afterDispatchScan exposes the empty-scan boundary for lost-wakeup tests.
115 afterDispatchScan func(found bool)
116 // beforeDispatchSubmit injects a transient owner-level dispatch failure.
117 beforeDispatchSubmit func(itemID string) error
118 // scheduleDispatchRetry replaces the production timer in deterministic tests.
119 scheduleDispatchRetry func(delay time.Duration, retry func())
120 }
121
122 func (s *inboxState) trackActive(id string) {
123 if s == nil || id == "" {
124 return
125 }
126 if s.activeItemIDs == nil {
127 s.activeItemIDs = make(map[string]struct{})
128 }
129 s.activeOwnership.Store(id, struct{}{})
130 s.activeItemIDs[id] = struct{}{}
131 }
132
133 func (s *inboxState) untrackActive(id string) {
134 if s == nil || id == "" {
135 return
136 }
137 if s.activeItemIDs != nil {
138 delete(s.activeItemIDs, id)
139 }
140 s.activeOwnership.Delete(id)
141 }
142
143 func (s *inboxState) untrackActiveSet(ids []string) {
144 if s == nil {
145 return
146 }
147 for _, id := range ids {
148 if s.activeItemIDs != nil {
149 delete(s.activeItemIDs, id)
150 }
151 s.activeOwnership.Delete(id)
152 }
153 }
154
155 func (s *inboxState) clearActive() {
156 if s == nil {
157 return
158 }
159 s.activeItemIDs = nil
160 s.activeOwnership.Clear()
161 }
162
163 func (s *inboxState) trackAdmission(id string) {
164 if s != nil && id != "" {
165 s.admittingOwnership.Store(id, struct{}{})
166 }
167 }
168
169 func (s *inboxState) untrackAdmission(id string) {
170 if s != nil && id != "" {
171 s.admittingOwnership.Delete(id)
172 }
173 }
174
175 // ownsItem is intentionally lock-free: Store recovery calls it while holding
176 // its own transaction lock, and no Store -> Controller lock edge is allowed.
177 func (s *inboxState) ownsItem(id string) bool {
178 if s == nil || id == "" {
179 return false
180 }
181 if _, ok := s.admittingOwnership.Load(id); ok {
182 return true
183 }
184 _, ok := s.activeOwnership.Load(id)
185 return ok
186 }
187
188 func (s *inboxState) activeIDs() []string {
189 if s == nil || len(s.activeItemIDs) == 0 {
190 return nil
191 }
192 out := make([]string, 0, len(s.activeItemIDs))
193 for id := range s.activeItemIDs {
194 out = append(out, id)
195 }
196 return out
197 }
198
199 func (c *Controller) bindInboxStoreNotifications(st *sessioninbox.Store) {
200 if c == nil || st == nil {
201 return
202 }
203 st.OnChange(func(snap sessioninbox.InboxSnapshot) {
204 notifyInboxChanged(c.sink, snap)
205 })
206 }
207
208 // inboxBinding selects identity and storage ownership from the same runtime.
209 // An import/display path can coexist with a canonical binding during startup;
210 // it must not turn the canonical locator into a filesystem path.
211 func (c *Controller) inboxBinding() (locator string, temporary bool) {
212 if _, runtime, exclusive := c.v3Binding(); exclusive {
213 if runtime == nil {
214 return "", true
215 }
216 return "session-id:" + runtime.Ref().SessionID, true
217 }
218 return c.SessionPath(), false
219 }
220
221 // inboxDirectoryLocked retains the existing process-local storage lifetime.
222 // c.inbox.mu must be held by the caller.
223 func (c *Controller) inboxDirectoryLocked(path string, temporary bool) (string, error) {
224 if temporary {
225 if c.inbox.tempLease == nil {
226 lease, err := c.sessionTemp.Acquire()
227 if err != nil {
228 return "", fmt.Errorf("open runtime inbox: %w", err)
229 }
230 c.inbox.tempLease = lease
231 }
232 return store.SessionInboxDir(filepath.Join(c.inbox.tempLease.Dir(), "runtime-inbox.jsonl")), nil
233 }
234 return store.SessionInboxDir(path), nil
235 }
236
237 func (c *Controller) ensureInbox() (*sessioninbox.Store, error) {
238 path, temporary := c.inboxBinding()
239 c.inbox.mu.Lock()
240 defer c.inbox.mu.Unlock()
241 if path == "" {
242 return nil, fmt.Errorf("inbox requires a session identity")
243 }
244 if c.inbox.store != nil && c.inbox.store.SessionPath() == path {
245 return c.inbox.store, nil
246 }
247 if c.inbox.closed {
248 return nil, fmt.Errorf("controller inbox is closed")
249 }
250 if c.inbox.store != nil {
251 c.inbox.store.Close()
252 c.inbox.store = nil
253 }
254 dir, err := c.inboxDirectoryLocked(path, temporary)
255 if err != nil {
256 return nil, err
257 }
258 st, err := sessioninbox.OpenAt(path, dir, sessioninbox.Limits{})
259 if err != nil {
260 return nil, err
261 }
262 c.bindInboxStoreNotifications(st)
263 c.inbox.store = st
264 snap := st.Snapshot()
265 if snap.Recovered && snap.RecoveredN > 0 {
266 c.sink.Emit(event.Event{
267 Kind: event.Notice,
268 Level: event.LevelWarn,
269 Code: "inbox_recovered",
270 Text: fmt.Sprintf("Recovered %d pending instruction(s). Inbox is paused — review with /queue before resuming.", snap.RecoveredN),
271 })
272 sessioninbox.NoteRecovered(snap.RecoveredN)
273 }
274 return st, nil
275 }
276
277 // rebindInbox opens the inbox for the current session path. Safe across
278 // NewSession/Resume/SetSessionPath; does not copy items on fork.
279 func (c *Controller) rebindInbox() {
280 path, temporary := c.inboxBinding()
281 c.inbox.mu.Lock()
282 defer c.inbox.mu.Unlock()
283 if c.inbox.closed {
284 return
285 }
286 if c.inbox.store != nil {
287 if path != "" && c.inbox.store.SessionPath() == path {
288 return
289 }
290 // Pending work must remain inspectable if this session is reopened.
291 _ = c.inbox.store.PauseIfPending()
292 c.inbox.store.Close()
293 c.inbox.store = nil
294 c.inbox.clearActive()
295 }
296 if c.inbox.tempLease != nil {
297 c.inbox.tempLease.Release()
298 c.inbox.tempLease = nil
299 }
300 if path == "" {
301 return
302 }
303 dir, err := c.inboxDirectoryLocked(path, temporary)
304 if err != nil {
305 slog.Warn("controller: open runtime inbox", "err", err)
306 return
307 }
308 st, err := sessioninbox.OpenAt(path, dir, sessioninbox.Limits{})
309 if err != nil {
310 slog.Warn("controller: open session inbox", "err", err, "path", path)
311 return
312 }
313 c.bindInboxStoreNotifications(st)
314 c.inbox.store = st
315 snap := st.Snapshot()
316 if snap.Recovered && snap.RecoveredN > 0 {
317 // Emit after unlock via deferred sink call would race; emit here.
318 go func(n int) {
319 c.sink.Emit(event.Event{
320 Kind: event.Notice,
321 Level: event.LevelWarn,
322 Code: "inbox_recovered",
323 Text: fmt.Sprintf("Recovered %d pending instruction(s). Inbox is paused — review with /queue before resuming.", n),
324 })
325 }(snap.RecoveredN)
326 sessioninbox.NoteRecovered(snap.RecoveredN)
327 }
328 }
329
330 func (c *Controller) pauseInboxOnRotate() {
331 c.inbox.mu.Lock()
332 st := c.inbox.store
333 c.inbox.mu.Unlock()
334 if st != nil {
335 _ = st.PauseIfPending()
336 }
337 }
338
339 func (c *Controller) InboxSnapshot() sessioninbox.InboxSnapshot {
340 st, err := c.ensureInbox()
341 if err != nil {
342 return sessioninbox.InboxSnapshot{}
343 }
344 if recovered, recoverErr := st.RecoverOrphanedInFlightOwnedBy(c.inbox.ownsItem); recoverErr != nil {
345 slog.Warn("controller: recover orphaned inbox items", "err", recoverErr)
346 } else if recovered > 0 {
347 sessioninbox.NoteRecovered(recovered)
348 }
349 c.inbox.mu.Lock()
350 beforeSnapshotRead := c.inbox.beforeSnapshotRead
351 c.inbox.mu.Unlock()
352 if beforeSnapshotRead != nil {
353 beforeSnapshotRead()
354 }
355 return st.Snapshot()
356 }
357
358 func (c *Controller) ReadInboxItem(id string) (sessioninbox.InboxItemMeta, sessioninbox.PromptEnvelope, error) {
359 st, err := c.ensureInbox()
360 if err != nil {
361 return sessioninbox.InboxItemMeta{}, sessioninbox.PromptEnvelope{}, err
362 }
363 return st.ReadItem(id)
364 }
365
366 func (c *Controller) UpdateInboxItem(id, display, raw, submit string) (sessioninbox.InboxItemMeta, error) {
367 st, err := c.ensureInbox()
368 if err != nil {
369 return sessioninbox.InboxItemMeta{}, err
370 }
371 submit = strings.TrimSpace(firstNonEmptyStr(submit, raw, display))
372 display = firstNonEmptyStr(display, submit)
373 raw = firstNonEmptyStr(raw, submit)
374 meta, previous, err := st.ReadItem(id)
375 if err != nil {
376 return sessioninbox.InboxItemMeta{}, err
377 }
378 env := previous
379 env.DisplayText, env.RawText, env.SubmitText = display, raw, submit
380 if err := c.freezeInboxEnvelopeReferences(context.Background(), &env, submit, env.ExplicitRefs); err != nil {
381 return sessioninbox.InboxItemMeta{}, err
382 }
383 updated, err := st.UpdateItemIfVersion(id, env, sessioninbox.ContentVersion(meta))
384 if err != nil {
385 return sessioninbox.InboxItemMeta{}, err
386 }
387 return updated, nil
388 }
389
390 // AppendInboxItem atomically merges collect-mode text and binds the inbound
391 // platform message ID as an idempotency alias for the existing durable item.
392 func (c *Controller) AppendInboxItem(id, text, idempotency string, extra map[string]string) (sessioninbox.InboxItemMeta, error) {
393 st, err := c.ensureInbox()
394 if err != nil {
395 return sessioninbox.InboxItemMeta{}, err
396 }
397 meta, previous, err := st.ReadItem(id)
398 if err != nil {
399 return sessioninbox.InboxItemMeta{}, err
400 }
401 text = strings.TrimSpace(text)
402 if text == "" {
403 return sessioninbox.InboxItemMeta{}, sessioninbox.ErrEmpty
404 }
405 merged := strings.TrimSpace(previous.SubmitText)
406 if merged != "" {
407 merged += "\n" + text
408 } else {
409 merged = text
410 }
411 env := previous
412 env.DisplayText = merged
413 env.RawText = merged
414 env.SubmitText = merged
415 if len(extra) > 0 {
416 env.Extra = maps.Clone(extra)
417 }
418 if err := c.freezeInboxEnvelopeReferences(context.Background(), &env, merged, env.ExplicitRefs); err != nil {
419 return sessioninbox.InboxItemMeta{}, err
420 }
421 aliasEnv := sessioninbox.PromptEnvelope{
422 DisplayText: text,
423 RawText: text,
424 SubmitText: text,
425 Source: previous.Source,
426 Extra: maps.Clone(extra),
427 }
428 return st.UpdateItemWithIdempotencyIfVersion(id, env, idempotency, aliasEnv, sessioninbox.ContentVersion(meta))
429 }
430
431 func (c *Controller) DeleteInboxItem(id string) error {
432 c.inbox.admissionMu.Lock()
433 defer c.inbox.admissionMu.Unlock()
434 st, err := c.ensureInbox()
435 if err != nil {
436 return err
437 }
438 if _, recoverErr := st.RecoverOrphanedInFlightOwnedBy(c.inbox.ownsItem); recoverErr != nil {
439 slog.Warn("controller: recover inbox item before delete", "err", recoverErr, "id", id)
440 }
441 err = st.DeletePendingOrAcceptedItem(id)
442 if err == nil || errors.Is(err, sessioninbox.ErrNotFound) {
443 return nil
444 }
445 return err
446 }
447
448 func (c *Controller) MoveInboxItem(id string, toIndex int) error {
449 st, err := c.ensureInbox()
450 if err != nil {
451 return err
452 }
453 return st.MoveItem(id, toIndex)
454 }
455
456 func (c *Controller) SetInboxPaused(paused bool) error {
457 return c.setInboxPaused(paused, true)
458 }
459
460 // SetInboxPausedPassive changes pause state without starting a background turn.
461 // Blocking transports such as Bot own their render sink and drain explicitly.
462 func (c *Controller) SetInboxPausedPassive(paused bool) error {
463 return c.setInboxPaused(paused, false)
464 }
465
466 func (c *Controller) setInboxPaused(paused, dispatch bool) error {
467 st, err := c.ensureInbox()
468 if err != nil {
469 return err
470 }
471 if err := st.SetPaused(paused); err != nil {
472 return err
473 }
474 if paused {
475 sessioninbox.NotePaused()
476 } else if dispatch {
477 // On resume, try to dispatch if idle.
478 c.maybeDispatchInbox()
479 }
480 return nil
481 }
482
483 func (c *Controller) RetryInboxItem(id string) error {
484 return c.retryInboxItem(id, true)
485 }
486
487 // RetryInboxItemPassive requeues an item without detached background dispatch.
488 func (c *Controller) RetryInboxItemPassive(id string) error {
489 return c.retryInboxItem(id, false)
490 }
491
492 func (c *Controller) retryInboxItem(id string, dispatch bool) error {
493 st, err := c.ensureInbox()
494 if err != nil {
495 return err
496 }
497 if err := st.RetryItem(id); err != nil {
498 return err
499 }
500 if dispatch {
501 c.maybeDispatchInbox()
502 }
503 return nil
504 }
505
506 // TrySubmitInboxItem admits a queued item as a new turn when the session is idle.
507 func (c *Controller) TrySubmitInboxItem(id string) (sessioninbox.InboxReceipt, error) {
508 c.mu.Lock()
509 beforeDispatch := c.modelSettings.beforeInboxDispatch
510 c.mu.Unlock()
511 if beforeDispatch != nil {
512 release, err := beforeDispatch(c)
513 if err != nil {
514 return sessioninbox.InboxReceipt{}, err
515 }
516 if release != nil {
517 defer release()
518 }
519 }
520 c.inbox.admissionMu.Lock()
521 defer c.inbox.admissionMu.Unlock()
522 st, err := c.ensureInbox()
523 if err != nil {
524 return sessioninbox.InboxReceipt{}, err
525 }
526 meta, env, err := st.ReadItem(id)
527 if err != nil {
528 return sessioninbox.InboxReceipt{}, err
529 }
530 if meta.State != sessioninbox.StateQueued {
531 return sessioninbox.InboxReceipt{}, sessioninbox.ErrInvalidState
532 }
533 if st.Snapshot().Paused {
534 return sessioninbox.InboxReceipt{}, sessioninbox.ErrPaused
535 }
536 run, block, materializeErr := c.prepareInboxRun(env)
537 if materializeErr != nil {
538 return sessioninbox.InboxReceipt{}, materializeErr
539 }
540 if block != "" {
541 if err := st.TransitionPrepared(id, sessioninbox.ContentVersion(meta), sessioninbox.StateBlocked, block, true); err != nil {
542 return sessioninbox.InboxReceipt{}, err
543 }
544 return sessioninbox.InboxReceipt{}, fmt.Errorf("%w: %s", sessioninbox.ErrInvalidState, block)
545 }
546 // Persist the in-flight state before admission. Active tracking is installed
547 // only after Controller admission is reserved and before the turn can finish.
548 c.inbox.trackAdmission(id)
549 defer c.inbox.untrackAdmission(id)
550 if err := st.TransitionPrepared(id, sessioninbox.ContentVersion(meta), sessioninbox.StateRunning, "", true); err != nil {
551 return sessioninbox.InboxReceipt{}, err
552 }
553 c.inbox.mu.Lock()
554 beforeAdmission := c.inbox.beforePreparedAdmission
555 c.inbox.mu.Unlock()
556 if beforeAdmission != nil {
557 beforeAdmission()
558 }
559 // Start the classified envelope directly. Submit would parse @tokens again
560 // and mix live workspace bytes with the enqueue-time snapshot.
561 result := c.submitPreparedInboxTurn(id, run)
562 if result != turnStarted {
563 if err := st.SetState(id, sessioninbox.StateQueued, ""); err != nil {
564 _ = st.ForcePause(true, 1)
565 return sessioninbox.InboxReceipt{}, err
566 }
567 return c.receiptForAdmissionResult(id, st, result), nil
568 }
569 return sessioninbox.InboxReceipt{
570 ItemID: id,
571 Disposition: sessioninbox.DispositionStarted,
572 Capacity: st.Snapshot().Capacity,
573 }, nil
574 }
575
576 func (c *Controller) receiptForAdmissionResult(id string, st *sessioninbox.Store, result admissionResult) sessioninbox.InboxReceipt {
577 disposition := sessioninbox.DispositionRejectedBusy
578 switch result {
579 case turnDroppedClosed:
580 disposition = sessioninbox.DispositionRejectedClosed
581 case turnDroppedRotating:
582 disposition = sessioninbox.DispositionRejectedRotating
583 }
584 return sessioninbox.InboxReceipt{ItemID: id, Disposition: disposition, Capacity: st.Snapshot().Capacity}
585 }
586
587 // onInboxTurnDone acknowledges durable completion of every active inbox item
588 // (running follow-up + all steers accepted this turn). Dispatch of the next
589 // item is deferred until the finishing window closes so admission is not
590 // rejected as busy.
591 func (c *Controller) onInboxTurnDone() {
592 c.inbox.mu.Lock()
593 // Keep these IDs published as live ownership while SnapshotActivity runs.
594 // Inbox recovery can therefore proceed without waiting on extension hooks,
595 // transcript I/O, or the session file lock and will preserve this turn.
596 ids := c.inbox.activeIDs()
597 st := c.inbox.store
598 beforeSnapshot := c.inbox.beforeCompletionSnapshot
599 beforeAck := c.inbox.beforeCompletionAck
600 c.inbox.mu.Unlock()
601 if st == nil || len(ids) == 0 {
602 return
603 }
604 if beforeSnapshot != nil {
605 beforeSnapshot()
606 }
607 // Transcript snapshot is the durable receipt boundary for the whole set.
608 if err := c.SnapshotActivity(); err != nil {
609 slog.Warn("controller: inbox turn snapshot", "err", err)
610 for _, id := range ids {
611 _ = st.SetState(id, sessioninbox.StateUncertain, "turn completed but transcript snapshot failed")
612 }
613 _ = st.SetPaused(true)
614 c.inbox.mu.Lock()
615 c.inbox.untrackActiveSet(ids)
616 c.inbox.mu.Unlock()
617 sessioninbox.NoteUncertain()
618 return
619 }
620 // Keep ownership published through every durable acknowledgement. Recovery
621 // can run concurrently, sees these IDs as live without a Controller lock,
622 // and ownership is removed only after dequeue or uncertain state is durable.
623 if beforeAck != nil {
624 beforeAck()
625 }
626 ackFailed := false
627 for _, id := range ids {
628 if err := st.AckDequeue(id); err != nil {
629 if errors.Is(err, sessioninbox.ErrNotFound) {
630 continue
631 }
632 slog.Warn("controller: inbox ack dequeue", "err", err, "id", id)
633 _ = st.SetState(id, sessioninbox.StateUncertain, "turn completed but inbox acknowledgement failed")
634 ackFailed = true
635 }
636 }
637 if ackFailed {
638 _ = st.SetPaused(true)
639 sessioninbox.NoteUncertain()
640 }
641 c.inbox.mu.Lock()
642 c.inbox.untrackActiveSet(ids)
643 c.inbox.mu.Unlock()
644 }
645
646 // onInboxUnappliedSteer keeps accepted-but-unapplied steers for inspection.
647 func (c *Controller) onInboxUnappliedSteer(itemID string) {
648 if itemID == "" {
649 return
650 }
651 st, err := c.ensureInbox()
652 if err != nil {
653 return
654 }
655 if err := st.MarkAcceptedSteerUncertain(itemID, "steer accepted but unapplied before turn exit"); err != nil {
656 if errors.Is(err, sessioninbox.ErrNotFound) {
657 c.inbox.mu.Lock()
658 c.inbox.untrackActive(itemID)
659 c.inbox.mu.Unlock()
660 }
661 return
662 }
663 _ = st.SetPaused(true)
664 c.inbox.mu.Lock()
665 c.inbox.untrackActive(itemID)
666 c.inbox.mu.Unlock()
667 sessioninbox.NoteUncertain()
668 }
669
670 // TryEnqueueAndSteer is a convenience for frontends: durable steer then TrySteer.
671 func (c *Controller) TryEnqueueAndSteer(req InboxRequest) (sessioninbox.InboxReceipt, error) {
672 return c.tryEnqueueAndSteerForTurn("", req)
673 }
674
675 // TryEnqueueAndSteerForTurn preserves the durable fallback semantics while
676 // fencing the mid-turn steer against the exact lifecycle turn observed by the
677 // caller. If that turn has already ended, the instruction remains a queued
678 // follow-up and is never injected into a replacement turn.
679 func (c *Controller) TryEnqueueAndSteerForTurn(turnID string, req InboxRequest) (sessioninbox.InboxReceipt, error) {
680 turnID = strings.TrimSpace(turnID)
681 if turnID == "" {
682 return sessioninbox.InboxReceipt{}, fmt.Errorf("turnId is required")
683 }
684 return c.tryEnqueueAndSteerForTurn(turnID, req)
685 }
686
687 func (c *Controller) tryEnqueueAndSteerForTurn(turnID string, req InboxRequest) (sessioninbox.InboxReceipt, error) {
688 req.Intent = sessioninbox.IntentSteer
689 rec, err := c.EnqueueInbox(req)
690 if err != nil {
691 return rec, err
692 }
693 steered, err := c.trySteerInboxItem(rec.ItemID, turnID)
694 if errors.Is(err, sessioninbox.ErrPaused) {
695 rec.Disposition = sessioninbox.DispositionQueuedFollowup
696 rec.Paused = true
697 return rec, nil
698 }
699 if err != nil {
700 return rec, err
701 }
702 return steered, nil
703 }
704
705 // TryEnqueueFollowup durably queues a follow-up and may dispatch if idle.
706 func (c *Controller) TryEnqueueFollowup(req InboxRequest) (sessioninbox.InboxReceipt, error) {
707 return c.TryEnqueueFollowupContext(c.attachmentContext(), req)
708 }
709
710 func (c *Controller) TryEnqueueFollowupContext(ctx context.Context, req InboxRequest) (sessioninbox.InboxReceipt, error) {
711 req.Intent = sessioninbox.IntentFollowup
712 rec, err := c.EnqueueInboxContext(ctx, req)
713 if err != nil {
714 return rec, err
715 }
716 if !c.Running() {
717 c.maybeDispatchInbox()
718 }
719 return rec, nil
720 }
721
722 func firstNonEmptyStr(vals ...string) string {
723 for _, v := range vals {
724 if strings.TrimSpace(v) != "" {
725 return strings.TrimSpace(v)
726 }
727 }
728 return ""
729 }
730
730 lines GO