| 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 |