返回 DeepSeek-Reasonix
inbox_reconciliation_test.go
根目录 / internal / control / inbox_reconciliation_test.go
1 package control
2
3 import (
4 "errors"
5 "os"
6 "path/filepath"
7 "testing"
8 "time"
9
10 "reasonix/internal/agent"
11 "reasonix/internal/event"
12 "reasonix/internal/sessioninbox"
13 "reasonix/internal/tool"
14 "reasonix/internal/transcript"
15 )
16
17 func TestSteerEventFollowsDurableConsumedTransition(t *testing.T) {
18 dir := t.TempDir()
19 prov := &inboxSteerProvider{started: make(chan struct{}), release: make(chan struct{})}
20 exec := agent.New(prov, tool.NewRegistry(), agent.NewSession("sys"), agent.Options{}, event.Discard)
21 observed := make(chan sessioninbox.InboxState, 1)
22 steerEvents := make(chan event.Event, 1)
23 done := make(chan struct{})
24 var c *Controller
25 sink := event.FuncSink(func(e event.Event) {
26 if e.Kind == event.Steer && c != nil {
27 state := sessioninbox.InboxState("")
28 for _, item := range c.InboxSnapshot().Items {
29 if item.ID == e.ItemID {
30 state = item.State
31 break
32 }
33 }
34 observed <- state
35 steerEvents <- e
36 }
37 if e.Kind == event.TurnDone {
38 select {
39 case <-done:
40 default:
41 close(done)
42 }
43 }
44 })
45 c = newOwnedTestController(t, Options{
46 Runner: exec,
47 Executor: exec,
48 Sink: sink,
49 SessionDir: dir,
50 SessionPath: filepath.Join(dir, "s.jsonl"),
51 })
52 t.Cleanup(func() {
53 c.Close()
54 c.autosaveWG.Wait()
55 })
56 c.Submit("initial turn")
57 prov.awaitStarted(t, c)
58 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "durable steer"})
59 if err != nil {
60 t.Fatal(err)
61 }
62 if _, err := c.TrySteerInboxItem(rec.ItemID); err != nil {
63 t.Fatal(err)
64 }
65 close(prov.release)
66 select {
67 case state := <-observed:
68 if state != sessioninbox.StateSteerConsumed {
69 t.Fatalf("state at steer event = %q, want %q", state, sessioninbox.StateSteerConsumed)
70 }
71 case <-time.After(5 * time.Second):
72 t.Fatal("timed out waiting for steer event")
73 }
74 select {
75 case <-done:
76 case <-time.After(5 * time.Second):
77 t.Fatal("timed out waiting for turn completion")
78 }
79 e := <-steerEvents
80 if e.MessageID == "" || e.ItemID != rec.ItemID {
81 t.Fatalf("steer receipt lacks message/inbox identity: %+v", e)
82 }
83 var count int
84 for _, row := range transcript.History(exec.Session().Snapshot(), transcript.HistoryOptions{}) {
85 if row.MessageID == e.MessageID {
86 count++
87 if row.RecordID != "m:"+e.MessageID || row.Role != "notice" || row.Content != "↪ durable steer" {
88 t.Fatalf("steer receipt and history disagree: %+v", row)
89 }
90 }
91 }
92 if count != 1 {
93 t.Fatalf("steer receipt owns %d canonical rows, want 1", count)
94 }
95 }
96
97 func TestCancelWithInboxItemsResultRestoresOnlyUnconsumedItems(t *testing.T) {
98 dir := t.TempDir()
99 session := filepath.Join(dir, "s.jsonl")
100 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
101 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
102 accepted, err := c.EnqueueInbox(InboxRequest{Submit: "accepted", Source: "desktop"})
103 if err != nil {
104 t.Fatal(err)
105 }
106 consumed, err := c.EnqueueInbox(InboxRequest{Submit: "consumed", Source: "desktop"})
107 if err != nil {
108 t.Fatal(err)
109 }
110 st, err := c.ensureInbox()
111 if err != nil {
112 t.Fatal(err)
113 }
114 if err := st.SetState(accepted.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
115 t.Fatal(err)
116 }
117 if err := st.SetState(consumed.ItemID, sessioninbox.StateSteerConsumed, ""); err != nil {
118 t.Fatal(err)
119 }
120
121 result, err := c.CancelWithInboxItemsResult([]string{accepted.ItemID, consumed.ItemID}, "desktop")
122 if err != nil {
123 t.Fatal(err)
124 }
125 if len(result.DiscardedItemIDs) != 1 || result.DiscardedItemIDs[0] != accepted.ItemID {
126 t.Fatalf("discarded ids = %v", result.DiscardedItemIDs)
127 }
128 items := c.InboxSnapshot().Items
129 if len(items) != 1 || items[0].ID != consumed.ItemID {
130 t.Fatalf("remaining items = %+v", items)
131 }
132 }
133
134 func TestDeleteInboxItemDoesNotOverwriteConsumedSteer(t *testing.T) {
135 dir := t.TempDir()
136 session := filepath.Join(dir, "s.jsonl")
137 _ = os.WriteFile(session, []byte("{}\n"), 0o644)
138 c := newOwnedTestController(t, Options{SessionPath: session, SessionDir: dir, Sink: event.Discard})
139 rec, err := c.EnqueueInbox(InboxRequest{Intent: sessioninbox.IntentSteer, Submit: "consumed"})
140 if err != nil {
141 t.Fatal(err)
142 }
143 st, err := c.ensureInbox()
144 if err != nil {
145 t.Fatal(err)
146 }
147 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
148 t.Fatal(err)
149 }
150 c.inbox.mu.Lock()
151 c.inbox.trackActive(rec.ItemID)
152 c.inbox.mu.Unlock()
153 if err := st.MarkSteerConsumed(rec.ItemID); err != nil {
154 t.Fatal(err)
155 }
156 if err := c.DeleteInboxItem(rec.ItemID); !errors.Is(err, sessioninbox.ErrInvalidState) {
157 t.Fatalf("delete consumed steer = %v, want ErrInvalidState", err)
158 }
159 meta, _, err := c.ReadInboxItem(rec.ItemID)
160 if err != nil {
161 t.Fatal(err)
162 }
163 if meta.State != sessioninbox.StateSteerConsumed {
164 t.Fatalf("state = %q, want %q", meta.State, sessioninbox.StateSteerConsumed)
165 }
166 }
167
168 type inboxChangedCapture struct {
169 changed chan sessioninbox.InboxSnapshot
170 }
171
172 func (s *inboxChangedCapture) Emit(event.Event) {}
173
174 func (s *inboxChangedCapture) InboxChanged(snap sessioninbox.InboxSnapshot) {
175 s.changed <- snap
176 }
177
178 func TestInboxStoreChangesReachOptionalSink(t *testing.T) {
179 dir := t.TempDir()
180 sink := &inboxChangedCapture{changed: make(chan sessioninbox.InboxSnapshot, 1)}
181 c := newOwnedTestController(t, Options{
182 SessionPath: filepath.Join(dir, "s.jsonl"),
183 SessionDir: dir,
184 Sink: sink,
185 })
186 rec, err := c.EnqueueInbox(InboxRequest{Submit: "notify", Source: "desktop"})
187 if err != nil {
188 t.Fatal(err)
189 }
190 awaitState := func(want sessioninbox.InboxState) sessioninbox.InboxSnapshot {
191 t.Helper()
192 select {
193 case snap := <-sink.changed:
194 if len(snap.Items) != 1 || snap.Items[0].State != want {
195 t.Fatalf("notification = %+v, want state %q", snap.Items, want)
196 }
197 return snap
198 case <-time.After(time.Second):
199 t.Fatalf("timed out waiting for %q notification", want)
200 return sessioninbox.InboxSnapshot{}
201 }
202 }
203 queued := awaitState(sessioninbox.StateQueued)
204 st, err := c.ensureInbox()
205 if err != nil {
206 t.Fatal(err)
207 }
208 if err := st.SetState(rec.ItemID, sessioninbox.StateSteerAccepted, ""); err != nil {
209 t.Fatal(err)
210 }
211 accepted := awaitState(sessioninbox.StateSteerAccepted)
212 if err := st.MarkSteerConsumed(rec.ItemID); err != nil {
213 t.Fatal(err)
214 }
215 consumed := awaitState(sessioninbox.StateSteerConsumed)
216 if !(queued.Revision < accepted.Revision && accepted.Revision < consumed.Revision) {
217 t.Fatalf("revisions did not increase: %d, %d, %d", queued.Revision, accepted.Revision, consumed.Revision)
218 }
219 }
220
220 lines GO