This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / db / db.go
10 kB 393 lines
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}