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