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 id integer primary key autoincrement,
83 did text not null,
84 rkey text not null,
85 instance text not null,
86 subject text not null,
87 created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
88 unique (did, rkey)
89 );
90
91 create table if not exists events (
92 rkey text not null,
93 nsid text not null,
94 event text not null,
95 created integer not null
96 );
97
98 create table if not exists nixos_toplevel_cache (
99 config_key text primary key,
100 toplevel text not null,
101 updated_at text not null
102 );
103
104 create table if not exists pipelines (
105 id text primary key,
106 repo_did text not null,
107 commit_id text not null
108 );
109 create table if not exists jobs (
110 id integer primary key autoincrement,
111 repo_did text not null,
112 pipeline_id_knot text not null,
113 pipeline_id_rkey text not null,
114 source_repo text,
115 tpl text not null,
116 created_at integer not null default (strftime('%s', 'now'))
117 );
118
119 create table if not exists workflows (
120 id integer primary key autoincrement,
121 pipeline_id text not null,
122 name text not null,
123 status text not null default 'pending',
124
125 unique(pipeline_id, id),
126 foreign key (pipeline_id) references pipelines(id) on delete cascade
127 );
128
129 create table if not exists mill_executors (
130 name text primary key,
131 token_hash text not null unique,
132 created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
133 expires_at text,
134 labels text,
135 quarantine_reason text,
136 quarantined_at text
137 );
138
139 create table if not exists mill_leases (
140 lease_id text primary key,
141 node_id text not null,
142 epoch text not null,
143 engine text not null,
144 knot text not null,
145 rkey text not null,
146 workflow text not null,
147 state text not null
148 );
149
150 create table if not exists mill_executor_cursors (
151 node_id text not null,
152 epoch text not null,
153 acked_seqno integer not null,
154 primary key (node_id, epoch)
155 );
156
157 create table if not exists mill_outbox_state (
158 epoch text not null,
159 next_seqno integer not null,
160 primary key (epoch)
161 );
162
163 create table if not exists mill_outbox_rows (
164 epoch text not null,
165 seqno integer not null,
166 payload blob not null,
167 byte_size integer not null,
168 control integer not null,
169 primary key (epoch, seqno),
170 foreign key (epoch) references mill_outbox_state(epoch) on delete cascade
171 );
172
173 create table if not exists mill_artifacts (
174 id integer primary key autoincrement,
175 lease_id text not null,
176 workflow text not null,
177 ref text not null,
178 hash text not null
179 );
180
181 create table if not exists executor_pending_artifacts (
182 lease_id text primary key,
183 workflow text not null,
184 status text not null,
185 error text not null default '',
186 exit_code integer not null default 0,
187 ref text not null,
188 hash text not null
189 );
190
191 create table if not exists migrations (
192 id integer primary key autoincrement,
193 name text unique
194 );
195 `)
196 if err := runMigrations(ctx, conn, logger); err != nil {
197 return nil, err
198 }
199
200 return &DB{db}, nil
201}
202
203func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error {
204 if err := orm.RunMigration(conn, logger, "repos-to-repo-did", func(tx *sql.Tx) error {
205 var hasName int
206 if err := tx.QueryRow(
207 `select count(*) from pragma_table_info('repos') where name = 'name'`,
208 ).Scan(&hasName); err != nil {
209 return err
210 }
211
212 if hasName > 0 {
213 var totalRows, copiedRows int
214 if err := tx.QueryRow(`select count(*) from repos`).Scan(&totalRows); err != nil {
215 return err
216 }
217 if err := tx.QueryRow(`select count(*) from repos where coalesce(name, '') <> ''`).Scan(&copiedRows); err != nil {
218 return err
219 }
220 if dropped := totalRows - copiedRows; dropped > 0 {
221 logger.Warn("dropping repo rows with empty name during migration", "dropped", dropped, "kept", copiedRows)
222 }
223
224 if _, err := tx.Exec(`
225 create table if not exists repos_new (
226 id integer primary key autoincrement,
227 knot text not null,
228 owner text not null,
229 rkey text not null,
230 repo_did text,
231 created_at text,
232 addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
233
234 unique(owner, rkey)
235 );
236
237 insert into repos_new (id, knot, owner, rkey, addedAt)
238 select id, knot, owner, name, addedAt from repos where coalesce(name, '') <> '';
239
240 drop table repos;
241 alter table repos_new rename to repos;
242 `); err != nil {
243 return err
244 }
245 }
246
247 _, err := tx.Exec(`
248 create index if not exists idx_repos_repo_did on repos(repo_did);
249 create index if not exists idx_repos_owner_repo_did on repos(owner, repo_did);
250 create index if not exists idx_repo_collaborators_repo_did
251 on repo_collaborators(repo_did);
252 `)
253 return err
254 }); err != nil {
255 return err
256 }
257
258 if err := orm.RunMigration(conn, logger, "spindle-members-unique-on-rkey", func(tx *sql.Tx) error {
259 hasTarget, err := hasUniqueIndex(tx, "spindle_members", []string{"did", "rkey"})
260 if err != nil {
261 return err
262 }
263 if hasTarget {
264 return nil
265 }
266
267 var totalRows, distinctRows int
268 if err := tx.QueryRow(`select count(*) from spindle_members`).Scan(&totalRows); err != nil {
269 return err
270 }
271 if err := tx.QueryRow(`select count(*) from (select 1 from spindle_members group by did, rkey)`).Scan(&distinctRows); err != nil {
272 return err
273 }
274 if dropped := totalRows - distinctRows; dropped > 0 {
275 logger.Warn("dropping duplicate (did, rkey) rows during spindle_members rebuild", "dropped", dropped, "kept", distinctRows)
276 }
277
278 _, err = tx.Exec(`
279 create table spindle_members_new (
280 id integer primary key autoincrement,
281 did text not null,
282 rkey text not null,
283 instance text not null,
284 subject text not null,
285 created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
286 unique (did, rkey)
287 );
288
289 insert into spindle_members_new (id, did, rkey, instance, subject, created)
290 select id, did, rkey, instance, subject, created
291 from spindle_members sm
292 where id = (
293 select max(id) from spindle_members
294 where did = sm.did and rkey = sm.rkey
295 );
296
297 drop table spindle_members;
298 alter table spindle_members_new rename to spindle_members;
299 `)
300 return err
301 }); err != nil {
302 return err
303 }
304
305 if err := orm.RunMigration(conn, logger, "events-pipeline-index", func(tx *sql.Tx) error {
306 _, err := tx.Exec(`
307 create index if not exists idx_events_pipeline_lookup on events(
308 coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'),
309 json_extract(event, '$.triggerMetadata.repo.did')),
310 coalesce(json_extract(event, '$.triggerMetadata.push.newSha'),
311 json_extract(event, '$.triggerMetadata.pullRequest.sourceSha'),
312 json_extract(event, '$.triggerMetadata.manual.sha'))
313 ) where nsid = 'sh.tangled.pipeline';
314
315 create index if not exists idx_events_pipeline_status on events(
316 json_extract(event, '$.pipeline'),
317 json_extract(event, '$.workflow')
318 ) where nsid = 'sh.tangled.pipeline.status';
319 `)
320 return err
321 }); err != nil {
322 return err
323 }
324
325 return nil
326}
327
328func hasUniqueIndex(tx *sql.Tx, table string, cols []string) (bool, error) {
329 rows, err := tx.Query(
330 `select name from pragma_index_list(?) where "unique" = 1`,
331 table,
332 )
333 if err != nil {
334 return false, err
335 }
336 defer rows.Close()
337
338 var indexNames []string
339 for rows.Next() {
340 var name string
341 if err := rows.Scan(&name); err != nil {
342 return false, err
343 }
344 indexNames = append(indexNames, name)
345 }
346 if err := rows.Err(); err != nil {
347 return false, err
348 }
349
350 wantSorted := slices.Clone(cols)
351 slices.Sort(wantSorted)
352
353 for _, name := range indexNames {
354 colRows, err := tx.Query(
355 `select name from pragma_index_info(?) order by seqno`,
356 name,
357 )
358 if err != nil {
359 return false, err
360 }
361 var got []string
362 for colRows.Next() {
363 var c string
364 if err := colRows.Scan(&c); err != nil {
365 colRows.Close()
366 return false, err
367 }
368 got = append(got, c)
369 }
370 colRows.Close()
371 slices.Sort(got)
372 if slices.Equal(got, wantSorted) {
373 return true, nil
374 }
375 }
376 return false, nil
377}
378
379func (d *DB) SaveLastTimeUs(lastTimeUs int64) error {
380 _, err := d.Exec(`
381 insert into _jetstream (id, last_time_us)
382 values (1, ?)
383 on conflict(id) do update set last_time_us = excluded.last_time_us
384 `, lastTimeUs)
385 return err
386}
387
388func (d *DB) GetLastTimeUs() (int64, error) {
389 var lastTimeUs int64
390 row := d.QueryRow(`select last_time_us from _jetstream where id = 1;`)
391 err := row.Scan(&lastTimeUs)
392 return lastTimeUs, err
393}