| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "path/filepath" |
| 8 | "strings" |
| 9 | "sync" |
| 10 | "sync/atomic" |
| 11 | "testing" |
| 12 | "time" |
| 13 | |
| 14 | "reasonix/internal/agent" |
| 15 | "reasonix/internal/control" |
| 16 | "reasonix/internal/event" |
| 17 | "reasonix/internal/provider" |
| 18 | ) |
| 19 | |
| 20 | type compactReceiptProvider struct { |
| 21 | started chan struct{} |
| 22 | release chan struct{} |
| 23 | calls atomic.Int32 |
| 24 | } |
| 25 | |
| 26 | func (p *compactReceiptProvider) Name() string { return "compact-receipt" } |
| 27 | func (p *compactReceiptProvider) Stream(ctx context.Context, _ provider.Request) (<-chan provider.Chunk, error) { |
| 28 | p.calls.Add(1) |
| 29 | p.started <- struct{}{} |
| 30 | select { |
| 31 | case <-p.release: |
| 32 | case <-ctx.Done(): |
| 33 | return nil, ctx.Err() |
| 34 | } |
| 35 | chunks := make(chan provider.Chunk, 2) |
| 36 | chunks <- provider.Chunk{Type: provider.ChunkText, Text: "Preserve the user's task and the completed work."} |
| 37 | chunks <- provider.Chunk{Type: provider.ChunkDone} |
| 38 | close(chunks) |
| 39 | return chunks, nil |
| 40 | } |
| 41 | |
| 42 | func TestCompactDuplicateReceiptKeepsExistingOperation(t *testing.T) { |
| 43 | prov := &compactReceiptProvider{started: make(chan struct{}, 2), release: make(chan struct{})} |
| 44 | defer close(prov.release) |
| 45 | sess := agent.NewSession("system") |
| 46 | for range 8 { |
| 47 | sess.Add(provider.Message{Role: provider.RoleUser, Content: strings.Repeat("question ", 300)}) |
| 48 | sess.Add(provider.Message{Role: provider.RoleAssistant, Content: strings.Repeat("answer ", 300)}) |
| 49 | } |
| 50 | exec := agent.New(prov, nil, sess, agent.Options{ContextWindow: 32_000}, event.Discard) |
| 51 | dir := t.TempDir() |
| 52 | ctrl := control.New(control.Options{Executor: exec, Sink: event.Discard, SessionDir: dir, SessionPath: filepath.Join(dir, "session.jsonl")}) |
| 53 | cleanupExactTurnController(t, ctrl) |
| 54 | tab := &WorkspaceTab{ID: "source", Scope: "global", Ready: true, Ctrl: ctrl} |
| 55 | app := &App{tabs: map[string]*WorkspaceTab{tab.ID: tab}, activeTabID: tab.ID} |
| 56 | first, err := app.StartTurnForTab(tab.ID, "/compact", "first") |
| 57 | if err != nil || first.OperationID == "" || first.ManagementErrorCode != "" { |
| 58 | t.Fatalf("first compact = %+v, %v", first, err) |
| 59 | } |
| 60 | select { |
| 61 | case <-prov.started: |
| 62 | case <-time.After(5 * time.Second): |
| 63 | t.Fatal("compaction did not enter provider") |
| 64 | } |
| 65 | // A late response still belongs to the source tab after the user navigates away. |
| 66 | app.activeTabID = "other" |
| 67 | for _, input := range []string{"/compact", "/compact preserve decisions"} { |
| 68 | duplicate, err := app.StartTurnForTab(tab.ID, input, "duplicate:"+input) |
| 69 | if err != nil || duplicate.Disposition != control.SubmitManagementHandled || duplicate.ManagementErrorCode != "maintenance_busy" || |
| 70 | duplicate.OperationID != first.OperationID || duplicate.TurnID != "" { |
| 71 | t.Fatalf("duplicate compact = %+v, %v", duplicate, err) |
| 72 | } |
| 73 | } |
| 74 | if _, err := app.StartTurnForTab(tab.ID, "ordinary question", "ordinary"); !errors.Is(err, control.ErrTurnRunning) { |
| 75 | t.Fatalf("ordinary turn bypassed maintenance guard: %v", err) |
| 76 | } |
| 77 | if prov.calls.Load() != 1 || ctrl.ActiveMaintenanceOperationID() != first.OperationID { |
| 78 | t.Fatal("duplicate request replaced or restarted the existing operation") |
| 79 | } |
| 80 | } |
| 81 | |
| 82 | type maintenanceReceiptController struct { |
| 83 | *tabScopedActionController |
| 84 | maintenance *event.MaintenanceState |
| 85 | admissionMu *sync.Mutex |
| 86 | readOutsideAdmission bool |
| 87 | } |
| 88 | |
| 89 | func (c *maintenanceReceiptController) ClassifySubmitRoute(string) control.SubmitDisposition { |
| 90 | return control.SubmitManagementHandled |
| 91 | } |
| 92 | func (c *maintenanceReceiptController) RuntimeStateSnapshot() event.RuntimeStateSnapshot { |
| 93 | if c.admissionMu != nil && c.admissionMu.TryLock() { |
| 94 | c.readOutsideAdmission = true |
| 95 | c.admissionMu.Unlock() |
| 96 | } |
| 97 | return event.RuntimeStateSnapshot{Maintenance: c.maintenance} |
| 98 | } |
| 99 | |
| 100 | func TestCompactRecoveryReceiptAndLegacyShape(t *testing.T) { |
| 101 | ctrl := &maintenanceReceiptController{tabScopedActionController: newTabScopedActionController(), |
| 102 | maintenance: &event.MaintenanceState{OperationID: "recover-operation", Kind: "compact", Activity: "recovery_required"}} |
| 103 | tab := &WorkspaceTab{ID: "source", Scope: "global", Ready: true, Ctrl: ctrl} |
| 104 | ctrl.admissionMu = &tab.turnStartMu |
| 105 | app := &App{tabs: map[string]*WorkspaceTab{tab.ID: tab}} |
| 106 | receipt, err := app.StartTurnForTab(tab.ID, "/compact", "retry") |
| 107 | if err != nil || receipt.ManagementErrorCode != "maintenance_recovery_required" || receipt.OperationID != "recover-operation" { |
| 108 | t.Fatalf("recovery refusal was lost: %+v, %v", receipt, err) |
| 109 | } |
| 110 | if ctrl.readOutsideAdmission { |
| 111 | t.Fatal("management conflict read its owner outside the admission lock") |
| 112 | } |
| 113 | // Existing readers retain the management disposition and operation identity. |
| 114 | encoded, err := json.Marshal(receipt) |
| 115 | if err != nil { |
| 116 | t.Fatal(err) |
| 117 | } |
| 118 | var legacy struct { |
| 119 | Disposition control.SubmitDisposition `json:"disposition"` |
| 120 | OperationID string `json:"operationId"` |
| 121 | } |
| 122 | if err := json.Unmarshal(encoded, &legacy); err != nil || legacy.Disposition != control.SubmitManagementHandled || legacy.OperationID != receipt.OperationID { |
| 123 | t.Fatalf("legacy receipt decode = %+v, %v", legacy, err) |
| 124 | } |
| 125 | var old TurnStartView |
| 126 | if err := json.Unmarshal([]byte(`{"disposition":"management_handled"}`), &old); err != nil || old.ManagementErrorCode != "" { |
| 127 | t.Fatalf("old successful receipt changed meaning: %+v, %v", old, err) |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | type classifiedCompactController struct { |
| 132 | *control.Controller |
| 133 | classified chan struct{} |
| 134 | } |
| 135 | |
| 136 | func (c *classifiedCompactController) ClassifySubmitRoute(input string) control.SubmitDisposition { |
| 137 | result := c.Controller.ClassifySubmitRoute(input) |
| 138 | close(c.classified) |
| 139 | return result |
| 140 | } |
| 141 | |
| 142 | func TestCompactDuplicateWaitsForAdmissionOwner(t *testing.T) { |
| 143 | prov := &compactReceiptProvider{started: make(chan struct{}, 2), release: make(chan struct{})} |
| 144 | defer close(prov.release) |
| 145 | sess := agent.NewSession("system") |
| 146 | for range 8 { |
| 147 | sess.Add(provider.Message{Role: provider.RoleUser, Content: strings.Repeat("question ", 300)}) |
| 148 | sess.Add(provider.Message{Role: provider.RoleAssistant, Content: strings.Repeat("answer ", 300)}) |
| 149 | } |
| 150 | exec := agent.New(prov, nil, sess, agent.Options{ContextWindow: 32_000}, event.Discard) |
| 151 | dir := t.TempDir() |
| 152 | ctrl := control.New(control.Options{Executor: exec, Sink: event.Discard, SessionDir: dir, SessionPath: filepath.Join(dir, "session.jsonl")}) |
| 153 | cleanupExactTurnController(t, ctrl) |
| 154 | wrapper := &classifiedCompactController{Controller: ctrl, classified: make(chan struct{})} |
| 155 | tab := &WorkspaceTab{ID: "source", Scope: "global", Ready: true, Ctrl: wrapper} |
| 156 | app := &App{tabs: map[string]*WorkspaceTab{tab.ID: tab}} |
| 157 | type response struct { |
| 158 | receipt TurnStartView |
| 159 | err error |
| 160 | } |
| 161 | done := make(chan response, 1) |
| 162 | // The first submit owns the tab lock while the duplicate classifies its route. |
| 163 | tab.turnStartMu.Lock() |
| 164 | go func() { |
| 165 | receipt, err := app.StartTurnForTab(tab.ID, "/compact preserve decisions", "duplicate") |
| 166 | done <- response{receipt, err} |
| 167 | }() |
| 168 | select { |
| 169 | case <-wrapper.classified: |
| 170 | case <-time.After(5 * time.Second): |
| 171 | tab.turnStartMu.Unlock() |
| 172 | t.Fatal("duplicate did not classify") |
| 173 | } |
| 174 | first := ctrl.SubmitDisplayWithResult("/compact", "/compact") |
| 175 | tab.turnStartMu.Unlock() |
| 176 | select { |
| 177 | case got := <-done: |
| 178 | if got.err != nil || got.receipt.ManagementErrorCode != "maintenance_busy" || got.receipt.OperationID != first.OperationID || first.OperationID == "" { |
| 179 | t.Fatalf("duplicate lost the admission owner's identity: %+v, %v; first=%+v", got.receipt, got.err, first) |
| 180 | } |
| 181 | case <-time.After(5 * time.Second): |
| 182 | t.Fatal("duplicate remained blocked") |
| 183 | } |
| 184 | select { |
| 185 | case <-prov.started: |
| 186 | case <-time.After(5 * time.Second): |
| 187 | t.Fatal("first compaction did not enter provider") |
| 188 | } |
| 189 | if prov.calls.Load() != 1 { |
| 190 | t.Fatalf("provider called %d times", prov.calls.Load()) |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | type compactDuringStatusController struct{ *maintenanceReceiptController } |
| 195 | |
| 196 | func (c *compactDuringStatusController) RuntimeStatus() control.RuntimeStatus { |
| 197 | // The synchronous Compact button also owns controller maintenance. Model its |
| 198 | // registration while the slash request is entering desktop admission. |
| 199 | c.maintenance = &event.MaintenanceState{OperationID: "button-compact", Kind: "compact", Activity: "running"} |
| 200 | return control.RuntimeStatus{Running: true} |
| 201 | } |
| 202 | |
| 203 | func TestCompactButtonConflictUsesAdmissionReceipt(t *testing.T) { |
| 204 | ctrl := &compactDuringStatusController{&maintenanceReceiptController{tabScopedActionController: newTabScopedActionController()}} |
| 205 | tab := &WorkspaceTab{ID: "source", Scope: "global", Ready: true, Ctrl: ctrl} |
| 206 | app := &App{tabs: map[string]*WorkspaceTab{tab.ID: tab}} |
| 207 | receipt, err := app.StartTurnForTab(tab.ID, "/compact", "duplicate") |
| 208 | if err != nil || receipt.ManagementErrorCode != "maintenance_busy" || receipt.OperationID != "button-compact" { |
| 209 | t.Fatalf("button conflict lost its identity: %+v, %v", receipt, err) |
| 210 | } |
| 211 | } |
| 212 |