| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "encoding/json" |
| 6 | "errors" |
| 7 | "os" |
| 8 | "path/filepath" |
| 9 | "strings" |
| 10 | "testing" |
| 11 | "time" |
| 12 | |
| 13 | "reasonix/internal/agent" |
| 14 | "reasonix/internal/agent/testutil" |
| 15 | "reasonix/internal/event" |
| 16 | "reasonix/internal/provider" |
| 17 | "reasonix/internal/session" |
| 18 | "reasonix/internal/store" |
| 19 | "reasonix/internal/tool" |
| 20 | ) |
| 21 | |
| 22 | func loadDurableSessionProjection(t *testing.T, legacyPath string) session.Projection { |
| 23 | t.Helper() |
| 24 | commits, err := session.Replay(sessionDirectory(legacyPath), nil) |
| 25 | if err != nil { |
| 26 | t.Fatal(err) |
| 27 | } |
| 28 | projection, err := session.Project(commits) |
| 29 | if err != nil { |
| 30 | t.Fatal(err) |
| 31 | } |
| 32 | return projection |
| 33 | } |
| 34 | |
| 35 | func TestRuntimeSnapshotRecoversRecoveryRequiredFromV3Alone(t *testing.T) { |
| 36 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 37 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 38 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 39 | t.Cleanup(func() { c.Close() }) |
| 40 | recovery := event.RecoveryStatus{State: "recovery_required", Reason: "uncooperative_tool", Phase: "tool"} |
| 41 | payload, err := json.Marshal(recovery) |
| 42 | if err != nil { |
| 43 | t.Fatal(err) |
| 44 | } |
| 45 | store := c.sessionEventStore() |
| 46 | if _, err := store.Append(t.Context(), session.Batch{OperationID: "v3-only-recovery", Events: []session.Event{{Kind: "runtime/recovery", Payload: payload}}}); err != nil { |
| 47 | t.Fatal(err) |
| 48 | } |
| 49 | c.refreshRuntimeState(event.Event{}) |
| 50 | state := c.RuntimeStateSnapshot() |
| 51 | if state.Phase != "recovery_required" || state.TurnStatus != event.TurnRecoveryRequired || state.Recovery == nil || state.Recovery.Reason != recovery.Reason { |
| 52 | t.Fatalf("runtime did not adopt v3 recovery: %+v", state) |
| 53 | } |
| 54 | } |
| 55 | |
| 56 | func TestExclusiveControllerUsesBoundSessionIdentityAndWritesNoLegacyTranscript(t *testing.T) { |
| 57 | root := t.TempDir() |
| 58 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(root, "sessions-v4"))) |
| 59 | if err != nil { |
| 60 | t.Fatal(err) |
| 61 | } |
| 62 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 63 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "stable-session"}) |
| 64 | if err != nil { |
| 65 | t.Fatal(err) |
| 66 | } |
| 67 | legacyPath := filepath.Join(root, "sessions", "legacy-name.jsonl") |
| 68 | providerMock := testutil.NewMock("test", testutil.Turn{Text: "answer"}) |
| 69 | exec := agent.New(providerMock, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 70 | c := newOwnedTestController(t, Options{ |
| 71 | Runner: exec, Executor: exec, Sink: event.Discard, |
| 72 | SessionPath: legacyPath, SessionDir: filepath.Dir(legacyPath), |
| 73 | SessionService: service, SessionRuntime: runtime, ExclusiveSession: true, |
| 74 | }) |
| 75 | if err := c.RunTurn(t.Context(), "question"); err != nil { |
| 76 | t.Fatal(err) |
| 77 | } |
| 78 | if err := c.Snapshot(); err != nil { |
| 79 | t.Fatal(err) |
| 80 | } |
| 81 | ref, ok := c.SessionRef() |
| 82 | if !ok || ref.HostID != "desktop" || ref.SessionID != "stable-session" { |
| 83 | t.Fatalf("session ref = %+v, %v", ref, ok) |
| 84 | } |
| 85 | state := c.RuntimeStateSnapshot() |
| 86 | if state.HostID != ref.HostID || state.SessionID != ref.SessionID || state.SessionCodec != session.Codec { |
| 87 | t.Fatalf("runtime identity = %+v", state) |
| 88 | } |
| 89 | if _, err := os.Stat(legacyPath); !os.IsNotExist(err) { |
| 90 | t.Fatalf("exclusive controller wrote legacy transcript: %v", err) |
| 91 | } |
| 92 | if got := c.History(); len(got) < 3 || got[len(got)-1].Content != "answer" { |
| 93 | t.Fatalf("event-derived history = %#v", got) |
| 94 | } |
| 95 | c.Close() |
| 96 | if cached, ok := service.Runtime(ref); !ok || cached != runtime { |
| 97 | t.Fatal("terminal controller close did not retain the idle runtime") |
| 98 | } |
| 99 | if err := service.Close(t.Context(), ref); err != nil { |
| 100 | t.Fatal(err) |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | func TestOpenFailureDoesNotCloseAnotherControllersRuntime(t *testing.T) { |
| 105 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4"))) |
| 106 | if err != nil { |
| 107 | t.Fatal(err) |
| 108 | } |
| 109 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 110 | first, err := service.Create(t.Context(), session.CreateOptions{SessionID: "first"}) |
| 111 | if err != nil { |
| 112 | t.Fatal(err) |
| 113 | } |
| 114 | second, err := service.Create(t.Context(), session.CreateOptions{SessionID: "second"}) |
| 115 | if err != nil { |
| 116 | t.Fatal(err) |
| 117 | } |
| 118 | newController := func(runtime *session.Runtime) *Controller { |
| 119 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 120 | return newOwnedTestController(t, Options{ |
| 121 | Executor: exec, Sink: event.Discard, SessionService: service, |
| 122 | SessionRuntime: runtime, ExclusiveSession: true, |
| 123 | }) |
| 124 | } |
| 125 | owner := newController(second) |
| 126 | t.Cleanup(owner.Close) |
| 127 | caller := newController(first) |
| 128 | t.Cleanup(caller.Close) |
| 129 | |
| 130 | if _, err := second.Session().AppendBatch(t.Context(), "bad-plan", []session.Event{{Kind: "plan/state", Payload: []byte(`{"enabled":"invalid"}`)}}); err != nil { |
| 131 | t.Fatal(err) |
| 132 | } |
| 133 | if _, err := caller.OpenSession(t.Context(), second.Ref()); err == nil { |
| 134 | t.Fatal("OpenSession accepted an invalid domain projection") |
| 135 | } |
| 136 | if got, ok := service.Runtime(second.Ref()); !ok || got != second { |
| 137 | t.Fatal("failed client publication disposed another controller's runtime") |
| 138 | } |
| 139 | if _, err := second.Session().AppendBatch(t.Context(), "still-owned", []session.Event{{Kind: "diagnostic", Optional: true}}); err != nil { |
| 140 | t.Fatalf("shared runtime is unusable after another controller failed to bind: %v", err) |
| 141 | } |
| 142 | } |
| 143 | |
| 144 | func TestExclusiveControllerRuntimeSnapshotAndCancelUseExactV3Instance(t *testing.T) { |
| 145 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4"))) |
| 146 | if err != nil { |
| 147 | t.Fatal(err) |
| 148 | } |
| 149 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 150 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "runtime-owned"}) |
| 151 | if err != nil { |
| 152 | t.Fatal(err) |
| 153 | } |
| 154 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 155 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 156 | t.Cleanup(c.Close) |
| 157 | |
| 158 | initial := c.RuntimeStateSnapshot() |
| 159 | runtimeInitial := runtime.Snapshot() |
| 160 | if initial.RuntimeEpoch != runtimeInitial.Epoch || initial.ActivityRevision != runtimeInitial.ActivityRevision || initial.SessionID != runtime.Ref().SessionID || initial.HeadID != "" { |
| 161 | t.Fatalf("controller snapshot = %+v, runtime = %+v", initial, runtimeInitial) |
| 162 | } |
| 163 | |
| 164 | started := make(chan struct{}) |
| 165 | release := make(chan struct{}) |
| 166 | if got := c.runGuarded(func(ctx context.Context) error { |
| 167 | close(started) |
| 168 | select { |
| 169 | case <-ctx.Done(): |
| 170 | return ctx.Err() |
| 171 | case <-release: |
| 172 | return nil |
| 173 | } |
| 174 | }); got != turnStarted { |
| 175 | t.Fatalf("turn admission = %v", got) |
| 176 | } |
| 177 | <-started |
| 178 | receipt := c.CancelSession() |
| 179 | if !receipt.Accepted || receipt.SessionRef != runtime.Ref().SessionID || receipt.HeadID != "" || receipt.RuntimeEpoch != runtimeInitial.Epoch { |
| 180 | t.Fatalf("cancel receipt = %+v", receipt) |
| 181 | } |
| 182 | close(release) |
| 183 | for deadline := time.Now().Add(5 * time.Second); c.Running() && time.Now().Before(deadline); { |
| 184 | time.Sleep(time.Millisecond) |
| 185 | } |
| 186 | if c.Running() { |
| 187 | t.Fatal("cancelled v3 activity did not settle") |
| 188 | } |
| 189 | } |
| 190 | |
| 191 | func TestExclusiveSessionSwitchRestoresPlanAndGoalWithoutTodo(t *testing.T) { |
| 192 | service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4"))) |
| 193 | if err != nil { |
| 194 | t.Fatal(err) |
| 195 | } |
| 196 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 197 | first, err := service.Create(t.Context(), session.CreateOptions{SessionID: "first-domain"}) |
| 198 | if err != nil { |
| 199 | t.Fatal(err) |
| 200 | } |
| 201 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 202 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: first, ExclusiveSession: true}) |
| 203 | t.Cleanup(c.Close) |
| 204 | c.SetPlanMode(true) |
| 205 | c.SetGoal("finish the migration") |
| 206 | if _, err := first.Session().Flush(t.Context()); err != nil { |
| 207 | t.Fatal(err) |
| 208 | } |
| 209 | goalRaw := first.Session().Snapshot().Projection.GoalState |
| 210 | if string(goalRaw) == "" || strings.Contains(string(goalRaw), `"todos"`) { |
| 211 | t.Fatalf("goal event retained todo state: %s", goalRaw) |
| 212 | } |
| 213 | |
| 214 | second, err := service.Create(t.Context(), session.CreateOptions{SessionID: "second-domain"}) |
| 215 | if err != nil { |
| 216 | t.Fatal(err) |
| 217 | } |
| 218 | if _, err := second.Session().Flush(t.Context()); err != nil { |
| 219 | t.Fatal(err) |
| 220 | } |
| 221 | if err := service.Close(t.Context(), second.Ref()); err != nil { |
| 222 | t.Fatal(err) |
| 223 | } |
| 224 | if _, err := c.OpenSession(t.Context(), second.Ref()); err != nil { |
| 225 | t.Fatal(err) |
| 226 | } |
| 227 | if c.PlanMode() || c.Goal() != "" || c.GoalStatus() == GoalStatusRunning { |
| 228 | t.Fatalf("empty target inherited plan/goal: plan=%v goal=%q status=%q", c.PlanMode(), c.Goal(), c.GoalStatus()) |
| 229 | } |
| 230 | if _, err := c.OpenSession(t.Context(), first.Ref()); err != nil { |
| 231 | t.Fatal(err) |
| 232 | } |
| 233 | if !c.PlanMode() || c.Goal() != "finish the migration" || c.GoalStatus() == GoalStatusRunning { |
| 234 | t.Fatalf("restored domain state: plan=%v goal=%q status=%q", c.PlanMode(), c.Goal(), c.GoalStatus()) |
| 235 | } |
| 236 | } |
| 237 | |
| 238 | func TestExclusiveControllerNewPublishesFreshIdentityAndKeepsOldHistory(t *testing.T) { |
| 239 | root := t.TempDir() |
| 240 | persistence := session.NewFilesystemPersistence(filepath.Join(root, "sessions-v4")) |
| 241 | service, err := session.NewService("desktop", persistence) |
| 242 | if err != nil { |
| 243 | t.Fatal(err) |
| 244 | } |
| 245 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 246 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "old-session"}) |
| 247 | if err != nil { |
| 248 | t.Fatal(err) |
| 249 | } |
| 250 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 251 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionDir: filepath.Join(root, "sessions"), SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 252 | if err := c.NewSession(); err != nil { |
| 253 | t.Fatal(err) |
| 254 | } |
| 255 | ref, ok := c.SessionRef() |
| 256 | if !ok || ref.SessionID == "old-session" { |
| 257 | t.Fatalf("new identity = %+v, ok=%v", ref, ok) |
| 258 | } |
| 259 | oldRef := session.SessionRef{HostID: "desktop", SessionID: "old-session"} |
| 260 | if cached, ok := service.Runtime(oldRef); !ok || cached != runtime { |
| 261 | t.Fatal("old runtime was not retained for quick switching") |
| 262 | } |
| 263 | read, err := persistence.Open("old-session", session.ReadOnly) |
| 264 | if err != nil { |
| 265 | t.Fatalf("old history was removed by NewSession: %v", err) |
| 266 | } |
| 267 | _ = read.Close(t.Context()) |
| 268 | if err := service.Close(t.Context(), oldRef); err != nil { |
| 269 | t.Fatal(err) |
| 270 | } |
| 271 | c.Close() |
| 272 | if err := service.Close(t.Context(), ref); err != nil { |
| 273 | t.Fatal(err) |
| 274 | } |
| 275 | } |
| 276 | |
| 277 | func TestExclusiveControllerOpenMissingKeepsCurrentExactRuntime(t *testing.T) { |
| 278 | root := t.TempDir() |
| 279 | persistence := session.NewFilesystemPersistence(filepath.Join(root, "sessions-v4")) |
| 280 | service, err := session.NewService("desktop", persistence) |
| 281 | if err != nil { |
| 282 | t.Fatal(err) |
| 283 | } |
| 284 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 285 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "current"}) |
| 286 | if err != nil { |
| 287 | t.Fatal(err) |
| 288 | } |
| 289 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 290 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 291 | if _, err := c.OpenSession(t.Context(), session.SessionRef{HostID: "desktop", SessionID: "missing"}); !errors.Is(err, session.ErrSessionNotFound) { |
| 292 | t.Fatalf("OpenSession missing error = %v", err) |
| 293 | } |
| 294 | ref, ok := c.SessionRef() |
| 295 | if !ok || ref != runtime.Ref() { |
| 296 | t.Fatalf("failed open changed binding to %+v", ref) |
| 297 | } |
| 298 | if _, err := persistence.Stat(t.Context(), "missing"); !errors.Is(err, session.ErrSessionNotFound) { |
| 299 | t.Fatalf("failed open created missing session: %v", err) |
| 300 | } |
| 301 | c.Close() |
| 302 | } |
| 303 | |
| 304 | func TestExclusiveControllerClearDeletesClosedSourceAfterPublishingFreshIdentity(t *testing.T) { |
| 305 | root := t.TempDir() |
| 306 | persistence := session.NewFilesystemPersistence(filepath.Join(root, "sessions-v4")) |
| 307 | service, err := session.NewService("desktop", persistence) |
| 308 | if err != nil { |
| 309 | t.Fatal(err) |
| 310 | } |
| 311 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 312 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "clear-source"}) |
| 313 | if err != nil { |
| 314 | t.Fatal(err) |
| 315 | } |
| 316 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 317 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionDir: filepath.Join(root, "sessions"), SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 318 | if err := c.ClearSession(); err != nil { |
| 319 | t.Fatal(err) |
| 320 | } |
| 321 | ref, ok := c.SessionRef() |
| 322 | if !ok || ref.SessionID == "clear-source" { |
| 323 | t.Fatalf("clear identity = %+v, ok=%v", ref, ok) |
| 324 | } |
| 325 | if _, err := persistence.Stat(t.Context(), "clear-source"); !errors.Is(err, session.ErrSessionNotFound) { |
| 326 | t.Fatalf("cleared source remains in catalog: %v", err) |
| 327 | } |
| 328 | c.Close() |
| 329 | } |
| 330 | |
| 331 | func TestExclusiveControllerDelegatesClearPersistenceToIdentityHost(t *testing.T) { |
| 332 | root := t.TempDir() |
| 333 | persistence := session.NewFilesystemPersistence(filepath.Join(root, "sessions-v5")) |
| 334 | service, err := session.NewService("desktop", persistence) |
| 335 | if err != nil { |
| 336 | t.Fatal(err) |
| 337 | } |
| 338 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "clear-source", CWD: root, Origin: session.SessionOriginNew}) |
| 339 | if err != nil { |
| 340 | t.Fatal(err) |
| 341 | } |
| 342 | var committed SessionRotationRequest |
| 343 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 344 | c := newOwnedTestController(t, Options{ |
| 345 | Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true, |
| 346 | OnSessionRotation: func(_ context.Context, request SessionRotationRequest) (SessionRotationPlan, error) { |
| 347 | committed = request |
| 348 | return SessionRotationPlan{ |
| 349 | CreateOptions: session.CreateOptions{SessionID: "replacement", CWD: root, Origin: session.SessionOriginNew}, |
| 350 | Commit: func(_ context.Context, ref session.SessionRef) error { |
| 351 | if ref.SessionID != "replacement" { |
| 352 | t.Fatalf("replacement ref = %+v", ref) |
| 353 | } |
| 354 | return nil |
| 355 | }, |
| 356 | }, nil |
| 357 | }, |
| 358 | }) |
| 359 | if err := c.ClearSession(); err != nil { |
| 360 | t.Fatal(err) |
| 361 | } |
| 362 | if committed.Source.SessionID != "clear-source" || committed.Reason != "clear" { |
| 363 | t.Fatalf("rotation request = %+v", committed) |
| 364 | } |
| 365 | if _, err := persistence.Stat(t.Context(), "clear-source"); err != nil { |
| 366 | t.Fatalf("host-owned clear permanently deleted source: %v", err) |
| 367 | } |
| 368 | info, err := persistence.Stat(t.Context(), "replacement") |
| 369 | if err != nil { |
| 370 | t.Fatal(err) |
| 371 | } |
| 372 | if info.CWD != root || info.Origin != session.SessionOriginNew { |
| 373 | t.Fatalf("replacement header = %+v", info) |
| 374 | } |
| 375 | c.Close() |
| 376 | } |
| 377 | |
| 378 | func TestExclusiveControllerForkUsesTypedTurnBoundaryAndNoLegacyTranscript(t *testing.T) { |
| 379 | root := t.TempDir() |
| 380 | persistence := session.NewFilesystemPersistence(filepath.Join(root, "sessions-v4")) |
| 381 | service, err := session.NewService("desktop", persistence) |
| 382 | if err != nil { |
| 383 | t.Fatal(err) |
| 384 | } |
| 385 | t.Cleanup(func() { _ = service.CloseAll(context.Background()) }) |
| 386 | runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "fork-parent", CWD: root, Origin: session.SessionOriginNew}) |
| 387 | if err != nil { |
| 388 | t.Fatal(err) |
| 389 | } |
| 390 | prov := testutil.NewMock("fork", testutil.Turn{Text: "answer one"}, testutil.Turn{Text: "answer two"}) |
| 391 | exec := agent.New(prov, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 392 | c := newOwnedTestController(t, Options{Runner: exec, Executor: exec, Sink: event.Discard, SessionDir: filepath.Join(root, "legacy"), SessionService: service, SessionRuntime: runtime, ExclusiveSession: true}) |
| 393 | if err := c.RunTurn(t.Context(), "one"); err != nil { |
| 394 | t.Fatal(err) |
| 395 | } |
| 396 | if err := c.RunTurn(t.Context(), "two"); err != nil { |
| 397 | t.Fatal(err) |
| 398 | } |
| 399 | childID, err := c.ForkNamed(1, "first turn") |
| 400 | if err != nil { |
| 401 | t.Fatal(err) |
| 402 | } |
| 403 | ref, ok := c.SessionRef() |
| 404 | if !ok || ref.SessionID != childID || ref.SessionID == "fork-parent" { |
| 405 | t.Fatalf("fork ref = %+v, child id=%q", ref, childID) |
| 406 | } |
| 407 | snapshot := runtimeSnapshotForTest(t, c) |
| 408 | if got := snapshot.Projection.Title; got != "first turn" { |
| 409 | t.Fatalf("child title = %q", got) |
| 410 | } |
| 411 | if got := snapshot.Projection.Messages; len(got) != 3 || got[1].Content != "one" || got[2].Content != "answer one" { |
| 412 | t.Fatalf("child history = %#v", got) |
| 413 | } |
| 414 | childInfo, err := persistence.Stat(t.Context(), childID) |
| 415 | if err != nil { |
| 416 | t.Fatal(err) |
| 417 | } |
| 418 | if childInfo.CWD != root || childInfo.ParentSessionID != "fork-parent" || childInfo.Origin != session.SessionOriginFork { |
| 419 | t.Fatalf("child header = %+v", childInfo) |
| 420 | } |
| 421 | entries, err := os.ReadDir(filepath.Join(root, "legacy")) |
| 422 | if err != nil && !os.IsNotExist(err) { |
| 423 | t.Fatal(err) |
| 424 | } |
| 425 | if len(entries) != 0 { |
| 426 | t.Fatalf("v3 fork created legacy artifacts: %+v", entries) |
| 427 | } |
| 428 | c.Close() |
| 429 | } |
| 430 | |
| 431 | func runtimeSnapshotForTest(t *testing.T, c *Controller) session.Snapshot { |
| 432 | t.Helper() |
| 433 | _, runtime, ok := c.SessionBinding() |
| 434 | if !ok { |
| 435 | t.Fatal("controller has no v3 runtime") |
| 436 | } |
| 437 | return runtime.Session().Snapshot() |
| 438 | } |
| 439 | |
| 440 | func TestSessionPathBindingSeedsTranscriptBeforeLaterStateEvents(t *testing.T) { |
| 441 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 442 | session := agent.NewSession("system") |
| 443 | session.Add(provider.Message{Role: provider.RoleUser, Content: "question"}) |
| 444 | session.Add(provider.Message{Role: provider.RoleAssistant, Content: "answer"}) |
| 445 | exec := agent.New(nil, tool.NewRegistry(), session, agent.Options{}, event.Discard) |
| 446 | c := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard}) |
| 447 | t.Cleanup(func() { c.Close() }) |
| 448 | |
| 449 | c.SetSessionPath(legacyPath) |
| 450 | snapshot, ok := c.sessionEventSnapshot() |
| 451 | if !ok || len(snapshot.Projection.ModelMessages) != 3 { |
| 452 | t.Fatalf("bound projection = %#v, available=%v", snapshot.Projection.ModelMessages, ok) |
| 453 | } |
| 454 | for i, message := range snapshot.Projection.ModelMessages { |
| 455 | if message.ID == "" { |
| 456 | t.Fatalf("bound projection message %d has no stable id", i) |
| 457 | } |
| 458 | } |
| 459 | if snapshot.Projection.ModelMessages[0].Role != provider.RoleSystem { |
| 460 | t.Fatalf("bound projection lost leading system message: %#v", snapshot.Projection.ModelMessages) |
| 461 | } |
| 462 | |
| 463 | // A later state-only event must not become the first authoritative v3 |
| 464 | // commit and hide the transcript from History. |
| 465 | c.SetPlanMode(false) |
| 466 | if got := c.History(); len(got) != 3 || got[0].Role != provider.RoleSystem { |
| 467 | t.Fatalf("history after state event = %#v", got) |
| 468 | } |
| 469 | } |
| 470 | |
| 471 | func TestLateManagedHostBindingFencesCandidateUntilLeaseActivation(t *testing.T) { |
| 472 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 473 | initial := agent.NewSession("system") |
| 474 | initial.Add(provider.Message{Role: provider.RoleUser, Content: "old"}) |
| 475 | exec := agent.New(nil, tool.NewRegistry(), initial, agent.Options{}, event.Discard) |
| 476 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 477 | t.Cleanup(func() { c.Close() }) |
| 478 | |
| 479 | // Serve and other embedders may install their transition owner after the |
| 480 | // controller is built. From that point onward a private replacement without |
| 481 | // the final lease may observe, but may not publish, the shared projection. |
| 482 | c.SetOnSessionTransition(func(SessionTransitionInfo) error { return nil }) |
| 483 | before, _ := c.sessionEventSnapshot() |
| 484 | candidate := agent.NewSession("system") |
| 485 | candidate.Add(provider.Message{Role: provider.RoleUser, Content: "old"}) |
| 486 | candidate.Add(provider.Message{Role: provider.RoleAssistant, Content: "candidate"}) |
| 487 | c.executor.SetSession(candidate) |
| 488 | if err := c.RecordSessionMessages(t.Context(), "unpublished-candidate", []provider.Message{{ID: agent.NewMessageID(), Role: provider.RoleAssistant, Content: "must-not-publish"}}); err != nil { |
| 489 | t.Fatal(err) |
| 490 | } |
| 491 | after, _ := c.sessionEventSnapshot() |
| 492 | if after.EventSequence != before.EventSequence { |
| 493 | t.Fatalf("unleased candidate mutated v3: before=%d after=%d", before.EventSequence, after.EventSequence) |
| 494 | } |
| 495 | |
| 496 | lease, err := agent.TryAcquireSessionLease(legacyPath) |
| 497 | if err != nil { |
| 498 | t.Fatal(err) |
| 499 | } |
| 500 | t.Cleanup(lease.Release) |
| 501 | if err := c.BindSessionWriteAuthority(lease); err != nil { |
| 502 | t.Fatal(err) |
| 503 | } |
| 504 | activated, _ := c.sessionEventSnapshot() |
| 505 | if got := activated.Projection.ModelMessages; len(got) != 3 || got[2].Content != "candidate" { |
| 506 | t.Fatalf("activated projection = %#v", got) |
| 507 | } |
| 508 | } |
| 509 | |
| 510 | func TestControllerUsesV3AsBusinessEventStore(t *testing.T) { |
| 511 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 512 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 513 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 514 | t.Cleanup(func() { c.Close() }) |
| 515 | |
| 516 | admitted := c.prepareTurnAdmission(func(context.Context) error { return nil }) |
| 517 | if err := admitted(context.Background()); err != nil { |
| 518 | t.Fatal(err) |
| 519 | } |
| 520 | todos := []event.Todo{{Content: "one", Status: "in_progress"}} |
| 521 | toolMessage := provider.Message{ID: "tool-message-1", Role: provider.RoleTool, ToolCallID: "todo-1", Name: "todo_write", Content: "updated"} |
| 522 | if err := c.emitTurnEventChecked(event.Event{Kind: event.ToolResult, Tool: event.Tool{ |
| 523 | ID: "todo-1", Name: "todo_write", TodoWritten: true, Todos: todos, |
| 524 | PresentedFiles: []provider.PresentedFile{{Path: "report.txt", Description: "report"}}, |
| 525 | Execution: &event.ShellExecution{Kind: "shell", State: "completed"}, |
| 526 | }, CommittedMessage: &toolMessage}); err != nil { |
| 527 | t.Fatal(err) |
| 528 | } |
| 529 | |
| 530 | snapshot, ok := c.sessionEventSnapshot() |
| 531 | if !ok || snapshot.EventSequence == 0 { |
| 532 | t.Fatalf("v3 snapshot = %#v, %v", snapshot, ok) |
| 533 | } |
| 534 | if !snapshot.Projection.TodoWritten || len(snapshot.Projection.Todos) != 1 { |
| 535 | t.Fatalf("v3 todos = %#v, written=%v", snapshot.Projection.Todos, snapshot.Projection.TodoWritten) |
| 536 | } |
| 537 | if _, err := os.Stat(store.SessionTurnEventLog(legacyPath)); !os.IsNotExist(err) { |
| 538 | t.Fatalf("legacy turn ledger was written: %v", err) |
| 539 | } |
| 540 | if err := c.CheckpointSession(context.Background(), agent.CheckpointBeforeModel); err != nil { |
| 541 | t.Fatal(err) |
| 542 | } |
| 543 | snapshot, _ = c.sessionEventSnapshot() |
| 544 | if snapshot.DurableSequence != snapshot.EventSequence { |
| 545 | t.Fatalf("durable=%d event=%d", snapshot.DurableSequence, snapshot.EventSequence) |
| 546 | } |
| 547 | commits, err := session.Replay(sessionDirectory(legacyPath), nil) |
| 548 | if err != nil || len(commits) == 0 { |
| 549 | t.Fatalf("Replay = %d commits, %v", len(commits), err) |
| 550 | } |
| 551 | foundAtomicTodo := false |
| 552 | for _, commit := range commits { |
| 553 | if len(commit.Events) == 3 && commit.Events[0].Kind == "message/complete" && commit.Events[1].Kind == "todo/write" && commit.Events[2].Kind == "tool/result" { |
| 554 | foundAtomicTodo = true |
| 555 | var payload struct { |
| 556 | Todos []event.Todo `json:"todos"` |
| 557 | TodoWritten bool `json:"todoWritten"` |
| 558 | PresentedFiles []provider.PresentedFile `json:"presentedFiles"` |
| 559 | Execution event.ShellExecution `json:"execution"` |
| 560 | } |
| 561 | if err := json.Unmarshal(commit.Events[2].Payload, &payload); err != nil { |
| 562 | t.Fatal(err) |
| 563 | } |
| 564 | if !payload.TodoWritten || len(payload.Todos) != 1 || len(payload.PresentedFiles) != 1 || payload.Execution.State != "completed" { |
| 565 | t.Fatalf("structured tool result metadata was lost: %#v", payload) |
| 566 | } |
| 567 | break |
| 568 | } |
| 569 | } |
| 570 | if !foundAtomicTodo { |
| 571 | t.Fatalf("todo result was not one ordered atomic batch: %#v", commits) |
| 572 | } |
| 573 | } |
| 574 | |
| 575 | func TestCheckpointDoesNotInferMessagesFromLegacyTranscript(t *testing.T) { |
| 576 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 577 | session := agent.NewSession("system") |
| 578 | exec := agent.New(nil, tool.NewRegistry(), session, agent.Options{}, event.Discard) |
| 579 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 580 | t.Cleanup(func() { c.Close() }) |
| 581 | recorded := provider.Message{ID: "recorded", Role: provider.RoleUser, Content: "recorded"} |
| 582 | if err := c.RecordSessionMessages(context.Background(), "test", []provider.Message{recorded}); err != nil { |
| 583 | t.Fatal(err) |
| 584 | } |
| 585 | session.Add(recorded) |
| 586 | // Simulate an obsolete caller mutating the display transcript directly. |
| 587 | // The checkpoint must not mirror or infer that mutation into v3. |
| 588 | session.Add(provider.Message{ID: "legacy-only", Role: provider.RoleAssistant, Content: "legacy"}) |
| 589 | if err := c.CheckpointSession(context.Background(), agent.CheckpointBeforeModel); err != nil { |
| 590 | t.Fatal(err) |
| 591 | } |
| 592 | messages := c.sessionEventStore().Snapshot().Projection.Messages |
| 593 | for _, message := range messages { |
| 594 | if message.ID == "legacy-only" { |
| 595 | t.Fatalf("checkpoint inferred an unrecorded message: %#v", messages) |
| 596 | } |
| 597 | } |
| 598 | } |
| 599 | |
| 600 | func TestTurnEndUsesExplicitFinalMessageCommit(t *testing.T) { |
| 601 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 602 | session := agent.NewSession("system") |
| 603 | exec := agent.New(nil, tool.NewRegistry(), session, agent.Options{}, event.Discard) |
| 604 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 605 | t.Cleanup(func() { c.Close() }) |
| 606 | if err := c.prepareTurnAdmission(func(context.Context) error { return nil })(context.Background()); err != nil { |
| 607 | t.Fatal(err) |
| 608 | } |
| 609 | assistant := provider.Message{ID: "a1", Role: provider.RoleAssistant, Content: "done"} |
| 610 | if err := c.RecordSessionMessages(context.Background(), "test-final", []provider.Message{assistant}); err != nil { |
| 611 | t.Fatal(err) |
| 612 | } |
| 613 | session.Add(assistant) |
| 614 | if err := c.emitTurnEventChecked(event.Event{Kind: event.TurnDone, Status: event.TurnCompleted}); err != nil { |
| 615 | t.Fatal(err) |
| 616 | } |
| 617 | snapshot, _ := c.sessionEventSnapshot() |
| 618 | if snapshot.Projection.TurnID != "" || len(snapshot.Projection.ModelMessages) == 0 { |
| 619 | t.Fatalf("terminal projection = %#v", snapshot.Projection) |
| 620 | } |
| 621 | } |
| 622 | |
| 623 | func TestTurnEndClosesPendingInteractionsInTheSameBatch(t *testing.T) { |
| 624 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 625 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 626 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 627 | t.Cleanup(func() { c.Close() }) |
| 628 | if err := c.prepareTurnAdmission(func(context.Context) error { return nil })(context.Background()); err != nil { |
| 629 | t.Fatal(err) |
| 630 | } |
| 631 | if err := c.emitTurnEventChecked(event.Event{Kind: event.AskRequest, ItemID: "ask-1", PromptKind: string(PromptAsk)}); err != nil { |
| 632 | t.Fatal(err) |
| 633 | } |
| 634 | if err := c.emitTurnEventChecked(event.Event{Kind: event.TurnDone, Status: event.TurnInterrupted, Cancelled: true}); err != nil { |
| 635 | t.Fatal(err) |
| 636 | } |
| 637 | snapshot, _ := c.sessionEventSnapshot() |
| 638 | if len(snapshot.Projection.Interactions) != 0 { |
| 639 | t.Fatalf("terminal projection kept pending interactions: %#v", snapshot.Projection.Interactions) |
| 640 | } |
| 641 | if _, err := c.flushSessionEvents(context.Background()); err != nil { |
| 642 | t.Fatal(err) |
| 643 | } |
| 644 | commits, err := session.Replay(sessionDirectory(legacyPath), nil) |
| 645 | if err != nil { |
| 646 | t.Fatal(err) |
| 647 | } |
| 648 | last := commits[len(commits)-1] |
| 649 | found := false |
| 650 | for _, item := range last.Events { |
| 651 | if item.Kind == "interaction/resolved" && string(item.Payload) == `{"id":"ask-1","state":"cancelled"}` { |
| 652 | found = true |
| 653 | } |
| 654 | } |
| 655 | if !found { |
| 656 | t.Fatalf("turn-end batch did not cancel the interaction: %#v", last.Events) |
| 657 | } |
| 658 | } |
| 659 | |
| 660 | func TestPromptResolutionAndPlanStateShareOneAtomicBatch(t *testing.T) { |
| 661 | legacyPath := filepath.Join(t.TempDir(), "sessions", "chat.jsonl") |
| 662 | exec := agent.New(nil, tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard) |
| 663 | c := newOwnedTestController(t, Options{Executor: exec, SessionPath: legacyPath, Sink: event.Discard}) |
| 664 | t.Cleanup(func() { c.Close() }) |
| 665 | if err := c.prepareTurnAdmission(func(context.Context) error { return nil })(context.Background()); err != nil { |
| 666 | t.Fatal(err) |
| 667 | } |
| 668 | if err := c.emitTurnEventChecked(event.Event{Kind: event.ApprovalRequest, ItemID: "plan-1", PromptKind: string(PromptPlan)}); err != nil { |
| 669 | t.Fatal(err) |
| 670 | } |
| 671 | payload := []byte(`{"requestId":"plan-1","decision":"start_execution"}`) |
| 672 | if err := c.emitTurnEventChecked(event.Event{ |
| 673 | Kind: event.PromptAnswered, ItemID: "plan-1", InteractionState: string(PromptAnswered), |
| 674 | DomainKind: "plan/state", DomainPayload: payload, |
| 675 | }); err != nil { |
| 676 | t.Fatal(err) |
| 677 | } |
| 678 | if _, err := c.flushSessionEvents(context.Background()); err != nil { |
| 679 | t.Fatal(err) |
| 680 | } |
| 681 | commits, err := session.Replay(sessionDirectory(legacyPath), nil) |
| 682 | if err != nil { |
| 683 | t.Fatal(err) |
| 684 | } |
| 685 | last := commits[len(commits)-1] |
| 686 | if len(last.Events) != 2 || last.Events[0].Kind != "interaction/resolved" || last.Events[1].Kind != "plan/state" { |
| 687 | t.Fatalf("resolution split across batches: %#v", last.Events) |
| 688 | } |
| 689 | } |
| 690 |