Monorepo for Tangled forked from tangled.org/core
Something went wrong. Try again.
Go
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298package db
import ( "context" "database/sql" "fmt" "log/slog" "reflect" "strings"
_ "github.com/mattn/go-sqlite3" "tangled.org/core/log")
type DB struct { *sql.DB logger *slog.Logger}
type Execer interface { Query(query string, args ...any) (*sql.Rows, error) QueryContext(ctx context.Context, query string, args ...any) (*sql.Rows, error) QueryRow(query string, args ...any) *sql.Row QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row Exec(query string, args ...any) (sql.Result, error) ExecContext(ctx context.Context, query string, args ...any) (sql.Result, error) Prepare(query string) (*sql.Stmt, error) PrepareContext(ctx context.Context, query string) (*sql.Stmt, error)}
func Make(ctx context.Context, dbPath string) (*DB, error) { // https://github.com/mattn/go-sqlite3#connection-string opts := []string{ "_foreign_keys=1", "_journal_mode=WAL", "_synchronous=NORMAL", "_auto_vacuum=incremental", }
logger := log.FromContext(ctx) logger = log.SubLogger(logger, "db")
db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) 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 registrations ( id integer primary key autoincrement, domain text not null unique, did text not null, secret text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), registered text ); create table if not exists public_keys ( id integer primary key autoincrement, did text not null, name text not null, key text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(did, name, key) ); create table if not exists repos ( id integer primary key autoincrement, did text not null, name text not null, knot text not null, rkey text not null, at_uri text not null unique, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(did, name, knot, rkey) ); create table if not exists collaborators ( id integer primary key autoincrement, did text not null, repo integer not null, foreign key (repo) references repos(id) on delete cascade ); create table if not exists follows ( user_did text not null, subject_did text not null, rkey text not null, followed_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), primary key (user_did, subject_did), check (user_did <> subject_did) ); create table if not exists issues ( id integer primary key autoincrement, owner_did text not null, repo_at text not null, issue_id integer not null, title text not null, body text not null, open integer not null default 1, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), issue_at text, unique(repo_at, issue_id), foreign key (repo_at) references repos(at_uri) on delete cascade ); create table if not exists comments ( id integer primary key autoincrement, owner_did text not null, issue_id integer not null, repo_at text not null, comment_id integer not null, comment_at text not null, body text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(issue_id, comment_id), foreign key (repo_at, issue_id) references issues(repo_at, issue_id) on delete cascade ); create table if not exists pulls ( -- identifiers id integer primary key autoincrement, pull_id integer not null,
-- at identifiers repo_at text not null, owner_did text not null, rkey text not null, pull_at text,
-- content title text not null, body text not null, target_branch text not null, state integer not null default 0 check (state in (0, 1, 2)), -- open, merged, closed
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique(repo_at, pull_id), foreign key (repo_at) references repos(at_uri) on delete cascade );
-- every pull must have atleast 1 submission: the initial submission create table if not exists pull_submissions ( -- identifiers id integer primary key autoincrement, pull_id integer not null,
-- at identifiers repo_at text not null,
-- content, these are immutable, and require a resubmission to update round_number integer not null default 0, patch text,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique(repo_at, pull_id, round_number), foreign key (repo_at, pull_id) references pulls(repo_at, pull_id) on delete cascade );
create table if not exists pull_comments ( -- identifiers id integer primary key autoincrement, pull_id integer not null, submission_id integer not null,
-- at identifiers repo_at text not null, owner_did text not null, comment_at text not null,
-- content body text not null,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints foreign key (repo_at, pull_id) references pulls(repo_at, pull_id) on delete cascade, foreign key (submission_id) references pull_submissions(id) on delete cascade );
create table if not exists _jetstream ( id integer primary key autoincrement, last_time_us integer not null );
create table if not exists repo_issue_seqs ( repo_at text primary key, next_issue_id integer not null default 1 );
create table if not exists repo_pull_seqs ( repo_at text primary key, next_pull_id integer not null default 1 );
create table if not exists stars ( id integer primary key autoincrement, starred_by_did text not null, repo_at text not null, rkey text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), foreign key (repo_at) references repos(at_uri) on delete cascade, unique(starred_by_did, repo_at) );
create table if not exists reactions ( id integer primary key autoincrement, reacted_by_did text not null, thread_at text not null, kind text not null, rkey text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(reacted_by_did, thread_at, kind) );
create table if not exists emails ( id integer primary key autoincrement, did text not null, email text not null, verified integer not null default 0, verification_code text not null, last_sent text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), is_primary integer not null default 0, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(did, email) );
create table if not exists artifacts ( -- id id integer primary key autoincrement, did text not null, rkey text not null,
-- meta repo_at text not null, tag binary(20) not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- data blob_cid text not null, name text not null, size integer not null default 0, mimetype string not null default "*/*",
-- constraints unique(did, rkey), -- record must be unique unique(repo_at, tag, name), -- for a given tag object, each file must be unique foreign key (repo_at) references repos(at_uri) on delete cascade );
create table if not exists profile ( -- id id integer primary key autoincrement, did text not null,
-- data description text not null, include_bluesky integer not null default 0, location text,
-- constraints unique(did) ); create table if not exists profile_links ( -- id id integer primary key autoincrement, did text not null,
-- data link text not null,
-- constraints foreign key (did) references profile(did) on delete cascade ); create table if not exists profile_stats ( -- id id integer primary key autoincrement, did text not null,
-- data kind text not null check (kind in ( "merged-pull-request-count", "closed-pull-request-count", "open-pull-request-count", "open-issue-count", "closed-issue-count", "repository-count" )),
-- constraints foreign key (did) references profile(did) on delete cascade ); create table if not exists profile_pinned_repositories ( -- id id integer primary key autoincrement, did text not null,
-- data at_uri text not null,
-- constraints unique(did, at_uri), foreign key (did) references profile(did) on delete cascade, foreign key (at_uri) references repos(at_uri) on delete cascade );
create table if not exists oauth_requests ( id integer primary key autoincrement, auth_server_iss text not null, state text not null, did text not null, handle text not null, pds_url text not null, pkce_verifier text not null, dpop_auth_server_nonce text not null, dpop_private_jwk text not null );
create table if not exists oauth_sessions ( id integer primary key autoincrement, did text not null, handle text not null, pds_url text not null, auth_server_iss text not null, access_jwt text not null, refresh_jwt text not null, dpop_pds_nonce text, dpop_auth_server_nonce text not null, dpop_private_jwk text not null, expiry text not null );
create table if not exists punchcard ( did text not null, date text not null, -- yyyy-mm-dd count integer, primary key (did, date) );
create table if not exists spindles ( id integer primary key autoincrement, owner text not null, instance text not null, verified text, -- time of verification created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
unique(owner, instance) );
create table if not exists spindle_members ( -- identifiers for the record id integer primary key autoincrement, did text not null, rkey text not null,
-- data instance text not null, subject text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique (did, instance, subject) );
create table if not exists pipelines ( -- identifiers id integer primary key autoincrement, knot text not null, rkey text not null,
repo_owner text not null, repo_name text not null,
-- every pipeline must be associated with exactly one commit sha text not null check (length(sha) = 40), created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- trigger data trigger_id integer not null,
unique(knot, rkey), foreign key (trigger_id) references triggers(id) on delete cascade );
create table if not exists triggers ( -- primary key id integer primary key autoincrement,
-- top-level fields kind text not null,
-- pushTriggerData fields push_ref text, push_new_sha text check (length(push_new_sha) = 40), push_old_sha text check (length(push_old_sha) = 40),
-- pullRequestTriggerData fields pr_source_branch text, pr_target_branch text, pr_source_sha text check (length(pr_source_sha) = 40), pr_action text );
create table if not exists pipeline_statuses ( -- identifiers id integer primary key autoincrement, spindle text not null, rkey text not null,
-- referenced pipeline. these form the (did, rkey) pair pipeline_knot text not null, pipeline_rkey text not null,
-- content created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), workflow text not null, status text not null, error text, exit_code integer not null default 0,
unique (spindle, rkey), foreign key (pipeline_knot, pipeline_rkey) references pipelines (knot, rkey) on delete cascade );
create table if not exists repo_languages ( -- identifiers id integer primary key autoincrement,
-- repo identifiers repo_at text not null, ref text not null, is_default_ref integer not null default 0,
-- language breakdown language text not null, bytes integer not null check (bytes >= 0),
unique(repo_at, ref, language) );
create table if not exists signups_inflight ( id integer primary key autoincrement, email text not null unique, invite_code text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) );
create table if not exists strings ( -- identifiers did text not null, rkey text not null,
-- content filename text not null, description text, content text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), edited text,
primary key (did, rkey) );
create table if not exists label_definitions ( -- identifiers id integer primary key autoincrement, did text not null, rkey text not null, at_uri text generated always as ('at://' || did || '/' || 'sh.tangled.label.definition' || '/' || rkey) stored,
-- content name text not null, value_type text not null check (value_type in ( "null", "boolean", "integer", "string" )), value_format text not null default "any", value_enum text, -- comma separated list scope text not null, -- comma separated list of nsid color text, multiple integer not null default 0, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique (did, rkey) unique (at_uri) );
-- ops are flattened, a record may contain several additions and deletions, but the table will include one row per add/del create table if not exists label_ops ( -- identifiers id integer primary key autoincrement, did text not null, rkey text not null, at_uri text generated always as ('at://' || did || '/' || 'sh.tangled.label.op' || '/' || rkey) stored,
-- content subject text not null, operation text not null check (operation in ("add", "del")), operand_key text not null, operand_value text not null, -- we need two time values: performed is declared by the user, indexed is calculated by the av performed text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), indexed text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints -- traditionally (did, rkey) pair should be unique, but not in this case -- operand_key should reference a label definition foreign key (operand_key) references label_definitions (at_uri) on delete cascade, unique (did, rkey, subject, operand_key, operand_value) );
create table if not exists repo_labels ( -- identifiers id integer primary key autoincrement,
-- repo identifiers repo_at text not null,
-- label to subscribe to label_at text not null,
unique (repo_at, label_at) );
create table if not exists notifications ( id integer primary key autoincrement, recipient_did text not null, actor_did text not null, type text not null, entity_type text not null, entity_id text not null, read integer not null default 0, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), repo_id integer references repos(id), issue_id integer references issues(id), pull_id integer references pulls(id) );
create table if not exists notification_preferences ( id integer primary key autoincrement, user_did text not null unique, repo_starred integer not null default 1, issue_created integer not null default 1, issue_commented integer not null default 1, pull_created integer not null default 1, pull_commented integer not null default 1, followed integer not null default 1, pull_merged integer not null default 1, issue_closed integer not null default 1, email_notifications integer not null default 0 );
create table if not exists reference_links ( id integer primary key autoincrement, from_at text not null, to_at text not null, unique (from_at, to_at) );
create table if not exists migrations ( id integer primary key autoincrement, name text unique );
-- indexes for better performance create index if not exists idx_notifications_recipient_created on notifications(recipient_did, created desc); create index if not exists idx_notifications_recipient_read on notifications(recipient_did, read); create index if not exists idx_references_from_at on reference_links(from_at); create index if not exists idx_references_to_at on reference_links(to_at); `) if err != nil { return nil, err }
// run migrations runMigration(conn, logger, "add-description-to-repos", func(tx *sql.Tx) error { tx.Exec(` alter table repos add column description text check (length(description) <= 200); `) return nil })
runMigration(conn, logger, "add-rkey-to-pubkeys", func(tx *sql.Tx) error { // add unconstrained column _, err := tx.Exec(` alter table public_keys add column rkey text; `) if err != nil { return err }
// backfill _, err = tx.Exec(` update public_keys set rkey = '' where rkey is null; `) if err != nil { return err }
return nil })
runMigration(conn, logger, "add-rkey-to-comments", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table comments drop column comment_at; alter table comments add column rkey text; `) return err })
runMigration(conn, logger, "add-deleted-and-edited-to-issue-comments", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table comments add column deleted text; -- timestamp alter table comments add column edited text; -- timestamp `) return err })
runMigration(conn, logger, "add-source-info-to-pulls-and-submissions", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table pulls add column source_branch text; alter table pulls add column source_repo_at text; alter table pull_submissions add column source_rev text; `) return err })
runMigration(conn, logger, "add-source-to-repos", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table repos add column source text; `) return err })
// disable foreign-keys for the next migration // NOTE: this cannot be done in a transaction, so it is run outside [0] // // [0]: https://sqlite.org/pragma.html#pragma_foreign_keys conn.ExecContext(ctx, "pragma foreign_keys = off;") runMigration(conn, logger, "recreate-pulls-column-for-stacking-support", func(tx *sql.Tx) error { _, err := tx.Exec(` create table pulls_new ( -- identifiers id integer primary key autoincrement, pull_id integer not null,
-- at identifiers repo_at text not null, owner_did text not null, rkey text not null,
-- content title text not null, body text not null, target_branch text not null, state integer not null default 0 check (state in (0, 1, 2, 3)), -- closed, open, merged, deleted
-- source info source_branch text, source_repo_at text,
-- stacking stack_id text, change_id text, parent_change_id text,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique(repo_at, pull_id), foreign key (repo_at) references repos(at_uri) on delete cascade );
insert into pulls_new ( id, pull_id, repo_at, owner_did, rkey, title, body, target_branch, state, source_branch, source_repo_at, created ) select id, pull_id, repo_at, owner_did, rkey, title, body, target_branch, state, source_branch, source_repo_at, created FROM pulls;
drop table pulls; alter table pulls_new rename to pulls; `) return err }) conn.ExecContext(ctx, "pragma foreign_keys = on;")
runMigration(conn, logger, "add-spindle-to-repos", func(tx *sql.Tx) error { tx.Exec(` alter table repos add column spindle text; `) return nil })
// drop all knot secrets, add unique constraint to knots // // knots will henceforth use service auth for signed requests runMigration(conn, logger, "no-more-secrets", func(tx *sql.Tx) error { _, err := tx.Exec(` create table registrations_new ( id integer primary key autoincrement, domain text not null, did text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), registered text, read_only integer not null default 0, unique(domain, did) );
insert into registrations_new (id, domain, did, created, registered, read_only) select id, domain, did, created, registered, 1 from registrations where registered is not null;
drop table registrations; alter table registrations_new rename to registrations; `) return err })
// recreate and add rkey + created columns with default constraint runMigration(conn, logger, "rework-collaborators-table", func(tx *sql.Tx) error { // create new table // - repo_at instead of repo integer // - rkey field // - created field _, err := tx.Exec(` create table collaborators_new ( -- identifiers for the record id integer primary key autoincrement, did text not null, rkey text,
-- content subject_did text not null, repo_at text not null,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints foreign key (repo_at) references repos(at_uri) on delete cascade ) `) if err != nil { return err }
// copy data _, err = tx.Exec(` insert into collaborators_new (id, did, rkey, subject_did, repo_at) select c.id, r.did, '', c.did, r.at_uri from collaborators c join repos r on c.repo = r.id `) if err != nil { return err }
// drop old table _, err = tx.Exec(`drop table collaborators`) if err != nil { return err }
// rename new table _, err = tx.Exec(`alter table collaborators_new rename to collaborators`) return err })
runMigration(conn, logger, "add-rkey-to-issues", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table issues add column rkey text not null default '';
-- get last url section from issue_at and save to rkey column update issues set rkey = replace(issue_at, rtrim(issue_at, replace(issue_at, '/', '')), ''); `) return err })
// repurpose the read-only column to "needs-upgrade" runMigration(conn, logger, "rename-registrations-read-only-to-needs-upgrade", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table registrations rename column read_only to needs_upgrade; `) return err })
// require all knots to upgrade after the release of total xrpc runMigration(conn, logger, "migrate-knots-to-total-xrpc", func(tx *sql.Tx) error { _, err := tx.Exec(` update registrations set needs_upgrade = 1; `) return err })
// require all knots to upgrade after the release of total xrpc runMigration(conn, logger, "migrate-spindles-to-xrpc-owner", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table spindles add column needs_upgrade integer not null default 0; `) return err })
// remove issue_at from issues and replace with generated column // // this requires a full table recreation because stored columns // cannot be added via alter // // couple other changes: // - columns renamed to be more consistent // - adds edited and deleted fields // // disable foreign-keys for the next migration conn.ExecContext(ctx, "pragma foreign_keys = off;") runMigration(conn, logger, "remove-issue-at-from-issues", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists issues_new ( -- identifiers id integer primary key autoincrement, did text not null, rkey text not null, at_uri text generated always as ('at://' || did || '/' || 'sh.tangled.repo.issue' || '/' || rkey) stored,
-- at identifiers repo_at text not null,
-- content issue_id integer not null, title text not null, body text not null, open integer not null default 1, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), edited text, -- timestamp deleted text, -- timestamp
unique(did, rkey), unique(repo_at, issue_id), unique(at_uri), foreign key (repo_at) references repos(at_uri) on delete cascade ); `) if err != nil { return err }
// transfer data _, err = tx.Exec(` insert into issues_new (id, did, rkey, repo_at, issue_id, title, body, open, created) select i.id, i.owner_did, i.rkey, i.repo_at, i.issue_id, i.title, i.body, i.open, i.created from issues i; `) if err != nil { return err }
// drop old table _, err = tx.Exec(`drop table issues`) if err != nil { return err }
// rename new table _, err = tx.Exec(`alter table issues_new rename to issues`) return err }) conn.ExecContext(ctx, "pragma foreign_keys = on;")
// - renames the comments table to 'issue_comments' // - rework issue comments to update constraints: // * unique(did, rkey) // * remove comment-id and just use the global ID // * foreign key (repo_at, issue_id) // - new columns // * column "reply_to" which can be any other comment // * column "at-uri" which is a generated column runMigration(conn, logger, "rework-issue-comments", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists issue_comments ( -- identifiers id integer primary key autoincrement, did text not null, rkey text, at_uri text generated always as ('at://' || did || '/' || 'sh.tangled.repo.issue.comment' || '/' || rkey) stored,
-- at identifiers issue_at text not null, reply_to text, -- at_uri of parent comment
-- content body text not null, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), edited text, deleted text,
-- constraints unique(did, rkey), unique(at_uri), foreign key (issue_at) references issues(at_uri) on delete cascade ); `) if err != nil { return err }
// transfer data _, err = tx.Exec(` insert into issue_comments (id, did, rkey, issue_at, body, created, edited, deleted) select c.id, c.owner_did, c.rkey, i.at_uri, -- get at_uri from issues table c.body, c.created, c.edited, c.deleted from comments c join issues i on c.repo_at = i.repo_at and c.issue_id = i.issue_id; `) if err != nil { return err }
// drop old table _, err = tx.Exec(`drop table comments`) return err })
// add generated at_uri column to pulls table // // this requires a full table recreation because stored columns // cannot be added via alter // // disable foreign-keys for the next migration conn.ExecContext(ctx, "pragma foreign_keys = off;") runMigration(conn, logger, "add-at-uri-to-pulls", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists pulls_new ( -- identifiers id integer primary key autoincrement, pull_id integer not null, at_uri text generated always as ('at://' || owner_did || '/' || 'sh.tangled.repo.pull' || '/' || rkey) stored,
-- at identifiers repo_at text not null, owner_did text not null, rkey text not null,
-- content title text not null, body text not null, target_branch text not null, state integer not null default 0 check (state in (0, 1, 2, 3)), -- closed, open, merged, deleted
-- source info source_branch text, source_repo_at text,
-- stacking stack_id text, change_id text, parent_change_id text,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique(repo_at, pull_id), unique(at_uri), foreign key (repo_at) references repos(at_uri) on delete cascade ); `) if err != nil { return err }
// transfer data _, err = tx.Exec(` insert into pulls_new ( id, pull_id, repo_at, owner_did, rkey, title, body, target_branch, state, source_branch, source_repo_at, stack_id, change_id, parent_change_id, created ) select id, pull_id, repo_at, owner_did, rkey, title, body, target_branch, state, source_branch, source_repo_at, stack_id, change_id, parent_change_id, created from pulls; `) if err != nil { return err }
// drop old table _, err = tx.Exec(`drop table pulls`) if err != nil { return err }
// rename new table _, err = tx.Exec(`alter table pulls_new rename to pulls`) return err }) conn.ExecContext(ctx, "pragma foreign_keys = on;")
// remove repo_at and pull_id from pull_submissions and replace with pull_at // // this requires a full table recreation because stored columns // cannot be added via alter // // disable foreign-keys for the next migration conn.ExecContext(ctx, "pragma foreign_keys = off;") runMigration(conn, logger, "remove-repo-at-pull-id-from-pull-submissions", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists pull_submissions_new ( -- identifiers id integer primary key autoincrement, pull_at text not null,
-- content, these are immutable, and require a resubmission to update round_number integer not null default 0, patch text, source_rev text,
-- meta created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')),
-- constraints unique(pull_at, round_number), foreign key (pull_at) references pulls(at_uri) on delete cascade ); `) if err != nil { return err }
// transfer data, constructing pull_at from pulls table _, err = tx.Exec(` insert into pull_submissions_new (id, pull_at, round_number, patch, created) select ps.id, 'at://' || p.owner_did || '/sh.tangled.repo.pull/' || p.rkey, ps.round_number, ps.patch, ps.created from pull_submissions ps join pulls p on ps.repo_at = p.repo_at and ps.pull_id = p.pull_id; `) if err != nil { return err }
// drop old table _, err = tx.Exec(`drop table pull_submissions`) if err != nil { return err }
// rename new table _, err = tx.Exec(`alter table pull_submissions_new rename to pull_submissions`) return err }) conn.ExecContext(ctx, "pragma foreign_keys = on;")
// knots may report the combined patch for a comparison, we can store that on the appview side // (but not on the pds record), because calculating the combined patch requires a git index runMigration(conn, logger, "add-combined-column-submissions", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table pull_submissions add column combined text; `) return err })
runMigration(conn, logger, "add-pronouns-profile", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table profile add column pronouns text; `) return err })
runMigration(conn, logger, "add-meta-column-repos", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table repos add column website text; alter table repos add column topics text; `) return err })
runMigration(conn, logger, "add-usermentioned-preference", func(tx *sql.Tx) error { _, err := tx.Exec(` alter table notification_preferences add column user_mentioned integer not null default 1; `) return err })
// remove the foreign key constraints from stars. runMigration(conn, logger, "generalize-stars-subject", func(tx *sql.Tx) error { _, err := tx.Exec(` create table stars_new ( id integer primary key autoincrement, did text not null, rkey text not null,
subject_at text not null,
created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(did, rkey), unique(did, subject_at) );
insert into stars_new ( id, did, rkey, subject_at, created ) select id, starred_by_did, rkey, repo_at, created from stars;
drop table stars; alter table stars_new rename to stars;
create index if not exists idx_stars_created on stars(created); create index if not exists idx_stars_subject_at_created on stars(subject_at, created); `) return err })
return &DB{ db, logger, }, nil}
type migrationFn = func(*sql.Tx) error
func runMigration(c *sql.Conn, logger *slog.Logger, name string, migrationFn migrationFn) error { logger = logger.With("migration", name)
tx, err := c.BeginTx(context.Background(), nil) if err != nil { return err } defer tx.Rollback()
var exists bool err = tx.QueryRow("select exists (select 1 from migrations where name = ?)", name).Scan(&exists) if err != nil { return err }
if !exists { // run migration err = migrationFn(tx) if err != nil { logger.Error("failed to run migration", "err", err) return err }
// mark migration as complete _, err = tx.Exec("insert into migrations (name) values (?)", name) if err != nil { logger.Error("failed to mark migration as complete", "err", err) return err }
// commit the transaction if err := tx.Commit(); err != nil { return err }
logger.Info("migration applied successfully") } else { logger.Warn("skipped migration, already applied") }
return nil}
func (d *DB) Close() error { return d.DB.Close()}
type filter struct { key string arg any cmp string}
func newFilter(key, cmp string, arg any) filter { return filter{ key: key, arg: arg, cmp: cmp, }}
func FilterEq(key string, arg any) filter { return newFilter(key, "=", arg) }func FilterNotEq(key string, arg any) filter { return newFilter(key, "<>", arg) }func FilterGte(key string, arg any) filter { return newFilter(key, ">=", arg) }func FilterLte(key string, arg any) filter { return newFilter(key, "<=", arg) }func FilterIs(key string, arg any) filter { return newFilter(key, "is", arg) }func FilterIsNot(key string, arg any) filter { return newFilter(key, "is not", arg) }func FilterIn(key string, arg any) filter { return newFilter(key, "in", arg) }func FilterLike(key string, arg any) filter { return newFilter(key, "like", arg) }func FilterNotLike(key string, arg any) filter { return newFilter(key, "not like", arg) }func FilterContains(key string, arg any) filter { return newFilter(key, "like", fmt.Sprintf("%%%v%%", arg))}
func (f filter) Condition() string { rv := reflect.ValueOf(f.arg) kind := rv.Kind()
// if we have `FilterIn(k, [1, 2, 3])`, compile it down to `k in (?, ?, ?)` if (kind == reflect.Slice && rv.Type().Elem().Kind() != reflect.Uint8) || kind == reflect.Array { if rv.Len() == 0 { // always false return "1 = 0" }
placeholders := make([]string, rv.Len()) for i := range placeholders { placeholders[i] = "?" }
return fmt.Sprintf("%s %s (%s)", f.key, f.cmp, strings.Join(placeholders, ", ")) }
return fmt.Sprintf("%s %s ?", f.key, f.cmp)}
func (f filter) Arg() []any { rv := reflect.ValueOf(f.arg) kind := rv.Kind() if (kind == reflect.Slice && rv.Type().Elem().Kind() != reflect.Uint8) || kind == reflect.Array { if rv.Len() == 0 { return nil }
out := make([]any, rv.Len()) for i := range rv.Len() { out[i] = rv.Index(i).Interface() } return out }
return []any{f.arg}}