返回 DeepSeek-Reasonix
history_index_boundary_test.go
根目录 / internal / session / history_index_boundary_test.go
1 package session
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "os"
8 "path/filepath"
9 "testing"
10
11 "reasonix/internal/projectiondb"
12 )
13
14 func historyBoundaryFixture(t *testing.T) (*Service, *Runtime, string, string) {
15 t.Helper()
16 root := t.TempDir()
17 s, err := NewService("local", NewFilesystemPersistence(root))
18 if err != nil {
19 t.Fatal(err)
20 }
21 t.Cleanup(func() { _ = s.CloseAll(context.Background()) })
22 r, err := s.Create(t.Context(), CreateOptions{SessionID: "boundary"})
23 if err != nil {
24 t.Fatal(err)
25 }
26 appendRecoveryTestMessage(t, r.Session(), "first", "first")
27 if _, err := r.Session().Flush(t.Context()); err != nil {
28 t.Fatal(err)
29 }
30 return s, r, filepath.Join(root, r.Ref().SessionID), historyIndexPath(root, r.Ref().SessionID)
31 }
32
33 func TestHistoryRebuildProgressMatchesScannedTransactions(t *testing.T) {
34 s, r, dir, path := historyBoundaryFixture(t)
35 // Force the startup ordering: stat, append, then scan. No timing dependency.
36 before, err := revisionOfLog(dir)
37 if err != nil {
38 t.Fatal(err)
39 }
40 appendRecoveryTestMessage(t, r.Session(), "second", "second")
41 if _, err := r.Session().Flush(t.Context()); err != nil {
42 t.Fatal(err)
43 }
44 if err := rebuildHistoryIndex(t.Context(), dir, path, r.Ref().SessionID, before); err != nil {
45 t.Fatal(err)
46 }
47 page := historyPageReady(t, s.Query(), r.Ref(), "", 100)
48 if len(page.Messages) != 2 {
49 t.Fatalf("messages = %d, want 2", len(page.Messages))
50 }
51 }
52
53 func TestHistoryRepairsPersistedMismatchedProgress(t *testing.T) {
54 s, r, dir, path := historyBoundaryFixture(t)
55 before, err := revisionOfLog(dir)
56 if err != nil {
57 t.Fatal(err)
58 }
59 appendRecoveryTestMessage(t, r.Session(), "second", "second")
60 if _, err := r.Session().Flush(t.Context()); err != nil {
61 t.Fatal(err)
62 }
63 _ = historyPageReady(t, s.Query(), r.Ref(), "", 100)
64 // Reproduce the shipped index: new sequence, old byte offset.
65 h, err := projectiondb.Open(t.Context(), projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true})
66 if err != nil {
67 t.Fatal(err)
68 }
69 _, err = h.DB.ExecContext(t.Context(), "UPDATE metadata SET value=? WHERE key='log_size'", fmt.Sprint(before.Size))
70 _ = h.DB.Close()
71 if err != nil {
72 t.Fatal(err)
73 }
74 // A new Query proves recovery survives process-local state loss.
75 restarted, err := NewService("local", s.persistence)
76 if err != nil {
77 t.Fatal(err)
78 }
79 t.Cleanup(func() { _ = restarted.CloseAll(context.Background()) })
80 page, err := waitHistoryWindow(t, restarted.Query(), r.Ref(), HistoryWindowRequest{Anchor: "newest", Limit: 32})
81 if err != nil || len(page.Messages) != 2 {
82 t.Fatalf("recovered window: status=%s messages=%d err=%v", page.Status, len(page.Messages), err)
83 }
84 }
85
86 func TestHistoryRebuildDoesNotPublishIncompleteTail(t *testing.T) {
87 s, r, dir, path := historyBoundaryFixture(t)
88 before, err := os.ReadFile(filepath.Join(dir, "events.frames"))
89 if err != nil {
90 t.Fatal(err)
91 }
92 appendRecoveryTestMessage(t, r.Session(), "second", "second")
93 if _, err := r.Session().Flush(t.Context()); err != nil {
94 t.Fatal(err)
95 }
96 logPath := filepath.Join(dir, "events.frames")
97 complete, err := os.ReadFile(logPath)
98 if err != nil {
99 t.Fatal(err)
100 }
101 if err := os.WriteFile(logPath, complete[:len(complete)-1], 0o600); err != nil {
102 t.Fatal(err)
103 }
104 revision, err := revisionOfLog(dir)
105 if err != nil {
106 t.Fatal(err)
107 }
108 if err := rebuildHistoryIndex(t.Context(), dir, path, r.Ref().SessionID, revision); err != nil {
109 t.Fatal(err)
110 }
111 h, err := projectiondb.Open(t.Context(), projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true})
112 if err != nil {
113 t.Fatal(err)
114 }
115 metadata, err := readHistoryIndexMetadata(t.Context(), h.DB)
116 _ = h.DB.Close()
117 if err != nil {
118 t.Fatal(err)
119 }
120 if metadata.logSize != int64(len(before)) {
121 t.Fatalf("published offset=%d, last complete transaction ends at %d", metadata.logSize, len(before))
122 }
123 if err := os.WriteFile(logPath, complete, 0o600); err != nil {
124 t.Fatal(err)
125 }
126 page := historyPageReady(t, s.Query(), r.Ref(), "", 100)
127 if len(page.Messages) != 2 {
128 t.Fatalf("completed tail messages=%d", len(page.Messages))
129 }
130 }
131
132 func TestHistoryRepairDoesNotHideDamagedLog(t *testing.T) {
133 s, r, dir, _ := historyBoundaryFixture(t)
134 _ = historyPageReady(t, s.Query(), r.Ref(), "", 100)
135 log, err := os.OpenFile(filepath.Join(dir, "events.frames"), os.O_WRONLY|os.O_APPEND, 0o600)
136 if err != nil {
137 t.Fatal(err)
138 }
139 _, err = log.Write([]byte("invalid frame header"))
140 _ = log.Close()
141 if err != nil {
142 t.Fatal(err)
143 }
144 _, err = waitHistoryWindow(t, s.Query(), r.Ref(), HistoryWindowRequest{Anchor: "newest"})
145 if !errors.Is(err, ErrDamagedStore) {
146 t.Fatalf("corruption error=%v", err)
147 }
148 }
149
150 func TestHistoryScanStopsAtCapturedSnapshot(t *testing.T) {
151 _, r, dir, _ := historyBoundaryFixture(t)
152 log, err := os.Open(filepath.Join(dir, "events.frames"))
153 if err != nil {
154 t.Fatal(err)
155 }
156 defer log.Close()
157 before, err := log.Stat()
158 if err != nil {
159 t.Fatal(err)
160 }
161 visited := 0
162 progress, err := scanHistoryLog(t.Context(), log, 0, 1, before.Size(), contentStoreForSessionDir(dir), func(commit Commit) bool {
163 visited++
164 appendRecoveryTestMessage(t, r.Session(), "during-scan", "during scan")
165 if _, err := r.Session().Flush(t.Context()); err != nil {
166 t.Fatal(err)
167 }
168 return true
169 })
170 if err != nil {
171 t.Fatal(err)
172 }
173 if visited != 1 || progress.end != before.Size() || progress.sequence != 1 {
174 t.Fatalf("scan followed concurrent append: visited=%d progress=%+v", visited, progress)
175 }
176 }
177
178 func TestHistoryScanRejectsReplacedFile(t *testing.T) {
179 _, _, dir, _ := historyBoundaryFixture(t)
180 path := filepath.Join(dir, "events.frames")
181 generation, err := historyProjectionGeneration(dir, 0)
182 if err != nil {
183 t.Fatal(err)
184 }
185 data, err := os.ReadFile(path)
186 if err != nil {
187 t.Fatal(err)
188 }
189 // Model the reader retained after a replacement: identical bytes, but a
190 // different file identity from the manifest's current log. Windows prevents
191 // renaming over an open file, so construct that state without a live rename.
192 previousPath := filepath.Join(t.TempDir(), "previous.frames")
193 if err := os.WriteFile(previousPath, data, 0o600); err != nil {
194 t.Fatal(err)
195 }
196 log, err := os.Open(previousPath)
197 if err != nil {
198 t.Fatal(err)
199 }
200 defer log.Close()
201 if err := validateHistoryLog(t.Context(), dir, log, generation); !errors.Is(err, ErrStaleGeneration) {
202 t.Fatalf("replacement accepted: %v", err)
203 }
204 }
205
205 lines GO