This repository has no description
1package db
2
3import (
4 "context"
5 "database/sql"
6 "log/slog"
7 "slices"
8 "strings"
9
10 _ "github.com/mattn/go-sqlite3"
11 "tangled.org/core/log"
12 "tangled.org/core/orm"
13)
14
15type DB struct {
16 *sql.DB
17}
18
19type DBTX interface {
20 QueryRow(query string, args ...any) *sql.Row
21 Exec(query string, args ...any) (sql.Result, error)
22}
23
24func Make(ctx context.Context, dbPath string) (*DB, error) {
25 // https://github.com/mattn/go-sqlite3#connection-string
26 opts := []string{
27 "_foreign_keys=1",
28 "_journal_mode=WAL",
29 "_synchronous=NORMAL",
30 "_auto_vacuum=incremental",
31 "_busy_timeout=5000",
32 }
33
34 logger := log.FromContext(ctx)
35 logger = log.SubLogger(logger, "db")
36
37 db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&"))
38 if err != nil {
39 return nil, err
40 }
41
42 conn, err := db.Conn(ctx)
43 if err != nil {
44 return nil, err
45 }
46 defer conn.Close()
47
48 _, err = conn.ExecContext(ctx, `
49 create table if not exists _jetstream (
50 id integer primary key autoincrement,
51 last_time_us integer not null
52 );
53
54 create table if not exists known_dids (
55 did text primary key
56 );
57
58 create table if not exists repos (
59 id integer primary key autoincrement,
60 knot text not null,
61 owner text not null,
62 rkey text not null,
63 repo_did text,
64 created_at text,
65 addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
66
67 unique(owner, rkey)
68 );
69
70 create table if not exists repo_collaborators (
71 id integer primary key autoincrement,
72 owner_did text not null,
73 rkey text not null,
74 subject text not null,
75 repo_did text not null,
76 addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
77
78 unique(owner_did, rkey)
79 );
80
81 create table if not exists spindle_members (
82 -- identifiers for the record
83 id integer primary key autoincrement,
84 did text not null,
85 rkey text not null,
86
87 -- data
88 instance text not null,
89 subject text not null,
90 created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
91
92 -- constraints
93 unique (did, rkey)
94 );
95
96 -- status event for a single workflow
97 create table if not exists events (
98 rkey text not null,
99 nsid text not null,
100 event text not null, -- json
101 created integer not null -- unix nanos
102 );
103
104 create table if not exists nixos_toplevel_cache (
105 config_key text primary key,
106 toplevel text not null,
107 updated_at text not null
108 );
109
110 create table if not exists pipelines (
111 id text primary key,
112 repo_did text not null,
113 commit_id text not null
114 );
115 create table if not exists jobs (
116 id integer primary key autoincrement,
117 repo_did text not null,
118 pipeline_id_knot text not null,
119 pipeline_id_rkey text not null,
120 source_repo text,
121 tpl text not null,
122 created_at integer not null default (strftime('%s', 'now'))
123 );
124
125 create table if not exists workflows (
126 id integer primary key autoincrement,
127 pipeline_id text not null,
128 name text not null,
129 status text not null default 'pending',
130
131 unique(pipeline_id, id),
132 foreign key (pipeline_id) references pipelines(id) on delete cascade
133 );
134
135 create table if not exists migrations (
136 id integer primary key autoincrement,
137 name text unique
138 );
139 `)
140 if err != nil {
141 return nil, err
142 }
143
144 if err := runMigrations(ctx, conn, logger); err != nil {
145 return nil, err
146 }
147
148 return &DB{db}, nil
149}
150
151func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error {
152 if err := orm.RunMigration(conn, logger, "repos-to-repo-did", func(tx *sql.Tx) error {
153 var hasName int
154 if err := tx.QueryRow(
155 `select count(*) from pragma_table_info('repos') where name = 'name'`,
156 ).Scan(&hasName); err != nil {
157 return err
158 }
159
160 if hasName > 0 {
161 var totalRows, copiedRows int
162 if err := tx.QueryRow(`select count(*) from repos`).Scan(&totalRows); err != nil {
163 return err
164 }
165 if err := tx.QueryRow(`select count(*) from repos where coalesce(name, '') <> ''`).Scan(&copiedRows); err != nil {
166 return err
167 }
168 if dropped := totalRows - copiedRows; dropped > 0 {
169 logger.Warn("dropping repo rows with empty name during migration", "dropped", dropped, "kept", copiedRows)
170 }
171
172 if _, err := tx.Exec(`
173 create table if not exists repos_new (
174 id integer primary key autoincrement,
175 knot text not null,
176 owner text not null,
177 rkey text not null,
178 repo_did text,
179 created_at text,
180 addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
181
182 unique(owner, rkey)
183 );
184
185 insert into repos_new (id, knot, owner, rkey, addedAt)
186 select id, knot, owner, name, addedAt from repos where coalesce(name, '') <> '';
187
188 drop table repos;
189 alter table repos_new rename to repos;
190 `); err != nil {
191 return err
192 }
193 }
194
195 _, err := tx.Exec(`
196 create index if not exists idx_repos_repo_did on repos(repo_did);
197 create index if not exists idx_repos_owner_repo_did on repos(owner, repo_did);
198 create index if not exists idx_repo_collaborators_repo_did
199 on repo_collaborators(repo_did);
200 `)
201 return err
202 }); err != nil {
203 return err
204 }
205
206 if err := orm.RunMigration(conn, logger, "spindle-members-unique-on-rkey", func(tx *sql.Tx) error {
207 hasTarget, err := hasUniqueIndex(tx, "spindle_members", []string{"did", "rkey"})
208 if err != nil {
209 return err
210 }
211 if hasTarget {
212 return nil
213 }
214
215 var totalRows, distinctRows int
216 if err := tx.QueryRow(`select count(*) from spindle_members`).Scan(&totalRows); err != nil {
217 return err
218 }
219 if err := tx.QueryRow(`select count(*) from (select 1 from spindle_members group by did, rkey)`).Scan(&distinctRows); err != nil {
220 return err
221 }
222 if dropped := totalRows - distinctRows; dropped > 0 {
223 logger.Warn("dropping duplicate (did, rkey) rows during spindle_members rebuild", "dropped", dropped, "kept", distinctRows)
224 }
225
226 _, err = tx.Exec(`
227 create table spindle_members_new (
228 id integer primary key autoincrement,
229 did text not null,
230 rkey text not null,
231 instance text not null,
232 subject text not null,
233 created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
234 unique (did, rkey)
235 );
236
237 insert into spindle_members_new (id, did, rkey, instance, subject, created)
238 select id, did, rkey, instance, subject, created
239 from spindle_members sm
240 where id = (
241 select max(id) from spindle_members
242 where did = sm.did and rkey = sm.rkey
243 );
244
245 drop table spindle_members;
246 alter table spindle_members_new rename to spindle_members;
247 `)
248 return err
249 }); err != nil {
250 return err
251 }
252
253 if err := orm.RunMigration(conn, logger, "events-pipeline-index", func(tx *sql.Tx) error {
254 _, err := tx.Exec(`
255 create index if not exists idx_events_pipeline_lookup on events(
256 coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'),
257 json_extract(event, '$.triggerMetadata.repo.did')),
258 coalesce(json_extract(event, '$.triggerMetadata.push.newSha'),
259 json_extract(event, '$.triggerMetadata.pullRequest.sourceSha'),
260 json_extract(event, '$.triggerMetadata.manual.sha'))
261 ) where nsid = 'sh.tangled.pipeline';
262
263 create index if not exists idx_events_pipeline_status on events(
264 json_extract(event, '$.pipeline'),
265 json_extract(event, '$.workflow')
266 ) where nsid = 'sh.tangled.pipeline.status';
267 `)
268 return err
269 }); err != nil {
270 return err
271 }
272
273 return nil
274}
275
276func hasUniqueIndex(tx *sql.Tx, table string, cols []string) (bool, error) {
277 rows, err := tx.Query(
278 `select name from pragma_index_list(?) where "unique" = 1`,
279 table,
280 )
281 if err != nil {
282 return false, err
283 }
284 defer rows.Close()
285
286 var indexNames []string
287 for rows.Next() {
288 var name string
289 if err := rows.Scan(&name); err != nil {
290 return false, err
291 }
292 indexNames = append(indexNames, name)
293 }
294 if err := rows.Err(); err != nil {
295 return false, err
296 }
297
298 wantSorted := slices.Clone(cols)
299 slices.Sort(wantSorted)
300
301 for _, name := range indexNames {
302 colRows, err := tx.Query(
303 `select name from pragma_index_info(?) order by seqno`,
304 name,
305 )
306 if err != nil {
307 return false, err
308 }
309 var got []string
310 for colRows.Next() {
311 var c string
312 if err := colRows.Scan(&c); err != nil {
313 colRows.Close()
314 return false, err
315 }
316 got = append(got, c)
317 }
318 colRows.Close()
319 slices.Sort(got)
320 if slices.Equal(got, wantSorted) {
321 return true, nil
322 }
323 }
324 return false, nil
325}
326
327func (d *DB) SaveLastTimeUs(lastTimeUs int64) error {
328 _, err := d.Exec(`
329 insert into _jetstream (id, last_time_us)
330 values (1, ?)
331 on conflict(id) do update set last_time_us = excluded.last_time_us
332 `, lastTimeUs)
333 return err
334}
335
336func (d *DB) GetLastTimeUs() (int64, error) {
337 var lastTimeUs int64
338 row := d.QueryRow(`select last_time_us from _jetstream where id = 1;`)
339 err := row.Scan(&lastTimeUs)
340 return lastTimeUs, err
341}