| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "net/http/httptest" |
| 8 | "os" |
| 9 | "path/filepath" |
| 10 | "testing" |
| 11 | "time" |
| 12 | |
| 13 | "reasonix/internal/config" |
| 14 | "reasonix/internal/control" |
| 15 | "reasonix/internal/event" |
| 16 | "reasonix/internal/serve" |
| 17 | ) |
| 18 | |
| 19 | // A foreground submit that lands while the previous turn is still running must |
| 20 | // fall back to the durable inbox follow-up the serve asks for, instead of |
| 21 | // surfacing "session is busy" to the user. |
| 22 | func TestRemoteBusySubmitFallsBackToDurableInbox(t *testing.T) { |
| 23 | isolateDesktopUserDirs(t) |
| 24 | dir := t.TempDir() |
| 25 | path := filepath.Join(dir, "remote.jsonl") |
| 26 | if err := os.WriteFile(path, nil, 0o600); err != nil { |
| 27 | t.Fatal(err) |
| 28 | } |
| 29 | runner := remoteInboxRunner{started: make(chan string, 4), release: make(chan struct{}, 4)} |
| 30 | sink := &remoteInboxSink{states: make(chan event.RuntimeStateSnapshot, 64)} |
| 31 | ctrl := control.New(control.Options{SessionDir: dir, SessionPath: path, Runner: runner, Sink: sink}) |
| 32 | t.Cleanup(func() { |
| 33 | runner.release <- struct{}{} |
| 34 | runner.release <- struct{}{} |
| 35 | closeRemoteTestController(t, ctrl) |
| 36 | }) |
| 37 | server := httptest.NewServer(operatorServeHandler(serve.New(ctrl, nil, config.ServeConfig{}))) |
| 38 | defer server.Close() |
| 39 | a, tab := remoteRuntimeTestApp(server.Client()) |
| 40 | tab.base, tab.routing.currentPath, tab.session.path = server.URL, path, path |
| 41 | // emitInboxChanged requires a context; production App always has one. |
| 42 | a.ctx = context.Background() |
| 43 | |
| 44 | // Turn 1 occupies the foreground. |
| 45 | ctrl.Send("hold the foreground") |
| 46 | select { |
| 47 | case <-runner.started: |
| 48 | case <-time.After(5 * time.Second): |
| 49 | t.Fatal("first turn did not start") |
| 50 | } |
| 51 | |
| 52 | // A submit during the running turn must queue durably and report the receipt |
| 53 | // as a queued outcome, so the renderer can show the visible queue entry |
| 54 | // immediately instead of a failure. |
| 55 | events := make([]string, 0, 4) |
| 56 | previousEmit := runtimeEventsEmitFallback |
| 57 | runtimeEventsEmitFallback = func(_ context.Context, name string, payload ...any) { |
| 58 | if name == "InboxChanged" { |
| 59 | events = append(events, name) |
| 60 | } |
| 61 | } |
| 62 | t.Cleanup(func() { runtimeEventsEmitFallback = previousEmit }) |
| 63 | |
| 64 | err := a.SubmitRemoteTabWithSubmission(tab.id, "follow-up message", "submission-busy-1") |
| 65 | var queued interface{ RPCErrorData() map[string]any } |
| 66 | if !errors.As(err, &queued) { |
| 67 | t.Fatalf("busy submit = %v, want a queuedFollowupError", err) |
| 68 | } |
| 69 | data := queued.RPCErrorData() |
| 70 | var receipt InboxReceiptView |
| 71 | raw, _ := json.Marshal(data["queuedFollowup"]) |
| 72 | if json.Unmarshal(raw, &receipt) != nil || receipt.ItemID == "" || receipt.Disposition == "" { |
| 73 | t.Fatalf("queued outcome carries no valid receipt: %v", data) |
| 74 | } |
| 75 | if len(events) != 1 { |
| 76 | t.Fatalf("remote enqueue emitted %d InboxChanged events, want 1", len(events)) |
| 77 | } |
| 78 | snapshot, err := a.InboxSnapshot(tab.id) |
| 79 | if err != nil || len(snapshot.Items) != 1 || snapshot.Items[0].Preview != "follow-up message" || snapshot.Items[0].ID != receipt.ItemID { |
| 80 | t.Fatalf("queued receipt did not resolve to a durable item: %+v, %v", snapshot, err) |
| 81 | } |
| 82 | |
| 83 | // The idempotency key rides the submission id: a retried identical submit |
| 84 | // must not double-queue. |
| 85 | if err := a.SubmitRemoteTabWithSubmission(tab.id, "follow-up message", "submission-busy-1"); err == nil { |
| 86 | t.Fatal("retried busy submit reported a started turn") |
| 87 | } else if !errors.As(err, &queued) { |
| 88 | t.Fatalf("retried busy submit = %v, want a queued outcome", err) |
| 89 | } |
| 90 | snapshot, err = a.InboxSnapshot(tab.id) |
| 91 | if err != nil || len(snapshot.Items) != 1 { |
| 92 | t.Fatalf("retried submit double-queued: %+v, %v", snapshot, err) |
| 93 | } |
| 94 | |
| 95 | // Release turn 1; the queued follow-up dispatches as turn 2. |
| 96 | runner.release <- struct{}{} |
| 97 | select { |
| 98 | case <-runner.started: |
| 99 | case <-time.After(5 * time.Second): |
| 100 | t.Fatal("queued follow-up was not dispatched") |
| 101 | } |
| 102 | runner.release <- struct{}{} |
| 103 | } |
| 104 |