| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "path/filepath" |
| 7 | "strings" |
| 8 | "sync" |
| 9 | "testing" |
| 10 | "time" |
| 11 | |
| 12 | "reasonix/internal/agent" |
| 13 | "reasonix/internal/event" |
| 14 | "reasonix/internal/provider" |
| 15 | ) |
| 16 | |
| 17 | type blockingMaintenanceProvider struct { |
| 18 | started chan struct{} |
| 19 | cancelled chan struct{} |
| 20 | once sync.Once |
| 21 | } |
| 22 | |
| 23 | type failNthMaintenanceEventSink struct { |
| 24 | mu sync.Mutex |
| 25 | n int |
| 26 | failAt int |
| 27 | failure error |
| 28 | } |
| 29 | |
| 30 | func (s *failNthMaintenanceEventSink) Emit(event.Event) {} |
| 31 | func (s *failNthMaintenanceEventSink) EmitChecked(e event.Event) error { |
| 32 | if e.Kind != event.SessionOperation { |
| 33 | return nil |
| 34 | } |
| 35 | s.mu.Lock() |
| 36 | defer s.mu.Unlock() |
| 37 | s.n++ |
| 38 | if s.n == s.failAt { |
| 39 | return s.failure |
| 40 | } |
| 41 | return nil |
| 42 | } |
| 43 | |
| 44 | func (p *blockingMaintenanceProvider) Name() string { return "blocking-maintenance" } |
| 45 | func (p *blockingMaintenanceProvider) Stream(ctx context.Context, _ provider.Request) (<-chan provider.Chunk, error) { |
| 46 | p.once.Do(func() { close(p.started) }) |
| 47 | <-ctx.Done() |
| 48 | close(p.cancelled) |
| 49 | return nil, ctx.Err() |
| 50 | } |
| 51 | |
| 52 | func maintenanceFixtureSession() *agent.Session { |
| 53 | sess := agent.NewSession("sys") |
| 54 | for range 8 { |
| 55 | sess.Add(provider.Message{Role: provider.RoleUser, Content: strings.Repeat("question ", 300)}) |
| 56 | sess.Add(provider.Message{Role: provider.RoleAssistant, Content: strings.Repeat("answer ", 300)}) |
| 57 | } |
| 58 | return sess |
| 59 | } |
| 60 | |
| 61 | func TestManualCompactIsCancellableForegroundMaintenance(t *testing.T) { |
| 62 | prov := &blockingMaintenanceProvider{started: make(chan struct{}), cancelled: make(chan struct{})} |
| 63 | exec := agent.New(prov, nil, maintenanceFixtureSession(), agent.Options{ContextWindow: 32_000}, event.Discard) |
| 64 | // Cancellation and foreground ownership do not require filesystem writes. |
| 65 | // TestMaintenanceLifecycleDurableOutsideTurn covers the durable boundary. |
| 66 | c := newOwnedTestController(t, Options{Executor: exec, SystemPrompt: "sys", Sink: event.Discard}) |
| 67 | |
| 68 | done := make(chan error, 1) |
| 69 | go func() { done <- c.Compact(context.Background(), "") }() |
| 70 | select { |
| 71 | case <-prov.started: |
| 72 | case <-time.After(3 * time.Second): |
| 73 | t.Fatal("summary provider was not called") |
| 74 | } |
| 75 | |
| 76 | snapshot := c.RuntimeStateSnapshot() |
| 77 | if !snapshot.Running || !snapshot.Cancellable || snapshot.Maintenance == nil || snapshot.Maintenance.Kind != "compact" { |
| 78 | t.Fatalf("runtime during compact = %+v", snapshot) |
| 79 | } |
| 80 | receipt := c.CancelSessionFrom("test") |
| 81 | if receipt.AlreadyIdle { |
| 82 | t.Fatalf("CancelSession reported idle during compaction: %+v", receipt) |
| 83 | } |
| 84 | select { |
| 85 | case <-prov.cancelled: |
| 86 | case <-time.After(time.Second): |
| 87 | t.Fatal("Stop did not cancel the summary request") |
| 88 | } |
| 89 | select { |
| 90 | case err := <-done: |
| 91 | if !errors.Is(err, context.Canceled) { |
| 92 | t.Fatalf("Compact error = %v, want context.Canceled", err) |
| 93 | } |
| 94 | case <-time.After(3 * time.Second): |
| 95 | t.Fatal("compaction did not finish after cancellation") |
| 96 | } |
| 97 | |
| 98 | if snapshot = c.RuntimeStateSnapshot(); snapshot.Running || snapshot.Maintenance != nil { |
| 99 | t.Fatalf("runtime after cancelled compact = %+v", snapshot) |
| 100 | } |
| 101 | } |
| 102 | |
| 103 | func TestCompactSubmitReceiptIncludesRegisteredOperation(t *testing.T) { |
| 104 | prov := &blockingMaintenanceProvider{started: make(chan struct{}), cancelled: make(chan struct{})} |
| 105 | exec := agent.New(prov, nil, maintenanceFixtureSession(), agent.Options{ContextWindow: 32_000}, event.Discard) |
| 106 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard}) |
| 107 | result := c.SubmitDisplayWithResult("/compact", "/compact") |
| 108 | if result.Disposition != SubmitManagementHandled || result.OperationID == "" { |
| 109 | t.Fatalf("compact receipt = %+v", result) |
| 110 | } |
| 111 | if got := c.ActiveMaintenanceOperationID(); got != result.OperationID { |
| 112 | t.Fatalf("active operation = %q, receipt = %q", got, result.OperationID) |
| 113 | } |
| 114 | select { |
| 115 | case <-prov.started: |
| 116 | case <-time.After(3 * time.Second): |
| 117 | t.Fatal("summary provider was not called") |
| 118 | } |
| 119 | c.CancelSessionFrom("test") |
| 120 | select { |
| 121 | case <-prov.cancelled: |
| 122 | case <-time.After(time.Second): |
| 123 | t.Fatal("Stop did not cancel the submitted compaction") |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | func TestMaintenanceBlocksRotationAndInboxDispatch(t *testing.T) { |
| 128 | c := newOwnedTestController(t, Options{Executor: agent.New(nil, nil, maintenanceFixtureSession(), agent.Options{}, event.Discard), Sink: event.Discard}) |
| 129 | op, _, err := c.beginMaintenance(context.Background(), "compact") |
| 130 | if err != nil { |
| 131 | t.Fatalf("beginMaintenance: %v", err) |
| 132 | } |
| 133 | if err := c.beginRotation(); !errors.Is(err, ErrMaintenanceBusy) { |
| 134 | t.Fatalf("beginRotation during maintenance = %v, want ErrMaintenanceBusy", err) |
| 135 | } |
| 136 | if _, _, err := c.beginMaintenance(context.Background(), "compact"); !errors.Is(err, ErrMaintenanceBusy) { |
| 137 | t.Fatalf("second maintenance = %v, want ErrMaintenanceBusy", err) |
| 138 | } |
| 139 | c.mu.Lock() |
| 140 | busy := c.maintenance == op |
| 141 | c.mu.Unlock() |
| 142 | if !busy || !c.Running() { |
| 143 | t.Fatal("maintenance was not retained as foreground work") |
| 144 | } |
| 145 | _, _, _ = c.signalMaintenanceCancel() |
| 146 | c.mu.Lock() |
| 147 | c.maintenance = nil |
| 148 | close(op.done) |
| 149 | c.mu.Unlock() |
| 150 | } |
| 151 | |
| 152 | func TestMaintenanceStopNeverWithdrawsQueuedItemsDuringFinalizing(t *testing.T) { |
| 153 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 154 | op, _, err := c.beginMaintenance(context.Background(), "compact") |
| 155 | if err != nil { |
| 156 | t.Fatalf("beginMaintenance: %v", err) |
| 157 | } |
| 158 | c.mu.Lock() |
| 159 | op.activity = "finalizing" |
| 160 | c.mu.Unlock() |
| 161 | result, err := c.CancelWithInboxItemsResult([]string{"queued-1"}, "test") |
| 162 | if err != nil { |
| 163 | t.Fatalf("CancelWithInboxItemsResult: %v", err) |
| 164 | } |
| 165 | if len(result.DiscardedItemIDs) != 0 { |
| 166 | t.Fatalf("discarded queued items during finalizing: %v", result.DiscardedItemIDs) |
| 167 | } |
| 168 | c.mu.Lock() |
| 169 | c.maintenance = nil |
| 170 | close(op.done) |
| 171 | c.mu.Unlock() |
| 172 | } |
| 173 | |
| 174 | func TestCloseWaitsForMaintenanceThatIgnoredCancellation(t *testing.T) { |
| 175 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 176 | c.testCancelGrace = 5 * time.Millisecond |
| 177 | op, workCtx, err := c.beginMaintenance(context.Background(), "compact") |
| 178 | if err != nil { |
| 179 | t.Fatalf("beginMaintenance: %v", err) |
| 180 | } |
| 181 | started := make(chan struct{}) |
| 182 | release := make(chan struct{}) |
| 183 | done := make(chan error, 1) |
| 184 | go func() { |
| 185 | done <- c.executeMaintenance(op, workCtx, func(context.Context) error { |
| 186 | close(started) |
| 187 | <-release // deliberately ignore cancellation until the owner releases us |
| 188 | return workCtx.Err() |
| 189 | }) |
| 190 | }() |
| 191 | select { |
| 192 | case <-started: |
| 193 | case <-time.After(3 * time.Second): |
| 194 | t.Fatal("maintenance worker did not start") |
| 195 | } |
| 196 | c.CancelSessionFrom("test") |
| 197 | deadline := time.Now().Add(time.Second) |
| 198 | recoveryObserved := false |
| 199 | for time.Now().Before(deadline) { |
| 200 | snapshot := c.RuntimeStateSnapshot() |
| 201 | if snapshot.Maintenance != nil && snapshot.Maintenance.Activity == "recovery_required" { |
| 202 | recoveryObserved = true |
| 203 | break |
| 204 | } |
| 205 | time.Sleep(time.Millisecond) |
| 206 | } |
| 207 | if !recoveryObserved { |
| 208 | t.Fatal("maintenance did not enter recovery after ignoring cancellation") |
| 209 | } |
| 210 | c.Close() |
| 211 | select { |
| 212 | case <-c.Closed(): |
| 213 | t.Fatal("controller resources closed while the maintenance worker was still running") |
| 214 | default: |
| 215 | } |
| 216 | close(release) |
| 217 | select { |
| 218 | case err := <-done: |
| 219 | if !errors.Is(err, context.Canceled) { |
| 220 | t.Fatalf("Compact error = %v, want context.Canceled", err) |
| 221 | } |
| 222 | case <-time.After(3 * time.Second): |
| 223 | t.Fatal("late maintenance worker did not settle") |
| 224 | } |
| 225 | select { |
| 226 | case <-c.Closed(): |
| 227 | case <-time.After(3 * time.Second): |
| 228 | t.Fatal("controller did not release resources after maintenance settled") |
| 229 | } |
| 230 | } |
| 231 | |
| 232 | func TestTerminalOperationPersistenceFailureRetainsRecovery(t *testing.T) { |
| 233 | want := errors.New("operation journal unavailable") |
| 234 | sink := &failNthMaintenanceEventSink{failAt: 2, failure: want} |
| 235 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 236 | op, workCtx, err := c.beginMaintenance(context.Background(), "compact") |
| 237 | if err != nil { |
| 238 | t.Fatalf("beginMaintenance: %v", err) |
| 239 | } |
| 240 | c.sink = sink // finalizing succeeds; the terminal record fails |
| 241 | err = c.executeMaintenance(op, workCtx, func(context.Context) error { return nil }) |
| 242 | if !errors.Is(err, want) { |
| 243 | t.Fatalf("executeMaintenance error = %v, want %v", err, want) |
| 244 | } |
| 245 | snapshot := c.RuntimeStateSnapshot() |
| 246 | if snapshot.Maintenance == nil || snapshot.Maintenance.Activity != "recovery_required" || !snapshot.Running || snapshot.Cancellable { |
| 247 | t.Fatalf("runtime after terminal persistence failure = %+v", snapshot) |
| 248 | } |
| 249 | } |
| 250 | |
| 251 | func TestMaintenanceDoesNotStartWhenOperationRecordCannotPersist(t *testing.T) { |
| 252 | want := errors.New("operation journal unavailable") |
| 253 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 254 | c.sink = &failNthMaintenanceEventSink{failAt: 1, failure: want} |
| 255 | if _, _, err := c.beginMaintenance(context.Background(), "compact"); !errors.Is(err, want) { |
| 256 | t.Fatalf("beginMaintenance error = %v, want %v", err, want) |
| 257 | } |
| 258 | if state := c.RuntimeStateSnapshot(); state.Maintenance == nil || state.Maintenance.Status != "recovery_required" || state.Cancellable { |
| 259 | t.Fatalf("failed operation registration must retain a recovery barrier: %+v", state) |
| 260 | } |
| 261 | } |
| 262 | |
| 263 | func TestEmptyManualCompactIsNoop(t *testing.T) { |
| 264 | c := newOwnedTestController(t, Options{ |
| 265 | Executor: agent.New(nil, nil, agent.NewSession("sys"), agent.Options{ContextWindow: 32_000}, event.Discard), |
| 266 | Sink: event.Discard, |
| 267 | }) |
| 268 | var terminal string |
| 269 | c.sink = event.FuncSink(func(e event.Event) { |
| 270 | if e.SessionOperation != nil { |
| 271 | terminal = e.SessionOperation.Status |
| 272 | } |
| 273 | }) |
| 274 | if err := c.Compact(context.Background(), ""); err != nil { |
| 275 | t.Fatalf("Compact(empty) = %v, want nil", err) |
| 276 | } |
| 277 | if terminal != "noop" { |
| 278 | t.Fatalf("terminal status = %q, want noop", terminal) |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | func TestMaintenanceCompletionDrainsParkedGuidance(t *testing.T) { |
| 283 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 284 | op, ctx, err := c.beginMaintenance(context.Background(), "compact") |
| 285 | if err != nil { |
| 286 | t.Fatal(err) |
| 287 | } |
| 288 | if got := c.submitSteerFallback("retain this guidance"); got != turnParked { |
| 289 | t.Fatalf("admission = %v, want turnParked", got) |
| 290 | } |
| 291 | if err := c.executeMaintenance(op, ctx, func(context.Context) error { return nil }); err != nil { |
| 292 | t.Fatal(err) |
| 293 | } |
| 294 | c.mu.Lock() |
| 295 | pending := len(c.turns.pending) |
| 296 | active := c.bodyActiveLocked() |
| 297 | c.mu.Unlock() |
| 298 | if pending != 0 && !active { |
| 299 | t.Fatalf("maintenance ended idle with %d stranded guidance item(s)", pending) |
| 300 | } |
| 301 | } |
| 302 | |
| 303 | func TestSynchronousTurnCannotEnterMaintenance(t *testing.T) { |
| 304 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 305 | op, ctx, err := c.beginMaintenance(context.Background(), "compact") |
| 306 | if err != nil { |
| 307 | t.Fatal(err) |
| 308 | } |
| 309 | ran := false |
| 310 | err = c.runSynchronousTurn(context.Background(), nil, func(context.Context) error { |
| 311 | ran = true |
| 312 | return nil |
| 313 | }) |
| 314 | if ran || !errors.Is(err, ErrMaintenanceBusy) || !errors.Is(err, ErrTurnRunning) { |
| 315 | t.Fatalf("synchronous turn during maintenance: ran=%v err=%v", ran, err) |
| 316 | } |
| 317 | if err := c.executeMaintenance(op, ctx, func(context.Context) error { return nil }); err != nil { |
| 318 | t.Fatal(err) |
| 319 | } |
| 320 | } |
| 321 | |
| 322 | func TestRunInboxTurnRemainsQueuedDuringMaintenance(t *testing.T) { |
| 323 | dir := t.TempDir() |
| 324 | c := newOwnedTestController(t, Options{ |
| 325 | SessionDir: dir, SessionPath: filepath.Join(dir, "session.jsonl"), Sink: event.Discard, |
| 326 | }) |
| 327 | op, ctx, err := c.beginMaintenance(context.Background(), "compact") |
| 328 | if err != nil { |
| 329 | t.Fatal(err) |
| 330 | } |
| 331 | receipt, err := c.EnqueueInbox(InboxRequest{Submit: "queued during compact", Idempotency: "maintenance-inbox"}) |
| 332 | if err != nil { |
| 333 | t.Fatal(err) |
| 334 | } |
| 335 | if err := c.RunInboxTurn(context.Background(), receipt.ItemID); !errors.Is(err, ErrMaintenanceBusy) || !errors.Is(err, ErrTurnRunning) { |
| 336 | t.Fatalf("RunInboxTurn error = %v, want retryable maintenance busy", err) |
| 337 | } |
| 338 | meta, _, err := c.ReadInboxItem(receipt.ItemID) |
| 339 | if err != nil { |
| 340 | t.Fatal(err) |
| 341 | } |
| 342 | if meta.State != "queued" { |
| 343 | t.Fatalf("inbox state = %q, want queued", meta.State) |
| 344 | } |
| 345 | if err := c.SetInboxPaused(true); err != nil { |
| 346 | t.Fatal(err) |
| 347 | } |
| 348 | if err := c.executeMaintenance(op, ctx, func(context.Context) error { return nil }); err != nil { |
| 349 | t.Fatal(err) |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | func TestTerminalMaintenanceOperationRejectsLateCancellingEvent(t *testing.T) { |
| 354 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 355 | op, ctx, err := c.beginMaintenance(context.Background(), "compact") |
| 356 | if err != nil { |
| 357 | t.Fatal(err) |
| 358 | } |
| 359 | entered := make(chan struct{}) |
| 360 | release := make(chan struct{}) |
| 361 | stopped := make(chan struct{}) |
| 362 | var mu sync.Mutex |
| 363 | var statuses []string |
| 364 | c.sink = event.FuncSink(func(e event.Event) { |
| 365 | if e.SessionOperation == nil { |
| 366 | return |
| 367 | } |
| 368 | status := e.SessionOperation.Status |
| 369 | if status == "cancelling" { |
| 370 | close(entered) |
| 371 | <-release |
| 372 | } |
| 373 | mu.Lock() |
| 374 | statuses = append(statuses, status) |
| 375 | mu.Unlock() |
| 376 | }) |
| 377 | go func() { |
| 378 | c.signalMaintenanceCancel() |
| 379 | close(stopped) |
| 380 | }() |
| 381 | <-entered |
| 382 | finished := make(chan error, 1) |
| 383 | go func() { |
| 384 | finished <- c.executeMaintenance(op, ctx, func(context.Context) error { return ctx.Err() }) |
| 385 | }() |
| 386 | close(release) |
| 387 | <-stopped |
| 388 | if err := <-finished; !errors.Is(err, context.Canceled) { |
| 389 | t.Fatalf("executeMaintenance = %v, want context.Canceled", err) |
| 390 | } |
| 391 | mu.Lock() |
| 392 | defer mu.Unlock() |
| 393 | if len(statuses) == 0 || statuses[len(statuses)-1] != "cancelled" { |
| 394 | t.Fatalf("operation statuses = %v, want terminal cancelled last", statuses) |
| 395 | } |
| 396 | } |
| 397 | |
| 398 | func TestMaintenanceOperationRevisionsAreMonotonicAndSnapshotIsLossless(t *testing.T) { |
| 399 | var mu sync.Mutex |
| 400 | var records []event.SessionOperationInfo |
| 401 | sink := event.FuncSink(func(e event.Event) { |
| 402 | if e.SessionOperation == nil { |
| 403 | return |
| 404 | } |
| 405 | mu.Lock() |
| 406 | records = append(records, *e.SessionOperation) |
| 407 | mu.Unlock() |
| 408 | }) |
| 409 | c := newOwnedTestController(t, Options{Sink: sink}) |
| 410 | op, ctx, err := c.beginMaintenance(context.Background(), "compact") |
| 411 | if err != nil { |
| 412 | t.Fatal(err) |
| 413 | } |
| 414 | running := c.RuntimeStateSnapshot().Maintenance |
| 415 | if running == nil || running.OperationID != op.id || running.OperationRevision == 0 || running.Status != "running" { |
| 416 | t.Fatalf("running maintenance snapshot = %+v", running) |
| 417 | } |
| 418 | if err := c.executeMaintenance(op, ctx, func(context.Context) error { return errors.New("summary unavailable") }); err == nil { |
| 419 | t.Fatal("executeMaintenance unexpectedly succeeded") |
| 420 | } |
| 421 | mu.Lock() |
| 422 | defer mu.Unlock() |
| 423 | if len(records) != 3 { |
| 424 | t.Fatalf("operation records = %+v, want running/finalizing/failed", records) |
| 425 | } |
| 426 | for i, record := range records { |
| 427 | if record.OperationRevision != uint64(i+1) { |
| 428 | t.Fatalf("record %d revision = %d, want %d", i, record.OperationRevision, i+1) |
| 429 | } |
| 430 | if record.RuntimeEpoch != records[0].RuntimeEpoch { |
| 431 | t.Fatalf("record %d runtime epoch = %q, want %q", i, record.RuntimeEpoch, records[0].RuntimeEpoch) |
| 432 | } |
| 433 | } |
| 434 | if records[2].Status != "failed" || records[2].ErrorCode != "summary_failed" || records[2].Detail == "" { |
| 435 | t.Fatalf("terminal record = %+v", records[2]) |
| 436 | } |
| 437 | } |
| 438 |