返回 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) Stream(ctx context.Context, req provider.Request) (<-chan provider.Chunk, error) {
271 p.requests = append(p.requests, req)
272 ch := make(chan provider.Chunk, 2)
273 if len(p.requests) == 1 {
274 close(p.started)
275 go func() {
276 defer close(ch)
277 select {
278 case <-p.release:
279 ch <- provider.Chunk{Type: provider.ChunkText, Text: "ready"}
280 ch <- provider.Chunk{Type: provider.ChunkDone}
281 case <-ctx.Done():
282 }
283 }()
284 return ch, nil
285 }
286 ch <- provider.Chunk{Type: provider.ChunkText, Text: "applied"}
287 ch <- provider.Chunk{Type: provider.ChunkDone}
288 close(ch)
289 return ch, nil
290 }
291
292 func TestThirtySteersApplyAndAckExactlyOnce(t *testing.T) {
293 dir := t.TempDir()
294 prov := &inboxSteerProvider{started: make(chan struct{}), release: make(chan struct{})}
295 sess := agent.NewSession("sys")
296 exec := agent.New(prov, tool.NewRegistry(), sess, agent.Options{}, event.Discard)
297 sink, done, _ := collectSink()
298 c := newOwnedTestController(t, Options{
299 Runner: exec,
300 Executor: exec,
301 Sink: sink,
302 SessionDir: dir,
303 SessionPath: filepath.Join(dir, "s.jsonl"),
304 })
305 defer c.autosaveWG.Wait()
306 c.Submit("initial turn")
307 select {
308 case <-prov.started:
309 case <-time.After(time.Second):
310 t.Fatal("initial provider turn did not start")
311 }
312
313 const steerCount = 30
314 for i := range steerCount {
315 body := fmt.Sprintf("durable-steer-%02d", i)
316 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: body})
317 if err != nil {
318 t.Fatal(err)
319 }
320 got, err := c.TrySteerInboxItem(rec.ItemID)
321 if err != nil {
322 t.Fatal(err)
323 }
324 if got.Disposition != sessioninbox.DispositionSteerAccepted {
325 t.Fatalf("steer %d disposition = %q", i, got.Disposition)
326 }
327 }
328 close(prov.release)
329 // Thirty durable round trips are real filesystem work; a loaded Windows
330 // runner spends most of the default five seconds before the turn is even
331 // released. This asserts exactly-once acknowledgement, not latency.
332 waitForDoneWithin(t, done, 60*time.Second)
333
334 if items := c.InboxSnapshot().Items; len(items) != 0 {
335 t.Fatalf("accepted steers were not all acknowledged: %+v", items)
336 }
337 if got := len(prov.requests); got != steerCount+1 {
338 t.Fatalf("provider requests = %d, want %d", got, steerCount+1)
339 }
340 messages := sess.Snapshot()
341 for i := range steerCount {
342 body := fmt.Sprintf("durable-steer-%02d", i)
343 count := 0
344 for _, message := range messages {
345 count += strings.Count(message.Content, body)
346 }
347 if count != 1 {
348 t.Fatalf("%q appears %d times in transcript, want exactly once", body, count)
349 }
350 }
351 }
352
353 func TestMultiSteerActiveSetAcksAll(t *testing.T) {
354 dir := t.TempDir()
355 session := filepath.Join(dir, "s.jsonl")
356 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
357 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
358
359 st, err := c.ensureInbox()
360 if err != nil {
361 t.Fatal(err)
362 }
363 var ids []string
364 for i := range 3 {
365 rec, err := c.EnqueueInbox(InboxRequest{Submit: "body-" + string(rune('a'+i))})
366 if err != nil {
367 t.Fatal(err)
368 }
369 ids = append(ids, rec.ItemID)
370 _ = st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, "")
371 }
372 c.inbox.mu.Lock()
373 c.inbox.clearActive()
374 for _, id := range ids {
375 c.inbox.trackActive(id)
376 }
377 c.inbox.mu.Unlock()
378
379 c.onInboxTurnDone()
380 if n := len(c.InboxSnapshot().Items); n != 0 {
381 t.Fatalf("want all 3 steers acked/dequeued, still have %d items", n)
382 }
383 }
384
385 func TestSubmitInboxUsesFrozenReferenceWithoutLiveReresolve(t *testing.T) {
386 dir := t.TempDir()
387 workspace := filepath.Join(dir, "workspace")
388 if err := os.MkdirAll(workspace, 0o755); err != nil {
389 t.Fatal(err)
390 }
391 refPath := filepath.Join(workspace, "note.txt")
392 if err := os.WriteFile(refPath, []byte("enqueue-time-body"), 0o600); err != nil {
393 t.Fatal(err)
394 }
395 sessionPath := filepath.Join(dir, "s.jsonl")
396 sess := agent.NewSession("sys")
397 exec := agent.New(nil, nil, sess, agent.Options{}, event.Discard)
398 sink, done, _ := collectSink()
399 c := newOwnedTestController(t, Options{
400 Runner: appendingRunner{session: sess},
401 Executor: exec,
402 Sink: sink,
403 SessionDir: dir,
404 SessionPath: sessionPath,
405 WorkspaceRoot: workspace,
406 })
407 defer c.autosaveWG.Wait()
408
409 rec, err := c.EnqueueInbox(InboxRequest{Submit: "review @note.txt"})
410 if err != nil {
411 t.Fatal(err)
412 }
413 if err := os.WriteFile(refPath, []byte("live-body-after-enqueue"), 0o600); err != nil {
414 t.Fatal(err)
415 }
416 got, err := c.TrySubmitInboxItem(rec.ItemID)
417 if err != nil {
418 t.Fatal(err)
419 }
420 if got.Disposition != sessioninbox.DispositionStarted {
421 t.Fatalf("disposition = %q, want started", got.Disposition)
422 }
423 waitForDone(t, done)
424
425 messages := sess.Snapshot()
426 if len(messages) < 2 {
427 t.Fatalf("messages = %+v", messages)
428 }
429 input := messages[len(messages)-1].Content
430 if !strings.Contains(input, "enqueue-time-body") {
431 t.Fatalf("prepared inbox turn omitted frozen body: %q", input)
432 }
433 if strings.Contains(input, "live-body-after-enqueue") {
434 t.Fatalf("prepared inbox turn re-resolved live reference: %q", input)
435 }
436 if strings.Count(input, "enqueue-time-body") != 1 {
437 t.Fatalf("frozen body injected more than once: %q", input)
438 }
439 }
440
441 func TestInboxFreezesTypedDirectoryAndPathInstructions(t *testing.T) {
442 dir := t.TempDir()
443 workspace := filepath.Join(dir, "workspace")
444 service := filepath.Join(workspace, "service")
445 if err := os.MkdirAll(service, 0o755); err != nil {
446 t.Fatal(err)
447 }
448 for path, body := range map[string]string{
449 filepath.Join(workspace, "AGENTS.md"): "ROOT RULE",
450 filepath.Join(service, "AGENTS.md"): "SERVICE RULE",
451 filepath.Join(service, "old.go"): "package service",
452 } {
453 if err := os.WriteFile(path, []byte(body), 0o644); err != nil {
454 t.Fatal(err)
455 }
456 }
457 sessionPath := filepath.Join(dir, "s.jsonl")
458 sess := agent.NewSession("sys")
459 exec := agent.New(nil, nil, sess, agent.Options{}, event.Discard)
460 sink, done, _ := collectSink()
461 c := newOwnedTestController(t, Options{
462 Runner: appendingRunner{session: sess},
463 Executor: exec,
464 Sink: sink,
465 SessionDir: dir,
466 SessionPath: sessionPath,
467 WorkspaceRoot: workspace,
468 Memory: memory.Load(memory.Options{CWD: workspace}),
469 })
470 defer c.autosaveWG.Wait()
471
472 rec, err := c.EnqueueInbox(InboxRequest{Submit: "review @service"})
473 if err != nil {
474 t.Fatal(err)
475 }
476 _, env, err := c.ReadInboxItem(rec.ItemID)
477 if err != nil {
478 t.Fatal(err)
479 }
480 for _, want := range []string{"<dir ", "old.go", "<path-instructions", "SERVICE RULE"} {
481 if !strings.Contains(env.FrozenRefBlock, want) {
482 t.Fatalf("frozen typed context missing %q:\n%s", want, env.FrozenRefBlock)
483 }
484 }
485 if err := os.WriteFile(filepath.Join(service, "new.go"), []byte("package changed"), 0o644); err != nil {
486 t.Fatal(err)
487 }
488 if _, err := c.TrySubmitInboxItem(rec.ItemID); err != nil {
489 t.Fatal(err)
490 }
491 waitForDone(t, done)
492 input := sess.Snapshot()[len(sess.Snapshot())-1].Content
493 if !strings.Contains(input, "old.go") || strings.Contains(input, "new.go") {
494 t.Fatalf("directory reference was re-resolved live: %q", input)
495 }
496 }
497
498 func TestInboxUsesFrozenImageBytesAfterWorkspaceChanges(t *testing.T) {
499 dir := t.TempDir()
500 workspace := filepath.Join(dir, "workspace")
501 if err := os.MkdirAll(workspace, 0o755); err != nil {
502 t.Fatal(err)
503 }
504 writeVisionTestConfig(t, workspace)
505 imagePath := filepath.Join(workspace, "diagram.png")
506 if err := os.WriteFile(imagePath, mustBase64(t, tinyPNG), 0o644); err != nil {
507 t.Fatal(err)
508 }
509 prov := &recordingProvider{streams: [][]provider.Chunk{{
510 {Type: provider.ChunkText, Text: "done"},
511 {Type: provider.ChunkDone},
512 }}}
513 sess := agent.NewSession("sys")
514 exec := agent.New(prov, tool.NewRegistry(), sess, agent.Options{}, event.Discard)
515 sink, done, _ := collectSink()
516 c := newOwnedTestController(t, Options{
517 Runner: exec,
518 Executor: exec,
519 Sink: sink,
520 SessionDir: dir,
521 SessionPath: filepath.Join(dir, "s.jsonl"),
522 WorkspaceRoot: workspace,
523 ModelRef: "custom/vision-pro",
524 })
525 defer c.autosaveWG.Wait()
526
527 rec, err := c.EnqueueInbox(InboxRequest{Submit: "inspect @diagram.png"})
528 if err != nil {
529 t.Fatal(err)
530 }
531 _, env, err := c.ReadInboxItem(rec.ItemID)
532 if err != nil || len(env.FrozenImages) != 1 {
533 t.Fatalf("frozen image envelope = %+v err=%v", env, err)
534 }
535 frozen := env.FrozenImages[0]
536 if err := os.WriteFile(imagePath, []byte("changed after enqueue"), 0o644); err != nil {
537 t.Fatal(err)
538 }
539 if _, err := c.TrySubmitInboxItem(rec.ItemID); err != nil {
540 t.Fatal(err)
541 }
542 waitForDone(t, done)
543 if len(prov.requests) != 1 {
544 t.Fatalf("provider requests = %d, want 1", len(prov.requests))
545 }
546 messages := prov.requests[0].Messages
547 if len(messages) == 0 || len(messages[len(messages)-1].Images) != 1 || messages[len(messages)-1].Images[0] != frozen {
548 t.Fatalf("provider did not receive the enqueue-time image snapshot: %+v", messages)
549 }
550 }
551
552 func TestTrySubmitInboxAdmissionRaceRestoresQueuedItem(t *testing.T) {
553 dir := t.TempDir()
554 session := filepath.Join(dir, "s.jsonl")
555 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
556 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
557 rec, err := c.EnqueueInbox(InboxRequest{Submit: "must remain durable"})
558 if err != nil {
559 t.Fatal(err)
560 }
561 competingStarted := make(chan struct{})
562 releaseCompeting := make(chan struct{})
563 c.inbox.mu.Lock()
564 c.inbox.beforePreparedAdmission = func() {
565 if result := c.runGuarded(func(context.Context) error {
566 close(competingStarted)
567 <-releaseCompeting
568 return nil
569 }); result != turnStarted {
570 t.Errorf("competing admission = %v, want turnStarted", result)
571 }
572 select {
573 case <-competingStarted:
574 case <-time.After(time.Second):
575 t.Error("competing turn did not start")
576 }
577 }
578 c.inbox.mu.Unlock()
579
580 receipt, err := c.TrySubmitInboxItem(rec.ItemID)
581 if err != nil {
582 t.Fatal(err)
583 }
584 if receipt.Disposition != sessioninbox.DispositionRejectedBusy {
585 t.Fatalf("race disposition = %q, want rejected_busy", receipt.Disposition)
586 }
587 meta, _, err := c.ReadInboxItem(rec.ItemID)
588 if err != nil || meta.State != sessioninbox.StateQueued {
589 t.Fatalf("raced item = %+v err=%v, want durable queued", meta, err)
590 }
591 if err := c.SetInboxPaused(true); err != nil {
592 t.Fatal(err)
593 }
594 close(releaseCompeting)
595 c.autosaveWG.Wait()
596 }
597
598 func TestCancelWithInboxItemsDiscardsOnlyOwnedPendingItems(t *testing.T) {
599 dir := t.TempDir()
600 session := filepath.Join(dir, "s.jsonl")
601 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
602 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
603 owned, err := c.EnqueueInbox(InboxRequest{Submit: "owned by composer", Source: "desktop"})
604 if err != nil {
605 t.Fatal(err)
606 }
607 unrelated, err := c.EnqueueInbox(InboxRequest{Submit: "owned by bot", Source: "bot"})
608 if err != nil {
609 t.Fatal(err)
610 }
611 if err := c.CancelWithInboxItems([]string{owned.ItemID, unrelated.ItemID}, "desktop"); err != nil {
612 t.Fatal(err)
613 }
614 snap := c.InboxSnapshot()
615 if snap.Paused {
616 t.Fatal("successful scoped cancel left inbox paused")
617 }
618 if len(snap.Items) != 1 || snap.Items[0].ID != unrelated.ItemID {
619 t.Fatalf("scoped cancel left items = %+v", snap.Items)
620 }
621 }
622
623 func TestRunTurnAcknowledgesAcceptedDurableItems(t *testing.T) {
624 dir := t.TempDir()
625 runner := &fakeTurnRunner{}
626 c := newOwnedTestController(t, Options{
627 Runner: runner,
628 SessionPath: filepath.Join(dir, "s.jsonl"),
629 SessionDir: dir,
630 Sink: event.Discard,
631 })
632 rec, err := c.EnqueueInbox(InboxRequest{Submit: "accepted steer", Idempotency: "steer-1"})
633 if err != nil {
634 t.Fatal(err)
635 }
636 st, err := c.ensureInbox()
637 if err != nil {
638 t.Fatal(err)
639 }
640 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerConsumed, ""); err != nil {
641 t.Fatal(err)
642 }
643 c.inbox.mu.Lock()
644 c.inbox.trackActive(rec.ItemID)
645 c.inbox.mu.Unlock()
646
647 if err := c.RunTurn(context.Background(), "foreground"); err != nil {
648 t.Fatal(err)
649 }
650 if got := c.InboxSnapshot().Items; len(got) != 0 {
651 t.Fatalf("synchronous completion left accepted item queued: %+v", got)
652 }
653 if len(runner.inputs) != 1 || runner.inputs[0] != "foreground" {
654 t.Fatalf("runner inputs = %q", runner.inputs)
655 }
656 }
657
658 func TestRunInboxTurnClaimsAndAcknowledgesFIFOItems(t *testing.T) {
659 dir := t.TempDir()
660 runner := &fakeTurnRunner{}
661 c := newOwnedTestController(t, Options{
662 Runner: runner,
663 SessionPath: filepath.Join(dir, "s.jsonl"),
664 SessionDir: dir,
665 Sink: event.Discard,
666 })
667 var ids []string
668 for _, input := range []string{"first", "second"} {
669 rec, err := c.EnqueueInbox(InboxRequest{Submit: input, Idempotency: "msg-" + input})
670 if err != nil {
671 t.Fatal(err)
672 }
673 ids = append(ids, rec.ItemID)
674 }
675 for _, id := range ids {
676 if err := c.RunInboxTurn(context.Background(), id); err != nil {
677 t.Fatal(err)
678 }
679 }
680 if got := runner.inputs; len(got) != 2 || got[0] != "first" || got[1] != "second" {
681 t.Fatalf("durable FIFO inputs = %q", got)
682 }
683 if got := c.InboxSnapshot().Items; len(got) != 0 {
684 t.Fatalf("completed FIFO items remain queued: %+v", got)
685 }
686 }
687
688 func TestStructuredInboxInvocationSurvivesReopenAndRunsSkill(t *testing.T) {
689 dir := t.TempDir()
690 path := filepath.Join(dir, "s.jsonl")
691 skills := []skill.Skill{{
692 Name: "init", Body: "INITIALIZE_FROM_DURABLE_INBOX", RunAs: skill.RunInline, Scope: skill.ScopeGlobal,
693 }}
694 first := newOwnedTestController(t, Options{SessionPath: path, SessionDir: dir, Skills: skills, Sink: event.Discard})
695 rec, err := first.EnqueueInbox(InboxRequest{
696 Display: "/init",
697 Idempotency: "desktop-submit-1",
698 Invocations: []InvocationRequest{{Name: "init", Kind: "skill", Offset: 0}},
699 })
700 if err != nil {
701 t.Fatal(err)
702 }
703 first.inbox.mu.Lock()
704 first.inbox.store.Close()
705 first.inbox.store = nil
706 first.inbox.mu.Unlock()
707
708 runner := &fakeTurnRunner{}
709 reopened := newOwnedTestController(t, Options{
710 Runner: runner, SessionPath: path, SessionDir: dir, Skills: skills, Sink: event.Discard,
711 })
712 if err := reopened.RunInboxTurn(context.Background(), rec.ItemID); err != nil {
713 t.Fatal(err)
714 }
715 if len(runner.inputs) != 1 || !strings.Contains(runner.inputs[0], "INITIALIZE_FROM_DURABLE_INBOX") {
716 t.Fatalf("reopened structured turn lost skill semantics: %q", runner.inputs)
717 }
718 if strings.Contains(runner.inputs[0], "/init") {
719 t.Fatalf("structured turn degraded to slash text: %q", runner.inputs[0])
720 }
721 if got := reopened.InboxSnapshot().Items; len(got) != 0 {
722 t.Fatalf("structured item was not acknowledged: %+v", got)
723 }
724 }
725
726 func TestLegacySingularInboxInvocationInfersSkillKind(t *testing.T) {
727 dir := t.TempDir()
728 runner := &fakeTurnRunner{}
729 c := newOwnedTestController(t, Options{
730 Runner: runner, SessionPath: filepath.Join(dir, "s.jsonl"), SessionDir: dir, Sink: event.Discard,
731 Skills: []skill.Skill{{Name: "legacy", Body: "LEGACY_SKILL_BODY", RunAs: skill.RunInline, Scope: skill.ScopeGlobal}},
732 })
733 st, err := c.ensureInbox()
734 if err != nil {
735 t.Fatal(err)
736 }
737 rec, err := st.Enqueue(sessioninbox.EnqueueRequest{Envelope: sessioninbox.PromptEnvelope{
738 DisplayText: "/legacy",
739 Invocation: &sessioninbox.StructuredInvocation{Name: "legacy"},
740 }})
741 if err != nil {
742 t.Fatal(err)
743 }
744 if err := c.RunInboxTurn(context.Background(), rec.ItemID); err != nil {
745 t.Fatal(err)
746 }
747 if len(runner.inputs) != 1 || !strings.Contains(runner.inputs[0], "LEGACY_SKILL_BODY") {
748 t.Fatalf("legacy structured input = %q", runner.inputs)
749 }
750 }
751
751 lines GO