返回 DeepSeek-Reasonix
inbox_session_fence_test.go
根目录 / internal / control / inbox_session_fence_test.go
1 package control
2
3 import (
4 "errors"
5 "path/filepath"
6 "testing"
7
8 "reasonix/internal/agent"
9 "reasonix/internal/event"
10 "reasonix/internal/session"
11 "reasonix/internal/sessioninbox"
12 "reasonix/internal/tool"
13 )
14
15 func TestTargetGuidanceDoesNotSteerSuccessorTurn(t *testing.T) {
16 dir := t.TempDir()
17 path := filepath.Join(dir, "session.jsonl")
18 prov := &inboxSteerProvider{started: make(chan struct{}), release: make(chan struct{})}
19 exec := agent.New(prov, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
20 c := newOwnedTestController(t, Options{Runner: exec, Executor: exec, SessionDir: dir, SessionPath: path, Sink: event.Discard})
21 c.Submit("successor turn")
22 prov.awaitStarted(t, c)
23 current := c.RuntimeStatus().TurnID
24 if current == "" {
25 t.Fatal("active turn has no identity")
26 }
27 stale, err := c.InboxQueue(path, InboxQueueRequest{Kind: "enqueue_steer", TurnID: "previous-turn", Text: "for the previous turn", IdempotencyKey: "stale"})
28 if err != nil || stale.Receipt == nil || stale.Receipt.Disposition != sessioninbox.DispositionQueuedFollowup {
29 t.Fatalf("stale guidance: %+v %v", stale, err)
30 }
31 meta, _, err := c.ReadInboxItem(stale.Receipt.ItemID)
32 if err != nil || meta.State != sessioninbox.StateQueued || meta.Intent != sessioninbox.IntentFollowup {
33 t.Fatalf("stale item: %+v %v", meta, err)
34 }
35 matching, err := c.InboxQueue(path, InboxQueueRequest{Kind: "enqueue_steer", TurnID: current, Text: "for the current turn", IdempotencyKey: "matching"})
36 if err != nil || matching.Receipt == nil || matching.Receipt.Disposition != sessioninbox.DispositionSteerAccepted {
37 t.Fatalf("matching guidance: %+v %v", matching, err)
38 }
39 if c.RuntimeStatus().TurnID != current {
40 t.Fatal("guidance replaced the running turn")
41 }
42 }
43
44 func TestInboxExpectedSessionCannotSubmitOrConfirmReplacement(t *testing.T) {
45 dir := t.TempDir()
46 first, second := filepath.Join(dir, "first.jsonl"), filepath.Join(dir, "second.jsonl")
47 c := newOwnedTestController(t, Options{SessionDir: dir, SessionPath: first, Sink: event.Discard})
48 defer c.Close()
49 if err := c.SetInboxPaused(true); err != nil {
50 t.Fatal(err)
51 }
52 request := InboxRequest{ExpectedSessionPath: first, Submit: "original", Idempotency: "original"}
53 receipt, err := c.TryEnqueueFollowup(request)
54 if err != nil {
55 t.Fatal(err)
56 }
57 confirmed, found, err := c.LookupInboxReceiptForSession(first, request.Idempotency)
58 if err != nil || !found || confirmed.ItemID != receipt.ItemID {
59 t.Fatalf("original confirmation = %+v, %v, %v", confirmed, found, err)
60 }
61 c.SetSessionPath(second)
62 if err := c.SetInboxPaused(true); err != nil {
63 t.Fatal(err)
64 }
65 if _, err := c.TryEnqueueFollowup(request); !errors.Is(err, ErrInboxSessionChanged) {
66 t.Fatalf("stale request = %v", err)
67 }
68 if _, _, err := c.LookupInboxReceiptForSession(first, request.Idempotency); !errors.Is(err, ErrInboxSessionChanged) {
69 t.Fatalf("stale lookup = %v", err)
70 }
71 if got := c.InboxSnapshot(); got.SessionPath != second || len(got.Items) != 0 {
72 t.Fatalf("replacement mutated: %+v", got)
73 }
74 request.ExpectedSessionPath = ""
75 if _, err := c.TryEnqueueFollowup(request); err != nil {
76 t.Fatalf("legacy request no longer works: %v", err)
77 }
78 }
79
80 func TestCanonicalInboxUsesSessionIdentityAcrossQueueOperations(t *testing.T) {
81 service, err := session.NewService("desktop", session.NewFilesystemPersistence(t.TempDir()))
82 if err != nil {
83 t.Fatal(err)
84 }
85 exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
86 c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service})
87 if _, err := c.BindFreshSession(t.Context(), "first"); err != nil {
88 t.Fatal(err)
89 }
90 const first = "session-id:first"
91 if err := c.SetInboxPaused(true); err != nil {
92 t.Fatal(err)
93 }
94 receipt, err := c.TryEnqueueFollowup(InboxRequest{ExpectedSessionPath: first, Submit: "follow up", Idempotency: "request"})
95 if err != nil {
96 t.Fatal(err)
97 }
98 confirmed, found, err := c.LookupInboxReceiptForSession(first, "request")
99 if err != nil || !found || confirmed.ItemID != receipt.ItemID {
100 t.Fatalf("canonical receipt: %+v %v %v", confirmed, found, err)
101 }
102 result, err := c.InboxQueue(first, InboxQueueRequest{Kind: "snapshot"})
103 if err != nil || result.Outcome != "unchanged" || result.Snapshot.SessionPath != first || len(result.Snapshot.Items) != 1 {
104 t.Fatalf("canonical queue: %+v %v", result, err)
105 }
106 result, err = c.InboxQueue(first, InboxQueueRequest{Kind: "read", ItemID: receipt.ItemID})
107 if err != nil || result.Edit == nil || result.Edit.Text != "follow up" {
108 t.Fatalf("canonical queue read: %+v %v", result, err)
109 }
110 if _, err := c.trySteerInboxItemForSession(receipt.ItemID, "old-turn", first); errors.Is(err, ErrInboxSessionChanged) {
111 t.Fatal("same-session steer rejected its identity")
112 }
113 if _, err := c.BindFreshSession(t.Context(), "second"); err != nil {
114 t.Fatal(err)
115 }
116 if _, err := c.EnqueueInbox(InboxRequest{ExpectedSessionPath: first, Submit: "stale"}); !errors.Is(err, ErrInboxSessionChanged) {
117 t.Fatalf("stale canonical enqueue: %v", err)
118 }
119 if _, _, err := c.LookupInboxReceiptForSession(first, "request"); !errors.Is(err, ErrInboxSessionChanged) {
120 t.Fatalf("stale canonical receipt: %v", err)
121 }
122 result, err = c.InboxQueue(first, InboxQueueRequest{Kind: "snapshot"})
123 if err != nil || result.Reason != "session_changed" {
124 t.Fatalf("stale canonical queue: %+v %v", result, err)
125 }
126 if got := c.InboxSnapshot(); got.SessionPath != "session-id:second" || len(got.Items) != 0 {
127 t.Fatalf("queue crossed the session boundary: %+v", got)
128 }
129 }
130
131 func TestCanonicalInboxAcceptsMessagesDuringTurnAndDispatchesFIFO(t *testing.T) {
132 service, err := session.NewService("desktop", session.NewFilesystemPersistence(t.TempDir()))
133 if err != nil {
134 t.Fatal(err)
135 }
136 runner := &gatedInboxDispatchRunner{inputs: make(chan string, 8), firstStarted: make(chan struct{}), releaseFirst: make(chan struct{})}
137 done := make(chan struct{}, 8)
138 exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
139 c := newOwnedTestController(t, Options{Executor: exec, Runner: runner, SessionService: service, Sink: event.FuncSink(func(e event.Event) {
140 if e.Kind == event.TurnDone {
141 done <- struct{}{}
142 }
143 })})
144 if _, err := c.BindFreshSession(t.Context(), "active-queue"); err != nil {
145 t.Fatal(err)
146 }
147 c.Submit("active turn")
148 inputs := &inboxDispatchRunner{inputs: runner.inputs}
149 if got := waitForInboxDispatch(t, c, inputs); got != "active turn" {
150 t.Fatalf("initial input: %q", got)
151 }
152 for _, text := range []string{"queued one", "queued two"} {
153 receipt, err := c.TryEnqueueFollowup(InboxRequest{ExpectedSessionPath: "session-id:active-queue", Submit: text, Idempotency: text})
154 if err != nil || receipt.ItemID == "" {
155 t.Fatalf("running canonical enqueue: %+v %v", receipt, err)
156 }
157 }
158 close(runner.releaseFirst)
159 waitForInboxTurnDone(t, c, done)
160 for _, want := range []string{"queued one", "queued two"} {
161 if got := waitForInboxDispatch(t, c, inputs); got != want {
162 t.Fatalf("canonical FIFO input = %q, want %q", got, want)
163 }
164 waitForInboxTurnDone(t, c, done)
165 }
166 if got := c.InboxSnapshot(); got.SessionPath != "session-id:active-queue" || len(got.Items) != 0 || got.Paused {
167 t.Fatalf("canonical completion left queued work: %+v", got)
168 }
169 }
170
170 lines GO