Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383package db
import ( "context" "database/sql" "encoding/json" "log/slog" "slices" "strings" "time"
_ "github.com/mattn/go-sqlite3" "tangled.org/core/api/tangled" "tangled.org/core/log" "tangled.org/core/orm" "tangled.org/core/spindle/models" pipelinecodec "tangled.org/core/spindle/pipeline" "tangled.org/core/sqlite")
type DB struct { *sql.DB}
type DBTX interface { QueryRow(query string, args ...any) *sql.Row Exec(query string, args ...any) (sql.Result, error)}
func Make(ctx context.Context, dbPath string) (*DB, error) { logger := log.FromContext(ctx) logger = log.SubLogger(logger, "db")
db, err := sqlite.Open(dbPath) if err != nil { return nil, err }
conn, err := db.Conn(ctx) if err != nil { return nil, err } defer conn.Close()
_, err = conn.ExecContext(ctx, ` create table if not exists _jetstream ( id integer primary key autoincrement, last_time_us integer not null );
create table if not exists known_dids ( did text primary key );
create table if not exists repos ( id integer primary key autoincrement, knot text not null, owner text not null, rkey text not null, repo_did text, created_at text, addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
unique(owner, rkey) );
create table if not exists feed_cursors ( knot text primary key, seq integer not null, feed text not null default 'atproto' );
create table if not exists feed_refs ( repo_did text not null, rkey text not null, sha text not null,
primary key (repo_did, rkey) );
create table if not exists repo_collaborators ( id integer primary key autoincrement, owner_did text not null, rkey text not null, subject text not null, repo_did text not null, addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
unique(owner_did, rkey) );
create table if not exists spindle_members ( id integer primary key autoincrement, did text not null, rkey text not null, instance text not null, subject text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique (did, rkey) );
create table if not exists nixos_toplevel_cache ( config_key text primary key, toplevel text not null, updated_at text not null );
create table if not exists cache_entries ( id text primary key, storage_key text unique not null, owner_did text not null, repo_did text not null, engine text not null, cache_key text not null, cache_hash text not null, checksum text not null default '', size_bytes integer not null default 0, restore_count integer not null default 0 check (restore_count >= 0), state text not null check (state in ('pending', 'ready', 'deleting')), created_at integer not null, last_used_at integer not null );
create index if not exists cache_entries_lookup on cache_entries (repo_did, engine, cache_key, cache_hash, created_at desc) where state = 'ready'; create index if not exists cache_entries_ready_expiry on cache_entries (last_used_at) where state = 'ready'; create index if not exists cache_entries_pending_deadline on cache_entries (last_used_at) where state in ('pending', 'deleting'); create index if not exists cache_entries_owner_state on cache_entries (owner_did, state);
create table if not exists cache_object_deletions ( storage_key text primary key, created_at integer not null );
create table if not exists pipelines ( id text primary key, repo_did text not null, commit_id text not null ); create table if not exists jobs ( id integer primary key autoincrement, repo_did text not null, pipeline_id_knot text not null, pipeline_id_rkey text not null, source_repo text, tpl text not null, traceparent text not null default '', tracestate text not null default '', created_at integer not null default (strftime('%s', 'now')), created_at_ns integer not null default 0 );
create table if not exists schedule_runs ( repo_did text not null, workflow text not null, scheduled_at integer not null, pipeline_id text,
primary key (repo_did, workflow, scheduled_at) );
create table if not exists scheduled_repos ( repo_did text primary key, branch text not null default '', sha text not null default '', refreshed_at integer not null ); create table if not exists workflow_schedules ( repo_did text not null, workflow text not null, expression text not null, timezone text not null default 'UTC',
primary key (repo_did, workflow, expression, timezone), foreign key (repo_did) references scheduled_repos(repo_did) on delete cascade );
create table if not exists workflows ( id integer primary key autoincrement, pipeline_id text not null, name text not null, status text not null default 'pending',
unique(pipeline_id, id), foreign key (pipeline_id) references pipelines(id) on delete cascade );
create table if not exists mill_executor_tokens ( token_hash text primary key, created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), expires_at text );
create table if not exists mill_executors ( name text primary key, token_hash text not null unique, labels text, quarantine_reason text, quarantined_at text, foreign key (token_hash) references mill_executor_tokens(token_hash) on delete cascade );
create table if not exists mill_leases ( lease_id text primary key, node_id text not null, epoch text not null, engine text not null, knot text not null, rkey text not null, workflow text not null, state text not null, quota_reservation_id text, owner_did text, repo_did text, mill_records_terminal_metrics integer not null default 0 );
create table if not exists mill_cache_capabilities ( lease_id text not null, action text not null check (action in ('restore', 'save')), cache_id text not null, storage_key text not null, primary key (lease_id, action, cache_id), foreign key (lease_id) references mill_leases(lease_id) on delete cascade );
create table if not exists mill_executor_cursors ( node_id text not null, epoch text not null, acked_seqno integer not null, primary key (node_id, epoch) );
create table if not exists mill_outbox_state ( epoch text not null, next_seqno integer not null, primary key (epoch) );
create table if not exists mill_outbox_rows ( epoch text not null, seqno integer not null, payload blob not null, byte_size integer not null, control integer not null, primary key (epoch, seqno), foreign key (epoch) references mill_outbox_state(epoch) on delete cascade );
create table if not exists mill_artifacts ( id integer primary key autoincrement, lease_id text not null, repo_did text not null, knot text not null, rkey text not null, workflow text not null, ref text not null, hash text not null );
create table if not exists executor_pending_artifacts ( lease_id text primary key, knot text not null default '', rkey text not null default '', workflow text not null, status text not null, error text not null default '', exit_code integer not null default 0, ref text not null, hash text not null, failure_class text not null default '', failure_reason text not null default '', mill_records_terminal_metrics integer not null default 0 );
create table if not exists webhooks ( id integer primary key autoincrement, repo_did text not null, url text not null, secret text, active integer not null default 1, events text not null, -- comma-separated event types created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), updated_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) ); create index if not exists idx_webhooks_repo_did on webhooks(repo_did);
create table if not exists webhook_deliveries ( id integer primary key autoincrement, webhook_id integer not null references webhooks(id) on delete cascade, event text not null, delivery_id text not null, url text not null, request_body text, response_code integer, response_body text, success integer not null default 0, created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) ); create index if not exists idx_webhook_deliveries_webhook_id on webhook_deliveries(webhook_id); create unique index if not exists idx_webhook_deliveries_delivery_id on webhook_deliveries(delivery_id);
create table if not exists pull_rounds ( repo_did text not null, rkey text not null, rounds integer not null, updated_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), primary key (repo_did, rkey) );
create table if not exists migrations ( id integer primary key autoincrement, name text unique ); `) if err := runMigrations(ctx, conn, logger); err != nil { return nil, err }
return &DB{DB: db}, nil}
func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error { if err := orm.RunMigration(conn, logger, "repos-to-repo-did", func(tx *sql.Tx) error { var hasName int if err := tx.QueryRow( `select count(*) from pragma_table_info('repos') where name = 'name'`, ).Scan(&hasName); err != nil { return err }
if hasName > 0 { var totalRows, copiedRows int if err := tx.QueryRow(`select count(*) from repos`).Scan(&totalRows); err != nil { return err } if err := tx.QueryRow(`select count(*) from repos where coalesce(name, '') <> ''`).Scan(&copiedRows); err != nil { return err } if dropped := totalRows - copiedRows; dropped > 0 { logger.Warn("dropping repo rows with empty name during migration", "dropped", dropped, "kept", copiedRows) }
if _, err := tx.Exec(` create table if not exists repos_new ( id integer primary key autoincrement, knot text not null, owner text not null, rkey text not null, repo_did text, created_at text, addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
unique(owner, rkey) );
insert into repos_new (id, knot, owner, rkey, addedAt) select id, knot, owner, name, addedAt from repos where coalesce(name, '') <> '';
drop table repos; alter table repos_new rename to repos; `); err != nil { return err } }
_, err := tx.Exec(` create index if not exists idx_repos_repo_did on repos(repo_did); create index if not exists idx_repos_owner_repo_did on repos(owner, repo_did); create index if not exists idx_repo_collaborators_repo_did on repo_collaborators(repo_did); `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "spindle-members-unique-on-rkey", func(tx *sql.Tx) error { hasTarget, err := hasUniqueIndex(tx, "spindle_members", []string{"did", "rkey"}) if err != nil { return err } if hasTarget { return nil }
var totalRows, distinctRows int if err := tx.QueryRow(`select count(*) from spindle_members`).Scan(&totalRows); err != nil { return err } if err := tx.QueryRow(`select count(*) from (select 1 from spindle_members group by did, rkey)`).Scan(&distinctRows); err != nil { return err } if dropped := totalRows - distinctRows; dropped > 0 { logger.Warn("dropping duplicate (did, rkey) rows during spindle_members rebuild", "dropped", dropped, "kept", distinctRows) }
_, err = tx.Exec(` create table spindle_members_new ( id integer primary key autoincrement, did text not null, rkey text not null, instance text not null, subject text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique (did, rkey) );
insert into spindle_members_new (id, did, rkey, instance, subject, created) select id, did, rkey, instance, subject, created from spindle_members sm where id = ( select max(id) from spindle_members where did = sm.did and rkey = sm.rkey );
drop table spindle_members; alter table spindle_members_new rename to spindle_members; `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "events-pipeline-index", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists events ( rkey text not null, nsid text not null, event text not null, created integer not null );
create index if not exists idx_events_pipeline_lookup on events( coalesce(json_extract(event, '$.triggerMetadata.repo.repoDid'), json_extract(event, '$.triggerMetadata.repo.did')), coalesce(json_extract(event, '$.triggerMetadata.push.newSha'), json_extract(event, '$.triggerMetadata.pullRequest.sourceSha'), json_extract(event, '$.triggerMetadata.manual.sha')) ) where nsid = 'sh.tangled.pipeline';
create index if not exists idx_events_pipeline_status on events( json_extract(event, '$.pipeline'), json_extract(event, '$.workflow') ) where nsid = 'sh.tangled.pipeline.status'; `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "mill-executor-credentials", func(tx *sql.Tx) error { var legacySchema int if err := tx.QueryRow( `select count(*) from pragma_table_info('mill_executors') where name = 'expires_at'`, ).Scan(&legacySchema); err != nil { return err } if legacySchema == 0 { return nil }
_, err := tx.Exec(` insert into mill_executor_tokens (token_hash, created_at, expires_at) select token_hash, created_at, expires_at from mill_executors;
create table mill_executors_new ( name text primary key, token_hash text not null unique, labels text, quarantine_reason text, quarantined_at text, foreign key (token_hash) references mill_executor_tokens(token_hash) on delete cascade );
insert into mill_executors_new ( name, token_hash, labels, quarantine_reason, quarantined_at ) select name, token_hash, labels, quarantine_reason, quarantined_at from mill_executors;
drop table mill_executors; alter table mill_executors_new rename to mill_executors; `) return err }); err != nil { return err }
// older generations select these fields during authentication, so clear rather than drop them if err := orm.RunMigration(conn, logger, "mill-clear-executor-quarantine", func(tx *sql.Tx) error { _, err := tx.Exec(`update mill_executors set quarantine_reason = null, quarantined_at = null`) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "jobs-trace-context", func(tx *sql.Tx) error { for _, column := range []string{"traceparent", "tracestate"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('jobs') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec(`alter table jobs add column ` + column + ` text not null default ''`); err != nil { return err } } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "cache-entry-integrity", func(tx *sql.Tx) error { for _, column := range []struct{ name, definition string }{ {"checksum", "text not null default ''"}, {"restore_count", "integer not null default 0"}, } { var present int if err := tx.QueryRow(`select count(*) from pragma_table_info('cache_entries') where name = ?`, column.name).Scan(&present); err != nil { return err } if present == 0 { if _, err := tx.Exec(`alter table cache_entries add column ` + column.name + ` ` + column.definition); err != nil { return err } } } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "cache-quotas-schema", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists cache_objects ( repo_did text not null, object_key text not null, byte_size integer not null check (byte_size >= 0), primary key (repo_did, object_key) ); create index if not exists idx_cache_objects_key on cache_objects(object_key);
create table if not exists cache_repo_owners ( repo_did text primary key, owner_did text not null ); create index if not exists idx_cache_repo_owners_owner on cache_repo_owners(owner_did);
create table if not exists cache_reservations ( id text primary key, object_key text not null, byte_size integer not null check (byte_size >= 0), repo_did text not null, owner_did text not null, phase text not null check (phase in ('reserved', 'publishing')), created_at integer not null ); create index if not exists idx_cache_reservations_repo on cache_reservations(repo_did); create index if not exists idx_cache_reservations_owner on cache_reservations(owner_did);
create table if not exists cache_overrides ( scope text not null check (scope in ('user', 'repo')), did text not null, max_bytes integer check (max_bytes >= 0 or max_bytes is null), primary key (scope, did) ); create index if not exists idx_cache_overrides_scope on cache_overrides(scope); `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "quota-schema", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists quota_limits ( did text not null, resource text not null, max_amount integer check (max_amount >= 0 or max_amount is null), primary key (did, resource) );
create table if not exists quota_reservations ( id text not null, resource text not null, kind text not null, key text not null, amount integer not null check (amount >= 0), repo_did text not null, owner_did text not null, phase text not null check (phase in ('reserved', 'publishing', 'active')), created_at integer not null, primary key (id, resource) ); create index if not exists idx_quota_reservations_repo on quota_reservations(repo_did); create index if not exists idx_quota_reservations_owner on quota_reservations(owner_did);
create table if not exists quota_allocations ( repo_did text not null, resource text not null, kind text not null, key text not null, amount integer not null check (amount >= 0), primary key (repo_did, resource, kind, key) ); create index if not exists idx_quota_allocations_key on quota_allocations(key);
create table if not exists quota_repo_owners ( repo_did text primary key, owner_did text not null ); create index if not exists idx_quota_repo_owners_owner on quota_repo_owners(owner_did); `) if err != nil { return err }
var hasOverrides int err = tx.QueryRow(`select count(*) from sqlite_master where type='table' and name='cache_overrides'`).Scan(&hasOverrides) if err != nil { return err } if hasOverrides > 0 { _, err = tx.Exec(` insert or ignore into quota_limits (did, resource, max_amount) select did, 'cache_storage_bytes', max_bytes from cache_overrides `) if err != nil { return err } }
var hasReservations int err = tx.QueryRow(`select count(*) from sqlite_master where type='table' and name='cache_reservations'`).Scan(&hasReservations) if err != nil { return err } if hasReservations > 0 { _, err = tx.Exec(` insert or ignore into quota_reservations (id, resource, kind, key, amount, repo_did, owner_did, phase, created_at) select id, 'cache_storage_bytes', 'nix_cache', object_key, byte_size, repo_did, owner_did, phase, created_at from cache_reservations `) if err != nil { return err } }
var hasObjects int err = tx.QueryRow(`select count(*) from sqlite_master where type='table' and name='cache_objects'`).Scan(&hasObjects) if err != nil { return err } if hasObjects > 0 { _, err = tx.Exec(` insert or ignore into quota_allocations (repo_did, resource, kind, key, amount) select repo_did, 'cache_storage_bytes', 'nix_cache', object_key, byte_size from cache_objects `) if err != nil { return err } }
var hasRepoOwners int err = tx.QueryRow(`select count(*) from sqlite_master where type='table' and name='cache_repo_owners'`).Scan(&hasRepoOwners) if err != nil { return err } if hasRepoOwners > 0 { _, err = tx.Exec(` insert or ignore into quota_repo_owners (repo_did, owner_did) select repo_did, owner_did from cache_repo_owners `) if err != nil { return err } }
_, err = tx.Exec(` drop index if exists idx_cache_objects_key; drop index if exists idx_cache_repo_owners_owner; drop index if exists idx_cache_reservations_repo; drop index if exists idx_cache_reservations_owner; drop index if exists idx_cache_overrides_scope; drop table if exists cache_objects; drop table if exists cache_repo_owners; drop table if exists cache_reservations; drop table if exists cache_overrides; `) if err != nil { return err }
return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "quota-zero-means-unlimited", func(tx *sql.Tx) error { _, err := tx.Exec(`update quota_limits set max_amount = null where max_amount = 0`) return err }); err != nil { return err }
// persist charged subjects so restarts can distinguish live reservations if err := orm.RunMigration(conn, logger, "mill-leases-quota-column", func(tx *sql.Tx) error { for _, column := range []string{"quota_reservation_id", "owner_did", "repo_did"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('mill_leases') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec(`alter table mill_leases add column ` + column + ` text`); err != nil { return err } } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "mill-artifacts-repo-did", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('mill_artifacts') where name = 'repo_did'`, ).Scan(&present); err != nil { return err } if present == 0 { if _, err := tx.Exec( `alter table mill_artifacts add column repo_did text not null default ''`, ); err != nil { return err } }
_, err := tx.Exec(` update mill_artifacts set repo_did = coalesce(( select repo_did from mill_leases where mill_leases.lease_id = mill_artifacts.lease_id ), repo_did, '') where coalesce(repo_did, '') = ''; `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "mill-artifacts-workflow-identity", func(tx *sql.Tx) error { for _, column := range []string{"knot", "rkey"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('mill_artifacts') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec( `alter table mill_artifacts add column ` + column + ` text not null default ''`, ); err != nil { return err } }
_, err := tx.Exec(` create index if not exists idx_mill_artifacts_workflow_identity on mill_artifacts (knot, rkey, workflow, id desc) `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "bans-schema", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists bans ( subject_did text primary key, created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) ); `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "jobs-created-at", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('jobs') where name = 'created_at'`, ).Scan(&present); err != nil { return err } if present != 0 { return nil } _, err := tx.Exec(`alter table jobs add column created_at integer not null default 0`) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "jobs-created-at-ns", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('jobs') where name = 'created_at_ns'`, ).Scan(&present); err != nil { return err } if present == 0 { if _, err := tx.Exec(`alter table jobs add column created_at_ns integer not null default 0`); err != nil { return err } } _, err := tx.Exec(` update jobs set created_at_ns = created_at * 1000000000 where created_at_ns = 0 and created_at between 1 and 9999999999 `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "pending-artifact-failure-attribution", func(tx *sql.Tx) error { for _, column := range []string{"failure_class", "failure_reason"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('executor_pending_artifacts') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec( `alter table executor_pending_artifacts add column ` + column + ` text not null default ''`, ); err != nil { return err } } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "pending-artifact-workflow-identity", func(tx *sql.Tx) error { for _, column := range []string{"knot", "rkey"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('executor_pending_artifacts') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec( `alter table executor_pending_artifacts add column ` + column + ` text not null default ''`, ); err != nil { return err } } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "pending-artifact-terminal-metric-authority", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('executor_pending_artifacts') where name = 'mill_records_terminal_metrics'`, ).Scan(&present); err != nil { return err } if present != 0 { return nil } _, err := tx.Exec(` alter table executor_pending_artifacts add column mill_records_terminal_metrics integer not null default 0 `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "mill-lease-terminal-metric-authority", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('mill_leases') where name = 'mill_records_terminal_metrics'`, ).Scan(&present); err != nil { return err } if present != 0 { return nil } _, err := tx.Exec(` alter table mill_leases add column mill_records_terminal_metrics integer not null default 0 `) return err }); err != nil { return err }
// re-introduce a repo display name, needed to build repository:renamed // webhook payloads. the earlier repos-to-repo-did migration dropped the // legacy name column; this adds it back as nullable metadata. if err := orm.RunMigration(conn, logger, "repos-add-name-column", func(tx *sql.Tx) error { var hasName int if err := tx.QueryRow( `select count(*) from pragma_table_info('repos') where name = 'name'`, ).Scan(&hasName); err != nil { return err } if hasName == 0 { if _, err := tx.Exec(`alter table repos add column name text`); err != nil { return err } } _, err := tx.Exec(`update repos set name = rkey where name is null`) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "pipelines-and-workflow-statuses", func(tx *sql.Tx) error { if _, err := tx.Exec(` drop table if exists workflows; drop table if exists pipelines; `); err != nil { return err }
if _, err := tx.Exec(` create table pipelines ( id integer primary key autoincrement, rkey text not null unique, knot text not null, repo_did text not null, commit_sha text not null, kind text not null, payload text not null ); create index idx_pipelines_repo_id on pipelines(repo_did, id); create index idx_pipelines_repo_commit on pipelines(repo_did, commit_sha);
create table workflow_statuses ( id integer primary key autoincrement, rkey text not null, workflow text not null, status text not null, error text, exit_code integer, created_at text not null ); create index idx_workflow_statuses_lookup on workflow_statuses(rkey, workflow, id); `); err != nil { return err }
if err := migratePipelines(tx, logger); err != nil { return err } if err := migrateWorkflowStatuses(tx, logger); err != nil { return err }
_, err := tx.Exec(` drop index if exists idx_events_pipeline_lookup; drop index if exists idx_events_pipeline_status; alter table events rename to events_legacy; `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "pipeline-id-identity", func(tx *sql.Tx) error { if _, err := tx.Exec(` create table pipeline_log_renames ( knot text not null, pipeline_id text not null, workflow text not null, primary key (knot, pipeline_id, workflow) ); `); err != nil { return err } if err := stagePipelineLogRenames(tx); err != nil { return err } if err := migratePipelineIdentityColumns(tx); err != nil { return err } if _, err := tx.Exec(` alter table jobs drop column pipeline_id_knot; alter table jobs rename column pipeline_id_rkey to pipeline_id;
alter table mill_leases drop column knot; alter table mill_leases rename column rkey to pipeline_id;
drop index if exists idx_mill_artifacts_workflow_identity; alter table mill_artifacts drop column knot; alter table mill_artifacts rename column rkey to pipeline_id; create index idx_mill_artifacts_workflow_identity on mill_artifacts (pipeline_id, workflow, id desc);
alter table executor_pending_artifacts drop column knot; alter table executor_pending_artifacts rename column rkey to pipeline_id; `); err != nil { return err } return nil }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "drop-legacy-pipeline-events", func(tx *sql.Tx) error { _, err := tx.Exec(` drop table if exists events; drop table if exists events_legacy; `) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "add-feed-to-feed-cursors", func(tx *sql.Tx) error { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('feed_cursors') where name = 'feed'`, ).Scan(&present); err != nil { return err } if present != 0 { return nil } _, err := tx.Exec( `alter table feed_cursors add column feed text not null default 'atproto'`, ) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "schedule-runs-prune-index", func(tx *sql.Tx) error { _, err := tx.Exec(`create index if not exists idx_schedule_runs_prune on schedule_runs(scheduled_at)`) return err }); err != nil { return err }
if err := orm.RunMigration(conn, logger, "workflow-schedules-timezone", func(tx *sql.Tx) error { var hasTimezone int if err := tx.QueryRow( `select count(*) from pragma_table_info('workflow_schedules') where name = 'timezone'`, ).Scan(&hasTimezone); err != nil { return err } if hasTimezone > 0 { return nil } _, err := tx.Exec(` create table workflow_schedules_new ( repo_did text not null, workflow text not null, expression text not null, timezone text not null default 'UTC',
primary key (repo_did, workflow, expression, timezone), foreign key (repo_did) references scheduled_repos(repo_did) on delete cascade ); insert into workflow_schedules_new (repo_did, workflow, expression, timezone) select repo_did, workflow, expression, 'UTC' from workflow_schedules; drop table workflow_schedules; alter table workflow_schedules_new rename to workflow_schedules; `) return err }); err != nil { return err } if err := orm.RunMigration(conn, logger, "scheduled-repos-revision", func(tx *sql.Tx) error { for _, column := range []string{"branch", "sha"} { var present int if err := tx.QueryRow( `select count(*) from pragma_table_info('scheduled_repos') where name = ?`, column, ).Scan(&present); err != nil { return err } if present != 0 { continue } if _, err := tx.Exec( `alter table scheduled_repos add column ` + column + ` text not null default ''`, ); err != nil { return err } } return nil }); err != nil { return err } return nil}
func tableHasColumn(tx *sql.Tx, table, column string) (bool, error) { var count int if err := tx.QueryRow( `select count(*) from pragma_table_info(?) where name = ?`, table, column, ).Scan(&count); err != nil { return false, err } return count != 0, nil}
func migratePipelineIdentityColumns(tx *sql.Tx) error { hasRkey, err := tableHasColumn(tx, "pipelines", "rkey") if err != nil { return err } if hasRkey { if _, err := tx.Exec(` alter table pipelines drop column knot; alter table pipelines rename column rkey to pipeline_id; `); err != nil { return err } }
hasStatusRkey, err := tableHasColumn(tx, "workflow_statuses", "rkey") if err != nil { return err } if hasStatusRkey { if _, err := tx.Exec(` drop index if exists idx_workflow_statuses_lookup; alter table workflow_statuses rename column rkey to pipeline_id; create index idx_workflow_statuses_lookup on workflow_statuses (pipeline_id, workflow, id); `); err != nil { return err } } return nil}
func hasUniqueIndex(tx *sql.Tx, table string, cols []string) (bool, error) { rows, err := tx.Query( `select name from pragma_index_list(?) where "unique" = 1`, table, ) if err != nil { return false, err } defer rows.Close()
var indexNames []string for rows.Next() { var name string if err := rows.Scan(&name); err != nil { return false, err } indexNames = append(indexNames, name) } if err := rows.Err(); err != nil { return false, err }
wantSorted := slices.Clone(cols) slices.Sort(wantSorted)
for _, name := range indexNames { colRows, err := tx.Query( `select name from pragma_index_info(?) order by seqno`, name, ) if err != nil { return false, err } var got []string for colRows.Next() { var c string if err := colRows.Scan(&c); err != nil { colRows.Close() return false, err } got = append(got, c) } colRows.Close() slices.Sort(got) if slices.Equal(got, wantSorted) { return true, nil } } return false, nil}
func (d *DB) SaveLastTimeUs(lastTimeUs int64) error { _, err := d.Exec(` insert into _jetstream (id, last_time_us) values (1, ?) on conflict(id) do update set last_time_us = excluded.last_time_us `, lastTimeUs) return err}
func (d *DB) GetLastTimeUs() (int64, error) { var lastTimeUs int64 row := d.QueryRow(`select last_time_us from _jetstream where id = 1;`) err := row.Scan(&lastTimeUs) return lastTimeUs, err}func (d *DB) GetRepoOwnerAndDid(knot, rkey string) (owner string, repoDid string, err error) { err = d.QueryRow("SELECT owner, repo_did FROM repos WHERE knot = ? AND rkey = ?", knot, rkey).Scan(&owner, &repoDid) return}
// migratePipelines orders converted rows by their original creation time so// pagination keeps the historical order.func migratePipelines(tx *sql.Tx, logger *slog.Logger) error { rows, err := tx.Query( `select rkey, event, created from events where nsid = ? order by created asc`, tangled.PipelineNSID, ) if err != nil { return err } defer rows.Close()
type converted struct { rkey, knot, repoDid, commitSha, kind string payload []byte } var out []converted var skipped int
for rows.Next() { var rkey, eventJson string var created int64 if err := rows.Scan(&rkey, &eventJson, &created); err != nil { return err }
var raw tangled.Pipeline if err := json.Unmarshal([]byte(eventJson), &raw); err != nil { skipped++ continue } if raw.TriggerMetadata == nil { skipped++ continue }
record, err := pipelinecodec.FromTangled(models.PipelineId(rkey), time.Unix(0, created), raw) if err != nil { skipped++ continue } payload, err := json.Marshal(record) if err != nil { skipped++ continue } var knot string if raw.TriggerMetadata.Repo != nil { knot = raw.TriggerMetadata.Repo.Knot } out = append(out, converted{rkey, knot, record.RepoDID, record.Commit, record.Trigger.Kind, payload}) } if err := rows.Err(); err != nil { return err }
for _, c := range out { if _, err := tx.Exec( `insert into pipelines (rkey, knot, repo_did, commit_sha, kind, payload) values (?, ?, ?, ?, ?, ?)`, c.rkey, c.knot, c.repoDid, c.commitSha, c.kind, string(c.payload), ); err != nil { return err } }
logger.Info("backfilled pipelines", "converted", len(out), "skipped", skipped) return nil}
// migrateWorkflowStatuses flattens the legacy status event log.func migrateWorkflowStatuses(tx *sql.Tx, logger *slog.Logger) error { rows, err := tx.Query( `select event from events where nsid = ? order by created asc`, tangled.PipelineStatusNSID, ) if err != nil { return err } defer rows.Close()
type converted struct { rkey, workflow, status, createdAt string wfError *string exitCode *int64 } var out []converted var skipped int
for rows.Next() { var eventJson string if err := rows.Scan(&eventJson); err != nil { return err }
var st tangled.PipelineStatus if err := json.Unmarshal([]byte(eventJson), &st); err != nil { skipped++ continue }
idx := strings.LastIndex(st.Pipeline, "/") if idx < 0 || idx == len(st.Pipeline)-1 { skipped++ continue } rkey := st.Pipeline[idx+1:]
out = append(out, converted{rkey, st.Workflow, st.Status, st.CreatedAt, st.Error, st.ExitCode}) } if err := rows.Err(); err != nil { return err }
convertedCount := 0 for _, c := range out { result, err := tx.Exec( `insert into workflow_statuses (rkey, workflow, status, error, exit_code, created_at) select ?, ?, ?, ?, ?, ? where exists (select 1 from pipelines where rkey = ?)`, c.rkey, c.workflow, c.status, c.wfError, c.exitCode, c.createdAt, c.rkey, ) if err != nil { return err } inserted, err := result.RowsAffected() if err != nil { return err } if inserted == 0 { skipped++ continue } convertedCount++ }
logger.Info("backfilled workflow statuses", "converted", convertedCount, "skipped", skipped) return nil}