返回 DeepSeek-Reasonix
inbox_steer.go
根目录 / internal / control / inbox_steer.go
1 package control
2
3 import (
4 "errors"
5 "fmt"
6 "strings"
7
8 "reasonix/internal/agent"
9 "reasonix/internal/sessioninbox"
10 )
11
12 func steerAlreadyAdmitted(state sessioninbox.InboxState) bool {
13 switch state {
14 case sessioninbox.StateRunning, sessioninbox.StateSteerAccepted, sessioninbox.StateSteerConsumed:
15 return true
16 default:
17 return false
18 }
19 }
20
21 func (c *Controller) readSteerCandidate(st *sessioninbox.Store, id string) (sessioninbox.InboxItemMeta, sessioninbox.PromptEnvelope, error) {
22 meta, env, err := st.ReadItem(id)
23 if err != nil || (meta.State != sessioninbox.StateRunning && meta.State != sessioninbox.StateSteerAccepted && meta.State != sessioninbox.StateSteerConsumed) {
24 return meta, env, err
25 }
26 recovered, err := st.RecoverOrphanedInFlightOwnedBy(c.inbox.ownsItem)
27 if err != nil {
28 return sessioninbox.InboxItemMeta{}, sessioninbox.PromptEnvelope{}, err
29 }
30 if recovered > 0 {
31 sessioninbox.NoteRecovered(recovered)
32 }
33 return st.ReadItem(id)
34 }
35
36 func (c *Controller) unlockInboxSteerAdmission(dispatch *bool) {
37 c.inbox.admissionMu.Unlock()
38 if *dispatch {
39 c.maybeDispatchInbox()
40 }
41 }
42
43 func inboxSteerLoader(st *sessioninbox.Store, itemID string) func() (string, error) {
44 return func() (string, error) {
45 _, env, err := st.ReadItem(itemID)
46 if err != nil {
47 if errors.Is(err, sessioninbox.ErrNotFound) {
48 return "", agent.ErrSteerWithdrawn
49 }
50 return "", err
51 }
52 text := strings.TrimSpace(env.SubmitText)
53 if text == "" {
54 text = strings.TrimSpace(env.DisplayText)
55 }
56 if text == "" {
57 return "", fmt.Errorf("inbox item %s has empty body", itemID)
58 }
59 materialized, images, block, materializeErr := applyInboxReferences(env)
60 if materializeErr != nil {
61 return "", materializeErr
62 }
63 if block != "" {
64 return "", fmt.Errorf("frozen reference unavailable: %s", block)
65 }
66 if len(images) > 0 {
67 return "", fmt.Errorf("image guidance requires a follow-up turn")
68 }
69 // This compare-and-transition is the durable hand-off boundary and
70 // closes the loader-vs-cancel gap after TrySteerInboxItem returns.
71 if err := st.MarkSteerConsumed(itemID); err != nil {
72 if errors.Is(err, sessioninbox.ErrNotFound) {
73 return "", agent.ErrSteerWithdrawn
74 }
75 return "", err
76 }
77 return firstNonEmptyStr(materialized, text), nil
78 }
79 }
80
81 // TrySteerInboxItem persists intent=steer (if needed) and attempts mid-turn
82 // admission. Rejected steers stay queued as follow-up.
83 //
84 // The agent loader only captures the item ID and re-reads the blob on consume
85 // so large steer bodies do not accumulate in the agent heap.
86 func (c *Controller) TrySteerInboxItem(id string) (sessioninbox.InboxReceipt, error) {
87 return c.trySteerInboxItem(id, "")
88 }
89
90 // TrySteerInboxItemForTurn applies an existing durable item only to the exact
91 // active turn. A stale target falls back to queued-follow-up semantics.
92 func (c *Controller) TrySteerInboxItemForTurn(turnID, id string) (sessioninbox.InboxReceipt, error) {
93 turnID = strings.TrimSpace(turnID)
94 if turnID == "" {
95 return sessioninbox.InboxReceipt{}, fmt.Errorf("turnId is required")
96 }
97 return c.trySteerInboxItem(id, turnID)
98 }
99
100 func (c *Controller) trySteerInboxItem(id, expectedTurnID string) (sessioninbox.InboxReceipt, error) {
101 return c.trySteerInboxItemForSession(id, expectedTurnID, "")
102 }
103
104 func (c *Controller) trySteerInboxItemForSession(id, expectedTurnID, expectedPath string) (sessioninbox.InboxReceipt, error) {
105 c.inbox.admissionMu.Lock()
106 dispatchAfterUnlock := false
107 defer c.unlockInboxSteerAdmission(&dispatchAfterUnlock)
108 st, err := c.ensureInbox()
109 if err != nil {
110 return sessioninbox.InboxReceipt{}, err
111 }
112 if expectedPath != "" && st.SessionPath() != expectedPath {
113 return sessioninbox.InboxReceipt{}, ErrInboxSessionChanged
114 }
115 meta, env, err := c.readSteerCandidate(st, id)
116 if err != nil {
117 return sessioninbox.InboxReceipt{}, err
118 }
119 // RetryInboxItem may start this item while the frontend holds stale running=true.
120 // Treat the follow-up Steer as idempotent: the current turn already owns the
121 // durable body, so it must not be applied twice or reported as a false failure.
122 if steerAlreadyAdmitted(meta.State) {
123 return sessioninbox.InboxReceipt{
124 ItemID: id,
125 Disposition: sessioninbox.DispositionSteerAccepted,
126 Paused: st.Snapshot().Paused,
127 Capacity: st.Snapshot().Capacity,
128 Idempotent: true,
129 }, nil
130 }
131 if meta.State != sessioninbox.StateQueued && meta.State != sessioninbox.StateUncertain {
132 return sessioninbox.InboxReceipt{}, sessioninbox.ErrInvalidState
133 }
134 snapshot := st.Snapshot()
135 if snapshot.Paused {
136 return sessioninbox.InboxReceipt{}, sessioninbox.ErrPaused
137 }
138 if meta.State == sessioninbox.StateUncertain {
139 if err := st.SetState(id, sessioninbox.StateQueued, ""); err != nil {
140 return sessioninbox.InboxReceipt{}, err
141 }
142 }
143 cap := snapshot.Capacity
144 c.mu.Lock()
145 rotating := c.rotating
146 closed := c.closed
147 c.mu.Unlock()
148 if closed {
149 return sessioninbox.InboxReceipt{ItemID: id, Disposition: sessioninbox.DispositionRejectedClosed, Capacity: cap}, nil
150 }
151 if rotating {
152 dispatchAfterUnlock = true
153 return sessioninbox.InboxReceipt{ItemID: id, Disposition: sessioninbox.DispositionRejectedRotating, Capacity: cap}, nil
154 }
155 // Capture only the store pointer + item id. Load body from disk at consume.
156 loader := inboxSteerLoader(st, id)
157 // Persist the admission boundary before exposing the loader to the agent.
158 // Holding c.mu for the short in-memory enqueue serializes active tracking
159 // with finishGuardedTurn, so TurnDone cannot overtake an accepted steer.
160 c.inbox.trackAdmission(id)
161 defer c.inbox.untrackAdmission(id)
162 if len(env.FrozenImages) == 0 {
163 if err := st.TransitionPrepared(id, sessioninbox.ContentVersion(meta), sessioninbox.StateSteerAccepted, "", false); err != nil {
164 return sessioninbox.InboxReceipt{}, err
165 }
166 }
167 c.mu.Lock()
168 turnMatches := true
169 if expectedTurnID != "" {
170 turnMatches = false
171 if ledger := c.turnEventLedger(); ledger != nil {
172 turnMatches = ledger.ActiveTurnID() == expectedTurnID
173 }
174 }
175 accepted := turnMatches && !c.closed && !c.rotating && c.bodyActiveLocked() && c.executor != nil && len(env.FrozenImages) == 0 && c.executor.SteerItem(id, loader)
176 if accepted {
177 c.inbox.mu.Lock()
178 c.inbox.trackActive(id)
179 c.inbox.mu.Unlock()
180 }
181 c.mu.Unlock()
182 if accepted {
183 sessioninbox.NoteSteerAccepted()
184 return sessioninbox.InboxReceipt{
185 ItemID: id,
186 Disposition: sessioninbox.DispositionSteerAccepted,
187 Paused: st.Snapshot().Paused,
188 Capacity: cap,
189 }, nil
190 }
191 // Rejected: keep as follow-up.
192 if len(env.FrozenImages) == 0 {
193 if err := st.SetState(id, sessioninbox.StateQueued, ""); err != nil {
194 _ = st.ForcePause(true, 1)
195 return sessioninbox.InboxReceipt{}, err
196 }
197 }
198 if err := st.ConvertIntent(id, sessioninbox.IntentFollowup); err != nil {
199 return sessioninbox.InboxReceipt{}, err
200 }
201 sessioninbox.NoteSteerRejected()
202 dispatchAfterUnlock = true
203 return sessioninbox.InboxReceipt{
204 ItemID: id,
205 Disposition: sessioninbox.DispositionQueuedFollowup,
206 Paused: st.Snapshot().Paused,
207 Capacity: cap,
208 }, nil
209 }
210
210 lines GO