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