返回 DeepSeek-Reasonix
session_events_test.go
根目录 / internal / agent / session_events_test.go
1 package agent
2
3 import (
4 "errors"
5 "os"
6 "path/filepath"
7 "strings"
8 "sync"
9 "testing"
10 "time"
11
12 "reasonix/internal/provider"
13 "reasonix/internal/store"
14 )
15
16 func sessionWithTurns(t *testing.T, path string, turns int) *Session {
17 t.Helper()
18 s := NewSession("sys")
19 for i := range turns {
20 s.Add(provider.Message{Role: provider.RoleUser, Content: "prompt " + strings.Repeat("x", 8)})
21 s.Add(provider.Message{Role: provider.RoleAssistant, Content: "reply"})
22 if err := s.SaveSnapshot(path); err != nil {
23 t.Fatalf("SaveSnapshot turn %d: %v", i, err)
24 }
25 }
26 return s
27 }
28
29 func TestLoadSessionToleratesTornEventLogTail(t *testing.T) {
30 path := filepath.Join(t.TempDir(), "session.jsonl")
31 s := sessionWithTurns(t, path, 2)
32 want := len(s.Snapshot())
33
34 logPath := store.SessionEventLog(path)
35 f, err := os.OpenFile(logPath, os.O_WRONLY|os.O_APPEND, 0o644)
36 if err != nil {
37 t.Fatalf("open log: %v", err)
38 }
39 if _, err := f.Write([]byte(`{"schema_version":1,"type":"append","message_index":5,"mess`)); err != nil {
40 t.Fatalf("write torn tail: %v", err)
41 }
42 f.Close()
43
44 loaded, err := LoadSession(path)
45 if err != nil {
46 t.Fatalf("LoadSession with torn tail: %v", err)
47 }
48 if len(loaded.Messages) != want {
49 t.Fatalf("messages with torn tail = %d, want %d", len(loaded.Messages), want)
50 }
51 if !loaded.eventLogDamaged {
52 t.Fatal("eventLogDamaged = false, want true for torn tail")
53 }
54 }
55
56 func TestSaveSnapshotHealsTornEventLogTail(t *testing.T) {
57 path := filepath.Join(t.TempDir(), "session.jsonl")
58 sessionWithTurns(t, path, 2)
59
60 logPath := store.SessionEventLog(path)
61 intact, err := os.ReadFile(logPath)
62 if err != nil {
63 t.Fatalf("read log: %v", err)
64 }
65 if err := os.WriteFile(logPath, append(intact, []byte(`{"schema_version":1,"type":"ap`)...), 0o644); err != nil {
66 t.Fatalf("write torn log: %v", err)
67 }
68
69 // A fresh runtime resumes the session and keeps chatting: the save must
70 // truncate the torn tail (not bury it) and the appended turn must replay.
71 resumed, err := LoadSession(path)
72 if err != nil {
73 t.Fatalf("LoadSession: %v", err)
74 }
75 resumed.Add(provider.Message{Role: provider.RoleUser, Content: "after crash"})
76 if err := resumed.SaveSnapshot(path); err != nil {
77 t.Fatalf("SaveSnapshot after torn tail: %v", err)
78 }
79
80 reloaded, err := LoadSession(path)
81 if err != nil {
82 t.Fatalf("LoadSession after heal: %v", err)
83 }
84 if reloaded.eventLogDamaged {
85 t.Fatal("event log still damaged after healing save")
86 }
87 tail := reloaded.Messages[len(reloaded.Messages)-1]
88 if tail.Content != "after crash" {
89 t.Fatalf("healed tail = %q, want %q", tail.Content, "after crash")
90 }
91 }
92
93 func TestLoadSessionIgnoresForeignEventLog(t *testing.T) {
94 path := filepath.Join(t.TempDir(), "session.jsonl")
95 sessionWithTurns(t, path, 1)
96
97 // A file this build cannot own squats the native log path: undecodable
98 // bytes, or a legacy v0.x Claude-style event transcript left in place by
99 // migration. It must be read-ignored and never written to.
100 logPath := store.SessionEventLog(path)
101 foreign := "garbage that is not json\n"
102 if err := os.WriteFile(logPath, []byte(foreign), 0o644); err != nil {
103 t.Fatalf("write foreign log: %v", err)
104 }
105
106 loaded, err := LoadSession(path)
107 if err != nil {
108 t.Fatalf("LoadSession with foreign log: %v", err)
109 }
110 if len(loaded.Messages) == 0 {
111 t.Fatal("expected checkpoint transcript, got empty session")
112 }
113 if loaded.eventLogDamaged {
114 t.Fatal("foreign log misreported as damaged native log")
115 }
116
117 // Saves keep working checkpoint-only and never touch the foreign file.
118 loaded.Add(provider.Message{Role: provider.RoleUser, Content: "still works"})
119 if err := loaded.SaveSnapshot(path); err != nil {
120 t.Fatalf("SaveSnapshot with foreign log: %v", err)
121 }
122 got, err := os.ReadFile(logPath)
123 if err != nil {
124 t.Fatalf("read foreign log: %v", err)
125 }
126 if string(got) != foreign {
127 t.Fatalf("foreign log was modified:\nbefore=%q\nafter=%q", foreign, got)
128 }
129 reloaded, err := LoadSession(path)
130 if err != nil {
131 t.Fatalf("LoadSession after checkpoint-only save: %v", err)
132 }
133 if got := reloaded.Messages[len(reloaded.Messages)-1].Content; got != "still works" {
134 t.Fatalf("checkpoint-only tail = %q, want %q", got, "still works")
135 }
136 }
137
138 func TestSaveLeavesLegacyEventTranscriptUntouched(t *testing.T) {
139 // The v0.x migration reconstructs sessions from legacy Claude-style
140 // ".events.jsonl" transcripts that live in the SAME directory the native
141 // session is imported into — i.e. exactly at the native event-log path.
142 // The import must succeed, and the user's original file must survive
143 // byte-for-byte.
144 dir := t.TempDir()
145 path := filepath.Join(dir, "chat-1.jsonl")
146 legacy := `{"type":"user.message","id":1,"ts":"t","turn":0,"text":"hello from v0.x"}` + "\n" +
147 `{"type":"model.final","id":2,"ts":"t","turn":0,"content":"hi","toolCalls":[],"usage":{},"costUsd":0}` + "\n"
148 logPath := store.SessionEventLog(path)
149 if err := os.WriteFile(logPath, []byte(legacy), 0o644); err != nil {
150 t.Fatalf("write legacy transcript: %v", err)
151 }
152
153 s := NewSession("")
154 s.Add(provider.Message{Role: provider.RoleUser, Content: "hello from v0.x"})
155 s.Add(provider.Message{Role: provider.RoleAssistant, Content: "hi"})
156 if err := s.Save(path); err != nil {
157 t.Fatalf("Save beside legacy transcript: %v", err)
158 }
159
160 got, err := os.ReadFile(logPath)
161 if err != nil {
162 t.Fatalf("read legacy transcript: %v", err)
163 }
164 if string(got) != legacy {
165 t.Fatalf("legacy transcript modified:\nbefore=%q\nafter=%q", legacy, got)
166 }
167 loaded, err := LoadSession(path)
168 if err != nil {
169 t.Fatalf("LoadSession of imported session: %v", err)
170 }
171 if len(loaded.Messages) != 2 || loaded.Messages[1].Content != "hi" {
172 t.Fatalf("imported transcript = %+v", loaded.Messages)
173 }
174 }
175
176 // TestRepairPreservesDamagedTailBytes locks in the #6607 content-loss fix:
177 // bytes the tail repair truncates away — including intact turns buried behind
178 // an out-of-order record, the dual-writer shape — must be salvaged to the
179 // .damaged sidecar instead of discarded forever.
180 func TestRepairPreservesDamagedTailBytes(t *testing.T) {
181 path := schemaOneSessionPath(t, "session.jsonl")
182 sessionWithTurns(t, path, 2)
183
184 logPath := store.SessionEventLog(path)
185 // An out-of-order append (as an interleaving second writer would leave)
186 // followed by a torn partial record: replay stops before both, and repair
187 // truncates both away.
188 buried := `{"schema_version":1,"type":"append","message_index":99,"messages":[{"role":"user","content":"buried turn"}],"created_at":"2026-01-01T00:00:00Z"}` + "\n"
189 torn := `{"schema_version":1,"type":"ap`
190 f, err := os.OpenFile(logPath, os.O_WRONLY|os.O_APPEND, 0o644)
191 if err != nil {
192 t.Fatalf("open log: %v", err)
193 }
194 if _, err := f.Write([]byte(buried + torn)); err != nil {
195 t.Fatalf("write damaged tail: %v", err)
196 }
197 f.Close()
198
199 // The healing save truncates the tail; preservation must run first.
200 resumed, err := LoadSession(path)
201 if err != nil {
202 t.Fatalf("LoadSession: %v", err)
203 }
204 resumed.Add(provider.Message{Role: provider.RoleUser, Content: "after damage"})
205 if err := resumed.SaveSnapshot(path); err != nil {
206 t.Fatalf("SaveSnapshot after damage: %v", err)
207 }
208
209 salvaged, err := os.ReadFile(store.SessionEventLogDamaged(path))
210 if err != nil {
211 t.Fatalf("read damaged sidecar: %v", err)
212 }
213 if !strings.Contains(string(salvaged), "buried turn") {
214 t.Fatalf("salvage sidecar is missing the buried turn:\n%s", salvaged)
215 }
216 if !strings.Contains(string(salvaged), `"damaged_tail":true`) {
217 t.Fatalf("salvage sidecar is missing the incident header:\n%s", salvaged)
218 }
219 // The log itself must be healed: replay to the end, no damage, and the
220 // buried record must not leak back into the transcript.
221 reloaded, err := LoadSession(path)
222 if err != nil {
223 t.Fatalf("LoadSession after heal: %v", err)
224 }
225 if reloaded.eventLogDamaged {
226 t.Fatal("event log still damaged after healing save")
227 }
228 for _, m := range reloaded.Messages {
229 if m.Content == "buried turn" {
230 t.Fatal("truncated record leaked back into the transcript")
231 }
232 }
233
234 // A second incident appends to the sidecar instead of overwriting the
235 // first salvage.
236 f, err = os.OpenFile(store.SessionEventLog(path), os.O_WRONLY|os.O_APPEND, 0o644)
237 if err != nil {
238 t.Fatalf("reopen log: %v", err)
239 }
240 if _, err := f.Write([]byte(`{"schema_version":1,"type":"append","message_index":77,"messages":[{"role":"user","content":"second incident"}],"created_at":"2026-01-01T00:00:00Z"}` + "\n")); err != nil {
241 t.Fatalf("write second damage: %v", err)
242 }
243 f.Close()
244 resumed2, err := LoadSession(path)
245 if err != nil {
246 t.Fatalf("LoadSession second incident: %v", err)
247 }
248 resumed2.Add(provider.Message{Role: provider.RoleUser, Content: "after second"})
249 if err := resumed2.SaveSnapshot(path); err != nil {
250 t.Fatalf("SaveSnapshot second incident: %v", err)
251 }
252 salvaged, err = os.ReadFile(store.SessionEventLogDamaged(path))
253 if err != nil {
254 t.Fatalf("read damaged sidecar after second incident: %v", err)
255 }
256 if !strings.Contains(string(salvaged), "buried turn") || !strings.Contains(string(salvaged), "second incident") {
257 t.Fatalf("salvage sidecar must keep both incidents:\n%s", salvaged)
258 }
259 }
260
261 func TestReplayStopsAtBrokenAppendChain(t *testing.T) {
262 path := schemaOneSessionPath(t, "session.jsonl")
263 sessionWithTurns(t, path, 1)
264 logPath := store.SessionEventLog(path)
265 // Append an event whose MessageIndex does not chain onto the transcript.
266 f, err := os.OpenFile(logPath, os.O_WRONLY|os.O_APPEND, 0o644)
267 if err != nil {
268 t.Fatalf("open log: %v", err)
269 }
270 if _, err := f.Write([]byte(`{"schema_version":1,"type":"append","message_index":99,"messages":[{"role":"user","content":"lost"}],"created_at":"2026-01-01T00:00:00Z"}` + "\n")); err != nil {
271 t.Fatalf("write broken chain: %v", err)
272 }
273 f.Close()
274
275 replay, err := replaySessionEventLog(logPath)
276 if err != nil {
277 t.Fatalf("replay: %v", err)
278 }
279 if !replay.damaged {
280 t.Fatal("replay.damaged = false, want true for broken chain")
281 }
282 for _, m := range replay.msgs {
283 if m.Content == "lost" {
284 t.Fatal("broken-chain record leaked into replayed transcript")
285 }
286 }
287 }
288
289 func TestReplayRejectsUnsupportedSchemaVersion(t *testing.T) {
290 dir := t.TempDir()
291 logPath := filepath.Join(dir, "future.events.jsonl")
292 if err := os.WriteFile(logPath, []byte(`{"schema_version":99,"type":"replace","messages":[]}`+"\n"), 0o644); err != nil {
293 t.Fatalf("write future log: %v", err)
294 }
295 if _, err := replaySessionEventLog(logPath); err == nil {
296 t.Fatal("replay of future schema succeeded, want hard error (no silent truncation of newer writers)")
297 }
298 }
299
300 func TestReplaySessionEventLogEnforcesResourceBudgetsWithoutMutation(t *testing.T) {
301 dir := t.TempDir()
302 logPath := filepath.Join(dir, "bounded.events.jsonl")
303 replace := `{"schema_version":1,"type":"replace","messages":[{"role":"system","content":"sys"},{"role":"user","content":"prompt"}]}` + "\n"
304 appendEvent := `{"schema_version":1,"type":"append","message_index":2,"messages":[{"role":"assistant","content":"reply"}]}` + "\n"
305 contents := replace + appendEvent
306 if err := os.WriteFile(logPath, []byte(contents), 0o600); err != nil {
307 t.Fatalf("write event log: %v", err)
308 }
309
310 tests := []struct {
311 name string
312 limits sessionReplayLimits
313 resource string
314 }{
315 {
316 name: "encoded bytes",
317 limits: sessionReplayLimits{
318 maxBytes: int64(len(contents) - 1), maxRecords: 10, maxMessages: 10,
319 },
320 resource: "encoded_bytes",
321 },
322 {
323 name: "event records",
324 limits: sessionReplayLimits{
325 maxBytes: int64(len(contents) + 1), maxRecords: 1, maxMessages: 10,
326 },
327 resource: "event_records",
328 },
329 {
330 name: "messages",
331 limits: sessionReplayLimits{
332 maxBytes: int64(len(contents) + 1), maxRecords: 10, maxMessages: 1,
333 },
334 resource: "messages",
335 },
336 {
337 name: "aggregate messages across records",
338 limits: sessionReplayLimits{
339 maxBytes: int64(len(contents) + 1), maxRecords: 10, maxMessages: 2,
340 },
341 resource: "messages",
342 },
343 }
344 for _, tc := range tests {
345 t.Run(tc.name, func(t *testing.T) {
346 if _, err := replaySessionEventLogWithLimits(logPath, tc.limits, nil); !errors.Is(err, ErrSessionReplayLimitExceeded) {
347 t.Fatalf("replay error = %v, want ErrSessionReplayLimitExceeded", err)
348 } else {
349 var limitErr *SessionReplayLimitError
350 if !errors.As(err, &limitErr) || limitErr.Resource != tc.resource {
351 t.Fatalf("limit error = %#v, want resource %q", limitErr, tc.resource)
352 }
353 if strings.Contains(err.Error(), logPath) {
354 t.Fatalf("user-facing error leaked local path: %q", err)
355 }
356 }
357 got, err := os.ReadFile(logPath)
358 if err != nil {
359 t.Fatalf("read event log after rejection: %v", err)
360 }
361 if string(got) != contents {
362 t.Fatal("resource-budget rejection modified the event log")
363 }
364 })
365 }
366 }
367
368 func TestReplayChecksMessageBudgetBeforeDecodingOverLimitElement(t *testing.T) {
369 logPath := filepath.Join(t.TempDir(), "bounded-before-decode.events.jsonl")
370 // The third element is valid JSON but cannot decode into provider.Message.
371 // A two-message budget must reject it before attempting that decode.
372 events := `{"schema_version":1,"type":"replace","messages":[{}, {}, 42]}` + "\n"
373 if err := os.WriteFile(logPath, []byte(events), 0o600); err != nil {
374 t.Fatalf("write event log: %v", err)
375 }
376
377 replay, err := replaySessionEventLogWithLimits(logPath, sessionReplayLimits{
378 maxBytes: int64(len(events) + 1), maxRecords: 10, maxMessages: 2,
379 }, nil)
380 if !errors.Is(err, ErrSessionReplayLimitExceeded) {
381 t.Fatalf("replay error = %v, want message replay limit", err)
382 }
383 var limitErr *SessionReplayLimitError
384 if !errors.As(err, &limitErr) || limitErr.Resource != "messages" || limitErr.Value != 3 {
385 t.Fatalf("limit error = %#v, want third message rejected before decode", limitErr)
386 }
387 if replay.records != 0 || len(replay.msgs) != 0 {
388 t.Fatalf("partial over-limit record was applied: records=%d messages=%d", replay.records, len(replay.msgs))
389 }
390 }
391
392 func TestReplayChecksCollectionBudgetBeforeDecodingOverLimitElement(t *testing.T) {
393 logPath := filepath.Join(t.TempDir(), "bounded-collection-before-decode.events.jsonl")
394 // The third tool call cannot decode into provider.ToolCall. A two-item
395 // collection budget must reject it before constructing the typed slice.
396 events := `{"schema_version":1,"type":"replace","messages":[{"role":"assistant","tool_calls":[{}, {}, 42]}]}` + "\n"
397 if err := os.WriteFile(logPath, []byte(events), 0o600); err != nil {
398 t.Fatalf("write event log: %v", err)
399 }
400
401 replay, err := replaySessionEventLogWithLimits(logPath, sessionReplayLimits{
402 maxBytes: int64(len(events) + 1), maxRecords: 10, maxMessages: 10, maxCollectionItems: 2,
403 }, nil)
404 if !errors.Is(err, ErrSessionReplayLimitExceeded) {
405 t.Fatalf("replay error = %v, want collection replay limit", err)
406 }
407 var limitErr *SessionReplayLimitError
408 if !errors.As(err, &limitErr) || limitErr.Resource != "message_collection_items" || limitErr.Value != 3 {
409 t.Fatalf("limit error = %#v, want third collection item rejected before decode", limitErr)
410 }
411 if replay.records != 0 || len(replay.msgs) != 0 {
412 t.Fatalf("partial over-limit record was applied: records=%d messages=%d", replay.records, len(replay.msgs))
413 }
414 if got, readErr := os.ReadFile(logPath); readErr != nil || string(got) != events {
415 t.Fatalf("event log changed after refusal: bytes=%q err=%v", got, readErr)
416 }
417 }
418
419 func TestReplayCollectionBudgetAggregatesAcrossRecords(t *testing.T) {
420 logPath := filepath.Join(t.TempDir(), "aggregate-collection-budget.events.jsonl")
421 replace := `{"schema_version":1,"type":"replace","messages":[{"role":"user","images":["one"]}]}` + "\n"
422 appendEvent := `{"schema_version":1,"type":"append","message_index":1,"messages":[{"role":"assistant","tool_calls":[{},{}]}]}` + "\n"
423 events := replace + appendEvent
424 if err := os.WriteFile(logPath, []byte(events), 0o600); err != nil {
425 t.Fatalf("write event log: %v", err)
426 }
427
428 replay, err := replaySessionEventLogWithLimits(logPath, sessionReplayLimits{
429 maxBytes: int64(len(events) + 1), maxRecords: 10, maxMessages: 10, maxCollectionItems: 2,
430 }, nil)
431 if !errors.Is(err, ErrSessionReplayLimitExceeded) {
432 t.Fatalf("replay error = %v, want aggregate collection replay limit", err)
433 }
434 var limitErr *SessionReplayLimitError
435 if !errors.As(err, &limitErr) || limitErr.Resource != "message_collection_items" || limitErr.Value != 3 {
436 t.Fatalf("limit error = %#v, want third live collection item rejected", limitErr)
437 }
438 if replay.records != 1 || len(replay.msgs) != 1 || replay.collectionItems != 1 {
439 t.Fatalf("replay prefix = records:%d messages:%d collections:%d, want first record only",
440 replay.records, len(replay.msgs), replay.collectionItems)
441 }
442 }
443
444 func TestLoadSessionMessagesDoesNotFallbackAfterEventReplayBudget(t *testing.T) {
445 path := filepath.Join(t.TempDir(), "session.jsonl")
446 checkpoint := `{"role":"system","content":"older checkpoint"}` + "\n"
447 if err := os.WriteFile(path, []byte(checkpoint), 0o600); err != nil {
448 t.Fatalf("write checkpoint: %v", err)
449 }
450 logPath := store.SessionEventLog(path)
451 events := `{"schema_version":1,"type":"replace","messages":[{"role":"system","content":"newer event state"},{"role":"user","content":"must not disappear"}]}` + "\n"
452 if err := os.WriteFile(logPath, []byte(events), 0o600); err != nil {
453 t.Fatalf("write event log: %v", err)
454 }
455
456 msgs, fromEvents, damaged, err := loadSessionMessagesWithLimits(path, sessionReplayLimits{
457 maxBytes: int64(len(events) + 1), maxRecords: 10, maxMessages: 1,
458 }, nil)
459 if !errors.Is(err, ErrSessionReplayLimitExceeded) {
460 t.Fatalf("load error = %v, want ErrSessionReplayLimitExceeded", err)
461 }
462 if msgs != nil || !fromEvents || damaged {
463 t.Fatalf("load result = msgs:%v fromEvents:%v damaged:%v, want hard event-log refusal", msgs, fromEvents, damaged)
464 }
465 if got, readErr := os.ReadFile(path); readErr != nil || string(got) != checkpoint {
466 t.Fatalf("checkpoint changed after refusal: bytes=%q err=%v", got, readErr)
467 }
468 if got, readErr := os.ReadFile(logPath); readErr != nil || string(got) != events {
469 t.Fatalf("event log changed after refusal: bytes=%q err=%v", got, readErr)
470 }
471 }
472
473 func TestMigrationReplayIsNotBoundByInteractiveHistoryBudget(t *testing.T) {
474 path := filepath.Join(t.TempDir(), "session.jsonl")
475 events := `{"schema_version":1,"type":"replace","messages":[{"role":"system","content":"system"},{"role":"user","content":"` +
476 strings.Repeat("x", 4096) + `"}]}` + "\n"
477 if err := os.WriteFile(store.SessionEventLog(path), []byte(events), 0o600); err != nil {
478 t.Fatalf("write event log: %v", err)
479 }
480
481 interactive := defaultSessionReplayLimits
482 interactive.maxBytes = 1024
483 if _, err := loadSessionUnlockedWithLimits(path, interactive); !errors.Is(err, ErrSessionReplayLimitExceeded) {
484 t.Fatalf("interactive load error = %v, want ErrSessionReplayLimitExceeded", err)
485 }
486 migrated, err := loadSessionUnlockedWithLimits(path, migrationSessionReplayLimits())
487 if err != nil {
488 t.Fatalf("migration load: %v", err)
489 }
490 if got := migrated.Snapshot(); len(got) != 2 || got[1].Content != strings.Repeat("x", 4096) {
491 t.Fatalf("migration transcript = %#v", got)
492 }
493 }
494
495 func TestProbeSessionEventLogReadsOnlyNativeHeader(t *testing.T) {
496 path := filepath.Join(t.TempDir(), "session.jsonl")
497 logPath := store.SessionEventLog(path)
498 largeRecord := `{"schema_version":1,"type":"replace","messages":[{"role":"user","content":"` +
499 strings.Repeat("x", int(sessionEventProbeMaxBytes*2)) + `"}]}` + "\n"
500 if err := os.WriteFile(logPath, []byte(largeRecord), 0o600); err != nil {
501 t.Fatalf("write event log: %v", err)
502 }
503 probe, err := probeSessionEventLogWithLimits(path, sessionReplayLimits{
504 maxBytes: 128, maxRecords: 10, maxMessages: 10,
505 })
506 if err != nil {
507 t.Fatalf("probe: %v", err)
508 }
509 if !probe.native || probe.size != int64(len(largeRecord)) {
510 t.Fatalf("probe = %+v, want native size %d", probe, len(largeRecord))
511 }
512 }
513
514 func TestLoadSessionMessagesAcceptsReorderedEventHeader(t *testing.T) {
515 path := filepath.Join(t.TempDir(), "session.jsonl")
516 checkpoint := `{"role":"system","content":"older checkpoint"}` + "\n"
517 if err := os.WriteFile(path, []byte(checkpoint), 0o600); err != nil {
518 t.Fatalf("write checkpoint: %v", err)
519 }
520 content := "newer event state " + strings.Repeat("x", int(sessionEventProbeMaxBytes*2))
521 // schema_version deliberately follows a messages payload larger than the
522 // prefix probe. Valid in-budget JSON must remain field-order independent.
523 events := `{"type":"replace","messages":[{"role":"system","content":"sys"},{"role":"user","content":"` +
524 content + `"}],"schema_version":1}` + "\n"
525 if err := os.WriteFile(store.SessionEventLog(path), []byte(events), 0o600); err != nil {
526 t.Fatalf("write reordered event log: %v", err)
527 }
528
529 msgs, fromEvents, damaged, err := loadSessionMessages(path)
530 if err != nil {
531 t.Fatalf("load reordered event log: %v", err)
532 }
533 if !fromEvents || damaged {
534 t.Fatalf("load result fromEvents=%v damaged=%v, want clean native replay", fromEvents, damaged)
535 }
536 if len(msgs) != 2 || msgs[1].Content != content {
537 t.Fatalf("loaded messages = %d, want reordered event state", len(msgs))
538 }
539 }
540
541 func TestDefaultSaveBootstrapsEventLog(t *testing.T) {
542 path := filepath.Join(t.TempDir(), "session.jsonl")
543 s := NewSession("sys")
544 s.Add(provider.Message{Role: provider.RoleUser, Content: "one-shot"})
545 if err := s.Save(path); err != nil {
546 t.Fatalf("Save: %v", err)
547 }
548 if _, err := os.Stat(store.SessionEventLog(path)); err != nil {
549 t.Fatalf("default save did not create an event log: %v", err)
550 }
551 if _, err := os.Stat(store.SessionEventIndex(path)); err != nil {
552 t.Fatalf("default save did not create an event index: %v", err)
553 }
554 }
555
556 func TestDefaultSaveRejectsDivergedOverwrite(t *testing.T) {
557 path := schemaOneSessionPath(t, "session.jsonl")
558 winner := sessionWithTurns(t, path, 3).Snapshot()
559
560 stale := NewSession("sys")
561 stale.Add(provider.Message{Role: provider.RoleUser, Content: "stale state"})
562 if err := stale.Save(path); !errors.Is(err, ErrSessionSnapshotConflict) {
563 t.Fatalf("diverged Save error = %v, want ErrSessionSnapshotConflict", err)
564 }
565 loaded, err := LoadSession(path)
566 if err != nil {
567 t.Fatalf("LoadSession: %v", err)
568 }
569 if !messagesEqualForStorageList(loaded.Messages, winner) {
570 t.Fatalf("default Save replaced newer transcript: got %d messages, want %d", len(loaded.Messages), len(winner))
571 }
572 }
573
574 func TestEventLogCompactionBoundsGrowth(t *testing.T) {
575 path := schemaOneSessionPath(t, "session.jsonl")
576 s := NewSession("sys")
577 filler := strings.Repeat("y", 8<<10)
578 // Repeated rewrites (each a full replace event) must not grow the log
579 // without bound: once past the threshold the log folds to one event.
580 for i := range 60 {
581 s.Replace([]provider.Message{
582 {Role: provider.RoleSystem, Content: "sys"},
583 {Role: provider.RoleUser, Content: filler},
584 {Role: provider.RoleAssistant, Content: strings.Repeat("z", i+1)},
585 })
586 if err := s.SaveRewrite(path); err != nil {
587 t.Fatalf("SaveRewrite %d: %v", i, err)
588 }
589 }
590 _, contentBytes, err := digestAndSizeSessionMessages(s.Snapshot())
591 if err != nil {
592 t.Fatalf("digest: %v", err)
593 }
594 logSize := sessionEventLogSize(path)
595 limit := sessionEventLogCompactFloor
596 if scaled := contentBytes * sessionEventLogCompactFactor; scaled > limit {
597 limit = scaled
598 }
599 // One post-compaction replace event of slack is allowed.
600 if logSize > limit+contentBytes+4096 {
601 t.Fatalf("event log grew unbounded: size=%d limit=%d content=%d", logSize, limit, contentBytes)
602 }
603 loaded, err := LoadSession(path)
604 if err != nil {
605 t.Fatalf("LoadSession after compaction churn: %v", err)
606 }
607 if got := loaded.Messages[len(loaded.Messages)-1].Content; got != strings.Repeat("z", 60) {
608 t.Fatalf("compacted transcript tail wrong: %q", got[:min(8, len(got))])
609 }
610 anchor, err := loadSessionMessagesFromJSONL(path, nil)
611 if err != nil {
612 t.Fatalf("read anchor: %v", err)
613 }
614 if got := anchor[len(anchor)-1].Content; got != strings.Repeat("z", 60) {
615 t.Fatal("anchor not refreshed by checkpoint compaction")
616 }
617 }
618
619 func TestMigrationLoadRejectsDamagedAuthoritativeLog(t *testing.T) {
620 path := filepath.Join(t.TempDir(), "session.jsonl")
621 s := NewSession("sys")
622 s.Add(provider.Message{Role: provider.RoleUser, Content: "checkpoint"})
623 if err := s.SaveSnapshot(path); err != nil {
624 t.Fatal(err)
625 }
626 file, err := os.OpenFile(SessionEventLogPath(path), os.O_APPEND|os.O_WRONLY, 0o600)
627 if err != nil {
628 t.Fatal(err)
629 }
630 if _, err := file.WriteString("{not-json}\n"); err != nil {
631 _ = file.Close()
632 t.Fatal(err)
633 }
634 if err := file.Close(); err != nil {
635 t.Fatal(err)
636 }
637 if _, err := LoadSessionForMigration(t.Context(), path); !errors.Is(err, ErrSessionHistoryDamaged) {
638 t.Fatalf("migration error = %v, want ErrSessionHistoryDamaged", err)
639 }
640 }
641
642 func TestConcurrentLoadDuringAppendsStaysConsistent(t *testing.T) {
643 path := filepath.Join(t.TempDir(), "session.jsonl")
644 s := NewSession("sys")
645 s.Add(provider.Message{Role: provider.RoleUser, Content: "seed"})
646 if err := s.SaveSnapshot(path); err != nil {
647 t.Fatalf("seed SaveSnapshot: %v", err)
648 }
649
650 stop := make(chan struct{})
651 var wg sync.WaitGroup
652 wg.Add(1)
653 errCh := make(chan error, 1)
654 go func() {
655 defer wg.Done()
656 for {
657 select {
658 case <-stop:
659 return
660 default:
661 }
662 loaded, err := LoadSession(path)
663 if err != nil {
664 select {
665 case errCh <- err:
666 default:
667 }
668 return
669 }
670 if loaded.eventLogDamaged {
671 select {
672 case errCh <- os.ErrInvalid:
673 default:
674 }
675 return
676 }
677 }
678 }()
679 for i := range 40 {
680 s.Add(provider.Message{Role: provider.RoleAssistant, Content: strings.Repeat("a", 512)})
681 if err := s.SaveSnapshot(path); err != nil {
682 t.Fatalf("SaveSnapshot %d: %v", i, err)
683 }
684 }
685 close(stop)
686 wg.Wait()
687 select {
688 case err := <-errCh:
689 t.Fatalf("concurrent LoadSession failed or saw damage: %v", err)
690 default:
691 }
692 }
693
694 func TestSessionsShareContentSeesEventLogDivergence(t *testing.T) {
695 dir := t.TempDir()
696 pathA := filepath.Join(dir, "a.jsonl")
697 pathB := filepath.Join(dir, "b.jsonl")
698 a := sessionWithTurns(t, pathA, 1)
699 _ = sessionWithTurns(t, pathB, 1)
700
701 same, err := SessionsShareContent(pathA, pathB)
702 if err != nil {
703 t.Fatalf("SessionsShareContent: %v", err)
704 }
705 if !same {
706 t.Fatal("identical transcripts reported as different")
707 }
708
709 // Grow A normally, then restore its old compatibility checkpoint to model a
710 // crash or an older Reasonix build that did not advance the display read
711 // model. Transcript equality must still follow the event log.
712 anchorB, err := os.ReadFile(pathB)
713 if err != nil {
714 t.Fatalf("read checkpoint B: %v", err)
715 }
716 a.Add(provider.Message{Role: provider.RoleUser, Content: "diverged"})
717 if err := a.SaveSnapshot(pathA); err != nil {
718 t.Fatalf("SaveSnapshot diverge: %v", err)
719 }
720 if err := os.WriteFile(pathA, anchorB, 0o600); err != nil {
721 t.Fatalf("restore stale checkpoint A: %v", err)
722 }
723 anchorA, _ := os.ReadFile(pathA)
724 if string(anchorA) != string(anchorB) {
725 t.Fatal("test setup did not preserve byte-identical checkpoints")
726 }
727 same, err = SessionsShareContent(pathA, pathB)
728 if err != nil {
729 t.Fatalf("SessionsShareContent after divergence: %v", err)
730 }
731 if same {
732 t.Fatal("diverged transcripts reported as identical")
733 }
734 }
735
736 func TestLoadSessionUserMessagesSeesEventLogTurns(t *testing.T) {
737 path := filepath.Join(t.TempDir(), "session.jsonl")
738 s := NewSession("sys")
739 s.Add(provider.Message{Role: provider.RoleUser, Content: "first prompt"})
740 if err := s.SaveSnapshot(path); err != nil {
741 t.Fatalf("SaveSnapshot: %v", err)
742 }
743 s.Add(provider.Message{Role: provider.RoleAssistant, Content: "reply"})
744 s.Add(provider.Message{Role: provider.RoleUser, Content: "second prompt"})
745 if err := s.SaveSnapshot(path); err != nil {
746 t.Fatalf("SaveSnapshot append: %v", err)
747 }
748
749 users, err := LoadSessionUserMessages(path)
750 if err != nil {
751 t.Fatalf("LoadSessionUserMessages: %v", err)
752 }
753 if len(users) != 2 || users[0].Message.Content != "first prompt" || users[1].Message.Content != "second prompt" {
754 t.Fatalf("user messages = %+v, want both prompts", users)
755 }
756 if users[1].At.IsZero() {
757 t.Fatal("appended prompt lost its event timestamp")
758 }
759 }
760
761 func TestLoadSessionUserMessagesDoesNotFallbackAfterEventReplayBudget(t *testing.T) {
762 path := filepath.Join(t.TempDir(), "session.jsonl")
763 checkpoint := `{"role":"user","content":"older checkpoint prompt"}` + "\n"
764 if err := os.WriteFile(path, []byte(checkpoint), 0o600); err != nil {
765 t.Fatalf("write checkpoint: %v", err)
766 }
767 logPath := store.SessionEventLog(path)
768 events := `{"schema_version":1,"type":"replace","messages":[{"role":"user","content":"newer prompt"},{"role":"assistant","content":"reply"}]}` + "\n"
769 if err := os.WriteFile(logPath, []byte(events), 0o600); err != nil {
770 t.Fatalf("write event log: %v", err)
771 }
772
773 users, err := loadSessionUserMessagesWithLimits(path, sessionReplayLimits{
774 maxBytes: int64(len(events) + 1), maxRecords: 10, maxMessages: 1, maxCollectionItems: 10,
775 })
776 if !errors.Is(err, ErrSessionReplayLimitExceeded) {
777 t.Fatalf("load error = %v, want ErrSessionReplayLimitExceeded", err)
778 }
779 if users != nil {
780 t.Fatalf("user messages = %+v, want no stale checkpoint fallback", users)
781 }
782 if got, readErr := os.ReadFile(path); readErr != nil || string(got) != checkpoint {
783 t.Fatalf("checkpoint changed after refusal: bytes=%q err=%v", got, readErr)
784 }
785 if got, readErr := os.ReadFile(logPath); readErr != nil || string(got) != events {
786 t.Fatalf("event log changed after refusal: bytes=%q err=%v", got, readErr)
787 }
788 }
789
790 func TestLoadSessionUserMessagesDoesNotFallbackFromFutureEventSchema(t *testing.T) {
791 path := filepath.Join(t.TempDir(), "session.jsonl")
792 checkpoint := `{"role":"user","content":"older checkpoint prompt"}` + "\n"
793 if err := os.WriteFile(path, []byte(checkpoint), 0o600); err != nil {
794 t.Fatalf("write checkpoint: %v", err)
795 }
796 events := `{"schema_version":99,"type":"replace","messages":[{"role":"user","content":"future prompt"}]}` + "\n"
797 if err := os.WriteFile(store.SessionEventLog(path), []byte(events), 0o600); err != nil {
798 t.Fatalf("write event log: %v", err)
799 }
800
801 users, err := LoadSessionUserMessages(path)
802 if err == nil || !strings.Contains(err.Error(), "uses schema 99") {
803 t.Fatalf("load error = %v, want future-schema refusal", err)
804 }
805 if users != nil {
806 t.Fatalf("user messages = %+v, want no stale checkpoint fallback", users)
807 }
808 }
809
810 func TestSessionEventLogPreservesUserCreatedAt(t *testing.T) {
811 path := filepath.Join(t.TempDir(), "session.jsonl")
812 want := time.Date(2026, 7, 16, 8, 30, 0, 0, time.UTC).UnixMilli()
813 s := NewSession("sys")
814 s.Add(provider.Message{Role: provider.RoleUser, Content: "timed prompt", CreatedAt: want})
815 if err := s.SaveSnapshot(path); err != nil {
816 t.Fatalf("SaveSnapshot: %v", err)
817 }
818
819 loaded, err := LoadSession(path)
820 if err != nil {
821 t.Fatalf("LoadSession: %v", err)
822 }
823 msgs := loaded.Snapshot()
824 if len(msgs) != 2 || msgs[1].CreatedAt != want {
825 t.Fatalf("loaded messages = %+v, want user createdAt %d", msgs, want)
826 }
827 }
828
829 func TestSessionDigestIgnoresUserCreatedAt(t *testing.T) {
830 withoutTime := []provider.Message{
831 {Role: provider.RoleSystem, Content: "sys"},
832 {Role: provider.RoleUser, Content: "prompt"},
833 }
834 withTime := append([]provider.Message(nil), withoutTime...)
835 withTime[1].CreatedAt = 1_718_000_000_000
836
837 withoutDigest, withoutSize, err := digestAndSizeSessionMessages(withoutTime)
838 if err != nil {
839 t.Fatalf("digest without createdAt: %v", err)
840 }
841 withDigest, withSize, err := digestAndSizeSessionMessages(withTime)
842 if err != nil {
843 t.Fatalf("digest with createdAt: %v", err)
844 }
845 if withoutDigest != withDigest || withoutSize != withSize {
846 t.Fatalf("local createdAt changed transcript identity: digestEqual=%v size=%d/%d", withoutDigest == withDigest, withoutSize, withSize)
847 }
848 if !messagesHavePrefix(withTime, withoutTime) || !messagesHavePrefix(withoutTime, withTime) {
849 t.Fatal("local createdAt broke cross-version transcript prefix matching")
850 }
851 if depth := messagesPrefixDigestDepth(withTime, withoutDigest); depth != len(withTime) {
852 t.Fatalf("createdAt-aware prefix digest depth = %d, want %d", depth, len(withTime))
853 }
854 }
855
856 func TestSessionContentModTimeTracksEventLog(t *testing.T) {
857 path := filepath.Join(t.TempDir(), "session.jsonl")
858 s := sessionWithTurns(t, path, 1)
859 s.Add(provider.Message{Role: provider.RoleUser, Content: "new"})
860 if err := s.SaveSnapshot(path); err != nil {
861 t.Fatalf("SaveSnapshot: %v", err)
862 }
863 // Filesystems such as macOS can expose coarse mtime resolution, so the
864 // checkpoint and event log written by one save are allowed to have the
865 // same timestamp. Make the ordering under test explicit after the save:
866 // the checkpoint is a stale anchor while the event log remains current.
867 old := time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC)
868 if err := os.Chtimes(path, old, old); err != nil {
869 t.Fatalf("age anchor: %v", err)
870 }
871 anchorInfo, err := os.Stat(path)
872 if err != nil {
873 t.Fatalf("stat anchor: %v", err)
874 }
875 if got := SessionContentModTime(path); !got.After(anchorInfo.ModTime()) {
876 t.Fatalf("SessionContentModTime = %v, want newer than stale anchor %v", got, anchorInfo.ModTime())
877 }
878 }
879
879 lines GO