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