返回 DeepSeek-Reasonix
read_snapshot_store_test.go
根目录 / desktop / read_snapshot_store_test.go
1 package main
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "strings"
9 "testing"
10 "time"
11 )
12
13 func mustListProjectTree(t *testing.T, a *App) []ProjectNode {
14 t.Helper()
15 nodes, err := a.ListProjectTree()
16 if err != nil {
17 t.Fatal(err)
18 }
19 return nodes
20 }
21
22 func TestReadSnapshotFrozenWindows(t *testing.T) {
23 for _, count := range []int{205, 405, 1000} {
24 t.Run(fmt.Sprint(count), func(t *testing.T) {
25 var store readSnapshotStore
26 defer store.close()
27 snap, err := store.build(t.Context(), "query", func(ctx context.Context, snap *readSnapshot) error {
28 for n := range count {
29 if err := store.append(ctx, snap, n); err != nil {
30 return err
31 }
32 }
33 return nil
34 })
35 if err != nil {
36 t.Fatal(err)
37 }
38 cursor, total := "", 0
39 for {
40 var values []int
41 next, id, _, _, err := store.page(t.Context(), "query", cursor, snap, 200, func(b []byte) error {
42 var n int
43 if err := json.Unmarshal(b, &n); err != nil {
44 return err
45 }
46 values = append(values, n)
47 return nil
48 })
49 if err != nil || id != snap.id {
50 t.Fatalf("page: %s %v", id, err)
51 }
52 for _, n := range values {
53 if n != total {
54 t.Fatalf("row=%d want %d", n, total)
55 }
56 total++
57 }
58 if next == "" {
59 break
60 }
61 cursor = next
62 }
63 if total != count {
64 t.Fatalf("rows=%d", total)
65 }
66 if _, _, _, _, err := store.page(t.Context(), "other", cursor, snap, 1, func([]byte) error { return nil }); err == nil {
67 t.Fatal("cross-query cursor accepted")
68 }
69 })
70 }
71 }
72
73 func TestReadSnapshotSpillExpiryAndRelease(t *testing.T) {
74 var store readSnapshotStore
75 defer store.close()
76 snap, err := store.build(t.Context(), "spill", func(ctx context.Context, snap *readSnapshot) error {
77 for range 600 {
78 if err := store.append(ctx, snap, strings.Repeat("x", 1024)); err != nil {
79 return err
80 }
81 }
82 return nil
83 })
84 if err != nil {
85 t.Fatal(err)
86 }
87 if snap.data.db == nil || snap.data.memory != 0 {
88 t.Fatal("snapshot did not spill")
89 }
90 next, _, _, _, err := store.page(t.Context(), "spill", "", snap, 200, func([]byte) error { return nil })
91 if err != nil {
92 t.Fatal(err)
93 }
94 store.mu.Lock()
95 snap.lifetime.used = time.Now().Add(-readSnapshotIdle)
96 store.mu.Unlock()
97 if _, _, _, _, err := store.page(t.Context(), "spill", next, nil, 200, func([]byte) error { return nil }); err == nil {
98 t.Fatal("expired snapshot accepted")
99 }
100 store.dispose(snap)
101 store.close()
102 if store.disk != 0 || store.memory != 0 {
103 t.Fatal("release leaked storage")
104 }
105 }
106
107 func TestReadSnapshotCancelledWaiterLeavesBuildAlive(t *testing.T) {
108 var store readSnapshotStore
109 defer store.close()
110 entered, release := make(chan struct{}), make(chan struct{})
111 ctx, cancel := context.WithCancel(t.Context())
112 done := make(chan error, 1)
113 go func() {
114 _, err := store.build(ctx, "same", func(ctx context.Context, snap *readSnapshot) error {
115 close(entered)
116 <-release
117 return store.append(ctx, snap, 42)
118 })
119 done <- err
120 }()
121 <-entered
122 cancel()
123 if err := <-done; !errors.Is(err, context.Canceled) {
124 t.Fatalf("cancel=%v", err)
125 }
126 store.mu.Lock()
127 job := store.building["same"]
128 store.mu.Unlock()
129 close(release)
130 <-job.done
131 if job.err != nil || job.snapshot == nil {
132 t.Fatalf("shared result=%v %v", job.snapshot, job.err)
133 }
134 }
135
136 func TestReadSnapshotSharedReleaseAndShutdown(t *testing.T) {
137 app := NewApp()
138 store := &app.desktopSessions.readSnapshots
139 entered, finish := make(chan struct{}), make(chan struct{})
140 done := make(chan *readSnapshot, 2)
141 fill := func(ctx context.Context, snap *readSnapshot) error {
142 close(entered)
143 <-finish
144 return store.append(ctx, snap, 42)
145 }
146 go func() {
147 snap, err := store.build(t.Context(), "shared", fill)
148 if err != nil {
149 t.Error(err)
150 }
151 done <- snap
152 }()
153 <-entered
154 // Done is consulted after joining the job and releasing the manager lock.
155 secondStarted := make(chan struct{})
156 go func() {
157 snap, err := store.build(snapshotWaitContext{t.Context(), secondStarted}, "shared", func(ctx context.Context, snap *readSnapshot) error { return store.append(ctx, snap, 42) })
158 if err != nil {
159 t.Error(err)
160 }
161 done <- snap
162 }()
163 <-secondStarted
164 close(finish)
165 a, b := <-done, <-done
166 if a == nil || b == nil || a.id == b.id || a.data != b.data {
167 t.Fatal("readers need independent handles")
168 }
169 app.ReleaseReadSnapshot(a.id)
170 app.ReleaseReadSnapshot(a.id)
171 if _, _, _, _, err := store.page(t.Context(), "shared", "", b, 1, func([]byte) error { return nil }); err != nil {
172 t.Fatalf("other reader was released: %v", err)
173 }
174 app.ReleaseReadSnapshot(b.id)
175 store.close()
176 store.close()
177 if store.memory != 0 || store.disk != 0 {
178 t.Fatal("shutdown leaked resources")
179 }
180 if _, err := store.build(t.Context(), "closed", fill); err == nil {
181 t.Fatal("closed owner restarted")
182 }
183 }
184
185 type snapshotWaitContext struct {
186 context.Context
187 joined chan struct{}
188 }
189
190 func (c snapshotWaitContext) Done() <-chan struct{} {
191 close(c.joined)
192 return c.Context.Done()
193 }
194
195 func TestReadSnapshotExpiredDiskReclaimedBeforeAdmission(t *testing.T) {
196 var store readSnapshotStore
197 defer store.close()
198 fill := func(ctx context.Context, snap *readSnapshot) error {
199 return store.append(ctx, snap, strings.Repeat("x", 600*1024))
200 }
201 for i := range 4 {
202 if _, err := store.build(t.Context(), fmt.Sprint(i), fill); err != nil {
203 t.Fatal(err)
204 }
205 }
206 store.mu.Lock()
207 reserved := store.disk
208 for _, snap := range store.entries {
209 snap.lifetime.used = time.Now().Add(-readSnapshotIdle)
210 }
211 store.mu.Unlock()
212 if reserved != readSnapshotDisk {
213 t.Fatalf("reservation=%d want %d", reserved, readSnapshotDisk)
214 }
215 if _, err := store.build(t.Context(), "after-expiry", fill); err != nil {
216 t.Fatalf("expired reservations blocked admission: %v", err)
217 }
218 store.close()
219 if store.disk != 0 || store.memory != 0 {
220 t.Fatal("expired storage leaked")
221 }
222 }
223
224 func TestReadSnapshotShutdownWaitsForActiveAndQueuedBuilders(t *testing.T) {
225 var store readSnapshotStore
226 defer store.close()
227 entered := make(chan struct{}, 2)
228 done := make(chan error, 3)
229 fill := func(ctx context.Context, snap *readSnapshot) error {
230 if err := store.append(ctx, snap, strings.Repeat("x", 600*1024)); err != nil {
231 return err
232 }
233 entered <- struct{}{}
234 <-ctx.Done()
235 return ctx.Err()
236 }
237 for i := range 2 {
238 go func() { _, err := store.build(t.Context(), fmt.Sprint(i), fill); done <- err }()
239 }
240 <-entered
241 <-entered
242 joined := make(chan struct{})
243 go func() { _, err := store.build(snapshotWaitContext{t.Context(), joined}, "queued", fill); done <- err }()
244 <-joined
245 store.close()
246 for range 3 {
247 if err := <-done; err == nil {
248 t.Fatal("shutdown published an unfinished snapshot")
249 }
250 }
251 if store.disk != 0 || store.memory != 0 || len(store.entries) != 0 || len(store.building) != 0 {
252 t.Fatal("shutdown returned before builder cleanup")
253 }
254 }
255
256 func TestReadSnapshotResourceFailureDoesNotPublishPrefix(t *testing.T) {
257 var store readSnapshotStore
258 defer store.close()
259 store.disk = readSnapshotDisk
260 _, err := store.build(t.Context(), "large", func(ctx context.Context, snap *readSnapshot) error {
261 for range 600 {
262 if err := store.append(ctx, snap, strings.Repeat("x", 1024)); err != nil {
263 return err
264 }
265 }
266 return nil
267 })
268 if err == nil || len(store.entries) != 0 || store.memory != 0 {
269 t.Fatalf("failed build leaked/published: err=%v memory=%d entries=%d", err, store.memory, len(store.entries))
270 }
271 store.disk = 0
272 }
273
274 func TestReadSnapshotEvictionKeepsBoundAndRejectsOldHandle(t *testing.T) {
275 var store readSnapshotStore
276 defer store.close()
277 var first *readSnapshot
278 for i := range 70 {
279 snap, err := store.build(t.Context(), fmt.Sprint(i), func(ctx context.Context, snap *readSnapshot) error { return store.append(ctx, snap, 42) })
280 if err != nil {
281 t.Fatal(err)
282 }
283 if i == 0 {
284 first = snap
285 }
286 if len(store.entries) > 64 {
287 t.Fatal("handle bound exceeded")
288 }
289 }
290 if _, _, _, _, err := store.page(t.Context(), "0", "", first, 1, func([]byte) error { return nil }); err == nil {
291 t.Fatal("evicted handle accepted")
292 }
293 store.close()
294 if store.memory != 0 || store.disk != 0 {
295 t.Fatal("eviction leaked storage")
296 }
297 }
298
299 func TestReadSnapshotEvictionOrdersEqualTimestampsAndPageAccess(t *testing.T) {
300 var store readSnapshotStore
301 defer store.close()
302 var ordered []*readSnapshot
303 for i := range 32 {
304 handle, err := store.build(t.Context(), fmt.Sprint(i), func(ctx context.Context, snap *readSnapshot) error {
305 return store.append(ctx, snap, i)
306 })
307 if err != nil {
308 t.Fatal(err)
309 }
310 ordered = append(ordered, handle.data, handle)
311 }
312 // Reading the oldest handle must refresh its position, even when every
313 // clock sample has the same value (as on a coarse platform clock).
314 first := ordered[1]
315 if _, _, _, _, err := store.page(t.Context(), "0", "", first, 1, func([]byte) error { return nil }); err != nil {
316 t.Fatal(err)
317 }
318 ordered = append(append(ordered[:1:1], ordered[2:]...), first)
319 sharedTime := time.Now()
320 store.mu.Lock()
321 for _, snap := range store.entries {
322 snap.lifetime.used = sharedTime
323 }
324 var victims []*readSnapshot
325 for range ordered {
326 victims = append(victims, store.evictOldestLocked())
327 }
328 store.mu.Unlock()
329 for i, victim := range victims {
330 store.dispose(victim)
331 if victim != ordered[i] {
332 t.Errorf("eviction %d did not follow insertion and page-access order", i)
333 }
334 }
335 }
336
336 lines GO