返回 DeepSeek-Reasonix
read_snapshot_store.go
根目录 / desktop / read_snapshot_store.go
1 package main
2
3 import (
4 "context"
5 "crypto/rand"
6 "crypto/sha256"
7 "database/sql"
8 "encoding/base64"
9 "encoding/hex"
10 "encoding/json"
11 "fmt"
12 "os"
13 "path/filepath"
14 "sync"
15 "time"
16 )
17
18 const (
19 readSnapshotMemory = 32 << 20
20 readSnapshotDisk = 256 << 20
21 readSnapshotMax = 64 << 20
22 readSnapshotIdle = 10 * time.Minute
23 readSnapshotLife = 30 * time.Minute
24 )
25
26 // A cursor owns a frozen result, never a position in a live ordering. The store
27 // is process-local and disposable; it does not change any durable user format.
28 type readSnapshotStore struct {
29 mu sync.Mutex
30 buildWG sync.WaitGroup
31 entries map[string]*readSnapshot
32 building map[string]*readSnapshotBuild
33 gate chan struct{}
34 ctx context.Context
35 cancel context.CancelFunc
36 closed bool
37 memory, disk int64
38 accessOrder uint64
39 }
40
41 type readSnapshotBuild struct {
42 done chan struct{}
43 snapshot *readSnapshot
44 err error
45 }
46
47 type readSnapshot struct {
48 data *readSnapshot
49 leases int
50 removed bool
51 mu sync.Mutex
52 id, binding string
53 lifetime readSnapshotLifetime
54 rows [][]byte
55 db *sql.DB
56 dir string
57 count int
58 size, memory, disk int64
59 overhead int64
60 metadata json.RawMessage
61 validate func() error
62 readPage func(context.Context, int, int) ([][]byte, bool, error)
63 closeRead func()
64 released bool
65 }
66
67 type readSnapshotLifetime struct {
68 created time.Time
69 used time.Time
70 order uint64
71 }
72
73 type readSnapshotCursor struct {
74 Version int `json:"v"`
75 ID string `json:"id"`
76 Binding string `json:"q"`
77 Offset int `json:"n"`
78 }
79
80 type ReadError struct {
81 Code string `json:"code"`
82 Reason string `json:"reason"`
83 Message string `json:"message"`
84 }
85
86 type ReadSnapshotDiagnostics struct {
87 Handles int `json:"handles"`
88 ActiveBuilders int `json:"activeBuilders"`
89 PendingBuilders int `json:"pendingBuilders"`
90 ResidentBytes int64 `json:"residentBytes"`
91 ReservedDiskBytes int64 `json:"reservedDiskBytes"`
92 }
93
94 func (s *readSnapshotStore) diagnostics() *ReadSnapshotDiagnostics {
95 s.mu.Lock()
96 defer s.mu.Unlock()
97 active := len(s.gate)
98 return &ReadSnapshotDiagnostics{Handles: len(s.entries), ActiveBuilders: active, PendingBuilders: max(0, len(s.building)-active), ResidentBytes: s.memory, ReservedDiskBytes: s.disk}
99 }
100
101 func snapshotBinding(kind string, query any) string {
102 b, _ := json.Marshal([]any{kind, query})
103 d := sha256.Sum256(b)
104 return hex.EncodeToString(d[:])
105 }
106
107 func snapshotStale(reason string) error {
108 return &SessionOperationError{Code: "stale_cursor", Message: "Read snapshot unavailable (" + reason + "). Reload it.", Retryable: true, ReadReason: reason}
109 }
110
111 func (s *readSnapshotStore) initLocked() {
112 if s.entries != nil {
113 return
114 }
115 s.entries = make(map[string]*readSnapshot)
116 s.building = make(map[string]*readSnapshotBuild)
117 s.gate = make(chan struct{}, 2)
118 s.ctx, s.cancel = context.WithCancel(context.Background())
119 }
120
121 // Cancellation belongs to the waiter. A shared build is cancelled only by
122 // shutdown, so one dismissed UI cannot interrupt another reader's build.
123 func (s *readSnapshotStore) build(ctx context.Context, binding string, fill func(context.Context, *readSnapshot) error) (*readSnapshot, error) {
124 job, err := s.joinBuild(binding, fill)
125 if err != nil {
126 return nil, err
127 }
128 select {
129 case <-ctx.Done():
130 return nil, ctx.Err()
131 case <-job.done:
132 if job.err != nil {
133 return nil, job.err
134 }
135 return s.acquireBuildHandle(binding, job.snapshot)
136 }
137 }
138
139 func (s *readSnapshotStore) joinBuild(binding string, fill func(context.Context, *readSnapshot) error) (*readSnapshotBuild, error) {
140 s.mu.Lock()
141 defer s.mu.Unlock()
142 if s.closed {
143 return nil, snapshotStale("restarted")
144 }
145 s.initLocked()
146 job := s.building[binding]
147 if job != nil {
148 return job, nil
149 }
150 if len(s.building) >= 64 {
151 return nil, fmt.Errorf("read snapshot resource limit: too many builds")
152 }
153 job = &readSnapshotBuild{done: make(chan struct{})}
154 s.building[binding] = job
155 s.buildWG.Add(1)
156 go s.runBuild(binding, job, fill)
157 return job, nil
158 }
159
160 func (s *readSnapshotStore) runBuild(binding string, job *readSnapshotBuild, fill func(context.Context, *readSnapshot) error) {
161 defer s.buildWG.Done()
162 select {
163 case s.gate <- struct{}{}:
164 defer func() { <-s.gate }()
165 case <-s.ctx.Done():
166 job.err = s.ctx.Err()
167 }
168 var snap *readSnapshot
169 if job.err == nil {
170 snap, job.err = s.executeBuild(binding, fill)
171 }
172 s.publishBuild(binding, job, snap)
173 }
174
175 func (s *readSnapshotStore) executeBuild(binding string, fill func(context.Context, *readSnapshot) error) (*readSnapshot, error) {
176 // Reclaim before allocating: expired disk reservations must not prevent the
177 // build that would otherwise reclaim them.
178 s.pruneExpired()
179 var id [24]byte
180 if _, err := rand.Read(id[:]); err != nil {
181 return nil, err
182 }
183 now := time.Now()
184 snap := &readSnapshot{id: hex.EncodeToString(id[:]), binding: binding, lifetime: readSnapshotLifetime{created: now, used: now}, leases: 1}
185 if err := fill(s.ctx, snap); err != nil {
186 return snap, err
187 }
188 if snap.validate != nil {
189 return snap, snap.validate()
190 }
191 return snap, nil
192 }
193
194 func (s *readSnapshotStore) publishBuild(binding string, job *readSnapshotBuild, snap *readSnapshot) {
195 s.mu.Lock()
196 var victims []*readSnapshot
197 if job.err == nil && !s.closed {
198 victims = s.evictExpiredLocked()
199 if len(s.entries) >= 64 {
200 victims = append(victims, s.evictOldestLocked())
201 }
202 s.touchLocked(snap)
203 s.entries[snap.id] = snap
204 job.snapshot = snap
205 } else if job.err == nil {
206 job.err = snapshotStale("restarted")
207 }
208 delete(s.building, binding)
209 s.mu.Unlock()
210 for _, victim := range victims {
211 s.dispose(victim)
212 }
213 if job.err != nil && snap != nil {
214 s.dispose(snap)
215 }
216 close(job.done)
217 }
218
219 func (s *readSnapshotStore) acquireBuildHandle(binding string, root *readSnapshot) (*readSnapshot, error) {
220 var id [24]byte
221 if _, err := rand.Read(id[:]); err != nil {
222 return nil, err
223 }
224 s.mu.Lock()
225 if s.closed || root.leases <= 0 {
226 s.mu.Unlock()
227 return nil, snapshotStale("evicted")
228 }
229 root.leases++
230 handle := &readSnapshot{id: hex.EncodeToString(id[:]), binding: binding, lifetime: readSnapshotLifetime{created: root.lifetime.created}, data: root}
231 s.touchLocked(handle)
232 var victim *readSnapshot
233 if len(s.entries) >= 64 {
234 victim = s.evictOldestLocked()
235 }
236 s.entries[handle.id] = handle
237 s.mu.Unlock()
238 if victim != nil {
239 s.dispose(victim)
240 }
241 return handle, nil
242
243 }
244
245 func (s *readSnapshotStore) evictExpiredLocked() []*readSnapshot {
246 var victims []*readSnapshot
247 for id, snap := range s.entries {
248 if time.Since(snap.lifetime.used) >= readSnapshotIdle || time.Since(snap.lifetime.created) >= readSnapshotLife {
249 delete(s.entries, id)
250 victims = append(victims, snap)
251 }
252 }
253 return victims
254 }
255
256 // Access order owns LRU eviction; wall-clock samples only own expiration.
257 func (s *readSnapshotStore) touchLocked(snap *readSnapshot) {
258 s.accessOrder++
259 snap.lifetime.used = time.Now()
260 snap.lifetime.order = s.accessOrder
261 }
262
263 func (s *readSnapshotStore) evictOldestLocked() *readSnapshot {
264 var oldest *readSnapshot
265 for _, snap := range s.entries {
266 if oldest == nil || snap.lifetime.order < oldest.lifetime.order {
267 oldest = snap
268 }
269 }
270 if oldest != nil {
271 delete(s.entries, oldest.id)
272 }
273 return oldest
274 }
275
276 func (s *readSnapshotStore) pruneExpired() {
277 s.mu.Lock()
278 victims := s.evictExpiredLocked()
279 s.mu.Unlock()
280 for _, snap := range victims {
281 s.dispose(snap)
282 }
283 }
284
285 // Append is builder-only. Encoded rows are accounted before storage, including
286 // unpublished builds. A spill reserves its maximum physical SQLite size.
287 func (s *readSnapshotStore) append(ctx context.Context, snap *readSnapshot, value any) error {
288 if err := ctx.Err(); err != nil {
289 return err
290 }
291 b, err := json.Marshal(value)
292 if err != nil {
293 return err
294 }
295 cost := int64(len(b) + 32)
296 if snap.size+cost > readSnapshotMax {
297 return fmt.Errorf("read snapshot resource limit: result exceeds 64 MiB")
298 }
299 if snap.db == nil {
300 s.mu.Lock()
301 resident := s.memory+cost <= readSnapshotMemory && snap.size+cost <= 512<<10
302 if resident {
303 s.memory += cost
304 snap.memory += cost
305 }
306 s.mu.Unlock()
307 if resident {
308 snap.rows = append(snap.rows, b)
309 snap.size += cost
310 snap.count++
311 return nil
312 }
313 s.mu.Lock()
314 if s.disk+readSnapshotMax > readSnapshotDisk {
315 s.mu.Unlock()
316 return fmt.Errorf("read snapshot resource limit: temporary storage exhausted")
317 }
318 s.disk += readSnapshotMax
319 snap.disk = readSnapshotMax
320 s.mu.Unlock()
321 snap.dir, err = os.MkdirTemp("", "reasonix-read-")
322 if err != nil {
323 return err
324 }
325 snap.db, err = sql.Open("sqlite", filepath.Join(snap.dir, "rows.sqlite"))
326 if err != nil {
327 return err
328 }
329 snap.db.SetMaxOpenConns(1)
330 if _, err = snap.db.ExecContext(ctx, `PRAGMA journal_mode=OFF; PRAGMA max_page_count=16384; CREATE TABLE rows (n INTEGER PRIMARY KEY, body BLOB NOT NULL)`); err != nil {
331 return err
332 }
333 for n, row := range snap.rows {
334 if _, err = snap.db.ExecContext(ctx, `INSERT INTO rows VALUES (?,?)`, n, row); err != nil {
335 return err
336 }
337 }
338 snap.rows = nil
339 s.mu.Lock()
340 s.memory -= snap.memory - snap.overhead
341 snap.memory = snap.overhead
342 s.mu.Unlock()
343 }
344 if _, err = snap.db.ExecContext(ctx, `INSERT INTO rows VALUES (?,?)`, snap.count, b); err != nil {
345 return err
346 }
347 snap.size += cost
348 snap.count++
349 return nil
350 }
351
352 // Retained dependency metadata cannot spill, so account for it separately from
353 // encoded rows. Charges are conservative and remain resident after a spill.
354 func (s *readSnapshotStore) reserve(snap *readSnapshot, bytes int64) error {
355 s.mu.Lock()
356 defer s.mu.Unlock()
357 if s.memory+bytes > readSnapshotMemory || snap.size+bytes > readSnapshotMax {
358 return fmt.Errorf("read snapshot resource limit: dependency metadata exhausted")
359 }
360 s.memory += bytes
361 snap.memory += bytes
362 snap.overhead += bytes
363 snap.size += bytes
364 return nil
365 }
366
367 func (s *readSnapshotStore) page(ctx context.Context, binding, cursor string, first *readSnapshot, limit int, consume func([]byte) error) (string, string, int64, json.RawMessage, error) {
368 offset := 0
369 id := ""
370 if first != nil {
371 id = first.id
372 }
373 if cursor != "" {
374 var c readSnapshotCursor
375 b, err := base64.RawURLEncoding.DecodeString(cursor)
376 if err != nil || json.Unmarshal(b, &c) != nil || c.Version != 1 || c.ID == "" || c.Offset < 0 {
377 return "", "", 0, nil, snapshotStale("invalid_cursor")
378 }
379 if c.Binding != binding {
380 return "", "", 0, nil, snapshotStale("query_mismatch")
381 }
382 id, offset = c.ID, c.Offset
383 }
384 s.mu.Lock()
385 snap := s.entries[id]
386 if snap == nil {
387 s.mu.Unlock()
388 return "", "", 0, nil, snapshotStale("expired_or_evicted")
389 }
390 if time.Since(snap.lifetime.used) >= readSnapshotIdle || time.Since(snap.lifetime.created) >= readSnapshotLife {
391 delete(s.entries, id)
392 s.mu.Unlock()
393 s.dispose(snap)
394 return "", "", 0, nil, snapshotStale("expired")
395 }
396 s.touchLocked(snap)
397 s.mu.Unlock()
398 if snap.data != nil {
399 snap = snap.data
400 }
401 // Per-snapshot lock is a read lease. Disposal cannot close a live DB read.
402 snap.mu.Lock()
403 defer snap.mu.Unlock()
404 if snap.released {
405 return "", "", 0, nil, snapshotStale("evicted")
406 }
407 if snap.binding != binding || snap.readPage == nil && offset > snap.count {
408 return "", "", 0, nil, snapshotStale("invalid_cursor")
409 }
410 if snap.validate != nil {
411 if err := snap.validate(); err != nil {
412 return "", "", 0, nil, err
413 }
414 }
415 if limit <= 0 {
416 limit = 50
417 }
418 limit = min(limit, 200)
419 if snap.readPage != nil {
420 rows, more, err := snap.readPage(ctx, offset, limit)
421 if err != nil {
422 return "", "", 0, nil, err
423 }
424 for _, row := range rows {
425 if err := consume(row); err != nil {
426 return "", "", 0, nil, err
427 }
428 }
429 next := ""
430 if more {
431 if len(rows) == 0 {
432 return "", "", 0, nil, fmt.Errorf("snapshot page did not advance")
433 }
434 b, _ := json.Marshal(readSnapshotCursor{1, id, binding, offset + len(rows)})
435 next = base64.RawURLEncoding.EncodeToString(b)
436 }
437 return next, id, snap.lifetime.created.Add(readSnapshotLife).UnixMilli(), snap.metadata, nil
438 }
439 end := min(offset+limit, snap.count)
440 if err := snap.consumeStoredRows(ctx, offset, end, consume); err != nil {
441 return "", "", 0, nil, err
442 }
443 next := ""
444 if end < snap.count {
445 b, _ := json.Marshal(readSnapshotCursor{1, id, binding, end})
446 next = base64.RawURLEncoding.EncodeToString(b)
447 }
448 return next, id, snap.lifetime.created.Add(readSnapshotLife).UnixMilli(), snap.metadata, nil
449 }
450
451 func (s *readSnapshotStore) dispose(snap *readSnapshot) {
452 s.mu.Lock()
453 if snap.removed {
454 s.mu.Unlock()
455 return
456 }
457 snap.removed = true
458 if snap.data != nil {
459 snap = snap.data
460 }
461 snap.leases--
462 if snap.leases > 0 {
463 s.mu.Unlock()
464 return
465 }
466 s.mu.Unlock()
467 snap.mu.Lock()
468 defer snap.mu.Unlock()
469 if snap.released {
470 return
471 }
472 snap.released = true
473 if snap.closeRead != nil {
474 snap.closeRead()
475 snap.closeRead = nil
476 }
477 snap.readPage = nil
478 if snap.db != nil {
479 _ = snap.db.Close()
480 }
481 if snap.dir != "" {
482 _ = os.RemoveAll(snap.dir)
483 }
484 snap.rows = nil
485 snap.metadata, snap.validate = nil, nil
486 s.mu.Lock()
487 s.memory -= snap.memory
488 s.disk -= snap.disk
489 s.mu.Unlock()
490 }
491
492 // walk is builder-only and permits streaming conversion of captured candidates
493 // after the source read transaction has closed.
494 func (snap *readSnapshot) walk(ctx context.Context, consume func([]byte) error) error {
495 if snap.db == nil {
496 for _, row := range snap.rows {
497 if err := consume(row); err != nil {
498 return err
499 }
500 }
501 return nil
502 }
503 rows, err := snap.db.QueryContext(ctx, `SELECT body FROM rows ORDER BY n`)
504 if err != nil {
505 return err
506 }
507 defer rows.Close()
508 for rows.Next() {
509 var row []byte
510 if err := rows.Scan(&row); err != nil {
511 return err
512 }
513 if err := consume(row); err != nil {
514 return err
515 }
516 }
517 return rows.Err()
518 }
519
520 func (a *App) ReleaseReadSnapshot(id string) {
521 s := &a.desktopSessions.readSnapshots
522 s.mu.Lock()
523 snap := s.entries[id]
524 delete(s.entries, id)
525 var cached *readSnapshot
526 if snap != nil && snap.data != nil {
527 cached = s.entries[snap.data.id]
528 delete(s.entries, snap.data.id)
529 }
530 s.mu.Unlock()
531 if cached != nil {
532 s.dispose(cached)
533 }
534 if snap != nil {
535 s.dispose(snap)
536 }
537 }
538
539 // The caller holds the snapshot read lease while reading memory or spilled rows.
540 func (snap *readSnapshot) consumeStoredRows(ctx context.Context, offset, end int, consume func([]byte) error) error {
541 if snap.db == nil {
542 for _, row := range snap.rows[offset:end] {
543 if err := consume(row); err != nil {
544 return err
545 }
546 }
547 } else {
548 rows, err := snap.db.QueryContext(ctx, `SELECT body FROM rows WHERE n>=? AND n<? ORDER BY n`, offset, end)
549 if err != nil {
550 return err
551 }
552 defer rows.Close()
553 for rows.Next() {
554 var b []byte
555 if err := rows.Scan(&b); err != nil {
556 return err
557 }
558 if err := consume(b); err != nil {
559 return err
560 }
561 }
562 if err := rows.Err(); err != nil {
563 return err
564 }
565 }
566 return nil
567 }
568
569 func (s *readSnapshotStore) close() {
570 s.mu.Lock()
571 s.closed = true
572 if s.cancel != nil {
573 s.cancel()
574 }
575 entries := s.entries
576 s.entries = nil
577 jobs := make([]*readSnapshotBuild, 0, len(s.building))
578 for _, job := range s.building {
579 jobs = append(jobs, job)
580 }
581 s.mu.Unlock()
582 s.buildWG.Wait()
583 for _, job := range jobs {
584 <-job.done
585 }
586 for _, snap := range entries {
587 s.dispose(snap)
588 }
589 }
590
590 lines GO