| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "os" |
| 8 | "path/filepath" |
| 9 | "reflect" |
| 10 | "sync/atomic" |
| 11 | "testing" |
| 12 | "time" |
| 13 | |
| 14 | "reasonix/internal/agent" |
| 15 | "reasonix/internal/event" |
| 16 | "reasonix/internal/extension" |
| 17 | "reasonix/internal/provider" |
| 18 | "reasonix/internal/session" |
| 19 | ) |
| 20 | |
| 21 | func TestMaintenanceRegressionMaintenanceRuntimeOwnership(t *testing.T) { |
| 22 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4"))) |
| 23 | if err != nil { |
| 24 | t.Fatal(err) |
| 25 | } |
| 26 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "maintenance-reaudit"}) |
| 27 | if err != nil { |
| 28 | t.Fatal(err) |
| 29 | } |
| 30 | newController := func() *Controller { |
| 31 | exec := agent.New(nil, nil, agent.NewSession("system"), agent.Options{}, event.Discard) |
| 32 | return newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 33 | } |
| 34 | first := newController() |
| 35 | second := newController() |
| 36 | before := runtime.Session().StateSnapshot().EventSequence |
| 37 | op, ctx, err := first.beginMaintenance(context.Background(), "compact") |
| 38 | if err != nil { |
| 39 | t.Fatal(err) |
| 40 | } |
| 41 | defer first.finishMaintenanceOperation(op) |
| 42 | if after := runtime.Session().StateSnapshot().EventSequence; after == before { |
| 43 | t.Errorf("maintenance start was never appended to session log: sequence stayed %d", after) |
| 44 | } |
| 45 | if got := runtime.StateSnapshot().Phase; got != session.RuntimeRunning { |
| 46 | t.Errorf("maintenance running, shared runtime phase = %s, want running", got) |
| 47 | } |
| 48 | if _, _, err := second.beginMaintenance(t.Context(), "compact"); !errors.Is(err, ErrMaintenanceBusy) { |
| 49 | t.Errorf("unpublished controller maintenance admission = %v", err) |
| 50 | } |
| 51 | if runtime.Cancel(); ctx.Err() == nil { |
| 52 | t.Error("shared runtime Cancel did not cancel maintenance context") |
| 53 | } |
| 54 | if err := ActivateControllerReplacement(first, second); err == nil { |
| 55 | t.Error("replacement stole execution ownership while maintenance was active") |
| 56 | } |
| 57 | } |
| 58 | |
| 59 | func TestMaintenanceLifecycleDurableOutsideTurn(t *testing.T) { |
| 60 | for _, status := range []string{"noop", "failed", "cancelled"} { |
| 61 | t.Run(status, func(t *testing.T) { |
| 62 | persistence := session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions")) |
| 63 | service, err := session.NewService("desktop", persistence) |
| 64 | if err != nil { |
| 65 | t.Fatal(err) |
| 66 | } |
| 67 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "durable-maintenance"}) |
| 68 | if err != nil { |
| 69 | t.Fatal(err) |
| 70 | } |
| 71 | exec := agent.New(nil, nil, agent.NewSession("system"), agent.Options{}, event.Discard) |
| 72 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 73 | before := runtime.Session().ExecutionSnapshot().Projection |
| 74 | op, ctx, err := c.beginMaintenance(t.Context(), "compact") |
| 75 | if err != nil { |
| 76 | t.Fatal(err) |
| 77 | } |
| 78 | if state := runtime.Session().StateSnapshot(); state.DurableSequence != state.EventSequence { |
| 79 | t.Fatal("start returned before durable checkpoint") |
| 80 | } |
| 81 | err = c.executeMaintenance(op, ctx, func(context.Context) error { |
| 82 | switch status { |
| 83 | case "failed": |
| 84 | return errors.New("summary test failure") |
| 85 | case "cancelled": |
| 86 | return context.Canceled |
| 87 | } |
| 88 | return nil |
| 89 | }) |
| 90 | if status == "noop" && err != nil { |
| 91 | t.Fatal(err) |
| 92 | } |
| 93 | after := runtime.Session().StateSnapshot() |
| 94 | if after.DurableSequence != after.EventSequence { |
| 95 | t.Fatal("maintenance released before durable terminal") |
| 96 | } |
| 97 | projection := runtime.Session().ExecutionSnapshot().Projection |
| 98 | if !reflect.DeepEqual(before.Messages, projection.Messages) || !reflect.DeepEqual(before.ModelMessages, projection.ModelMessages) || len(projection.Turns) != len(before.Turns) { |
| 99 | t.Fatal("display operation changed canonical/model history or turn count") |
| 100 | } |
| 101 | if got := runtime.StateSnapshot().Phase; got != session.RuntimeIdle { |
| 102 | t.Fatalf("settled runtime = %s", got) |
| 103 | } |
| 104 | c.Close() |
| 105 | <-c.Closed() |
| 106 | if err := service.CloseAll(t.Context()); err != nil { |
| 107 | t.Fatal(err) |
| 108 | } |
| 109 | cold, err := persistence.Open(runtime.Ref().SessionID, session.ReadOnly) |
| 110 | if err != nil { |
| 111 | t.Fatal(err) |
| 112 | } |
| 113 | defer cold.Close(context.Background()) |
| 114 | page, err := cold.Read(t.Context(), 0, 100) |
| 115 | if err != nil { |
| 116 | t.Fatal(err) |
| 117 | } |
| 118 | var records []event.SessionOperationInfo |
| 119 | for _, commit := range page.Commits { |
| 120 | for _, e := range commit.Events { |
| 121 | if e.Kind != "diagnostic" { |
| 122 | continue |
| 123 | } |
| 124 | var payload struct { |
| 125 | Type string `json:"type"` |
| 126 | Record provider.Message `json:"displayRecord"` |
| 127 | } |
| 128 | if err := json.Unmarshal(e.Payload, &payload); err != nil { |
| 129 | t.Fatal(err) |
| 130 | } |
| 131 | if payload.Type != "session-maintenance-v1" { |
| 132 | continue |
| 133 | } |
| 134 | if !e.Optional { |
| 135 | t.Fatal("maintenance record must be ignorable by older readers") |
| 136 | } |
| 137 | var info event.SessionOperationInfo |
| 138 | if err := json.Unmarshal([]byte(payload.Record.Content), &info); err != nil { |
| 139 | t.Fatal(err) |
| 140 | } |
| 141 | records = append(records, info) |
| 142 | } |
| 143 | } |
| 144 | if len(records) != 3 || records[0].Status != "running" || records[2].Status != status || records[2].OperationRevision != 3 { |
| 145 | t.Fatalf("reopened records = %+v", records) |
| 146 | } |
| 147 | if status == "failed" && records[2].Detail != "summary test failure" { |
| 148 | t.Fatal("failure details lost on reopen") |
| 149 | } |
| 150 | // The actual canonical history window collapses revisions into one row. |
| 151 | reader, err := session.NewService("desktop", persistence) |
| 152 | if err != nil { |
| 153 | t.Fatal(err) |
| 154 | } |
| 155 | defer reader.CloseAll(context.Background()) |
| 156 | deadline := time.Now().Add(5 * time.Second) |
| 157 | for { |
| 158 | window, err := reader.Query().ReadHistoryWindow(t.Context(), runtime.Ref(), session.HistoryWindowRequest{Anchor: "newest", Limit: 10}) |
| 159 | if err != nil { |
| 160 | t.Fatal(err) |
| 161 | } |
| 162 | if window.Status == "preparing" && time.Now().Before(deadline) { |
| 163 | time.Sleep(time.Millisecond) |
| 164 | continue |
| 165 | } |
| 166 | if window.Status != "ready" { |
| 167 | t.Fatalf("history status = %s", window.Status) |
| 168 | } |
| 169 | var count int |
| 170 | for _, row := range window.Messages { |
| 171 | if row.MessageID == "maintenance:"+op.id { |
| 172 | count++ |
| 173 | var message provider.Message |
| 174 | if err := json.Unmarshal(row.Inline, &message); err != nil { |
| 175 | t.Fatal(err) |
| 176 | } |
| 177 | var info event.SessionOperationInfo |
| 178 | if err := json.Unmarshal([]byte(message.Content), &info); err != nil { |
| 179 | t.Fatal(err) |
| 180 | } |
| 181 | if info.Status != status || info.OperationRevision != 3 { |
| 182 | t.Fatalf("history terminal = %+v", info) |
| 183 | } |
| 184 | } |
| 185 | } |
| 186 | if count != 1 { |
| 187 | t.Fatalf("maintenance history rows = %d", count) |
| 188 | } |
| 189 | break |
| 190 | } |
| 191 | }) |
| 192 | } |
| 193 | } |
| 194 | |
| 195 | func TestMaintenanceAdmissionRejectsUnpublishedAndDrainingOwners(t *testing.T) { |
| 196 | owner := extension.NewRuntimeOwner() |
| 197 | owner.Gate.Publish(2) |
| 198 | c := newOwnedTestController(t, Options{RuntimeOwner: owner, RuntimeGeneration: 1, Sink: event.Discard}) |
| 199 | if op, _, err := c.beginMaintenance(t.Context(), "compact"); !errors.Is(err, ErrRuntimeDraining) { |
| 200 | if op != nil { |
| 201 | c.finishMaintenanceOperation(op) |
| 202 | } |
| 203 | t.Fatalf("draining admission = %v", err) |
| 204 | } |
| 205 | } |
| 206 | |
| 207 | func TestMaintenanceTimeoutSettlesSharedRuntimeAndDrainsQueue(t *testing.T) { |
| 208 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions"))) |
| 209 | if err != nil { |
| 210 | t.Fatal(err) |
| 211 | } |
| 212 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "timeout"}) |
| 213 | if err != nil { |
| 214 | t.Fatal(err) |
| 215 | } |
| 216 | recovering := make(chan struct{}, 1) |
| 217 | sink := event.FuncSink(func(e event.Event) { |
| 218 | if e.SessionOperation != nil && e.SessionOperation.Status == "recovery_required" { |
| 219 | recovering <- struct{}{} |
| 220 | } |
| 221 | }) |
| 222 | exec := agent.New(nil, nil, agent.NewSession("system"), agent.Options{}, event.Discard) |
| 223 | c := newOwnedTestController(t, Options{Executor: exec, Sink: sink, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 224 | c.testCancelGrace = time.Millisecond |
| 225 | op, ctx, err := c.beginMaintenance(t.Context(), "compact") |
| 226 | if err != nil { |
| 227 | t.Fatal(err) |
| 228 | } |
| 229 | defer c.finishMaintenanceOperation(op) |
| 230 | queued := make(chan struct{}, 1) |
| 231 | if got := c.runGuardedOrPark(func(context.Context) error { queued <- struct{}{}; return nil }); got != turnParked { |
| 232 | t.Fatalf("queued admission = %v", got) |
| 233 | } |
| 234 | if !runtime.Cancel() { |
| 235 | t.Fatal("runtime did not acknowledge Stop") |
| 236 | } |
| 237 | select { |
| 238 | case <-recovering: |
| 239 | case <-time.After(3 * time.Second): |
| 240 | t.Fatal("timeout recovery was not published") |
| 241 | } |
| 242 | if runtime.StateSnapshot().Phase != session.RuntimeRecoveryRequired { |
| 243 | t.Fatal("runtime not fenced during timeout") |
| 244 | } |
| 245 | select { |
| 246 | case <-queued: |
| 247 | t.Fatal("queue ran during recovery") |
| 248 | default: |
| 249 | } |
| 250 | // The ignored worker has now exited with cancellation and no commit. The |
| 251 | // durable terminal is sufficient to release timeout recovery, not a model |
| 252 | // turn or controller replacement that could consume the queued request. |
| 253 | if err := c.executeMaintenance(op, ctx, func(context.Context) error { return ctx.Err() }); !errors.Is(err, context.Canceled) { |
| 254 | t.Fatalf("settled worker = %v", err) |
| 255 | } |
| 256 | select { |
| 257 | case <-queued: |
| 258 | case <-time.After(3 * time.Second): |
| 259 | t.Fatal("queue did not resume after safe settlement") |
| 260 | } |
| 261 | waitIdleAdmission(t, c) |
| 262 | if runtime.StateSnapshot().Phase != session.RuntimeIdle { |
| 263 | t.Fatal("shared runtime remained busy after queue drained") |
| 264 | } |
| 265 | } |
| 266 | |
| 267 | func TestMaintenanceDurabilityFailureRetainsOwner(t *testing.T) { |
| 268 | for _, boundary := range []string{"start", "terminal"} { |
| 269 | t.Run(boundary, func(t *testing.T) { |
| 270 | var fail atomic.Bool |
| 271 | store, err := session.CreateWithOptions(filepath.Join(t.TempDir(), "store"), "failure", session.OpenOptions{Sync: func(f *os.File) error { |
| 272 | if fail.Load() { |
| 273 | return errors.New("injected maintenance sync failure") |
| 274 | } |
| 275 | return f.Sync() |
| 276 | }}) |
| 277 | if err != nil { |
| 278 | t.Fatal(err) |
| 279 | } |
| 280 | service, err := session.NewService("desktop", failingFlushPersistence{session: store}) |
| 281 | if err != nil { |
| 282 | t.Fatal(err) |
| 283 | } |
| 284 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "failure"}) |
| 285 | if err != nil { |
| 286 | t.Fatal(err) |
| 287 | } |
| 288 | var completed atomic.Bool |
| 289 | sink := event.FuncSink(func(e event.Event) { |
| 290 | if e.SessionOperation != nil { |
| 291 | if e.SessionOperation.Status == "finalizing" && boundary == "terminal" { |
| 292 | fail.Store(true) |
| 293 | } |
| 294 | if maintenanceTerminalStatus(e.SessionOperation.Status) { |
| 295 | completed.Store(true) |
| 296 | } |
| 297 | } |
| 298 | }) |
| 299 | exec := agent.New(nil, nil, agent.NewSession("system"), agent.Options{}, event.Discard) |
| 300 | c := newOwnedTestController(t, Options{Executor: exec, Sink: sink, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 301 | if boundary == "start" { |
| 302 | fail.Store(true) |
| 303 | } |
| 304 | op, ctx, err := c.beginMaintenance(t.Context(), "compact") |
| 305 | if boundary == "terminal" { |
| 306 | if err != nil { |
| 307 | t.Fatal(err) |
| 308 | } |
| 309 | err = c.executeMaintenance(op, ctx, func(context.Context) error { return nil }) |
| 310 | } |
| 311 | if err == nil { |
| 312 | t.Fatal("failed sync reported success") |
| 313 | } |
| 314 | state := c.RuntimeStateSnapshot() |
| 315 | if state.Maintenance == nil || state.Maintenance.Status != "recovery_required" || state.Cancellable || completed.Load() { |
| 316 | t.Fatalf("persistence failure escaped recovery: %+v completed=%v", state, completed.Load()) |
| 317 | } |
| 318 | if runtime.StateSnapshot().Phase != session.RuntimeRecoveryRequired { |
| 319 | t.Fatal("shared runtime lost recovery barrier") |
| 320 | } |
| 321 | }) |
| 322 | } |
| 323 | } |
| 324 | |
| 325 | func TestMaintenanceRegressionEmptyPositionalCompression(t *testing.T) { |
| 326 | sess := agent.NewSession("system") |
| 327 | sess.Add(provider.Message{Role: provider.RoleUser, Content: "first user message"}) |
| 328 | exec := agent.New(nil, nil, sess, agent.Options{ContextWindow: 32000}, event.Discard) |
| 329 | if err := exec.SummarizeUpTo(t.Context(), 1); err != nil { |
| 330 | t.Errorf("empty range before first user should be noop, got %v", err) |
| 331 | } |
| 332 | } |
| 333 | |
| 334 | func TestMaintenanceRegressionLateCancelCannotReopenFinalizing(t *testing.T) { |
| 335 | c := newOwnedTestController(t, Options{Sink: event.Discard}) |
| 336 | op, _, err := c.beginMaintenance(context.Background(), "compact") |
| 337 | if err != nil { |
| 338 | t.Fatal(err) |
| 339 | } |
| 340 | defer c.finishMaintenanceOperation(op) |
| 341 | // Stop captures running, signals the worker, then is descheduled before |
| 342 | // publishing. The worker reaches finalization in that exact interval. |
| 343 | entered, resume, returned := make(chan struct{}), make(chan struct{}), make(chan struct{}) |
| 344 | cancel := op.cancel |
| 345 | op.cancel = func() { cancel(); close(entered); <-resume } |
| 346 | go func() { c.CancelSession(); close(returned) }() |
| 347 | <-entered |
| 348 | snapshot := c.ContextMaintenanceSnapshot() |
| 349 | if err := c.emitMaintenanceOperation(op, "finalizing", "", "", snapshot, false); err != nil { |
| 350 | t.Fatal(err) |
| 351 | } |
| 352 | close(resume) |
| 353 | <-returned |
| 354 | op.cancel = cancel |
| 355 | if got := c.RuntimeStateSnapshot(); got.Maintenance.Activity != "finalizing" || got.Cancellable { |
| 356 | t.Errorf("late cancellation reopened finalization: activity=%s cancellable=%v", got.Maintenance.Activity, got.Cancellable) |
| 357 | } |
| 358 | } |
| 359 |