返回 DeepSeek-Reasonix
service_test.go
根目录 / internal / session / service_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "os"
8 "path/filepath"
9 "strings"
10 "sync"
11 "testing"
12
13 "reasonix/internal/agent"
14 "reasonix/internal/provider"
15 )
16
17 func TestServiceForkAndRewindUsePersistedTurnBoundaries(t *testing.T) {
18 root := filepath.Join(t.TempDir(), "sessions-v4")
19 service, err := NewService("local", NewFilesystemPersistence(root))
20 if err != nil {
21 t.Fatal(err)
22 }
23 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
24 parent, err := service.Create(t.Context(), CreateOptions{SessionID: "parent"})
25 if err != nil {
26 t.Fatal(err)
27 }
28 appendTurn := func(operation, turnID, messageID string) {
29 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: messageID, Role: provider.RoleUser, Content: messageID}})
30 _, appendErr := parent.Session().Append(t.Context(), Batch{OperationID: operation, TurnID: turnID, Events: []Event{
31 {Kind: "turn/start"}, {Kind: "message/complete", Payload: payload}, {Kind: "turn/end", Payload: json.RawMessage(`{"status":"completed"}`)},
32 }})
33 if appendErr != nil {
34 t.Fatal(appendErr)
35 }
36 }
37 appendTurn("same-size-a", "turn-a", "message-a")
38 appendTurn("same-size-b", "turn-b", "message-b")
39
40 child, err := service.Fork(t.Context(), parent.Ref(), "turn-a", "child-after-a")
41 if err != nil {
42 t.Fatal(err)
43 }
44 childMessages := child.Session().Snapshot().Projection.ModelMessages
45 if len(childMessages) != 1 || childMessages[0].ID != "message-a" {
46 t.Fatalf("fork messages = %#v", childMessages)
47 }
48 if got := child.Session().Manifest().InheritedEvents; got != 3 {
49 t.Fatalf("inherited event count = %d", got)
50 }
51
52 rewound, err := service.Rewind(t.Context(), parent.Ref(), "turn-b", "child-before-b")
53 if err != nil {
54 t.Fatal(err)
55 }
56 rewoundMessages := rewound.Session().Snapshot().Projection.ModelMessages
57 if len(rewoundMessages) != 1 || rewoundMessages[0].ID != "message-a" {
58 t.Fatalf("rewind messages = %#v", rewoundMessages)
59 }
60
61 for _, runtime := range []*Runtime{child, rewound, parent} {
62 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
63 t.Fatal(err)
64 }
65 }
66 }
67
68 func TestServiceConcurrentOpenPublishesOneExactRuntime(t *testing.T) {
69 root := filepath.Join(t.TempDir(), "sessions-v4")
70 persistence := NewFilesystemPersistence(root)
71 handle, err := persistence.Create(CreateOptions{SessionID: "shared"})
72 if err != nil {
73 t.Fatal(err)
74 }
75 if err := handle.Close(t.Context()); err != nil {
76 t.Fatal(err)
77 }
78 service, err := NewService("local", persistence)
79 if err != nil {
80 t.Fatal(err)
81 }
82 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
83 ref := SessionRef{HostID: "local", SessionID: "shared"}
84
85 const callers = 16
86 results := make(chan *Runtime, callers)
87 errorsFound := make(chan error, callers)
88 start := make(chan struct{})
89 var group sync.WaitGroup
90 for range callers {
91 group.Go(func() {
92 <-start
93 runtime, openErr := service.Open(context.Background(), ref)
94 if openErr != nil {
95 errorsFound <- openErr
96 return
97 }
98 results <- runtime.Runtime()
99 })
100 }
101 close(start)
102 group.Wait()
103 close(results)
104 close(errorsFound)
105 for openErr := range errorsFound {
106 t.Fatalf("Open: %v", openErr)
107 }
108 var published *Runtime
109 for runtime := range results {
110 if published == nil {
111 published = runtime
112 } else if runtime != published {
113 t.Fatal("concurrent Open returned multiple runtime instances")
114 }
115 }
116 if published == nil {
117 t.Fatal("no runtime published")
118 }
119 if !service.Detach(published) {
120 t.Fatal("exact runtime did not detach")
121 }
122 if service.Detach(published) {
123 t.Fatal("stale callback detached twice")
124 }
125 if err := published.close(t.Context()); err != nil {
126 t.Fatal(err)
127 }
128 }
129
130 func TestServicePrepareCreateIsInvisibleUntilExactPublish(t *testing.T) {
131 root := filepath.Join(t.TempDir(), "sessions-v4")
132 service, err := NewService("local", NewFilesystemPersistence(root))
133 if err != nil {
134 t.Fatal(err)
135 }
136 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
137 prepared, err := service.PrepareCreate(t.Context(), CreateOptions{SessionID: "prepared"})
138 if err != nil {
139 t.Fatal(err)
140 }
141 ref := prepared.Runtime().Ref()
142 if _, ok := service.Runtime(ref); ok {
143 t.Fatal("prepared runtime was visible before publish")
144 }
145 if _, err := prepared.Runtime().Session().AppendBatch(t.Context(), "seed", []Event{{Kind: "diagnostic", Optional: true}}); err != nil {
146 t.Fatal(err)
147 }
148 if _, err := prepared.Runtime().Session().Flush(t.Context()); err != nil {
149 t.Fatal(err)
150 }
151 published, err := service.Publish(prepared)
152 if err != nil {
153 t.Fatal(err)
154 }
155 if current, ok := service.Runtime(ref); !ok || current != published.Runtime() {
156 t.Fatal("exact prepared runtime was not published")
157 }
158 if err := service.Discard(t.Context(), prepared); err == nil {
159 t.Fatal("published runtime was discarded outside service close")
160 }
161 if err := service.Close(t.Context(), ref); err != nil {
162 t.Fatal(err)
163 }
164 }
165
166 func TestQueryListProjectsEventBackedTitleAndCompletedTurns(t *testing.T) {
167 root := filepath.Join(t.TempDir(), "sessions-v4")
168 service, err := NewService("host-a", NewFilesystemPersistence(root))
169 if err != nil {
170 t.Fatal(err)
171 }
172 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
173 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "listed"})
174 if err != nil {
175 t.Fatal(err)
176 }
177 title, _ := json.Marshal(map[string]string{"title": "Event title"})
178 if _, err := runtime.Session().AppendBatch(t.Context(), "title", []Event{{Kind: "session/title", Payload: title}}); err != nil {
179 t.Fatal(err)
180 }
181 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "turn", TurnID: "turn-1", Events: []Event{
182 {Kind: "turn/start"}, {Kind: "turn/end", Payload: json.RawMessage(`{"status":"completed"}`)},
183 }}); err != nil {
184 t.Fatal(err)
185 }
186 if _, err := runtime.Session().Flush(t.Context()); err != nil {
187 t.Fatal(err)
188 }
189 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
190 t.Fatal(err)
191 }
192 page, err := service.Query().List(t.Context(), "", 50)
193 if err != nil {
194 t.Fatal(err)
195 }
196 if len(page.Sessions) != 1 {
197 t.Fatalf("sessions = %+v", page.Sessions)
198 }
199 got := page.Sessions[0]
200 if got.Ref != (SessionRef{HostID: "host-a", SessionID: "listed"}) || got.Codec != Codec || got.Title != "Event title" || got.Turns != 1 {
201 t.Fatalf("listed session = %+v", got)
202 }
203 }
204
205 func TestServiceDiscardPreparedCreateReleasesWriter(t *testing.T) {
206 root := filepath.Join(t.TempDir(), "sessions-v4")
207 service, err := NewService("local", NewFilesystemPersistence(root))
208 if err != nil {
209 t.Fatal(err)
210 }
211 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
212 prepared, err := service.PrepareCreate(t.Context(), CreateOptions{SessionID: "discarded"})
213 if err != nil {
214 t.Fatal(err)
215 }
216 if err := service.Discard(t.Context(), prepared); err != nil {
217 t.Fatal(err)
218 }
219 handle, err := NewFilesystemPersistence(root).Open("discarded", ReadWrite)
220 if err != nil {
221 t.Fatalf("writer lease remained held after discard: %v", err)
222 }
223 if err := handle.Close(t.Context()); err != nil {
224 t.Fatal(err)
225 }
226 }
227
228 func TestRuntimeCancellationUsesOwnedActivityWithoutTurnID(t *testing.T) {
229 root := filepath.Join(t.TempDir(), "sessions-v4")
230 service, err := NewService("local", NewFilesystemPersistence(root))
231 if err != nil {
232 t.Fatal(err)
233 }
234 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
235 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "cancel"})
236 if err != nil {
237 t.Fatal(err)
238 }
239 ctx, exec := bindTestExecution(t, runtime, "model")
240 ref := runtime.Ref()
241 snapshot, err := service.Cancel(ref)
242 if err != nil {
243 t.Fatal(err)
244 }
245 if snapshot.Phase != RuntimeCancelling {
246 t.Fatalf("phase after cancel = %s", snapshot.Phase)
247 }
248 select {
249 case <-ctx.Done():
250 default:
251 t.Fatal("cancel signal was not delivered")
252 }
253 exec.Finish()
254 if got := runtime.Snapshot().Phase; got != RuntimeIdle {
255 t.Fatalf("phase after activity exit = %s", got)
256 }
257 if err := service.Close(t.Context(), ref); err != nil {
258 t.Fatal(err)
259 }
260 if err := service.Close(t.Context(), ref); err != nil {
261 t.Fatalf("idempotent Close: %v", err)
262 }
263 }
264
265 func TestServiceObserveDistinguishesLiveAcceptedAndColdDurablePrefixes(t *testing.T) {
266 root := filepath.Join(t.TempDir(), "sessions-v4")
267 service, err := NewService("local", NewFilesystemPersistence(root))
268 if err != nil {
269 t.Fatal(err)
270 }
271 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
272 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "observe"})
273 if err != nil {
274 t.Fatal(err)
275 }
276 if _, err := runtime.Session().AppendBatch(t.Context(), "accepted", []Event{{Kind: "diagnostic", Optional: true}}); err != nil {
277 t.Fatal(err)
278 }
279 ref := runtime.Ref()
280 live, err := service.Observe(t.Context(), ref, 0, 10)
281 if err != nil {
282 t.Fatal(err)
283 }
284 if live.Runtime == nil || live.Runtime.Session.EventSequence != 1 || live.Runtime.Session.DurableSequence != 0 || len(live.Events.Commits) != 1 {
285 t.Fatalf("live observe = %+v", live)
286 }
287 if !service.Detach(runtime) {
288 t.Fatal("detach runtime")
289 }
290 // The detached writer still owns the unflushed prefix. A competing cold
291 // read sees the durable prefix without creating an Agent or taking a lease.
292 cold, err := service.Observe(t.Context(), ref, 0, 10)
293 if err != nil {
294 t.Fatal(err)
295 }
296 if cold.Runtime != nil || len(cold.Events.Commits) != 0 {
297 t.Fatalf("cold observe exposed unflushed events: %+v", cold)
298 }
299 if err := runtime.close(t.Context()); err != nil {
300 t.Fatal(err)
301 }
302 }
303
304 func TestServiceClosePreservesBusyRuntime(t *testing.T) {
305 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
306 if err != nil {
307 t.Fatal(err)
308 }
309 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
310 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "busy"})
311 if err != nil {
312 t.Fatal(err)
313 }
314 _, exec := bindTestExecution(t, runtime, "tool")
315 if err := service.Close(t.Context(), runtime.Ref()); !errors.Is(err, ErrRuntimeBusy) {
316 t.Fatalf("Close busy error = %v", err)
317 }
318 if current, ok := service.Runtime(runtime.Ref()); !ok || current != runtime {
319 t.Fatal("busy close released the exact runtime")
320 }
321 exec.Finish()
322 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
323 t.Fatal(err)
324 }
325 }
326
327 func TestServiceOpenClosesPersistedRuntimeWithoutRestoringAuthority(t *testing.T) {
328 root := filepath.Join(t.TempDir(), "sessions-v4")
329 persistence := NewFilesystemPersistence(root)
330 handle, err := persistence.Create(CreateOptions{SessionID: "restart"})
331 if err != nil {
332 t.Fatal(err)
333 }
334 store := handle
335 if _, err := store.Append(t.Context(), Batch{OperationID: "interrupted", TurnID: "turn", Events: []Event{
336 {Kind: "turn/start"},
337 {Kind: "tool/start", Payload: json.RawMessage(`{"id":"call","name":"bash"}`)},
338 {Kind: "interaction/created", Payload: json.RawMessage(`{"id":"approval","state":"pending"}`)},
339 }}); err != nil {
340 t.Fatal(err)
341 }
342 if err := store.Close(t.Context()); err != nil {
343 t.Fatal(err)
344 }
345 service, err := NewService("local", persistence)
346 if err != nil {
347 t.Fatal(err)
348 }
349 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
350 binding, err := service.Open(t.Context(), SessionRef{HostID: "local", SessionID: "restart"})
351 if err != nil {
352 t.Fatal(err)
353 }
354 runtime := binding.Runtime()
355 snapshot := runtime.Snapshot()
356 if snapshot.Phase != RuntimeIdle || snapshot.Session.Projection.TurnID != "" || len(snapshot.Session.Projection.ActiveTools) != 0 || len(snapshot.Session.Projection.Interactions) != 0 {
357 t.Fatalf("restored stale runtime authority: %+v", snapshot)
358 }
359 if snapshot.Session.Projection.TurnStatus != "interrupted" {
360 t.Fatalf("turn status = %q", snapshot.Session.Projection.TurnStatus)
361 }
362 if err := binding.Release(t.Context()); err != nil {
363 t.Fatal(err)
364 }
365 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
366 t.Fatal(err)
367 }
368 }
369
370 func TestSessionQueryColdReadDoesNotAcquireWriter(t *testing.T) {
371 root := filepath.Join(t.TempDir(), "sessions-v4")
372 persistence := NewFilesystemPersistence(root)
373 handle, err := persistence.Create(CreateOptions{SessionID: "cold"})
374 if err != nil {
375 t.Fatal(err)
376 }
377 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "m", Role: provider.RoleUser, Content: "hello"}})
378 if _, err := handle.Append(t.Context(), Batch{OperationID: "message", Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
379 t.Fatal(err)
380 }
381 if err := handle.Close(t.Context()); err != nil {
382 t.Fatal(err)
383 }
384 service, err := NewService("local", persistence)
385 if err != nil {
386 t.Fatal(err)
387 }
388 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
389 ref := SessionRef{HostID: "local", SessionID: "cold"}
390 history, err := service.Query().History(t.Context(), ref)
391 if err != nil || len(history) != 1 || history[0].ID != "m" {
392 t.Fatalf("cold history = %#v, %v", history, err)
393 }
394 // A query did not retain the writer lease; the execution runtime can attach.
395 binding, err := service.Open(t.Context(), ref)
396 if err != nil {
397 t.Fatal(err)
398 }
399 if err := binding.Release(t.Context()); err != nil {
400 t.Fatal(err)
401 }
402 if err := service.Close(t.Context(), ref); err != nil {
403 t.Fatal(err)
404 }
405 }
406
407 func TestServiceContinueLegacyPublishesNewIdentityAndLeavesSourceUnchanged(t *testing.T) {
408 root := t.TempDir()
409 legacy := filepath.Join(root, "sessions", "old.jsonl")
410 if err := os.MkdirAll(filepath.Dir(legacy), 0o700); err != nil {
411 t.Fatal(err)
412 }
413 source := agent.NewSession("system")
414 source.Add(provider.Message{ID: "user", Role: provider.RoleUser, Content: "hello"})
415 if err := source.Save(legacy); err != nil {
416 t.Fatal(err)
417 }
418 before, err := os.ReadFile(legacy)
419 if err != nil {
420 t.Fatal(err)
421 }
422 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(root, "sessions-v4")))
423 if err != nil {
424 t.Fatal(err)
425 }
426 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
427 runtime, migration, err := service.ContinueLegacy(t.Context(), legacy, "")
428 if err != nil {
429 t.Fatal(err)
430 }
431 if runtime.Ref().SessionID != migration.TargetID || runtime.Ref().SessionID == "old" {
432 t.Fatalf("runtime=%+v migration=%+v", runtime.Ref(), migration)
433 }
434 after, err := os.ReadFile(legacy)
435 if err != nil || string(after) != string(before) {
436 t.Fatal("legacy source changed during continue")
437 }
438 if got := runtime.Session().Snapshot().Projection.ModelMessages; len(got) != 2 || got[1].ID != "user" {
439 t.Fatalf("migrated history = %#v", got)
440 }
441 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
442 t.Fatal(err)
443 }
444 }
445
446 func TestServiceExportAndDeleteAreSessionDirectoryAtomic(t *testing.T) {
447 root := filepath.Join(t.TempDir(), "sessions-v4")
448 service, err := NewService("local", NewFilesystemPersistence(root))
449 if err != nil {
450 t.Fatal(err)
451 }
452 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
453 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "managed"})
454 if err != nil {
455 t.Fatal(err)
456 }
457 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "m", Role: provider.RoleUser, Content: strings.Repeat("exported", 20_000)}})
458 if _, err := runtime.Session().AppendBatch(t.Context(), "message", []Event{{Kind: "message/complete", Payload: payload}}); err != nil {
459 t.Fatal(err)
460 }
461 exported := filepath.Join(t.TempDir(), "exported-session")
462 if err := service.Export(t.Context(), runtime.Ref(), exported); err != nil {
463 t.Fatal(err)
464 }
465 if _, err := Replay(exported, nil); err != nil {
466 t.Fatalf("export is not self-contained: %v", err)
467 }
468 if _, err := os.Stat(filepath.Join(exported, "writer.lock")); !os.IsNotExist(err) {
469 t.Fatalf("export copied writer ownership artifact: %v", err)
470 }
471 ref := runtime.Ref()
472 if err := service.Delete(t.Context(), ref); err != nil {
473 t.Fatal(err)
474 }
475 if _, err := os.Stat(filepath.Join(root, ref.SessionID)); !os.IsNotExist(err) {
476 t.Fatalf("deleted session remains visible: %v", err)
477 }
478 if _, err := service.Query().Snapshot(t.Context(), ref); !errors.Is(err, os.ErrNotExist) && !errors.Is(err, ErrSessionNotFound) {
479 t.Fatalf("deleted session was revived by query cache: %v", err)
480 }
481 if err := os.RemoveAll(root); err != nil {
482 t.Fatal(err)
483 }
484 if _, err := Replay(exported, nil); err != nil {
485 t.Fatalf("export depended on the source content store: %v", err)
486 }
487 }
488
489 func TestServiceImportValidatesSelfContainedContentAndPublishesAtomically(t *testing.T) {
490 sourceRoot := filepath.Join(t.TempDir(), "source", "sessions-v4")
491 source, err := NewService("source", NewFilesystemPersistence(sourceRoot))
492 if err != nil {
493 t.Fatal(err)
494 }
495 t.Cleanup(func() { _ = source.Shutdown(context.Background()) })
496 runtime, err := source.Create(t.Context(), CreateOptions{SessionID: "portable"})
497 if err != nil {
498 t.Fatal(err)
499 }
500 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "large", Role: provider.RoleUser, Content: strings.Repeat("portable", 20_000)}})
501 if _, err := runtime.Session().AppendBatch(t.Context(), "large", []Event{{Kind: "message/complete", Payload: payload}}); err != nil {
502 t.Fatal(err)
503 }
504 bundle := filepath.Join(t.TempDir(), "bundle")
505 if err := source.Export(t.Context(), runtime.Ref(), bundle); err != nil {
506 t.Fatal(err)
507 }
508 if err := source.Close(t.Context(), runtime.Ref()); err != nil {
509 t.Fatal(err)
510 }
511 if err := os.RemoveAll(filepath.Dir(sourceRoot)); err != nil {
512 t.Fatal(err)
513 }
514
515 targetRoot := filepath.Join(t.TempDir(), "target", "sessions-v4")
516 target, err := NewService("target", NewFilesystemPersistence(targetRoot))
517 if err != nil {
518 t.Fatal(err)
519 }
520 t.Cleanup(func() { _ = target.Shutdown(context.Background()) })
521 ref, err := target.Import(t.Context(), bundle)
522 if err != nil {
523 t.Fatal(err)
524 }
525 page := historyPageReady(t, target.Query(), ref, "", 10)
526 if len(page.Messages) != 1 || page.Messages[0].ContentRef == nil {
527 t.Fatalf("imported history = %+v", page)
528 }
529 if _, err := target.Import(t.Context(), bundle); !errors.Is(err, ErrSessionExists) {
530 t.Fatalf("duplicate import = %v", err)
531 }
532 entries, err := os.ReadDir(targetRoot)
533 if err != nil {
534 t.Fatal(err)
535 }
536 for _, entry := range entries {
537 if strings.HasPrefix(entry.Name(), ".portable.import-") {
538 t.Fatalf("failed import left staging directory %q", entry.Name())
539 }
540 }
541 headerTarget, err := NewService("desktop", NewFilesystemPersistence(filepath.Join(t.TempDir(), "desktop-sessions-v5", "by-id")))
542 if err != nil {
543 t.Fatal(err)
544 }
545 t.Cleanup(func() { _ = headerTarget.Shutdown(context.Background()) })
546 headerRef, err := headerTarget.ImportWithHeader(t.Context(), bundle, CreateOptions{SessionID: "portable", CWD: "/workspace", Origin: SessionOriginCanonicalImport})
547 if err != nil {
548 t.Fatal(err)
549 }
550 headerInfo, err := headerTarget.persistence.Stat(t.Context(), headerRef.SessionID)
551 if err != nil {
552 t.Fatal(err)
553 }
554 // Headers record CWD in OS-native form; clean the expectation the same way.
555 if headerInfo.CWD != filepath.Clean("/workspace") || headerInfo.Origin != SessionOriginCanonicalImport {
556 t.Fatalf("imported header = %+v", headerInfo)
557 }
558 remapped, err := headerTarget.ImportWithHeader(t.Context(), bundle, CreateOptions{SessionID: "migr-conflict", CWD: "/workspace", Origin: SessionOriginCanonicalImport})
559 if err != nil {
560 t.Fatal(err)
561 }
562 if remapped.SessionID != "migr-conflict" {
563 t.Fatalf("remapped import = %+v", remapped)
564 }
565 remappedPage := historyPageReady(t, headerTarget.Query(), remapped, "", 10)
566 if len(remappedPage.Messages) != 1 || remappedPage.Messages[0].ContentRef == nil {
567 t.Fatalf("remapped history = %+v", remappedPage)
568 }
569 }
570
571 func TestDeleteRefusesAnOwnedSession(t *testing.T) {
572 root := filepath.Join(t.TempDir(), "sessions-v4")
573 persistence := NewFilesystemPersistence(root)
574 handle, err := persistence.Create(CreateOptions{SessionID: "owned"})
575 if err != nil {
576 t.Fatal(err)
577 }
578 ctx, cancel := context.WithCancel(t.Context())
579 cancel()
580 if err := persistence.Delete(ctx, "owned"); !errors.Is(err, context.Canceled) {
581 t.Fatalf("delete owned session error = %v", err)
582 }
583 if err := handle.Close(t.Context()); err != nil {
584 t.Fatal(err)
585 }
586 }
587
588 func TestCancelledActivityCannotCommitLateBusinessResult(t *testing.T) {
589 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
590 if err != nil {
591 t.Fatal(err)
592 }
593 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
594 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "fenced"})
595 if err != nil {
596 t.Fatal(err)
597 }
598 ctx, exec := bindTestExecution(t, runtime, "tool")
599 if receipt, err := service.CancelSession(runtime.Ref()); err != nil || !receipt.Accepted || receipt.Phase != RuntimeCancelling {
600 t.Fatalf("cancel receipt = %+v, %v", receipt, err)
601 }
602 if !errors.Is(ctx.Err(), context.Canceled) {
603 t.Fatalf("cancel context = %v", ctx.Err())
604 }
605 if _, err := runtime.Session().Append(t.Context(), Batch{
606 OperationID: "turn-end",
607 TurnID: "turn-1",
608 Events: []Event{{Kind: "turn/end", Payload: json.RawMessage(`{"status":"interrupted"}`)}},
609 }); err != nil {
610 t.Fatalf("terminal truth during cancel: %v", err)
611 }
612 exec.Finish()
613 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
614 t.Fatal(err)
615 }
616 }
617
618 func TestRecoveryOwnerCanCommitOnlyTerminalRecoveryFacts(t *testing.T) {
619 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
620 if err != nil {
621 t.Fatal(err)
622 }
623 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
624 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "recovery-terminal"})
625 if err != nil {
626 t.Fatal(err)
627 }
628 bindTestExecution(t, runtime, "tool")
629 runtime.RequireRecovery("tool")
630 if err := service.Close(t.Context(), runtime.Ref()); !errors.Is(err, ErrRuntimeBusy) {
631 t.Fatalf("close during recovery = %v", err)
632 }
633 terminal := Batch{OperationID: "recovery-terminal", TurnID: "turn-1", Events: []Event{
634 {Kind: "runtime/recovery", Payload: json.RawMessage(`{"state":"recovery_required","phase":"tool","reason":"cancellation_grace_expired","requires_user_decision":true}`)},
635 {Kind: "turn/end", Payload: json.RawMessage(`{"status":"recovery_required"}`)},
636 }}
637 if _, err := runtime.Session().Append(t.Context(), terminal); err != nil {
638 t.Fatalf("record recovery: %v", err)
639 }
640 projection := runtime.Session().Snapshot().Projection
641 if projection.Recovery == nil || projection.Recovery.State != "recovery_required" || projection.TurnStatus != "recovery_required" {
642 t.Fatalf("recovery projection = %#v, turn status = %q", projection.Recovery, projection.TurnStatus)
643 }
644 runtime.mu.Lock()
645 runtime.phase = RuntimeIdle
646 runtime.activity = ""
647 runtime.mu.Unlock()
648 if err := service.Close(t.Context(), runtime.Ref()); err != nil {
649 t.Fatal(err)
650 }
651 }
652
653 func TestCancelSessionIsIdempotentWithoutAttachedRuntime(t *testing.T) {
654 service, err := NewService("local", NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
655 if err != nil {
656 t.Fatal(err)
657 }
658 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
659 ref := SessionRef{HostID: "local", SessionID: "not-attached"}
660 receipt, err := service.CancelSession(ref)
661 if err != nil || !receipt.Accepted || receipt.Phase != RuntimeIdle {
662 t.Fatalf("idle cancel = %+v, %v", receipt, err)
663 }
664 }
665
665 lines GO