返回 DeepSeek-Reasonix
read_view.go
根目录 / internal / sessioncatalog / read_view.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "database/sql"
6 "errors"
7 "net/url"
8 "sync"
9 "time"
10
11 "reasonix/internal/sqliteuri"
12 )
13
14 type readViewKey struct{}
15 type readView struct {
16 owner *Catalog
17 tx *sql.Tx
18 revision uint64
19 now time.Time
20 }
21 type queryReader interface {
22 QueryContext(context.Context, string, ...any) (*sql.Rows, error)
23 QueryRowContext(context.Context, string, ...any) *sql.Row
24 }
25
26 // ReadLease retains an immutable WAL view for bounded, demand-driven paging.
27 // It uses a dedicated read-only connection so a long-lived UI cursor cannot
28 // exhaust the catalog writer's pool. The owner must release it with its cursor.
29 type ReadLease struct {
30 view *readView
31 db *sql.DB
32 cancel context.CancelFunc
33 stop func() bool
34 once sync.Once
35 }
36
37 var ErrReadLeaseUnavailable = errors.New("persistent catalog read lease unavailable")
38
39 func (c *Catalog) OpenReadLease(ctx context.Context) (_ *ReadLease, result error) {
40 defer func() { c.observeDatabaseError(result) }()
41 if err := c.readable(); err != nil {
42 return nil, err
43 }
44 status := c.Status()
45 if status.Mode != ModeDisk || status.Path == "" {
46 return nil, ErrReadLeaseUnavailable
47 }
48 dsn, err := sqliteuri.Disk(status.Path, url.Values{"mode": {"ro"}, "_pragma": {"busy_timeout(150)"}})
49 if err != nil {
50 return nil, err
51 }
52 db, err := sql.Open("sqlite", dsn)
53 if err != nil {
54 return nil, err
55 }
56 db.SetMaxOpenConns(1)
57 lifetime, cancel := context.WithCancel(ctx)
58 stop := context.AfterFunc(c.workerCtx, cancel)
59 tx, err := db.BeginTx(lifetime, &sql.TxOptions{ReadOnly: true})
60 if err != nil {
61 stop()
62 cancel()
63 db.Close()
64 return nil, err
65 }
66 var revision uint64
67 if err := tx.QueryRowContext(lifetime, `SELECT revision FROM catalog_state WHERE id=1`).Scan(&revision); err != nil {
68 _ = tx.Rollback()
69 stop()
70 cancel()
71 db.Close()
72 return nil, err
73 }
74 lease := &ReadLease{view: &readView{c, tx, revision, c.opts.Now()}, db: db, cancel: cancel, stop: stop}
75 if !c.registerReadLease(lease) {
76 lease.Close()
77 return nil, ErrCatalogInvalidated
78 }
79 return lease, nil
80 }
81
82 func (l *ReadLease) Context(ctx context.Context) context.Context {
83 return context.WithValue(ctx, readViewKey{}, l.view)
84 }
85
86 func (l *ReadLease) Close() {
87 l.once.Do(func() {
88 l.stop()
89 l.cancel()
90 _ = l.view.tx.Rollback()
91 l.db.Close()
92 c := l.view.owner
93 c.readLeasesMu.Lock()
94 delete(c.readLeases, l)
95 c.readLeasesMu.Unlock()
96 })
97 }
98
99 // WithReadView pins all nested list reads to one SQLite read transaction. No
100 // writer mutex is held and callers must finish materialization before return.
101 func (c *Catalog) WithReadView(ctx context.Context, visit func(context.Context) error) error {
102 if v, _ := ctx.Value(readViewKey{}).(*readView); v != nil && v.owner == c {
103 return visit(ctx)
104 }
105 tx, err := c.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
106 if err != nil {
107 return err
108 }
109 defer func() { _ = tx.Rollback() }()
110 // The first read establishes the WAL view, even for an empty result.
111 var revision uint64
112 if err := tx.QueryRowContext(ctx, `SELECT revision FROM catalog_state WHERE id=1`).Scan(&revision); err != nil {
113 return err
114 }
115 v := &readView{c, tx, revision, c.opts.Now()}
116 return visit(context.WithValue(ctx, readViewKey{}, v))
117 }
118
119 func (c *Catalog) readDB(ctx context.Context) catalogReader {
120 if v, _ := ctx.Value(readViewKey{}).(*readView); v != nil && v.owner == c {
121 return catalogReader{c, v.tx}
122 }
123 return catalogReader{c, c.db}
124 }
125 func (c *Catalog) readRevision(ctx context.Context) uint64 {
126 if v, _ := ctx.Value(readViewKey{}).(*readView); v != nil && v.owner == c {
127 return v.revision
128 }
129 return c.revision.Load()
130 }
131 func (c *Catalog) readTime(ctx context.Context) time.Time {
132 if v, _ := ctx.Value(readViewKey{}).(*readView); v != nil && v.owner == c {
133 return v.now
134 }
135 return c.opts.Now()
136 }
137
137 lines GO