返回 DeepSeek-Reasonix
projectiondb.go
根目录 / internal / projectiondb / projectiondb.go
1 // Package projectiondb owns the lifecycle of disposable SQLite projections.
2 // Business data must remain authoritative outside the database: callers are
3 // expected to be able to discard and rebuild every database opened here.
4 package projectiondb
5
6 import (
7 "context"
8 "database/sql"
9 "errors"
10 "fmt"
11 "net/url"
12 "os"
13 "path/filepath"
14 "runtime"
15 "strings"
16 "sync/atomic"
17 "time"
18
19 filelock "reasonix/internal/identitylock"
20 "reasonix/internal/sqliteuri"
21
22 moderncsqlite "modernc.org/sqlite"
23 sqlite3 "modernc.org/sqlite/lib"
24 )
25
26 type Mode string
27
28 var memoryDatabaseSequence atomic.Uint64
29
30 const (
31 ModeDisk Mode = "disk"
32 ModeMemory Mode = "memory"
33 )
34
35 type State string
36
37 const (
38 StateReady State = "ready"
39 StateDegraded State = "degraded"
40 )
41
42 type Status struct {
43 State State `json:"state"`
44 Mode Mode `json:"mode"`
45 Path string `json:"path,omitempty"`
46 Revision uint64 `json:"revision"`
47 Indexed int64 `json:"indexed"`
48 Total int64 `json:"total"`
49 Pending int64 `json:"pending"`
50 Failed int64 `json:"failed"`
51 LastError string `json:"lastError,omitempty"`
52 QuarantinedPath string `json:"quarantinedPath,omitempty"`
53 }
54
55 type Migration struct {
56 Version int
57 Apply func(context.Context, *sql.Tx) error
58 }
59
60 type OpenOptions struct {
61 Path string
62 MemoryName string
63 Migrations []Migration
64 InMemory bool
65 // RequireDisk disables the process-local memory fallback. Rebuild uses this
66 // so a failed temporary open reports the real disk error instead of
67 // "could not use disk storage" after silently opening :memory:.
68 RequireDisk bool
69 MaxOpenConns int
70 Now func() time.Time
71 SecureDelete bool
72 AutoVacuum bool
73 // RetainBackup keeps the previous database at its generated .replaced-
74 // timestamp path so disposable projections can offer a rollback point.
75 RetainBackup bool
76 // QuickCheck uses SQLite's quick_check for a disposable projection.
77 // Authoritative stores keep the full integrity_check default.
78 QuickCheck bool
79 // ResumeKey opts Rebuild into a single persistent staging database. The
80 // callback atomically commits data and position. Cancellation retains staging;
81 // a changed key starts fresh. Progress is disposable, never business authority.
82 ResumeKey string
83 }
84
85 // WALSizeLimit is the size a disk projection's WAL is truncated back to after
86 // a checkpoint.
87 const WALSizeLimit = 4 << 20
88
89 // CheckpointBeforeClose folds the WAL into the database and truncates it. It is
90 // best-effort: another process's reader leaves the WAL for its next checkpoint.
91 func CheckpointBeforeClose(ctx context.Context, db *sql.DB) {
92 if db != nil {
93 _, _ = db.ExecContext(ctx, `PRAGMA wal_checkpoint(TRUNCATE)`)
94 }
95 }
96
97 type Handle struct {
98 DB *sql.DB
99 Status Status
100 }
101
102 type FutureSchemaError struct {
103 Found int
104 Supported int
105 }
106
107 func (e *FutureSchemaError) Error() string {
108 return fmt.Sprintf("projection schema %d is newer than supported %d", e.Found, e.Supported)
109 }
110
111 func Open(ctx context.Context, opts OpenOptions) (*Handle, error) {
112 return openProjection(ctx, opts, true)
113 }
114
115 // OpenAdvisory opens disposable metadata before a full integrity audit. The
116 // owner must call CheckIntegrity in its maintenance lifecycle and invalidate
117 // all readers on corruption. This must never certify content or execution.
118 func OpenAdvisory(ctx context.Context, opts OpenOptions) (*Handle, error) {
119 return openProjection(ctx, opts, false)
120 }
121
122 func openProjection(ctx context.Context, opts OpenOptions, checkIntegrity bool) (*Handle, error) {
123 if opts.Now == nil {
124 opts.Now = time.Now
125 }
126 if opts.MemoryName == "" {
127 opts.MemoryName = "projection"
128 }
129 // Blank paths must never become relative "v1.sqlite" files under cwd.
130 if strings.TrimSpace(opts.Path) == "" {
131 opts.Path = ""
132 opts.InMemory = true
133 }
134 mode := ModeDisk
135 status := Status{State: StateReady, Mode: mode, Path: opts.Path}
136 useMemory := opts.InMemory || PathLooksRemote(opts.Path)
137 if !useMemory {
138 if err := os.MkdirAll(filepath.Dir(opts.Path), 0o700); err != nil {
139 if opts.RequireDisk {
140 return nil, fmt.Errorf("create projection directory: %w", err)
141 }
142 useMemory = true
143 status.LastError = err.Error()
144 } else {
145 _ = os.Chmod(filepath.Dir(opts.Path), 0o700)
146 if filesystemRemote(filepath.Dir(opts.Path)) {
147 if opts.RequireDisk {
148 return nil, errors.New("projection cache is on a remote filesystem")
149 }
150 useMemory = true
151 status.LastError = "projection cache is on a remote filesystem; using memory"
152 }
153 }
154 }
155 if useMemory {
156 if opts.RequireDisk {
157 return nil, errors.New("projection requires disk storage")
158 }
159 mode = ModeMemory
160 }
161
162 db, err := open(ctx, opts, mode, checkIntegrity)
163 if err != nil && mode == ModeDisk {
164 var future *FutureSchemaError
165 switch {
166 case errors.As(err, &future):
167 // A newer process wrote this projection. Keep the file intact and
168 // serve an empty memory projection so startup never quarantines a
169 // healthy future schema.
170 if opts.RequireDisk {
171 return nil, err
172 }
173 status.State = StateDegraded
174 status.LastError = err.Error()
175 mode = ModeMemory
176 db, err = open(ctx, opts, mode, checkIntegrity)
177 case isCorruptionError(err):
178 // Only integrity-level failures may rename the on-disk projection.
179 status.QuarantinedPath = Quarantine(opts.Path, opts.Now())
180 db, err = open(ctx, opts, mode, checkIntegrity)
181 if err != nil {
182 if opts.RequireDisk {
183 return nil, err
184 }
185 status.State = StateDegraded
186 status.LastError = err.Error()
187 mode = ModeMemory
188 db, err = open(ctx, opts, mode, checkIntegrity)
189 }
190 default:
191 // Busy, permission, IO, or transient open errors must never rename
192 // a healthy database. Fall back to memory for this process only.
193 if opts.RequireDisk {
194 return nil, err
195 }
196 status.State = StateDegraded
197 status.LastError = err.Error()
198 mode = ModeMemory
199 db, err = open(ctx, opts, mode, checkIntegrity)
200 }
201 }
202 if err != nil {
203 return nil, err
204 }
205 status.Mode = mode
206 if mode == ModeMemory {
207 status.Path = ""
208 }
209 return &Handle{DB: db, Status: status}, nil
210 }
211
212 func open(ctx context.Context, opts OpenOptions, mode Mode, checkIntegrity bool) (*sql.DB, error) {
213 var dsn string
214 if mode == ModeMemory {
215 // time.Now has coarse resolution on some platforms, notably Windows.
216 // A process-local sequence prevents concurrently opened projections with
217 // the same logical name from sharing one SQLite memory database by accident.
218 dsn = fmt.Sprintf("file:reasonix-%s-%d-%d?mode=memory&cache=shared", url.PathEscape(opts.MemoryName),
219 opts.Now().UnixNano(), memoryDatabaseSequence.Add(1))
220 } else {
221 var err error
222 // journal_size_limit is per connection, so it rides the DSN to reach
223 // every pooled one; without it a checkpointed WAL keeps its peak size.
224 dsn, err = sqliteuri.Disk(opts.Path, url.Values{
225 "_pragma": {"busy_timeout(150)", "foreign_keys(1)", fmt.Sprintf("journal_size_limit(%d)", WALSizeLimit)},
226 })
227 if err != nil {
228 return nil, err
229 }
230 }
231 db, err := sql.Open("sqlite", dsn)
232 if err != nil {
233 return nil, err
234 }
235 maxOpen := opts.MaxOpenConns
236 if maxOpen <= 0 {
237 maxOpen = 4
238 }
239 // Shared-cache memory databases cannot safely pool concurrent writers.
240 // Disk catalogs retain their requested pool and WAL read concurrency.
241 if mode == ModeMemory {
242 maxOpen = 1
243 }
244 db.SetMaxOpenConns(maxOpen)
245 db.SetMaxIdleConns(min(maxOpen, 2))
246 fail := func(err error) (*sql.DB, error) {
247 _ = db.Close()
248 return nil, err
249 }
250 if err := db.PingContext(ctx); err != nil {
251 return fail(err)
252 }
253 if mode == ModeDisk {
254 if _, err := db.ExecContext(ctx, `PRAGMA journal_mode=WAL`); err != nil {
255 return fail(err)
256 }
257 }
258 for _, pragma := range []string{`PRAGMA synchronous=NORMAL`, `PRAGMA foreign_keys=ON`, `PRAGMA busy_timeout=150`} {
259 if _, err := db.ExecContext(ctx, pragma); err != nil {
260 return fail(err)
261 }
262 }
263 if opts.SecureDelete {
264 // Best-effort: some builds/filesystems reject the pragma without making
265 // the projection unusable.
266 _, _ = db.ExecContext(ctx, `PRAGMA secure_delete=ON`)
267 }
268 if opts.AutoVacuum {
269 // auto_vacuum can only be changed on an empty database; ignore failures
270 // on already-initialized files so open does not degrade to memory.
271 _, _ = db.ExecContext(ctx, `PRAGMA auto_vacuum=INCREMENTAL`)
272 }
273 if checkIntegrity {
274 if err := checkDatabaseIntegrity(ctx, db, opts.QuickCheck); err != nil {
275 return fail(err)
276 }
277 }
278 if err := ApplyMigrations(ctx, db, opts.Migrations, opts.Now); err != nil {
279 return fail(err)
280 }
281 if mode == ModeDisk {
282 _ = os.Chmod(opts.Path, 0o600)
283 _ = os.Chmod(opts.Path+"-wal", 0o600)
284 _ = os.Chmod(opts.Path+"-shm", 0o600)
285 }
286 return db, nil
287 }
288
289 func ApplyMigrations(ctx context.Context, db *sql.DB, migrations []Migration, now func() time.Time) error {
290 if now == nil {
291 now = time.Now
292 }
293 if _, err := db.ExecContext(ctx, `CREATE TABLE IF NOT EXISTS schema_migrations (
294 version INTEGER PRIMARY KEY,
295 applied_at INTEGER NOT NULL
296 )`); err != nil {
297 return err
298 }
299 var current int
300 if err := db.QueryRowContext(ctx, `SELECT COALESCE(MAX(version), 0) FROM schema_migrations`).Scan(&current); err != nil {
301 return err
302 }
303 supported := 0
304 for _, migration := range migrations {
305 if migration.Version > supported {
306 supported = migration.Version
307 }
308 }
309 if current > supported {
310 return &FutureSchemaError{Found: current, Supported: supported}
311 }
312 for _, migration := range migrations {
313 if migration.Version <= current {
314 continue
315 }
316 if migration.Version != current+1 || migration.Apply == nil {
317 return fmt.Errorf("projection migration gap after version %d", current)
318 }
319 tx, err := db.BeginTx(ctx, nil)
320 if err != nil {
321 return err
322 }
323 if err := migration.Apply(ctx, tx); err != nil {
324 _ = tx.Rollback()
325 return fmt.Errorf("apply projection migration %d: %w", migration.Version, err)
326 }
327 if _, err := tx.ExecContext(ctx, `INSERT INTO schema_migrations(version, applied_at) VALUES(?, ?)`, migration.Version, now().UnixMilli()); err != nil {
328 _ = tx.Rollback()
329 return err
330 }
331 if err := tx.Commit(); err != nil {
332 return err
333 }
334 current = migration.Version
335 }
336 return nil
337 }
338
339 type Inspection struct {
340 Exists bool `json:"exists"`
341 Path string `json:"path"`
342 Schema int `json:"schema"`
343 Integrity string `json:"integrity,omitempty"`
344 Size int64 `json:"size"`
345 Error string `json:"error,omitempty"`
346 }
347
348 // Inspect is deliberately read-only: it never creates, migrates, repairs, or
349 // quarantines a database.
350 func Inspect(ctx context.Context, path string) Inspection {
351 out := Inspection{Path: path}
352 info, err := os.Stat(path)
353 if errors.Is(err, os.ErrNotExist) {
354 return out
355 }
356 if err != nil {
357 out.Error = err.Error()
358 return out
359 }
360 out.Exists = true
361 out.Size = info.Size()
362 // A live projection may hold its schema and latest commits only in WAL.
363 // immutable=1 would ignore that WAL and report a healthy database as broken.
364 dsn, err := sqliteuri.Disk(path, url.Values{
365 "_pragma": {"busy_timeout(150)", "foreign_keys(1)"},
366 "mode": {"ro"},
367 })
368 if err != nil {
369 out.Error = err.Error()
370 return out
371 }
372 db, err := sql.Open("sqlite", dsn)
373 if err != nil {
374 out.Error = err.Error()
375 return out
376 }
377 defer db.Close()
378 if err := db.QueryRowContext(ctx, `PRAGMA integrity_check`).Scan(&out.Integrity); err != nil {
379 out.Error = err.Error()
380 return out
381 }
382 if err := db.QueryRowContext(ctx, `SELECT COALESCE(MAX(version),0) FROM schema_migrations`).Scan(&out.Schema); err != nil {
383 out.Error = err.Error()
384 }
385 return out
386 }
387
388 func Quarantine(path string, now time.Time) string {
389 if strings.TrimSpace(path) == "" {
390 return ""
391 }
392 quarantined := fmt.Sprintf("%s.corrupt-%d", path, now.UnixMilli())
393 if err := os.Rename(path, quarantined); err != nil {
394 return ""
395 }
396 for _, suffix := range []string{"-wal", "-shm"} {
397 _ = os.Rename(path+suffix, quarantined+suffix)
398 }
399 return quarantined
400 }
401
402 // isCorruptionError reports whether err proves the on-disk projection is unsafe
403 // to keep open. Temporary busy/permission/IO failures must return false so a
404 // multi-process client never renames a healthy database out from under peers.
405 func isCorruptionError(err error) bool {
406 if err == nil {
407 return false
408 }
409 var future *FutureSchemaError
410 if errors.As(err, &future) {
411 return false
412 }
413 var se *moderncsqlite.Error
414 if errors.As(err, &se) {
415 switch se.Code() & 0xff {
416 case sqlite3.SQLITE_CORRUPT, sqlite3.SQLITE_NOTADB:
417 return true
418 default:
419 return false
420 }
421 }
422 msg := strings.ToLower(err.Error())
423 return strings.Contains(msg, "integrity check") ||
424 strings.Contains(msg, "malformed") ||
425 strings.Contains(msg, "file is not a database") ||
426 strings.Contains(msg, "not a database")
427 }
428
429 func PathLooksRemote(path string) bool {
430 clean := filepath.Clean(path)
431 if runtime.GOOS == "windows" && strings.HasPrefix(clean, `\\`) {
432 return true
433 }
434 slash := filepath.ToSlash(clean)
435 return strings.HasPrefix(slash, "/net/") || strings.HasPrefix(slash, "/nfs/") || strings.HasPrefix(slash, "/afs/")
436 }
437
438 // Rebuild constructs and validates a replacement beside the live database,
439 // then swaps it into place. The old projection remains untouched if building,
440 // validation, or the platform rename fails (notably an open database on
441 // Windows). Rebuild never touches authoritative business files.
442 func Rebuild(ctx context.Context, opts OpenOptions, populate func(context.Context, *sql.DB) error) error {
443 if strings.TrimSpace(opts.Path) == "" || opts.InMemory {
444 return errors.New("projection rebuild requires a disk path")
445 }
446 if err := os.MkdirAll(filepath.Dir(opts.Path), 0o700); err != nil {
447 return fmt.Errorf("create projection rebuild directory: %w", err)
448 }
449 release, err := filelock.Acquire(ctx, opts.Path+".rebuild.lock")
450 if err != nil {
451 return fmt.Errorf("lock projection rebuild: %w", err)
452 }
453 defer release()
454 if opts.Now == nil {
455 opts.Now = time.Now
456 }
457 temporary := fmt.Sprintf("%s.rebuild-%d", opts.Path, opts.Now().UnixNano())
458 if opts.ResumeKey != "" {
459 temporary = opts.Path + ".rebuild-pending"
460 }
461 replacement := opts
462 replacement.Path = temporary
463 replacement.InMemory = false
464 replacement.RequireDisk = true
465 handle, err := Open(ctx, replacement)
466 if err != nil {
467 return fmt.Errorf("open projection replacement: %w", err)
468 }
469 cleanupTemporary := func() {
470 _ = os.Remove(temporary)
471 _ = os.Remove(temporary + "-wal")
472 _ = os.Remove(temporary + "-shm")
473 }
474 if handle.Status.Mode != ModeDisk {
475 _ = handle.DB.Close()
476 cleanupTemporary()
477 detail := strings.TrimSpace(handle.Status.LastError)
478 if detail == "" {
479 detail = "unknown open fallback"
480 }
481 return fmt.Errorf("projection replacement could not use disk storage: %s", detail)
482 }
483 if opts.ResumeKey != "" {
484 handle, err = resumeRebuildReplacement(ctx, handle, replacement, cleanupTemporary)
485 if err != nil {
486 return err
487 }
488 }
489 if err := populateRebuildReplacement(ctx, opts, handle, populate, cleanupTemporary); err != nil {
490 return err
491 }
492
493 backup := fmt.Sprintf("%s.replaced-%d", opts.Path, opts.Now().UnixNano())
494 hadOld := false
495 if _, err := os.Stat(opts.Path); err == nil {
496 if err := os.Rename(opts.Path, backup); err != nil {
497 cleanupTemporary()
498 return fmt.Errorf("projection database is busy: %w", err)
499 }
500 hadOld = true
501 for _, suffix := range []string{"-wal", "-shm"} {
502 if err := os.Rename(opts.Path+suffix, backup+suffix); err != nil && !errors.Is(err, os.ErrNotExist) {
503 _ = os.Rename(backup, opts.Path)
504 for _, restored := range []string{"-wal", "-shm"} {
505 _ = os.Rename(backup+restored, opts.Path+restored)
506 }
507 cleanupTemporary()
508 return fmt.Errorf("projection database is busy: %w", err)
509 }
510 }
511 } else if !errors.Is(err, os.ErrNotExist) {
512 cleanupTemporary()
513 return err
514 }
515 if err := os.Rename(temporary, opts.Path); err != nil {
516 if hadOld {
517 _ = os.Rename(backup, opts.Path)
518 for _, suffix := range []string{"-wal", "-shm"} {
519 _ = os.Rename(backup+suffix, opts.Path+suffix)
520 }
521 }
522 cleanupTemporary()
523 return fmt.Errorf("install projection replacement: %w", err)
524 }
525 _ = os.Chmod(opts.Path, 0o600)
526 if hadOld && !opts.RetainBackup {
527 _ = os.Remove(backup)
528 _ = os.Remove(backup + "-wal")
529 _ = os.Remove(backup + "-shm")
530 }
531 cleanupTemporary()
532 return nil
533 }
534
534 lines GO