| 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 |