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
116 create table if not exists workflows (
117 id integer primary key autoincrement,
118 pipeline_id text not null,
119 name text not null,
120 status text not null default 'pending',
121
122 unique(pipeline_id, id),
123 foreign key (pipeline_id) references pipelines(id) on delete cascade
124 );
125
126 create table if not exists migrations (
127 id integer primary key autoincrement,
128 name text unique
129 );
130 `)
131 if err != nil {
132 return nil, err
133 }
134
135 if err := runMigrations(ctx, conn, logger); err != nil {
136 return nil, err
137 }
138
139 return &DB{db}, nil
140}
141
142func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error {
143 if err := orm.RunMigration(conn, logger, "repos-to-repo-did", func(tx *sql.Tx) error {
144 var hasName int
145 if err := tx.QueryRow(
146 `select count(*) from pragma_table_info('repos') where name = 'name'`,
147 ).Scan(&hasName); err != nil {
148 return err
149 }
150
151 if hasName > 0 {
152 var totalRows, copiedRows int
153 if err := tx.QueryRow(`select count(*) from repos`).Scan(&totalRows); err != nil {
154 return err
155 }
156 if err := tx.QueryRow(`select count(*) from repos where coalesce(name, '') <> ''`).Scan(&copiedRows); err != nil {
157 return err
158 }
159 if dropped := totalRows - copiedRows; dropped > 0 {
160 logger.Warn("dropping repo rows with empty name during migration", "dropped", dropped, "kept", copiedRows)
161 }
162
163 if _, err := tx.Exec(`
164 create table if not exists repos_new (
165 id integer primary key autoincrement,
166 knot text not null,
167 owner text not null,
168 rkey text not null,
169 repo_did text,
170 created_at text,
171 addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
172
173 unique(owner, rkey)
174 );
175
176 insert into repos_new (id, knot, owner, rkey, addedAt)
177 select id, knot, owner, name, addedAt from repos where coalesce(name, '') <> '';
178
179 drop table repos;
180 alter table repos_new rename to repos;
181 `); err != nil {
182 return err
183 }
184 }
185
186 _, err := tx.Exec(`
187 create index if not exists idx_repos_repo_did on repos(repo_did);
188 create index if not exists idx_repos_owner_repo_did on repos(owner, repo_did);
189 create index if not exists idx_repo_collaborators_repo_did
190 on repo_collaborators(repo_did);
191 `)
192 return err
193 }); err != nil {
194 return err
195 }
196
197 if err := orm.RunMigration(conn, logger, "spindle-members-unique-on-rkey", func(tx *sql.Tx) error {
198 hasTarget, err := hasUniqueIndex(tx, "spindle_members", []string{"did", "rkey"})
199 if err != nil {
200 return err
201 }
202 if hasTarget {
203 return nil
204 }
205
206 var totalRows, distinctRows int
207 if err := tx.QueryRow(`select count(*) from spindle_members`).Scan(&totalRows); err != nil {
208 return err
209 }
210 if err := tx.QueryRow(`select count(*) from (select 1 from spindle_members group by did, rkey)`).Scan(&distinctRows); err != nil {
211 return err
212 }
213 if dropped := totalRows - distinctRows; dropped > 0 {
214 logger.Warn("dropping duplicate (did, rkey) rows during spindle_members rebuild", "dropped", dropped, "kept", distinctRows)
215 }
216
217 _, err = tx.Exec(`
218 create table spindle_members_new (
219 id integer primary key autoincrement,
220 did text not null,
221 rkey text not null,
222 instance text not null,
223 subject text not null,
224 created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
225 unique (did, rkey)
226 );
227
228 insert into spindle_members_new (id, did, rkey, instance, subject, created)
229 select id, did, rkey, instance, subject, created
230 from spindle_members sm
231 where id = (
232 select max(id) from spindle_members
233 where did = sm.did and rkey = sm.rkey
234 );
235
236 drop table spindle_members;
237 alter table spindle_members_new rename to spindle_members;
238 `)
239 return err
240 }); err != nil {
241 return err
242 }
243
244 if err := orm.RunMigration(conn, logger, "events-pipeline-index", func(tx *sql.Tx) error {
245 _, err := tx.Exec(`
246 create index if not exists idx_events_pipeline_lookup on events(
247 coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'),
248 json_extract(event, '$.triggerMetadata.repo.did')),
249 coalesce(json_extract(event, '$.triggerMetadata.push.newSha'),
250 json_extract(event, '$.triggerMetadata.pullRequest.sourceSha'),
251 json_extract(event, '$.triggerMetadata.manual.sha'))
252 ) where nsid = 'sh.tangled.pipeline';
253
254 create index if not exists idx_events_pipeline_status on events(
255 json_extract(event, '$.pipeline'),
256 json_extract(event, '$.workflow')
257 ) where nsid = 'sh.tangled.pipeline.status';
258 `)
259 return err
260 }); err != nil {
261 return err
262 }
263
264 return nil
265}
266
267func hasUniqueIndex(tx *sql.Tx, table string, cols []string) (bool, error) {
268 rows, err := tx.Query(
269 `select name from pragma_index_list(?) where "unique" = 1`,
270 table,
271 )
272 if err != nil {
273 return false, err
274 }
275 defer rows.Close()
276
277 var indexNames []string
278 for rows.Next() {
279 var name string
280 if err := rows.Scan(&name); err != nil {
281 return false, err
282 }
283 indexNames = append(indexNames, name)
284 }
285 if err := rows.Err(); err != nil {
286 return false, err
287 }
288
289 wantSorted := slices.Clone(cols)
290 slices.Sort(wantSorted)
291
292 for _, name := range indexNames {
293 colRows, err := tx.Query(
294 `select name from pragma_index_info(?) order by seqno`,
295 name,
296 )
297 if err != nil {
298 return false, err
299 }
300 var got []string
301 for colRows.Next() {
302 var c string
303 if err := colRows.Scan(&c); err != nil {
304 colRows.Close()
305 return false, err
306 }
307 got = append(got, c)
308 }
309 colRows.Close()
310 slices.Sort(got)
311 if slices.Equal(got, wantSorted) {
312 return true, nil
313 }
314 }
315 return false, nil
316}
317
318func (d *DB) SaveLastTimeUs(lastTimeUs int64) error {
319 _, err := d.Exec(`
320 insert into _jetstream (id, last_time_us)
321 values (1, ?)
322 on conflict(id) do update set last_time_us = excluded.last_time_us
323 `, lastTimeUs)
324 return err
325}
326
327func (d *DB) GetLastTimeUs() (int64, error) {
328 var lastTimeUs int64
329 row := d.QueryRow(`select last_time_us from _jetstream where id = 1;`)
330 err := row.Scan(&lastTimeUs)
331 return lastTimeUs, err
332}