返回 DeepSeek-Reasonix
inbox_test.go
根目录 / internal / control / inbox_test.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "path/filepath"
9 "strings"
10 "testing"
11 "time"
12
13 "reasonix/internal/agent"
14 "reasonix/internal/event"
15 "reasonix/internal/memory"
16 "reasonix/internal/provider"
17 "reasonix/internal/sessioninbox"
18 "reasonix/internal/skill"
19 "reasonix/internal/tool"
20 )
21
22 func TestEnqueueInboxDurableAndSnapshot(t *testing.T) {
23 dir := t.TempDir()
24 session := filepath.Join(dir, "s.jsonl")
25 if err := os.WriteFile(session, []byte("{}\n"), 0o644); err != nil {
26 t.Fatal(err)
27 }
28 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
29 rec, err := c.EnqueueInbox(InboxRequest{
30 Intent: sessioninbox.IntentFollowup,
31 Display: "hello durable",
32 Submit: "hello durable",
33 Source: "test",
34 })
35 if err != nil {
36 t.Fatal(err)
37 }
38 if rec.ItemID == "" {
39 t.Fatal("empty item id")
40 }
41 snap := c.InboxSnapshot()
42 if len(snap.Items) != 1 || snap.Items[0].Preview == "" {
43 t.Fatalf("snapshot = %+v", snap)
44 }
45 if snap.SessionPath != session {
46 t.Fatalf("snapshot session path = %q, want %q", snap.SessionPath, session)
47 }
48 _, env, err := c.ReadInboxItem(rec.ItemID)
49 if err != nil || env.SubmitText != "hello durable" {
50 t.Fatalf("read = %+v err=%v", env, err)
51 }
52 }
53
54 func TestSessionRebindOnlyPausesInboxWithPendingWork(t *testing.T) {
55 for _, tc := range []struct {
56 name string
57 pending bool
58 }{
59 {name: "empty"},
60 {name: "pending", pending: true},
61 } {
62 t.Run(tc.name, func(t *testing.T) {
63 dir := t.TempDir()
64 oldPath := filepath.Join(dir, "old.jsonl")
65 c := newOwnedTestController(t, Options{SessionPath: oldPath, SessionDir: dir, Sink: event.Discard})
66 if tc.pending {
67 if _, err := c.EnqueueInbox(InboxRequest{Submit: "work"}); err != nil {
68 t.Fatal(err)
69 }
70 }
71
72 c.SetSessionPath(filepath.Join(dir, "new.jsonl"))
73 oldInbox, err := sessioninbox.Open(oldPath, sessioninbox.Limits{})
74 if err != nil {
75 t.Fatal(err)
76 }
77 defer oldInbox.Close()
78 if got := oldInbox.Snapshot().Paused; got != tc.pending {
79 t.Fatalf("paused = %v, want %v", got, tc.pending)
80 }
81 })
82 }
83 }
84
85 func TestTryEnqueueAndSteerWhenPausedKeepsQueuedFollowup(t *testing.T) {
86 dir := t.TempDir()
87 session := filepath.Join(dir, "s.jsonl")
88 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
89 if err := c.SetInboxPaused(true); err != nil {
90 t.Fatal(err)
91 }
92 got, err := c.TryEnqueueAndSteer(InboxRequest{Submit: "later"})
93 if err != nil {
94 t.Fatal(err)
95 }
96 if got.Disposition != sessioninbox.DispositionQueuedFollowup || !got.Paused || got.ItemID == "" {
97 t.Fatalf("receipt = %+v", got)
98 }
99 meta, _, err := c.ReadInboxItem(got.ItemID)
100 if err != nil {
101 t.Fatal(err)
102 }
103 if meta.State != sessioninbox.StateQueued {
104 t.Fatalf("meta = %+v", meta)
105 }
106 }
107
108 func TestDeleteInboxItemRecoversOrphanThenRemoves(t *testing.T) {
109 dir := t.TempDir()
110 session := filepath.Join(dir, "s.jsonl")
111 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
112 rec, err := c.EnqueueInbox(InboxRequest{Submit: "stuck"})
113 if err != nil {
114 t.Fatal(err)
115 }
116 st, err := c.ensureInbox()
117 if err != nil {
118 t.Fatal(err)
119 }
120 if err := st.ClaimItem(rec.ItemID); err != nil {
121 t.Fatal(err)
122 }
123 if err := c.DeleteInboxItem(rec.ItemID); err != nil {
124 t.Fatal(err)
125 }
126 if _, _, err := c.ReadInboxItem(rec.ItemID); !errors.Is(err, sessioninbox.ErrNotFound) {
127 t.Fatalf("item still present: %v", err)
128 }
129 if snap := c.InboxSnapshot(); snap.Paused || snap.Recovered || len(snap.Items) != 0 {
130 t.Fatalf("empty inbox stayed paused after deleting last orphan: %+v", snap)
131 }
132 }
133
134 func TestDeleteInboxItemWithdrawsUnconsumedSteer(t *testing.T) {
135 dir := t.TempDir()
136 session := filepath.Join(dir, "s.jsonl")
137 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
138 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "withdraw me"})
139 if err != nil {
140 t.Fatal(err)
141 }
142 st, err := c.ensureInbox()
143 if err != nil {
144 t.Fatal(err)
145 }
146 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
147 t.Fatal(err)
148 }
149 c.inbox.mu.Lock()
150 c.inbox.trackActive(rec.ItemID)
151 c.inbox.mu.Unlock()
152 if err := c.DeleteInboxItem(rec.ItemID); err != nil {
153 t.Fatal(err)
154 }
155 if _, _, err := c.ReadInboxItem(rec.ItemID); !errors.Is(err, sessioninbox.ErrNotFound) {
156 t.Fatalf("accepted steer still present: %v", err)
157 }
158 if snap := c.InboxSnapshot(); snap.Paused || len(snap.Items) != 0 {
159 t.Fatalf("withdrawing last steer left a paused empty inbox: %+v", snap)
160 }
161 }
162
163 func TestTrySteerRejectedBecomesFollowup(t *testing.T) {
164 dir := t.TempDir()
165 session := filepath.Join(dir, "s.jsonl")
166 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
167 runner := &gatedTurnRunner{started: make(chan struct{}), release: make(chan struct{})}
168 c := newOwnedTestController(t, Options{Runner: runner, SessionPath: session, SessionDir: dir, Sink: event.Discard})
169 defer c.autosaveWG.Wait()
170 defer close(runner.release)
171 rec, err := c.EnqueueInbox(InboxRequest{
172 Intent: sessioninbox.IntentSteer,
173 Submit: "mid-turn please",
174 })
175 if err != nil {
176 t.Fatal(err)
177 }
178 // No running turn → reject, keep as follow-up.
179 got, err := c.TrySteerInboxItem(rec.ItemID)
180 if err != nil {
181 t.Fatal(err)
182 }
183 if got.Disposition != sessioninbox.DispositionQueuedFollowup {
184 t.Fatalf("disposition = %s, want queued_followup", got.Disposition)
185 }
186 select {
187 case <-runner.started:
188 case <-time.After(time.Second):
189 t.Fatal("rejected idle steer did not dispatch as a follow-up")
190 }
191 meta, _, err := c.ReadInboxItem(rec.ItemID)
192 if err != nil {
193 t.Fatal(err)
194 }
195 if meta.State != sessioninbox.StateRunning || meta.Intent != sessioninbox.IntentFollowup {
196 t.Fatalf("meta = %+v", meta)
197 }
198 }
199
200 func TestIdempotentEnqueue(t *testing.T) {
201 dir := t.TempDir()
202 session := filepath.Join(dir, "s.jsonl")
203 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
204 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
205 a, err := c.EnqueueInbox(InboxRequest{Submit: "x", Idempotency: "k1"})
206 if err != nil {
207 t.Fatal(err)
208 }
209 b, err := c.EnqueueInbox(InboxRequest{Submit: "x", Idempotency: "k1"})
210 if err != nil {
211 t.Fatal(err)
212 }
213 if a.ItemID != b.ItemID || !b.Idempotent {
214 t.Fatalf("a=%+v b=%+v", a, b)
215 }
216 }
217
218 func TestIdempotentEnqueueDoesNotReclassifyExistingItem(t *testing.T) {
219 dir := t.TempDir()
220 workspace := filepath.Join(dir, "workspace")
221 if err := os.MkdirAll(workspace, 0o755); err != nil {
222 t.Fatal(err)
223 }
224 c := newOwnedTestController(t, Options{
225 SessionPath: filepath.Join(dir, "s.jsonl"),
226 SessionDir: dir,
227 WorkspaceRoot: workspace,
228 Sink: event.Discard,
229 })
230 first, err := c.EnqueueInbox(InboxRequest{Submit: "original", Idempotency: "same"})
231 if err != nil {
232 t.Fatal(err)
233 }
234 second, err := c.EnqueueInbox(InboxRequest{Submit: "original", Idempotency: "same"})
235 if err != nil {
236 t.Fatal(err)
237 }
238 if first.ItemID != second.ItemID || !second.Idempotent {
239 t.Fatalf("first=%+v second=%+v", first, second)
240 }
241 snapshot := c.InboxSnapshot()
242 if snapshot.Paused || len(snapshot.Items) != 1 || snapshot.Items[0].State != sessioninbox.StateQueued {
243 t.Fatalf("idempotent replay reclassified original item: %+v", snapshot)
244 }
245 }
246
247 func TestIdempotentEnqueueRejectsDifferentInput(t *testing.T) {
248 dir := t.TempDir()
249 c := newOwnedTestController(t, Options{
250 SessionPath: filepath.Join(dir, "s.jsonl"),
251 SessionDir: dir,
252 Sink: event.Discard,
253 })
254 if _, err := c.EnqueueInbox(InboxRequest{Submit: "original", Idempotency: "same"}); err != nil {
255 t.Fatal(err)
256 }
257 if _, err := c.EnqueueInbox(InboxRequest{Submit: "replacement", Idempotency: "same"}); !errors.Is(err, sessioninbox.ErrIdempotencyConflict) {
258 t.Fatalf("conflicting replay error = %v, want ErrIdempotencyConflict", err)
259 }
260 }
261
262 type inboxSteerProvider struct {
263 started chan struct{}
264 release chan struct{}
265 requests []provider.Request
266 }
267
268 func (p *inboxSteerProvider) Name() string { return "inbox-steer" }
269
270 func (p *inboxSteerProvider) awaitStarted(t *testing.T, c *Controller) {
271 t.Helper()
272 // Admission checkpoints use real durable I/O. These tests assert steer
273 // ordering and exactly-once consumption, not a one-second startup SLA.
274 select {
275 case <-p.started:
276 case <-time.After(inboxDispatchTestTimeout):
277 failInboxDispatchWait(t, c, "initial steer provider turn")
278 }
279 }
280
281 func (p *inboxSteerProvider) Stream(ctx context.Context, req provider.Request) (<-chan provider.Chunk, error) {
282 p.requests = append(p.requests, req)
283 ch := make(chan provider.Chunk, 2)
284 if len(p.requests) == 1 {
285 close(p.started)
286 go func() {
287 defer close(ch)
288 select {
289 case <-p.release:
290 ch <- provider.Chunk{Type: provider.ChunkText, Text: "ready"}
291 ch <- provider.Chunk{Type: provider.ChunkDone}
292 case <-ctx.Done():
293 }
294 }()
295 return ch, nil
296 }
297 ch <- provider.Chunk{Type: provider.ChunkText, Text: "applied"}
298 ch <- provider.Chunk{Type: provider.ChunkDone}
299 close(ch)
300 return ch, nil
301 }
302
303 func TestThirtySteersApplyAndAckExactlyOnce(t *testing.T) {
304 dir := t.TempDir()
305 prov := &inboxSteerProvider{started: make(chan struct{}), release: make(chan struct{})}
306 sess := agent.NewSession("sys")
307 exec := agent.New(prov, tool.NewRegistry(), sess, agent.Options{}, event.Discard)
308 sink, done, _ := collectSink()
309 c := newOwnedTestController(t, Options{
310 Runner: exec,
311 Executor: exec,
312 Sink: sink,
313 SessionDir: dir,
314 SessionPath: filepath.Join(dir, "s.jsonl"),
315 })
316 t.Cleanup(func() {
317 c.Close()
318 c.autosaveWG.Wait()
319 })
320 c.Submit("initial turn")
321 prov.awaitStarted(t, c)
322
323 const steerCount = 30
324 for i := range steerCount {
325 body := fmt.Sprintf("durable-steer-%02d", i)
326 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: body})
327 if err != nil {
328 t.Fatal(err)
329 }
330 got, err := c.TrySteerInboxItem(rec.ItemID)
331 if err != nil {
332 t.Fatal(err)
333 }
334 if got.Disposition != sessioninbox.DispositionSteerAccepted {
335 t.Fatalf("steer %d disposition = %q", i, got.Disposition)
336 }
337 }
338 close(prov.release)
339 // Thirty durable round trips are real filesystem work; a loaded Windows
340 // runner spends most of the default five seconds before the turn is even
341 // released. This asserts exactly-once acknowledgement, not latency.
342 waitForDoneWithin(t, done, 60*time.Second)
343
344 if items := c.InboxSnapshot().Items; len(items) != 0 {
345 t.Fatalf("accepted steers were not all acknowledged: %+v", items)
346 }
347 if got := len(prov.requests); got != steerCount+1 {
348 t.Fatalf("provider requests = %d, want %d", got, steerCount+1)
349 }
350 messages := sess.Snapshot()
351 for i := range steerCount {
352 body := fmt.Sprintf("durable-steer-%02d", i)
353 count := 0
354 for _, message := range messages {
355 count += strings.Count(message.Content, body)
356 }
357 if count != 1 {
358 t.Fatalf("%q appears %d times in transcript, want exactly once", body, count)
359 }
360 }
361 }
362
363 func TestMultiSteerActiveSetAcksAll(t *testing.T) {
364 dir := t.TempDir()
365 session := filepath.Join(dir, "s.jsonl")
366 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
367 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
368
369 st, err := c.ensureInbox()
370 if err != nil {
371 t.Fatal(err)
372 }
373 var ids []string
374 for i := range 3 {
375 rec, err := c.EnqueueInbox(InboxRequest{Submit: "body-" + string(rune('a'+i))})
376 if err != nil {
377 t.Fatal(err)
378 }
379 ids = append(ids, rec.ItemID)
380 _ = st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, "")
381 }
382 c.inbox.mu.Lock()
383 c.inbox.clearActive()
384 for _, id := range ids {
385 c.inbox.trackActive(id)
386 }
387 c.inbox.mu.Unlock()
388
389 c.onInboxTurnDone()
390 if n := len(c.InboxSnapshot().Items); n != 0 {
391 t.Fatalf("want all 3 steers acked/dequeued, still have %d items", n)
392 }
393 }
394
395 func TestSubmitInboxUsesFrozenReferenceWithoutLiveReresolve(t *testing.T) {
396 dir := t.TempDir()
397 workspace := filepath.Join(dir, "workspace")
398 if err := os.MkdirAll(workspace, 0o755); err != nil {
399 t.Fatal(err)
400 }
401 refPath := filepath.Join(workspace, "note.txt")
402 if err := os.WriteFile(refPath, []byte("enqueue-time-body"), 0o600); err != nil {
403 t.Fatal(err)
404 }
405 sessionPath := filepath.Join(dir, "s.jsonl")
406 sess := agent.NewSession("sys")
407 exec := agent.New(nil, nil, sess, agent.Options{}, event.Discard)
408 sink, done, _ := collectSink()
409 c := newOwnedTestController(t, Options{
410 Runner: appendingRunner{session: sess},
411 Executor: exec,
412 Sink: sink,
413 SessionDir: dir,
414 SessionPath: sessionPath,
415 WorkspaceRoot: workspace,
416 })
417 defer c.autosaveWG.Wait()
418
419 rec, err := c.EnqueueInbox(InboxRequest{Submit: "review @note.txt"})
420 if err != nil {
421 t.Fatal(err)
422 }
423 if err := os.WriteFile(refPath, []byte("live-body-after-enqueue"), 0o600); err != nil {
424 t.Fatal(err)
425 }
426 got, err := c.TrySubmitInboxItem(rec.ItemID)
427 if err != nil {
428 t.Fatal(err)
429 }
430 if got.Disposition != sessioninbox.DispositionStarted {
431 t.Fatalf("disposition = %q, want started", got.Disposition)
432 }
433 waitForDone(t, done)
434
435 messages := sess.Snapshot()
436 if len(messages) < 2 {
437 t.Fatalf("messages = %+v", messages)
438 }
439 input := messages[len(messages)-1].Content
440 if !strings.Contains(input, "enqueue-time-body") {
441 t.Fatalf("prepared inbox turn omitted frozen body: %q", input)
442 }
443 if strings.Contains(input, "live-body-after-enqueue") {
444 t.Fatalf("prepared inbox turn re-resolved live reference: %q", input)
445 }
446 if strings.Count(input, "enqueue-time-body") != 1 {
447 t.Fatalf("frozen body injected more than once: %q", input)
448 }
449 }
450
451 func TestInboxFreezesTypedDirectoryAndPathInstructions(t *testing.T) {
452 dir := t.TempDir()
453 workspace := filepath.Join(dir, "workspace")
454 service := filepath.Join(workspace, "service")
455 if err := os.MkdirAll(service, 0o755); err != nil {
456 t.Fatal(err)
457 }
458 for path, body := range map[string]string{
459 filepath.Join(workspace, "AGENTS.md"): "ROOT RULE",
460 filepath.Join(service, "AGENTS.md"): "SERVICE RULE",
461 filepath.Join(service, "old.go"): "package service",
462 } {
463 if err := os.WriteFile(path, []byte(body), 0o644); err != nil {
464 t.Fatal(err)
465 }
466 }
467 sessionPath := filepath.Join(dir, "s.jsonl")
468 sess := agent.NewSession("sys")
469 exec := agent.New(nil, nil, sess, agent.Options{}, event.Discard)
470 sink, done, _ := collectSink()
471 c := newOwnedTestController(t, Options{
472 Runner: appendingRunner{session: sess},
473 Executor: exec,
474 Sink: sink,
475 SessionDir: dir,
476 SessionPath: sessionPath,
477 WorkspaceRoot: workspace,
478 Memory: memory.Load(memory.Options{CWD: workspace}),
479 })
480 defer c.autosaveWG.Wait()
481
482 rec, err := c.EnqueueInbox(InboxRequest{Submit: "review @service"})
483 if err != nil {
484 t.Fatal(err)
485 }
486 _, env, err := c.ReadInboxItem(rec.ItemID)
487 if err != nil {
488 t.Fatal(err)
489 }
490 for _, want := range []string{"<dir ", "old.go", "<path-instructions", "SERVICE RULE"} {
491 if !strings.Contains(env.FrozenRefBlock, want) {
492 t.Fatalf("frozen typed context missing %q:\n%s", want, env.FrozenRefBlock)
493 }
494 }
495 if err := os.WriteFile(filepath.Join(service, "new.go"), []byte("package changed"), 0o644); err != nil {
496 t.Fatal(err)
497 }
498 if _, err := c.TrySubmitInboxItem(rec.ItemID); err != nil {
499 t.Fatal(err)
500 }
501 waitForDone(t, done)
502 input := sess.Snapshot()[len(sess.Snapshot())-1].Content
503 if !strings.Contains(input, "old.go") || strings.Contains(input, "new.go") {
504 t.Fatalf("directory reference was re-resolved live: %q", input)
505 }
506 }
507
508 func TestInboxUsesFrozenImageBytesAfterWorkspaceChanges(t *testing.T) {
509 dir := t.TempDir()
510 workspace := filepath.Join(dir, "workspace")
511 if err := os.MkdirAll(workspace, 0o755); err != nil {
512 t.Fatal(err)
513 }
514 writeVisionTestConfig(t, workspace)
515 imagePath := filepath.Join(workspace, "diagram.png")
516 if err := os.WriteFile(imagePath, mustBase64(t, tinyPNG), 0o644); err != nil {
517 t.Fatal(err)
518 }
519 prov := &recordingProvider{streams: [][]provider.Chunk{{
520 {Type: provider.ChunkText, Text: "done"},
521 {Type: provider.ChunkDone},
522 }}}
523 sess := agent.NewSession("sys")
524 exec := agent.New(prov, tool.NewRegistry(), sess, agent.Options{}, event.Discard)
525 sink, done, _ := collectSink()
526 c := newOwnedTestController(t, Options{
527 Runner: exec,
528 Executor: exec,
529 Sink: sink,
530 SessionDir: dir,
531 SessionPath: filepath.Join(dir, "s.jsonl"),
532 WorkspaceRoot: workspace,
533 ModelRef: "custom/vision-pro",
534 })
535 defer c.autosaveWG.Wait()
536
537 rec, err := c.EnqueueInbox(InboxRequest{Submit: "inspect @diagram.png"})
538 if err != nil {
539 t.Fatal(err)
540 }
541 _, env, err := c.ReadInboxItem(rec.ItemID)
542 if err != nil || len(env.FrozenImages) != 1 {
543 t.Fatalf("frozen image envelope = %+v err=%v", env, err)
544 }
545 frozen := env.FrozenImages[0]
546 if err := os.WriteFile(imagePath, []byte("changed after enqueue"), 0o644); err != nil {
547 t.Fatal(err)
548 }
549 if _, err := c.TrySubmitInboxItem(rec.ItemID); err != nil {
550 t.Fatal(err)
551 }
552 waitForDone(t, done)
553 if len(prov.requests) != 1 {
554 t.Fatalf("provider requests = %d, want 1", len(prov.requests))
555 }
556 messages := prov.requests[0].Messages
557 if len(messages) == 0 || len(messages[len(messages)-1].Images) != 1 || messages[len(messages)-1].Images[0] != frozen {
558 t.Fatalf("provider did not receive the enqueue-time image snapshot: %+v", messages)
559 }
560 }
561
562 func TestTrySubmitInboxAdmissionRaceRestoresQueuedItem(t *testing.T) {
563 dir := t.TempDir()
564 session := filepath.Join(dir, "s.jsonl")
565 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
566 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
567 rec, err := c.EnqueueInbox(InboxRequest{Submit: "must remain durable"})
568 if err != nil {
569 t.Fatal(err)
570 }
571 competingStarted := make(chan struct{})
572 releaseCompeting := make(chan struct{})
573 c.inbox.mu.Lock()
574 c.inbox.beforePreparedAdmission = func() {
575 if result := c.runGuarded(func(context.Context) error {
576 close(competingStarted)
577 <-releaseCompeting
578 return nil
579 }); result != turnStarted {
580 t.Errorf("competing admission = %v, want turnStarted", result)
581 }
582 select {
583 case <-competingStarted:
584 case <-time.After(time.Second):
585 t.Error("competing turn did not start")
586 }
587 }
588 c.inbox.mu.Unlock()
589
590 receipt, err := c.TrySubmitInboxItem(rec.ItemID)
591 if err != nil {
592 t.Fatal(err)
593 }
594 if receipt.Disposition != sessioninbox.DispositionRejectedBusy {
595 t.Fatalf("race disposition = %q, want rejected_busy", receipt.Disposition)
596 }
597 meta, _, err := c.ReadInboxItem(rec.ItemID)
598 if err != nil || meta.State != sessioninbox.StateQueued {
599 t.Fatalf("raced item = %+v err=%v, want durable queued", meta, err)
600 }
601 if err := c.SetInboxPaused(true); err != nil {
602 t.Fatal(err)
603 }
604 close(releaseCompeting)
605 c.autosaveWG.Wait()
606 }
607
608 func TestCancelWithInboxItemsDiscardsOnlyOwnedPendingItems(t *testing.T) {
609 dir := t.TempDir()
610 session := filepath.Join(dir, "s.jsonl")
611 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
612 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
613 owned, err := c.EnqueueInbox(InboxRequest{Submit: "owned by composer", Source: "desktop"})
614 if err != nil {
615 t.Fatal(err)
616 }
617 unrelated, err := c.EnqueueInbox(InboxRequest{Submit: "owned by bot", Source: "bot"})
618 if err != nil {
619 t.Fatal(err)
620 }
621 if err := c.CancelWithInboxItems([]string{owned.ItemID, unrelated.ItemID}, "desktop"); err != nil {
622 t.Fatal(err)
623 }
624 snap := c.InboxSnapshot()
625 if snap.Paused {
626 t.Fatal("successful scoped cancel left inbox paused")
627 }
628 if len(snap.Items) != 1 || snap.Items[0].ID != unrelated.ItemID {
629 t.Fatalf("scoped cancel left items = %+v", snap.Items)
630 }
631 }
632
633 func TestRunTurnAcknowledgesAcceptedDurableItems(t *testing.T) {
634 dir := t.TempDir()
635 runner := &fakeTurnRunner{}
636 c := newOwnedTestController(t, Options{
637 Runner: runner,
638 SessionPath: filepath.Join(dir, "s.jsonl"),
639 SessionDir: dir,
640 Sink: event.Discard,
641 })
642 rec, err := c.EnqueueInbox(InboxRequest{Submit: "accepted steer", Idempotency: "steer-1"})
643 if err != nil {
644 t.Fatal(err)
645 }
646 st, err := c.ensureInbox()
647 if err != nil {
648 t.Fatal(err)
649 }
650 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, ""); err != nil {
651 t.Fatal(err)
652 }
653 c.inbox.mu.Lock()
654 c.inbox.trackActive(rec.ItemID)
655 c.inbox.mu.Unlock()
656
657 if err := c.RunTurn(context.Background(), "foreground"); err != nil {
658 t.Fatal(err)
659 }
660 if got := c.InboxSnapshot().Items; len(got) != 0 {
661 t.Fatalf("synchronous completion left accepted item queued: %+v", got)
662 }
663 if len(runner.inputs) != 1 || runner.inputs[0] != "foreground" {
664 t.Fatalf("runner inputs = %q", runner.inputs)
665 }
666 }
667
668 func TestRunInboxTurnClaimsAndAcknowledgesFIFOItems(t *testing.T) {
669 dir := t.TempDir()
670 runner := &fakeTurnRunner{}
671 c := newOwnedTestController(t, Options{
672 Runner: runner,
673 SessionPath: filepath.Join(dir, "s.jsonl"),
674 SessionDir: dir,
675 Sink: event.Discard,
676 })
677 var ids []string
678 for _, input := range []string{"first", "second"} {
679 rec, err := c.EnqueueInbox(InboxRequest{Submit: input, Idempotency: "msg-" + input})
680 if err != nil {
681 t.Fatal(err)
682 }
683 ids = append(ids, rec.ItemID)
684 }
685 for _, id := range ids {
686 if err := c.RunInboxTurn(context.Background(), id); err != nil {
687 t.Fatal(err)
688 }
689 }
690 if got := runner.inputs; len(got) != 2 || got[0] != "first" || got[1] != "second" {
691 t.Fatalf("durable FIFO inputs = %q", got)
692 }
693 if got := c.InboxSnapshot().Items; len(got) != 0 {
694 t.Fatalf("completed FIFO items remain queued: %+v", got)
695 }
696 }
697
698 func TestStructuredInboxInvocationSurvivesReopenAndRunsSkill(t *testing.T) {
699 dir := t.TempDir()
700 path := filepath.Join(dir, "s.jsonl")
701 skills := []skill.Skill{{
702 Name: "init", Body: "INITIALIZE_FROM_DURABLE_INBOX", RunAs: skill.RunInline, Scope: skill.ScopeGlobal,
703 }}
704 first := newOwnedTestController(t, Options{SessionPath: path, SessionDir: dir, Skills: skills, Sink: event.Discard})
705 rec, err := first.EnqueueInbox(InboxRequest{
706 Display: "/init",
707 Idempotency: "desktop-submit-1",
708 Invocations: []InvocationRequest{{Name: "init", Kind: "skill", Offset: 0}},
709 })
710 if err != nil {
711 t.Fatal(err)
712 }
713 first.inbox.mu.Lock()
714 first.inbox.store.Close()
715 first.inbox.store = nil
716 first.inbox.mu.Unlock()
717
718 runner := &fakeTurnRunner{}
719 reopened := newOwnedTestController(t, Options{
720 Runner: runner, SessionPath: path, SessionDir: dir, Skills: skills, Sink: event.Discard,
721 })
722 if err := reopened.RunInboxTurn(context.Background(), rec.ItemID); err != nil {
723 t.Fatal(err)
724 }
725 if len(runner.inputs) != 1 || !strings.Contains(runner.inputs[0], "INITIALIZE_FROM_DURABLE_INBOX") {
726 t.Fatalf("reopened structured turn lost skill semantics: %q", runner.inputs)
727 }
728 if strings.Contains(runner.inputs[0], "/init") {
729 t.Fatalf("structured turn degraded to slash text: %q", runner.inputs[0])
730 }
731 if got := reopened.InboxSnapshot().Items; len(got) != 0 {
732 t.Fatalf("structured item was not acknowledged: %+v", got)
733 }
734 }
735
736 func TestLegacySingularInboxInvocationInfersSkillKind(t *testing.T) {
737 dir := t.TempDir()
738 runner := &fakeTurnRunner{}
739 c := newOwnedTestController(t, Options{
740 Runner: runner, SessionPath: filepath.Join(dir, "s.jsonl"), SessionDir: dir, Sink: event.Discard,
741 Skills: []skill.Skill{{Name: "legacy", Body: "LEGACY_SKILL_BODY", RunAs: skill.RunInline, Scope: skill.ScopeGlobal}},
742 })
743 st, err := c.ensureInbox()
744 if err != nil {
745 t.Fatal(err)
746 }
747 rec, err := st.Enqueue(sessioninbox.EnqueueRequest{Envelope: sessioninbox.PromptEnvelope{
748 DisplayText: "/legacy",
749 Invocation: &sessioninbox.StructuredInvocation{Name: "legacy"},
750 }})
751 if err != nil {
752 t.Fatal(err)
753 }
754 if err := c.RunInboxTurn(context.Background(), rec.ItemID); err != nil {
755 t.Fatal(err)
756 }
757 if len(runner.inputs) != 1 || !strings.Contains(runner.inputs[0], "LEGACY_SKILL_BODY") {
758 t.Fatalf("legacy structured input = %q", runner.inputs)
759 }
760 }
761
761 lines GO