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