返回 DeepSeek-Reasonix
store_test.go
根目录 / internal / session / store_test.go
1 package session
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "os"
10 "path/filepath"
11 "testing"
12 "time"
13
14 "reasonix/internal/agent"
15 "reasonix/internal/event"
16 "reasonix/internal/filelock"
17 "reasonix/internal/provider"
18 "reasonix/internal/store"
19 )
20
21 func TestStoredManifestRequiresExactFormatBoundary(t *testing.T) {
22 t.Parallel()
23 if !supportedStoredManifest(Manifest{SchemaVersion: SchemaVersion, Codec: Codec}) {
24 t.Fatal("in-place v4 session was rejected")
25 }
26 if !supportedStoredManifest(Manifest{SchemaVersion: SchemaVersion, Codec: Codec, StorageRevision: StorageRevision}) {
27 t.Fatal("final v4 manifest was rejected")
28 }
29 if !supportedStoredManifest(Manifest{SchemaVersion: 3, Codec: FinalV31Codec}) {
30 t.Fatal("frozen v3.1 input format was rejected")
31 }
32 if supportedStoredManifest(Manifest{SchemaVersion: SchemaVersion, Codec: FinalV31Codec}) {
33 t.Fatal("schema-4 data mislabeled as v3.1 was accepted")
34 }
35 }
36
37 func TestLegacyLinearStoreOpensAndAppendsWithoutCodecUpgrade(t *testing.T) {
38 for _, codec := range []string{LegacyLinearCodec, FinalV31Codec, PrototypeCodec} {
39 t.Run(codec, func(t *testing.T) { testNativeLinearStore(t, codec) })
40 }
41 }
42
43 func testNativeLinearStore(t *testing.T, codec string) {
44 root := t.TempDir()
45 dir := filepath.Join(root, "legacy-linear")
46 if err := os.MkdirAll(dir, 0o700); err != nil {
47 t.Fatal(err)
48 }
49 manifest := Manifest{SchemaVersion: 3, Codec: codec, SessionID: "legacy-linear", CreatedAt: time.Now().UTC(), WriterGeneration: 1}
50 if err := writeManifestFile(filepath.Join(dir, "manifest.json"), manifest); err != nil {
51 t.Fatal(err)
52 }
53 if err := os.WriteFile(filepath.Join(dir, legacyLogName), nil, 0o600); err != nil {
54 t.Fatal(err)
55 }
56 s, err := OpenWithOptions(dir, "legacy-linear", OpenOptions{ExternalHistory: true})
57 if err != nil {
58 t.Fatal(err)
59 }
60 t.Cleanup(func() { _ = s.Close(context.Background()) })
61 _, err = s.Append(t.Context(), Batch{OperationID: "continue-old", Events: []Event{{Kind: "turn/start"}}})
62 if err != nil {
63 t.Fatal(err)
64 }
65 if _, err := s.Flush(t.Context()); err != nil {
66 t.Fatal(err)
67 }
68 if err := s.Close(context.Background()); err != nil {
69 t.Fatal(err)
70 }
71 got, err := readStoredManifest(filepath.Join(dir, "manifest.json"))
72 if err != nil || got.Codec != codec || got.SchemaVersion != 3 {
73 t.Fatalf("manifest = %+v, err=%v", got, err)
74 }
75 commits, err := Replay(dir, nil)
76 if err != nil || len(commits) != 1 || commits[0].Codec != codec {
77 t.Fatalf("commits = %+v, err=%v", commits, err)
78 }
79 // Preserve the byte boundary of records produced by a different JSON
80 // encoder. Re-marshalling a record to derive its end would truncate this.
81 logPath := filepath.Join(dir, legacyLogName)
82 original, err := os.ReadFile(logPath)
83 if err != nil {
84 t.Fatal(err)
85 }
86 original = append([]byte(" "), original...)
87 if err := os.WriteFile(logPath, original, 0o600); err != nil {
88 t.Fatal(err)
89 }
90 s, err = OpenWithOptions(dir, "legacy-linear", OpenOptions{ExternalHistory: true})
91 if err != nil {
92 t.Fatal(err)
93 }
94 if got, err := os.ReadFile(logPath); err != nil || !bytes.Equal(got, original) {
95 t.Fatalf("opening changed original log bytes: %v", err)
96 }
97 if _, err := s.Append(t.Context(), Batch{OperationID: "finish-old", Events: []Event{{Kind: "turn/end", Payload: json.RawMessage(`{"status":"completed"}`)}}}); err != nil {
98 t.Fatal(err)
99 }
100 if err := s.Close(context.Background()); err != nil {
101 t.Fatal(err)
102 }
103 commits, err = Replay(dir, nil)
104 if err != nil || len(commits) != 2 || commits[1].FirstSequence != 2 {
105 t.Fatalf("reopened commits = %+v, err=%v", commits, err)
106 }
107 complete, err := os.ReadFile(logPath)
108 if err != nil {
109 t.Fatal(err)
110 }
111 if err := os.WriteFile(logPath, append(append([]byte(nil), complete...), []byte(`{"partial":`)...), 0o600); err != nil {
112 t.Fatal(err)
113 }
114 s, err = OpenWithOptions(dir, "legacy-linear", OpenOptions{ExternalHistory: true})
115 if err != nil {
116 t.Fatal(err)
117 }
118 if got, err := os.ReadFile(logPath); err != nil || !bytes.Equal(got, complete) {
119 t.Fatalf("torn repair changed complete prefix: %v", err)
120 }
121 backups, err := filepath.Glob(filepath.Join(dir, "events.torn-*.tail"))
122 if err != nil || len(backups) != 1 {
123 t.Fatalf("torn evidence missing: %v %v", backups, err)
124 }
125 if got, err := os.ReadFile(backups[0]); err != nil || string(got) != `{"partial":` {
126 t.Fatalf("torn evidence changed: %q %v", got, err)
127 }
128 }
129
130 func TestNativePrototypeHistoryRewriteRoundTrip(t *testing.T) {
131 dir := filepath.Join(t.TempDir(), "prototype")
132 payload, _ := json.Marshal(map[string]any{"messages": []provider.Message{{ID: "old", Role: provider.RoleUser, Content: "old"}}})
133 writePrototypeStore(t, dir, []Event{{Kind: "context/replace", Payload: payload}}, "")
134 s, err := Open(dir, "prototype")
135 if err != nil {
136 t.Fatal(err)
137 }
138 t.Cleanup(func() { _ = s.Close(context.Background()) })
139 if got := s.Snapshot().Projection.ModelMessages; len(got) != 1 || got[0].Content != "old" {
140 t.Fatalf("old projection: %+v", got)
141 }
142 payload, _ = json.Marshal(map[string]any{"messages": []provider.Message{{ID: "new", Role: provider.RoleUser, Content: "new"}}})
143 if _, err := s.Append(t.Context(), Batch{OperationID: "rewrite", Events: []Event{{Kind: "history/replace", Payload: payload}}}); err != nil {
144 t.Fatal(err)
145 }
146 if err := s.Close(context.Background()); err != nil {
147 t.Fatal(err)
148 }
149 raw, err := os.ReadFile(filepath.Join(dir, legacyLogName))
150 if err != nil || bytes.Contains(raw, []byte(`"kind":"history/replace"`)) {
151 t.Fatalf("prototype codec changed: %s %v", raw, err)
152 }
153 s, err = Open(dir, "prototype")
154 if err != nil {
155 t.Fatal(err)
156 }
157 if got := s.Snapshot().Projection.ModelMessages; len(got) != 1 || got[0].Content != "new" {
158 t.Fatalf("reopened projection: %+v", got)
159 }
160 }
161
162 func TestCommitBatchAndProjectionAreAtomic(t *testing.T) {
163 dir := filepath.Join(t.TempDir(), "s")
164 s, err := Open(dir, "s")
165 if err != nil {
166 t.Fatal(err)
167 }
168 t.Cleanup(func() { _ = s.Close(context.Background()) })
169 todos, _ := json.Marshal(map[string]any{"todos": []event.Todo{{Content: "A", Status: "in_progress"}, {Content: "B", Status: "in_progress"}}})
170 commit, err := s.Append(t.Context(), Batch{OperationID: "todo-1", TurnID: "turn-1", Events: []Event{
171 {Kind: "turn/start"},
172 {Kind: "todo/write", Payload: todos},
173 }})
174 if err != nil {
175 t.Fatal(err)
176 }
177 if commit.FirstSequence != 1 || commit.EventCount != 2 || commit.Events[1].Sequence != 2 {
178 t.Fatalf("commit boundary = %+v", commit)
179 }
180 projection := s.Snapshot().Projection
181 if !projection.TodoWritten || len(projection.Todos) != 2 || projection.Todos[1].Content != "B" {
182 t.Fatalf("live todo projection = %+v", projection)
183 }
184 if _, err := s.Flush(t.Context()); err != nil {
185 t.Fatal(err)
186 }
187 commits, err := Replay(dir, nil)
188 if err != nil {
189 t.Fatal(err)
190 }
191 replayed, err := Project(commits)
192 if err != nil {
193 t.Fatal(err)
194 }
195 if !replayed.TodoWritten || len(replayed.Todos) != 2 || replayed.Todos[1].Content != "B" {
196 t.Fatalf("replayed todo projection = %+v", replayed)
197 }
198 }
199
200 func TestProjectionRebuildsMessagesGoalPlanInteractionsAndRecovery(t *testing.T) {
201 dir := filepath.Join(t.TempDir(), "s")
202 s, err := Open(dir, "s")
203 if err != nil {
204 t.Fatal(err)
205 }
206 t.Cleanup(func() { _ = s.Close(context.Background()) })
207 message, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "m1", Role: provider.RoleUser, Content: "hello"}})
208 interaction, _ := json.Marshal(map[string]any{"id": "p1", "state": "pending"})
209 recovery, _ := json.Marshal(event.RecoveryStatus{State: "recovery_required", RequiresUserDecision: true, Reason: "worker did not exit"})
210 if _, err := s.Append(t.Context(), Batch{OperationID: "state-1", TurnID: "t1", Events: []Event{
211 {Kind: "turn/start"},
212 {Kind: "message/complete", Payload: message},
213 {Kind: "interaction/created", Payload: interaction},
214 {Kind: "plan/state", Payload: json.RawMessage(`{"status":"approved"}`)},
215 {Kind: "goal/state", Payload: json.RawMessage(`{"objective":"ship","status":"active"}`)},
216 {Kind: "runtime/recovery", Payload: recovery},
217 }}); err != nil {
218 t.Fatal(err)
219 }
220 projection := s.Snapshot().Projection
221 if len(projection.Messages) != 1 || projection.Messages[0].ID != "m1" || projection.Interactions["p1"] != "pending" {
222 t.Fatalf("message/interaction projection = %+v", projection)
223 }
224 if projection.Recovery == nil || !projection.Recovery.RequiresUserDecision || len(projection.PlanState) == 0 || len(projection.GoalState) == 0 {
225 t.Fatalf("state projection = %+v", projection)
226 }
227 }
228
229 func TestModelContextReplaceDoesNotRewriteUIHistory(t *testing.T) {
230 dir := filepath.Join(t.TempDir(), "replace")
231 store, err := Open(dir, "replace")
232 if err != nil {
233 t.Fatal(err)
234 }
235 t.Cleanup(func() { _ = store.Close(context.Background()) })
236 messages := []provider.Message{{ID: "u1", Role: provider.RoleUser, Content: "new context"}}
237 payload, _ := json.Marshal(map[string]any{"messages": messages, "reason": "compaction", "sourceSequences": []uint64{1}})
238 if _, err := store.Append(context.Background(), Batch{OperationID: "replace-1", Events: []Event{{Kind: "model/context-replace", Payload: payload}}}); err != nil {
239 t.Fatal(err)
240 }
241 got := store.Snapshot().Projection.Messages
242 if len(got) != 0 {
243 t.Fatalf("messages = %#v", got)
244 }
245 derived := store.DeriveMessages()
246 if len(derived) != 1 || derived[0].ID != "u1" {
247 t.Fatalf("model messages = %#v", derived)
248 }
249 }
250
251 func TestMessageUpsertPersistsLocalMetadataWithoutEnteringModelHistory(t *testing.T) {
252 store, err := Open(filepath.Join(t.TempDir(), "upsert"), "upsert")
253 if err != nil {
254 t.Fatal(err)
255 }
256 t.Cleanup(func() { _ = store.Close(context.Background()) })
257 local := provider.Message{ID: "receipt", Role: provider.RoleTool, Name: provider.LocalOnlyToolName, ToolCallID: provider.LocalOnlyToolID, LocalOnly: true, ProtocolRecovery: json.RawMessage(`{"version":1,"id":"r","state":"pending"}`)}
258 payload, _ := json.Marshal(map[string]any{"message": local})
259 if _, err := store.Append(t.Context(), Batch{OperationID: "upsert-local", Events: []Event{{Kind: "message/upsert", Payload: payload}}}); err != nil {
260 t.Fatal(err)
261 }
262 snapshot := store.Snapshot().Projection
263 if len(snapshot.Messages) != 1 || snapshot.Messages[0].ID != "receipt" {
264 t.Fatalf("ui messages = %#v", snapshot.Messages)
265 }
266 if len(snapshot.ModelMessages) != 0 {
267 t.Fatalf("local receipt leaked into model history: %#v", snapshot.ModelMessages)
268 }
269 local.ProtocolRecovery = json.RawMessage(`{"version":1,"id":"r","state":"consumed"}`)
270 payload, _ = json.Marshal(map[string]any{"message": local})
271 if _, err := store.Append(t.Context(), Batch{OperationID: "upsert-local-again", Events: []Event{{Kind: "message/upsert", Payload: payload}}}); err != nil {
272 t.Fatal(err)
273 }
274 if got := store.Snapshot().Projection; len(got.Messages) != 1 || len(got.ModelMessages) != 0 {
275 t.Fatalf("updated projections = %#v / %#v", got.Messages, got.ModelMessages)
276 }
277 }
278
279 func TestMessageCompleteRejectsDuplicateStableIDAtomically(t *testing.T) {
280 store, err := Open(filepath.Join(t.TempDir(), "duplicates"), "duplicates")
281 if err != nil {
282 t.Fatal(err)
283 }
284 t.Cleanup(func() { _ = store.Close(context.Background()) })
285 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "same", Role: provider.RoleUser, Content: "one"}})
286 if _, err := store.Append(t.Context(), Batch{OperationID: "one", Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
287 t.Fatal(err)
288 }
289 before := store.Snapshot()
290 payload, _ = json.Marshal(map[string]any{"message": provider.Message{ID: "same", Role: provider.RoleAssistant, Content: "two"}})
291 if _, err := store.Append(t.Context(), Batch{OperationID: "two", Events: []Event{{Kind: "message/complete", Payload: payload}}}); err == nil {
292 t.Fatal("duplicate stable id accepted")
293 }
294 after := store.Snapshot()
295 if after.EventSequence != before.EventSequence || len(after.Projection.Messages) != 1 {
296 t.Fatalf("failed batch mutated state: before=%+v after=%+v", before, after)
297 }
298 }
299
300 func TestCompactionReplacesOnlyModelProjection(t *testing.T) {
301 dir := filepath.Join(t.TempDir(), "compaction")
302 store, err := Open(dir, "compaction")
303 if err != nil {
304 t.Fatal(err)
305 }
306 t.Cleanup(func() { _ = store.Close(context.Background()) })
307 original := provider.Message{ID: "u1", Role: provider.RoleUser, Content: "original"}
308 messagePayload, _ := json.Marshal(map[string]any{"message": original})
309 projected := []provider.Message{{ID: "summary", Role: provider.RoleUser, Content: "summary"}}
310 compactionPayload, _ := json.Marshal(map[string]any{"messages": projected, "trigger": "manual"})
311 if _, err := store.Append(context.Background(), Batch{OperationID: "compact", Events: []Event{
312 {Kind: "message/complete", Payload: messagePayload},
313 {Kind: "compaction", Payload: compactionPayload},
314 }}); err != nil {
315 t.Fatal(err)
316 }
317 snapshot := store.Snapshot().Projection
318 if len(snapshot.Messages) != 1 || snapshot.Messages[0].ID != "u1" {
319 t.Fatalf("canonical messages changed: %#v", snapshot.Messages)
320 }
321 if len(snapshot.ModelMessages) != 1 || snapshot.ModelMessages[0].ID != "summary" {
322 t.Fatalf("model projection = %#v", snapshot.ModelMessages)
323 }
324 }
325
326 func TestRestartRecoveryClosesTurnToolsAndInteractionsWithoutReplay(t *testing.T) {
327 dir := filepath.Join(t.TempDir(), "recover")
328 first, err := Open(dir, "recover")
329 if err != nil {
330 t.Fatal(err)
331 }
332 call := json.RawMessage(`{"id":"call-1","name":"bash"}`)
333 interaction := json.RawMessage(`{"id":"approval-1","state":"pending"}`)
334 if _, err := first.Append(t.Context(), Batch{OperationID: "open-turn", TurnID: "turn-1", Events: []Event{
335 {Kind: "turn/start"}, {Kind: "tool/start", Payload: call}, {Kind: "interaction/created", Payload: interaction},
336 }}); err != nil {
337 t.Fatal(err)
338 }
339 if err := first.Close(t.Context()); err != nil {
340 t.Fatal(err)
341 }
342
343 second, err := Open(dir, "recover")
344 if err != nil {
345 t.Fatal(err)
346 }
347 defer second.Close(context.Background())
348 if _, recovered, err := second.RecoverInterrupted(t.Context()); err != nil || !recovered {
349 t.Fatalf("RecoverInterrupted = recovered %v, err %v", recovered, err)
350 }
351 projection := second.Snapshot().Projection
352 if projection.TurnID != "" || projection.TurnStatus != event.TurnInterrupted || len(projection.ActiveTools) != 0 || len(projection.Interactions) != 0 {
353 t.Fatalf("recovered projection = %+v", projection)
354 }
355 }
356
357 func TestAppendValidationDoesNotChangeLiveState(t *testing.T) {
358 dir := filepath.Join(t.TempDir(), "s")
359 s, err := Open(dir, "s")
360 if err != nil {
361 t.Fatal(err)
362 }
363 t.Cleanup(func() { _ = s.Close(context.Background()) })
364 tests := []Event{
365 {Kind: "todo/write", Payload: json.RawMessage(`{"todos":[{"content":" A ","status":"pending"}]}`)},
366 {Kind: "interaction/created", Payload: json.RawMessage(`{"id":"p","state":"answered"}`)},
367 {Kind: "interaction/resolved", Payload: json.RawMessage(`{"id":"p","state":"pending"}`)},
368 {Kind: "message/complete", Payload: json.RawMessage(`{"message":{"role":"user","content":"missing id"}}`)},
369 {Kind: "plan/state", Payload: json.RawMessage(`[]`)},
370 }
371 for i, invalid := range tests {
372 before := s.Snapshot()
373 _, err := s.Append(t.Context(), Batch{OperationID: "invalid-" + string(rune('a'+i)), Events: []Event{invalid}})
374 if !errors.Is(err, ErrDamagedStore) {
375 t.Fatalf("event %s error = %v", invalid.Kind, err)
376 }
377 after := s.Snapshot()
378 if after.EventSequence != before.EventSequence {
379 t.Fatalf("invalid %s advanced sequence: before=%d after=%d", invalid.Kind, before.EventSequence, after.EventSequence)
380 }
381 }
382 }
383
384 func TestSessionConfigIsProjectedFromTheEventStream(t *testing.T) {
385 dir := filepath.Join(t.TempDir(), "s")
386 s, err := Open(dir, "s")
387 if err != nil {
388 t.Fatal(err)
389 }
390 if _, err := s.Append(t.Context(), Batch{OperationID: "config", Events: []Event{{Kind: "session/config", Payload: json.RawMessage(`{"modelRef":"provider/model","modelIdentity":"revision-1"}`)}}}); err != nil {
391 t.Fatal(err)
392 }
393 if projection := s.Snapshot().Projection; projection.ModelRef != "provider/model" || projection.ModelIdentity != "revision-1" {
394 t.Fatalf("live model selection = %q / %q", projection.ModelRef, projection.ModelIdentity)
395 }
396 if err := s.Close(t.Context()); err != nil {
397 t.Fatal(err)
398 }
399 commits, err := Replay(dir, nil)
400 if err != nil {
401 t.Fatal(err)
402 }
403 projection, err := Project(commits)
404 if err != nil {
405 t.Fatal(err)
406 }
407 if projection.ModelRef != "provider/model" || projection.ModelIdentity != "revision-1" {
408 t.Fatalf("replayed model selection = %q / %q", projection.ModelRef, projection.ModelIdentity)
409 }
410 }
411
412 func TestReplayIgnoresTornTailAndWriterRecoversUnderLease(t *testing.T) {
413 dir := filepath.Join(t.TempDir(), "s")
414 s, err := Open(dir, "s")
415 if err != nil {
416 t.Fatal(err)
417 }
418 if _, err := s.Append(t.Context(), Batch{OperationID: "turn", Events: []Event{{Kind: "turn/start"}}}); err != nil {
419 t.Fatal(err)
420 }
421 if err := s.Close(t.Context()); err != nil {
422 t.Fatal(err)
423 }
424 path := filepath.Join(dir, currentLogName)
425 file, err := os.OpenFile(path, os.O_APPEND|os.O_WRONLY, 0o600)
426 if err != nil {
427 t.Fatal(err)
428 }
429 _, _ = file.WriteString(`{"torn`)
430 _ = file.Close()
431 commits, err := Replay(dir, nil)
432 if err != nil || len(commits) != 1 {
433 t.Fatalf("Replay partial tail = %d, %v", len(commits), err)
434 }
435 reopened, err := Open(dir, "s")
436 if err != nil {
437 t.Fatalf("writer recovery: %v", err)
438 }
439 backups, err := filepath.Glob(filepath.Join(dir, "events.torn-*.tail"))
440 if err != nil || len(backups) != 1 {
441 t.Fatalf("torn tail backups = %v, %v", backups, err)
442 }
443 tail, err := os.ReadFile(backups[0])
444 if err != nil || string(tail) != `{"torn` {
445 t.Fatalf("preserved tail = %q, %v", tail, err)
446 }
447 if _, err := reopened.Append(t.Context(), Batch{OperationID: "after-recovery", Events: []Event{{Kind: "diagnostic", Optional: true}}}); err != nil {
448 t.Fatal(err)
449 }
450 if err := reopened.Close(t.Context()); err != nil {
451 t.Fatal(err)
452 }
453 commits, err = Replay(dir, nil)
454 if err != nil || len(commits) != 2 {
455 t.Fatalf("Replay recovered log = %d, %v", len(commits), err)
456 }
457 }
458
459 func TestReplayRejectsUnknownRequiredEventAndAllowsOptionalDiagnostic(t *testing.T) {
460 dir := filepath.Join(t.TempDir(), "s")
461 s, err := Open(dir, "s")
462 if err != nil {
463 t.Fatal(err)
464 }
465 if _, err := s.Append(t.Context(), Batch{OperationID: "unknown", Events: []Event{{Kind: "future/required"}}}); !errors.Is(err, ErrUnsupportedVersion) {
466 t.Fatalf("unknown required append error = %v", err)
467 }
468 if _, err := s.Append(t.Context(), Batch{OperationID: "diagnostic", Events: []Event{{Kind: "future/diagnostic", Optional: true}}}); err != nil {
469 t.Fatal(err)
470 }
471 if err := s.Close(t.Context()); err != nil {
472 t.Fatal(err)
473 }
474 if commits, err := Replay(dir, nil); err != nil || len(commits) != 1 {
475 t.Fatalf("optional diagnostic replay = %d, %v", len(commits), err)
476 }
477 }
478
479 func TestOpenUsesSingleWriterLeaseAndAdvancesGeneration(t *testing.T) {
480 dir := filepath.Join(t.TempDir(), "s")
481 first, err := Open(dir, "s")
482 if err != nil {
483 t.Fatal(err)
484 }
485 if _, err := Open(dir, "s"); !errors.Is(err, ErrWriterOwned) {
486 t.Fatalf("second writer error = %v", err)
487 }
488 firstGeneration := first.manifest.WriterGeneration
489 if err := first.Close(t.Context()); err != nil {
490 t.Fatal(err)
491 }
492 second, err := Open(dir, "s")
493 if err != nil {
494 t.Fatal(err)
495 }
496 defer second.Close(context.Background())
497 if second.manifest.WriterGeneration != firstGeneration+1 {
498 t.Fatalf("writer generation = %d, want %d", second.manifest.WriterGeneration, firstGeneration+1)
499 }
500 }
501
502 func TestMigrateLegacyIsIdempotentAndUsesFrozenArtifacts(t *testing.T) {
503 root := t.TempDir()
504 legacyDir := filepath.Join(root, "sessions")
505 path := filepath.Join(legacyDir, "old.jsonl")
506 if err := os.MkdirAll(legacyDir, 0o700); err != nil {
507 t.Fatal(err)
508 }
509 session := agent.NewSession("sys")
510 session.Add(provider.Message{Role: provider.RoleUser, Content: "hello", ID: "user-stable-id"})
511 if err := session.Save(path); err != nil {
512 t.Fatal(err)
513 }
514 if err := agent.SetBranchModelSelectionPreserveUpdated(path, "provider/model", "connection-revision"); err != nil {
515 t.Fatal(err)
516 }
517 goal := `{"objective":"finish","status":"paused","token_budget":123,"todos":[{"content":"old","status":"in_progress"}],"auto_continue":true}`
518 if err := os.WriteFile(store.SessionGoalState(path), []byte(goal), 0o600); err != nil {
519 t.Fatal(err)
520 }
521 v3root := filepath.Join(root, "sessions-v4")
522 first, err := MigrateLegacy(context.Background(), path, v3root)
523 if err != nil {
524 t.Fatal(err)
525 }
526 second, err := MigrateLegacy(context.Background(), path, v3root)
527 if err != nil {
528 t.Fatal(err)
529 }
530 if first.TargetID != second.TargetID || !second.Reused {
531 t.Fatalf("migration idempotency: first=%+v second=%+v", first, second)
532 }
533 commits, err := Replay(first.TargetDir, nil)
534 if err != nil || len(commits) == 0 {
535 t.Fatalf("Replay migrated = %d, %v", len(commits), err)
536 }
537 projection, err := Project(commits)
538 if err != nil {
539 t.Fatal(err)
540 }
541 if len(projection.Messages) != 2 || projection.Messages[1].ID != "user-stable-id" {
542 t.Fatalf("migrated messages = %+v", projection.Messages)
543 }
544 if projection.ModelRef != "provider/model" || projection.ModelIdentity != "connection-revision" {
545 t.Fatalf("migrated model selection = %q / %q", projection.ModelRef, projection.ModelIdentity)
546 }
547 if containsJSONKey(projection.GoalState, "todos") || containsJSONKey(projection.GoalState, "auto_continue") {
548 t.Fatalf("migrated goal projection = %s", projection.GoalState)
549 }
550 legacyGoal := filepath.Join(first.TargetDir, "legacy", filepath.Base(store.SessionGoalState(path)))
551 if string(mustRead(t, legacyGoal)) != goal {
552 t.Fatal("raw goal sidecar was not preserved byte-for-byte")
553 }
554 }
555
556 func TestMigrateLegacyBeyondFormer128MiBReplayLimit(t *testing.T) {
557 if os.Getenv("REASONIX_LARGE_SESSION_TEST") != "1" {
558 t.Skip("set REASONIX_LARGE_SESSION_TEST=1 to run the exact 134,308,416-byte regression")
559 }
560 const (
561 totalBytes = int64(134_308_416)
562 messages = 8_192
563 )
564 root := t.TempDir()
565 legacy := filepath.Join(root, "oversized.jsonl")
566 file, err := os.OpenFile(legacy, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600)
567 if err != nil {
568 t.Fatal(err)
569 }
570 prefixBytes := int64(len(fmt.Sprintf(`{"role":"user","id":"m-%08d","content":"`, 0)))
571 suffix := []byte("\"}\n")
572 contentBytes := totalBytes - int64(messages)*(prefixBytes+int64(len(suffix)))
573 if contentBytes <= 0 {
574 t.Fatal("invalid oversized fixture dimensions")
575 }
576 base, extra := contentBytes/int64(messages), contentBytes%int64(messages)
577 chunk := bytes.Repeat([]byte{'x'}, 32<<10)
578 for i := range messages {
579 prefix := fmt.Sprintf(`{"role":"user","id":"m-%08d","content":"`, i)
580 if _, err := file.WriteString(prefix); err != nil {
581 t.Fatal(err)
582 }
583 n := base
584 if int64(i) < extra {
585 n++
586 }
587 for n > 0 {
588 part := min(n, int64(len(chunk)))
589 if _, err := file.Write(chunk[:part]); err != nil {
590 t.Fatal(err)
591 }
592 n -= part
593 }
594 if _, err := file.Write(suffix); err != nil {
595 t.Fatal(err)
596 }
597 }
598 if err := file.Sync(); err != nil {
599 t.Fatal(err)
600 }
601 if err := file.Close(); err != nil {
602 t.Fatal(err)
603 }
604 info, err := os.Stat(legacy)
605 if err != nil || info.Size() != totalBytes {
606 t.Fatalf("legacy fixture size = %d, %v", info.Size(), err)
607 }
608 targetRoot := filepath.Join(root, "sessions-v4")
609 result, err := MigrateLegacy(t.Context(), legacy, targetRoot)
610 if err != nil {
611 t.Fatal(err)
612 }
613 if result.Source.Size != totalBytes || result.MessageNum != messages {
614 t.Fatalf("migration result = %+v", result)
615 }
616 service, err := NewService("capacity", NewFilesystemPersistence(targetRoot))
617 if err != nil {
618 t.Fatal(err)
619 }
620 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
621 binding, err := service.Open(t.Context(), SessionRef{HostID: "capacity", SessionID: result.TargetID})
622 if err != nil {
623 t.Fatal(err)
624 }
625 if _, err := binding.Runtime().Session().AppendBatch(t.Context(), "continued-after-oversized-migration", []Event{{Kind: "session/title", Payload: json.RawMessage(`{"title":"continued"}`)}}); err != nil {
626 t.Fatal(err)
627 }
628 receipt, err := binding.Runtime().Session().Flush(t.Context())
629 if err != nil {
630 t.Fatal(err)
631 }
632 if receipt.DurableSequence != messages+1 {
633 t.Fatalf("continued durable sequence = %d, want %d", receipt.DurableSequence, messages+1)
634 }
635 if err := binding.Release(t.Context()); err != nil {
636 t.Fatal(err)
637 }
638 }
639
640 func TestMigrateEmptyLegacyTranscriptProducesValidExplicitEmptyHistory(t *testing.T) {
641 root := t.TempDir()
642 legacy := filepath.Join(root, "empty.jsonl")
643 if err := os.WriteFile(legacy, nil, 0o600); err != nil {
644 t.Fatal(err)
645 }
646 result, err := MigrateLegacy(t.Context(), legacy, filepath.Join(root, "sessions-v4"))
647 if err != nil {
648 t.Fatal(err)
649 }
650 store, err := Open(result.TargetDir, result.TargetID)
651 if err != nil {
652 t.Fatal(err)
653 }
654 defer store.Close(context.Background())
655 projection := store.Snapshot().Projection
656 if len(projection.Messages) != 0 || len(projection.ModelMessages) != 0 {
657 t.Fatalf("empty migration projection = messages %#v model %#v", projection.Messages, projection.ModelMessages)
658 }
659 }
660
661 func TestMigrateLegacyRefusesActiveSourceLease(t *testing.T) {
662 dir := t.TempDir()
663 path := filepath.Join(dir, "old.jsonl")
664 if err := agent.NewSession("sys").Save(path); err != nil {
665 t.Fatal(err)
666 }
667 lease, err := agent.TryAcquireSessionLease(path)
668 if err != nil {
669 t.Fatal(err)
670 }
671 defer lease.Release()
672 if _, err := MigrateLegacy(t.Context(), path, filepath.Join(dir, "sessions-v4")); !errors.Is(err, agent.ErrSessionLeaseHeld) {
673 t.Fatalf("migration with active source lease error = %v", err)
674 }
675 }
676
677 func TestMigrationMapLeaseWaitIsContextCancellable(t *testing.T) {
678 path := filepath.Join(t.TempDir(), "sessions-v4", "migration-map.json")
679 if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
680 t.Fatal(err)
681 }
682 release, err := filelock.TryAcquire(path + ".lock")
683 if err != nil {
684 t.Fatal(err)
685 }
686 defer release()
687 ctx, cancel := context.WithCancel(t.Context())
688 cancel()
689 if _, err := acquireMigrationMapLease(ctx, path); !errors.Is(err, context.Canceled) {
690 t.Fatalf("acquireMigrationMapLease error = %v", err)
691 }
692 }
693
694 func TestMigrateLegacyHeadsBecomeIndependentLinearSessions(t *testing.T) {
695 dir := t.TempDir()
696 path := filepath.Join(dir, "old.jsonl")
697 session := agent.NewSession("sys")
698 session.Add(provider.Message{Role: provider.RoleUser, Content: "root question"})
699 session.Add(provider.Message{Role: provider.RoleAssistant, Content: "root answer"})
700 if err := session.Save(path); err != nil {
701 t.Fatal(err)
702 }
703 childHead, err := session.ForkHead(path, session.Snapshot()[2].ID, agent.HeadKindFork, "child")
704 if err != nil {
705 t.Fatal(err)
706 }
707 session.Add(provider.Message{Role: provider.RoleUser, Content: "child only"})
708 if err := session.Save(path); err != nil {
709 t.Fatal(err)
710 }
711 heads, err := agent.ListSessionHeads(path)
712 if err != nil {
713 t.Fatal(err)
714 }
715 rootHead := ""
716 for _, head := range heads {
717 if head.ID != childHead {
718 rootHead = head.ID
719 break
720 }
721 }
722 if rootHead == "" {
723 t.Fatalf("legacy root head missing: %+v", heads)
724 }
725
726 v3root := filepath.Join(dir, "sessions-v4")
727 rootResult, err := MigrateLegacyHead(t.Context(), path, v3root, rootHead)
728 if err != nil {
729 t.Fatal(err)
730 }
731 childResult, err := MigrateLegacyHead(t.Context(), path, v3root, childHead)
732 if err != nil {
733 t.Fatal(err)
734 }
735 if rootResult.TargetID == childResult.TargetID {
736 t.Fatal("distinct legacy heads reused one v3 target")
737 }
738 rootCommits, err := Replay(rootResult.TargetDir, nil)
739 if err != nil {
740 t.Fatal(err)
741 }
742 childCommits, err := Replay(childResult.TargetDir, nil)
743 if err != nil {
744 t.Fatal(err)
745 }
746 rootProjection, err := Project(rootCommits)
747 if err != nil {
748 t.Fatal(err)
749 }
750 childProjection, err := Project(childCommits)
751 if err != nil {
752 t.Fatal(err)
753 }
754 if len(rootProjection.Messages) != 3 || len(childProjection.Messages) != 4 || childProjection.Messages[3].Content != "child only" {
755 t.Fatalf("head projections root=%+v child=%+v", rootProjection.Messages, childProjection.Messages)
756 }
757 }
758
759 func mustRead(t *testing.T, path string) []byte {
760 t.Helper()
761 b, err := os.ReadFile(path)
762 if err != nil {
763 t.Fatal(err)
764 }
765 return b
766 }
767
768 func containsJSONKey(raw []byte, key string) bool {
769 var value map[string]any
770 _ = json.Unmarshal(raw, &value)
771 _, ok := value[key]
772 return ok
773 }
774
774 lines GO