返回 DeepSeek-Reasonix
inbox_orphan_recovery_test.go
根目录 / internal / control / inbox_orphan_recovery_test.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "os"
7 "path/filepath"
8 "testing"
9 "time"
10
11 "reasonix/internal/event"
12 "reasonix/internal/filelock"
13 "reasonix/internal/sessioninbox"
14 )
15
16 func TestInboxSnapshotRecoversUnownedInFlightItem(t *testing.T) {
17 dir := t.TempDir()
18 session := filepath.Join(dir, "s.jsonl")
19 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
20 t.Fatal(err)
21 }
22 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
23 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "orphaned guidance"})
24 if err != nil {
25 t.Fatal(err)
26 }
27 st, err := c.ensureInbox()
28 if err != nil {
29 t.Fatal(err)
30 }
31 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
32 t.Fatal(err)
33 }
34
35 snap := c.InboxSnapshot()
36 if !snap.Paused || !snap.Recovered || snap.RecoveredN != 1 {
37 t.Fatalf("orphan recovery metadata = %+v", snap)
38 }
39 if len(snap.Items) != 1 || snap.Items[0].State != sessioninbox.StateUncertain {
40 t.Fatalf("orphan recovery items = %+v", snap.Items)
41 }
42 if err := c.DeleteInboxItem(rec.ItemID); err != nil {
43 t.Fatalf("delete recovered orphan: %v", err)
44 }
45 }
46
47 func TestInboxSnapshotPreservesActivelyOwnedSteer(t *testing.T) {
48 dir := t.TempDir()
49 session := filepath.Join(dir, "s.jsonl")
50 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
51 t.Fatal(err)
52 }
53 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
54 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "active guidance"})
55 if err != nil {
56 t.Fatal(err)
57 }
58 st, err := c.ensureInbox()
59 if err != nil {
60 t.Fatal(err)
61 }
62 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
63 t.Fatal(err)
64 }
65 c.inbox.mu.Lock()
66 c.inbox.trackActive(rec.ItemID)
67 c.inbox.mu.Unlock()
68
69 snap := c.InboxSnapshot()
70 if snap.Paused || snap.Recovered || len(snap.Items) != 1 || snap.Items[0].State != sessioninbox.StateSteerAccepted {
71 t.Fatalf("active steer was reclassified: %+v", snap)
72 }
73
74 c.inbox.mu.Lock()
75 c.inbox.untrackActive(rec.ItemID)
76 c.inbox.mu.Unlock()
77 snap = c.InboxSnapshot()
78 if !snap.Paused || len(snap.Items) != 1 || snap.Items[0].State != sessioninbox.StateUncertain {
79 t.Fatalf("unowned steer was not recovered: %+v", snap)
80 }
81 }
82
83 func TestTrySteerOrphanRequiresReviewBeforeExplicitRetry(t *testing.T) {
84 dir := t.TempDir()
85 session := filepath.Join(dir, "s.jsonl")
86 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
87 t.Fatal(err)
88 }
89 runner := &gatedTurnRunner{started: make(chan struct{}), release: make(chan struct{})}
90 c := newOwnedTestController(t, Options{Runner: runner, SessionPath: session, SessionDir: dir, Sink: event.Discard})
91 defer c.autosaveWG.Wait()
92 defer close(runner.release)
93 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "retry me"})
94 if err != nil {
95 t.Fatal(err)
96 }
97 st, err := c.ensureInbox()
98 if err != nil {
99 t.Fatal(err)
100 }
101 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
102 t.Fatal(err)
103 }
104
105 if _, err := c.TrySteerInboxItem(rec.ItemID); !errors.Is(err, sessioninbox.ErrPaused) {
106 t.Fatalf("first orphan retry error = %v, want ErrPaused", err)
107 }
108 snap := c.InboxSnapshot()
109 if len(snap.Items) != 1 || snap.Items[0].State != sessioninbox.StateUncertain {
110 t.Fatalf("first orphan retry state = %+v", snap)
111 }
112 if err := c.SetInboxPaused(false); err != nil {
113 t.Fatal(err)
114 }
115 receipt, err := c.TrySteerInboxItem(rec.ItemID)
116 if err != nil {
117 t.Fatal(err)
118 }
119 if receipt.Disposition != sessioninbox.DispositionQueuedFollowup {
120 t.Fatalf("explicit retry disposition = %q", receipt.Disposition)
121 }
122 select {
123 case <-runner.started:
124 case <-time.After(time.Second):
125 t.Fatal("explicit retry did not dispatch the recovered item")
126 }
127 meta, _, err := c.ReadInboxItem(rec.ItemID)
128 if err != nil {
129 t.Fatal(err)
130 }
131 if meta.State != sessioninbox.StateRunning || meta.Intent != sessioninbox.IntentFollowup {
132 t.Fatalf("explicit retry meta = %+v", meta)
133 }
134 }
135
136 func TestRetryThenStaleSteerTreatsAlreadyRunningItemAsIdempotent(t *testing.T) {
137 dir := t.TempDir()
138 session := filepath.Join(dir, "s.jsonl")
139 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
140 t.Fatal(err)
141 }
142 runner := &gatedTurnRunner{started: make(chan struct{}), release: make(chan struct{})}
143 c := newOwnedTestController(t, Options{
144 Runner: runner,
145 SessionPath: session,
146 SessionDir: dir,
147 Sink: event.Discard,
148 })
149 defer c.autosaveWG.Wait()
150 defer close(runner.release)
151 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentFollowup, Submit: "retry once"})
152 if err != nil {
153 t.Fatal(err)
154 }
155 st, err := c.ensureInbox()
156 if err != nil {
157 t.Fatal(err)
158 }
159 if err := st.SetState(rec.ItemID, sessioninbox.StateUncertain, "review retry"); err != nil {
160 t.Fatal(err)
161 }
162
163 if err := c.RetryInboxItem(rec.ItemID); err != nil {
164 t.Fatal(err)
165 }
166 select {
167 case <-runner.started:
168 case <-time.After(time.Second):
169 t.Fatal("retry did not start the recovered item")
170 }
171
172 receipt, err := c.TrySteerInboxItem(rec.ItemID)
173 if err != nil {
174 t.Fatalf("retry already started the item, but stale steer returned: %v", err)
175 }
176 if receipt.Disposition != sessioninbox.DispositionSteerAccepted || !receipt.Idempotent {
177 t.Fatalf("stale steer receipt = %+v, want idempotent accepted", receipt)
178 }
179 }
180
181 func TestInboxAdmissionOwnsClaimBeforeSnapshotRecovery(t *testing.T) {
182 dir := t.TempDir()
183 session := filepath.Join(dir, "s.jsonl")
184 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
185 t.Fatal(err)
186 }
187 c := newOwnedTestController(t, Options{
188 Runner: &fakeTurnRunner{},
189 SessionPath: session,
190 SessionDir: dir,
191 Sink: event.Discard,
192 })
193 rec, err := c.EnqueueInbox(InboxRequest{Submit: "claimed atomically"})
194 if err != nil {
195 t.Fatal(err)
196 }
197 claimed := make(chan struct{})
198 release := make(chan struct{})
199 c.inbox.mu.Lock()
200 c.inbox.beforePreparedAdmission = func() {
201 close(claimed)
202 <-release
203 }
204 c.inbox.mu.Unlock()
205 type result struct {
206 receipt sessioninbox.InboxReceipt
207 err error
208 }
209 resultCh := make(chan result, 1)
210 go func() {
211 receipt, submitErr := c.TrySubmitInboxItem(rec.ItemID)
212 resultCh <- result{receipt: receipt, err: submitErr}
213 }()
214 <-claimed
215 snapshotCh := make(chan sessioninbox.InboxSnapshot, 1)
216 go func() { snapshotCh <- c.InboxSnapshot() }()
217 var duringAdmission sessioninbox.InboxSnapshot
218 select {
219 case duringAdmission = <-snapshotCh:
220 case <-time.After(time.Second):
221 close(release)
222 <-resultCh
223 t.Fatal("snapshot recovery waited on the admission state machine")
224 }
225 if duringAdmission.Paused || len(duringAdmission.Items) != 1 || duringAdmission.Items[0].State != sessioninbox.StateRunning {
226 close(release)
227 <-resultCh
228 t.Fatalf("snapshot recovered a live admission: %+v", duringAdmission)
229 }
230 if c.inbox.admissionMu.TryLock() {
231 c.inbox.admissionMu.Unlock()
232 close(release)
233 <-resultCh
234 t.Fatal("admission hook did not hold the admission state machine")
235 }
236 close(release)
237 got := <-resultCh
238 if got.err != nil || got.receipt.Disposition != sessioninbox.DispositionStarted {
239 t.Fatalf("admission result = %+v, err=%v", got.receipt, got.err)
240 }
241 c.autosaveWG.Wait()
242 }
243
244 func TestInboxSnapshotDoesNotHoldAdmissionWhileDiskLocked(t *testing.T) {
245 dir := t.TempDir()
246 session := filepath.Join(dir, "s.jsonl")
247 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
248 t.Fatal(err)
249 }
250 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
251 st, err := c.ensureInbox()
252 if err != nil {
253 t.Fatal(err)
254 }
255 releaseDisk, err := filelock.Acquire(context.Background(), filepath.Join(st.Dir(), "transaction.lock"))
256 if err != nil {
257 t.Fatal(err)
258 }
259 reachedRead := make(chan struct{})
260 c.inbox.mu.Lock()
261 c.inbox.beforeSnapshotRead = func() { close(reachedRead) }
262 c.inbox.mu.Unlock()
263 done := make(chan struct{})
264 go func() {
265 _ = c.InboxSnapshot()
266 close(done)
267 }()
268 <-reachedRead
269 select {
270 case <-done:
271 releaseDisk()
272 t.Fatal("snapshot bypassed the held Store transaction lock")
273 default:
274 }
275 if !c.inbox.admissionMu.TryLock() {
276 releaseDisk()
277 <-done
278 t.Fatal("snapshot held admissionMu while waiting on transaction.lock")
279 }
280 c.inbox.admissionMu.Unlock()
281 releaseDisk()
282 <-done
283 }
284
285 func TestInboxCompletionKeepsOwnershipWithoutHoldingAdmissionDuringSnapshot(t *testing.T) {
286 dir := t.TempDir()
287 session := filepath.Join(dir, "s.jsonl")
288 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
289 t.Fatal(err)
290 }
291 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
292 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "complete atomically"})
293 if err != nil {
294 t.Fatal(err)
295 }
296 st, err := c.ensureInbox()
297 if err != nil {
298 t.Fatal(err)
299 }
300 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, ""); err != nil {
301 t.Fatal(err)
302 }
303 beforeSnapshot := make(chan struct{})
304 release := make(chan struct{})
305 c.inbox.mu.Lock()
306 c.inbox.trackActive(rec.ItemID)
307 c.inbox.beforeCompletionSnapshot = func() {
308 close(beforeSnapshot)
309 <-release
310 }
311 c.inbox.mu.Unlock()
312 done := make(chan struct{})
313 go func() {
314 c.onInboxTurnDone()
315 close(done)
316 }()
317 <-beforeSnapshot
318 if !c.inbox.admissionMu.TryLock() {
319 t.Fatal("completion held admission lock across transcript snapshot boundary")
320 }
321 c.inbox.admissionMu.Unlock()
322 whileSaving := c.InboxSnapshot()
323 if whileSaving.Paused || len(whileSaving.Items) != 1 || whileSaving.Items[0].State != sessioninbox.StateSteerConsumed {
324 t.Fatalf("snapshot recovery lost active ownership during transcript save: %+v", whileSaving)
325 }
326 close(release)
327 <-done
328 snap := c.InboxSnapshot()
329 if snap.Paused || len(snap.Items) != 0 {
330 t.Fatalf("completed item survived durable acknowledgement: %+v", snap)
331 }
332 }
333
334 func TestInboxCompletionOwnsItemWithoutHoldingAdmissionDuringDurableAck(t *testing.T) {
335 dir := t.TempDir()
336 session := filepath.Join(dir, "s.jsonl")
337 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
338 t.Fatal(err)
339 }
340 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
341 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "ack atomically"})
342 if err != nil {
343 t.Fatal(err)
344 }
345 st, err := c.ensureInbox()
346 if err != nil {
347 t.Fatal(err)
348 }
349 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, ""); err != nil {
350 t.Fatal(err)
351 }
352 beforeAck := make(chan struct{})
353 release := make(chan struct{})
354 c.inbox.mu.Lock()
355 c.inbox.trackActive(rec.ItemID)
356 c.inbox.beforeCompletionAck = func() {
357 close(beforeAck)
358 <-release
359 }
360 c.inbox.mu.Unlock()
361 done := make(chan struct{})
362 go func() {
363 c.onInboxTurnDone()
364 close(done)
365 }()
366 <-beforeAck
367 if !c.inbox.admissionMu.TryLock() {
368 close(release)
369 <-done
370 t.Fatal("completion held admission lock across durable acknowledgement")
371 }
372 c.inbox.admissionMu.Unlock()
373 whileAcking := c.InboxSnapshot()
374 if whileAcking.Paused || len(whileAcking.Items) != 1 || whileAcking.Items[0].State != sessioninbox.StateSteerConsumed {
375 close(release)
376 <-done
377 t.Fatalf("snapshot recovery lost active ownership during durable acknowledgement: %+v", whileAcking)
378 }
379 close(release)
380 <-done
381 snap := c.InboxSnapshot()
382 if snap.Paused || len(snap.Items) != 0 {
383 t.Fatalf("completed item survived durable acknowledgement: %+v", snap)
384 }
385 }
386
386 lines GO