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.6 kB 332 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 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}