返回 DeepSeek-Reasonix
session_events_test.go
根目录 / internal / control / session_events_test.go
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
690 lines GO