internal/store/store.go
306 lines · 9351 bytes
1// Package store owns SQLite access and schema migrations.
2package store
3
4import (
5 "context"
6 "database/sql"
7 "embed"
8 "errors"
9 "fmt"
10 "io/fs"
11 "log/slog"
12 "os"
13 "sort"
14 "strconv"
15 "strings"
16 "sync"
17
18 "gitbay.org/gitbay/internal/seal"
19 "modernc.org/sqlite"
20)
21
22//go:embed migrations/*.sql
23var migrationFS embed.FS
24
25type Store struct {
26 DB *sql.DB
27
28 // logWait holds one channel per build someone is following, closed
29 // by the next change to that build's row (BuildLogWait).
30 logMu sync.Mutex
31 logWait map[int64]chan struct{}
32
33 // onRevoke runs after each key revocation this process commits.
34 revokeMu sync.Mutex
35 onRevoke []func(Revoked)
36
37 // AuditJournal, when set, receives a copy of every audit row. The
38 // daemon sets it to its own logger, whose output the service
39 // journal keeps outside the database.
40 AuditJournal *slog.Logger
41 // secrets seals and opens the secret columns (secrets.go). Nil
42 // stores values as given and refuses to open sealed ones.
43 secrets *seal.Keyring
44}
45
46// Open opens (creating if needed) the database at path with WAL mode and
47// foreign keys enforced. Use ":memory:" in tests.
48//
49// _txlock=immediate is what serialises writers. Every transaction in this
50// package writes, and a deferred one takes the write lock only when it
51// reaches its first write — by which point another writer may hold it.
52// SQLite answers that with SQLITE_BUSY and does not invoke the busy
53// handler, because waiting would deadlock two transactions each holding a
54// read lock the other needs; busy_timeout cannot help. Measured with
55// eight concurrent read-then-write transactions, 44% of them failed.
56// Beginning immediate takes the write lock up front, where busy_timeout
57// does apply, so a second writer waits its turn: the same load runs with
58// no failures, and readers, which WAL keeps out of the way, are
59// unaffected (#121).
60func Open(path string) (*Store, error) {
61 // synchronous(FULL): a commit is durable before it returns, which
62 // secret key rotation needs before it drops the old keys (#273).
63 dsn := path + "?_txlock=immediate&_pragma=journal_mode(WAL)&_pragma=synchronous(FULL)&_pragma=foreign_keys(ON)&_pragma=busy_timeout(5000)"
64 if path == ":memory:" {
65 dsn = ":memory:?_txlock=immediate&_pragma=foreign_keys(ON)"
66 }
67 db, err := sql.Open("sqlite", dsn)
68 if err != nil {
69 return nil, err
70 }
71 if err := db.Ping(); err != nil {
72 db.Close()
73 return nil, err
74 }
75 // SQLite creates the file 0666&~umask, so it lands 0644 by default. The
76 // directory above it is the real boundary, but the file holds token
77 // hashes, addresses and private repo names and has no business being
78 // world-readable on its own.
79 if path != ":memory:" {
80 if err := os.Chmod(path, 0o640); err != nil && !errors.Is(err, fs.ErrNotExist) {
81 db.Close()
82 return nil, err
83 }
84 }
85 return &Store{DB: db}, nil
86}
87
88func (s *Store) Close() error { return s.DB.Close() }
89
90type migration struct {
91 version int
92 name string
93 up string
94 down string
95 // upFKOff and downFKOff are true when the up/down script's first line
96 // is the directive "-- foreign_keys: off".
97 upFKOff bool
98 downFKOff bool
99}
100
101// fkOffDirective, as the first line of a migration script, opts that
102// direction out of foreign-key enforcement for its step.
103const fkOffDirective = "-- foreign_keys: off"
104
105func loadMigrations() ([]migration, error) {
106 entries, err := fs.ReadDir(migrationFS, "migrations")
107 if err != nil {
108 return nil, err
109 }
110 byVersion := map[int]*migration{}
111 for _, e := range entries {
112 name := e.Name()
113 // <version>_<name>.<up|down>.sql
114 base, ok := strings.CutSuffix(name, ".sql")
115 if !ok {
116 return nil, fmt.Errorf("migration %q: not .sql", name)
117 }
118 var dir string
119 if b, ok := strings.CutSuffix(base, ".up"); ok {
120 base, dir = b, "up"
121 } else if b, ok := strings.CutSuffix(base, ".down"); ok {
122 base, dir = b, "down"
123 } else {
124 return nil, fmt.Errorf("migration %q: missing .up/.down", name)
125 }
126 verStr, rest, ok := strings.Cut(base, "_")
127 if !ok {
128 return nil, fmt.Errorf("migration %q: missing version prefix", name)
129 }
130 ver, err := strconv.Atoi(verStr)
131 if err != nil {
132 return nil, fmt.Errorf("migration %q: bad version: %w", name, err)
133 }
134 m := byVersion[ver]
135 if m == nil {
136 m = &migration{version: ver, name: rest}
137 byVersion[ver] = m
138 }
139 sqlBytes, err := migrationFS.ReadFile("migrations/" + name)
140 if err != nil {
141 return nil, err
142 }
143 text := string(sqlBytes)
144 firstLine, _, _ := strings.Cut(text, "\n")
145 fkOff := strings.TrimSpace(firstLine) == fkOffDirective
146 if dir == "up" {
147 m.up = text
148 m.upFKOff = fkOff
149 } else {
150 m.down = text
151 m.downFKOff = fkOff
152 }
153 }
154 var ms []migration
155 for _, m := range byVersion {
156 if m.up == "" || m.down == "" {
157 return nil, fmt.Errorf("migration %d %q: missing up or down file", m.version, m.name)
158 }
159 ms = append(ms, *m)
160 }
161 sort.Slice(ms, func(i, j int) bool { return ms[i].version < ms[j].version })
162 for i, m := range ms {
163 if m.version != i+1 {
164 return nil, fmt.Errorf("migration versions not contiguous at %d", m.version)
165 }
166 }
167 return ms, nil
168}
169
170// Version returns the current schema version (0 = empty database).
171func (s *Store) Version() (int, error) {
172 var v int
173 err := s.DB.QueryRow("PRAGMA user_version").Scan(&v)
174 return v, err
175}
176
177// MigrateUp applies all pending migrations.
178func (s *Store) MigrateUp() error { return s.migrateTo(-1) }
179
180// MigrateTo migrates up or down to the given version. 0 empties the schema.
181func (s *Store) MigrateTo(target int) error { return s.migrateTo(target) }
182
183func (s *Store) migrateTo(target int) error {
184 ms, err := loadMigrations()
185 if err != nil {
186 return err
187 }
188 if target < 0 {
189 target = len(ms)
190 }
191 if target > len(ms) {
192 return fmt.Errorf("no such schema version %d (max %d)", target, len(ms))
193 }
194 cur, err := s.Version()
195 if err != nil {
196 return err
197 }
198 for cur < target {
199 m := ms[cur]
200 if err := s.migrateStep(m.up, m.version, m.upFKOff); err != nil {
201 return fmt.Errorf("migration %d up: %w", m.version, err)
202 }
203 cur = m.version
204 }
205 for cur > target {
206 m := ms[cur-1]
207 if err := s.migrateStep(m.down, m.version-1, m.downFKOff); err != nil {
208 return fmt.Errorf("migration %d down: %w", m.version, err)
209 }
210 cur = m.version - 1
211 }
212 return nil
213}
214
215func (s *Store) migrateStep(sqlText string, newVersion int, fkOff bool) (retErr error) {
216 if !fkOff {
217 tx, err := s.DB.Begin()
218 if err != nil {
219 return err
220 }
221 defer tx.Rollback()
222 if _, err := tx.Exec(sqlText); err != nil {
223 return err
224 }
225 if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", newVersion)); err != nil {
226 return err
227 }
228 return tx.Commit()
229 }
230
231 // A script whose first line is "-- foreign_keys: off" rebuilds a
232 // table that other tables reference (labels, milestones): with
233 // foreign keys on, the rebuild-by-rename loses the children's
234 // rows. PRAGMA foreign_keys is a no-op inside a transaction, and
235 // the pool gives no guarantee that a pragma set on one connection
236 // is seen by the connection Begin() draws next, so the whole step
237 // — pragma off, transaction, foreign_key_check, commit, pragma on —
238 // runs on a single pinned connection. The check runs before commit:
239 // checking after would report a violation once the bad schema and
240 // user_version were already persisted.
241 ctx := context.Background()
242 conn, err := s.DB.Conn(ctx)
243 if err != nil {
244 return err
245 }
246 defer conn.Close()
247 if _, err := conn.ExecContext(ctx, "PRAGMA foreign_keys = OFF"); err != nil {
248 return err
249 }
250 // The connection goes back to the pool when this returns, so every
251 // path out of here has to put foreign keys back on first.
252 defer func() {
253 if _, err := conn.ExecContext(ctx, "PRAGMA foreign_keys = ON"); err != nil && retErr == nil {
254 retErr = err
255 }
256 }()
257 tx, err := conn.BeginTx(ctx, nil)
258 if err != nil {
259 return err
260 }
261 defer tx.Rollback()
262 if _, err := tx.Exec(sqlText); err != nil {
263 return err
264 }
265 if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", newVersion)); err != nil {
266 return err
267 }
268 // foreign_key_check works with enforcement off: it inspects the data
269 // directly rather than consulting the pragma. Running it here, inside
270 // the transaction, means a violation rolls back the whole rebuild
271 // (the deferred tx.Rollback fires) instead of leaving the bad schema
272 // and version committed.
273 rows, err := tx.QueryContext(ctx, "PRAGMA foreign_key_check")
274 if err != nil {
275 return err
276 }
277 if rows.Next() {
278 var table string
279 var rowid sql.NullInt64
280 var referredTable string
281 var fkid int
282 if err := rows.Scan(&table, &rowid, &referredTable, &fkid); err != nil {
283 rows.Close()
284 return err
285 }
286 rows.Close()
287 return fmt.Errorf("foreign_key_check failed after migration: %s row %v", table, rowid)
288 }
289 if err := rows.Err(); err != nil {
290 rows.Close()
291 return err
292 }
293 rows.Close()
294 return tx.Commit()
295}
296
297// IsInternal reports whether err is the database or the I/O beneath it
298// failing, as opposed to a sentinel or a message about the caller's
299// input. Callers map it to a failure exit rather than a usage error.
300func IsInternal(err error) bool {
301 var sqlErr *sqlite.Error
302 var pathErr *fs.PathError
303 return errors.As(err, &sqlErr) || errors.As(err, &pathErr) ||
304 errors.Is(err, sql.ErrTxDone) || errors.Is(err, sql.ErrConnDone) ||
305 errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled)
306}