返回 DeepSeek-Reasonix
rebuild_resume.go
根目录 / internal / projectiondb / rebuild_resume.go
1 package projectiondb
2
3 import (
4 "context"
5 "database/sql"
6 "errors"
7 )
8
9 // Called only while the existing cross-process rebuild lock is held. The
10 // stable staging name is never opened by ordinary readers or old Rebuild calls.
11 func matchRebuildResumeKey(ctx context.Context, db *sql.DB, key string) (bool, error) {
12 if _, err := db.ExecContext(ctx, `CREATE TABLE IF NOT EXISTS projection_rebuild_resume(
13 id INTEGER PRIMARY KEY CHECK(id=1), source_key TEXT NOT NULL)`); err != nil {
14 return false, err
15 }
16 var stored string
17 err := db.QueryRowContext(ctx, `SELECT source_key FROM projection_rebuild_resume WHERE id=1`).Scan(&stored)
18 if errors.Is(err, sql.ErrNoRows) {
19 _, err = db.ExecContext(ctx, `INSERT INTO projection_rebuild_resume VALUES(1,?)`, key)
20 // A missing fence cannot certify existing staging rows. Rebuild will
21 // discard them before installing the key in a fresh database.
22 return false, err
23 }
24 return stored == key, err
25 }
26
27 func resumeRebuildReplacement(ctx context.Context, handle *Handle, replacement OpenOptions, cleanupTemporary func()) (*Handle, error) {
28 matching, resumeErr := matchRebuildResumeKey(ctx, handle.DB, replacement.ResumeKey)
29 if resumeErr != nil {
30 _ = handle.DB.Close()
31 if ctx.Err() == nil {
32 cleanupTemporary()
33 }
34 return nil, resumeErr
35 }
36 if !matching {
37 _ = handle.DB.Close()
38 cleanupTemporary()
39 var err error
40 handle, err = Open(ctx, replacement)
41 if err != nil {
42 return nil, err
43 }
44 if _, err := matchRebuildResumeKey(ctx, handle.DB, replacement.ResumeKey); err != nil {
45 _ = handle.DB.Close()
46 cleanupTemporary()
47 return nil, err
48 }
49 }
50 return handle, nil
51 }
52
52 lines GO