Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240package db
import ( "context" "database/sql" "log/slog" "strings"
"tangled.org/core/log" "tangled.org/core/orm" "tangled.org/core/sqlite")
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 Exec(query string, args ...any) (sql.Result, error)}
func Make(ctx context.Context, dbPath string) (*DB, error) { logger := log.SubLogger(log.FromContext(ctx), "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, schema) if err != nil { return nil, err }
if err := runMigrations(conn, logger); err != nil { return nil, err }
return &DB{db, logger}, nil}
func (d *DB) Close() error { return d.DB.Close()}
const schema = `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 signups_inflight ( id integer primary key autoincrement, email text not null unique, invite_code text not null, verification_token text not null default '', created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), expires_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now', '+24 hours')), last_sent text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')));-- the verification_token index is built by the signup-token-index migration,-- never here: a pre-token database must gain the column before the index.
-- one materialized row per (recipient, source record). deliberi owns these;-- read/emailed are inline. unique(recipient_did, at_uri) dedupes fan-out.create table if not exists notifications ( id integer primary key autoincrement, recipient_did text not null, at_uri text not null, type text not null, actor_did text not null, repo_did text not null default '', knot_did text not null default '', entity_at text not null default '', entity_title text not null default '', read integer not null default 0, emailed integer not null default 0, created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), unique(recipient_did, at_uri));create index if not exists idx_deliberi_notifs_recipient on notifications(recipient_did, read, created);create index if not exists idx_deliberi_notifs_digest on notifications(recipient_did, emailed, read, created);
-- small denormalization caches the ingester fills from repo/issue/pull-- records so notifications render human names without extra lookups.create table if not exists repo_names ( repo_did text primary key, name text not null, owner_did text not null default '');-- repo_did lets a comment find its parent's repo.create table if not exists entity_titles ( at_uri text primary key, title text not null, repo_did text not null default '');
create table if not exists jetstream_cursor ( id integer primary key check (id = 0), last_time_us integer not null);
create table if not exists migrations ( name text primary key);
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, user_mentioned integer not null default 1, email_notifications integer not null default 0);`
func runMigrations(conn *sql.Conn, logger *slog.Logger) error { if err := orm.RunMigration(conn, logger, "add-owner-did-to-repo-names", func(tx *sql.Tx) error { _, err := tx.Exec(`alter table repo_names add column owner_did text not null default ''`) if err != nil && !isColumnExistsErr(err) { return err } return nil }); err != nil { return err } if err := orm.RunMigration(conn, logger, "email-canonical-unique", func(tx *sql.Tx) error { // legacy appview data could hold case-variant or repeated addresses; keep // one row per canonical address: verified wins, else primary, else earliest if _, err := tx.Exec(` delete from emails where id not in ( select id from ( select id, row_number() over ( partition by lower(trim(email)) order by verified desc, is_primary desc, id asc ) as rn from emails ) where rn = 1 ) `); err != nil { return err } // survivors keep legacy casing; normalize so handler lookups match if _, err := tx.Exec(`update emails set email = lower(trim(email))`); err != nil { return err } // legacy AddEmail could auto-promote an unverified first row; the // invariant is verified-only primary if _, err := tx.Exec(`update emails set is_primary = false where verified = false`); err != nil { return err } _, err := tx.Exec(`create unique index if not exists emails_email_lower_unique on emails (lower(email))`) return err }); err != nil { return err } if err := orm.RunMigration(conn, logger, "signup-expiry", func(tx *sql.Tx) error { if _, err := tx.Exec(`alter table signups_inflight add column expires_at text not null default ''`); err != nil && !isColumnExistsErr(err) { return err } // rows predating the column get a 24h lifetime from migration time; new // rows already default to +24h at insert (see the schema) _, err := tx.Exec(`update signups_inflight set expires_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now', '+24 hours') where expires_at = ''`) return err }); err != nil { return err } if err := orm.RunMigration(conn, logger, "signup-token-index", func(tx *sql.Tx) error { // '' rows from before the column existed get unique tokens so the index holds if _, err := tx.Exec(`alter table signups_inflight add column verification_token text not null default ''`); err != nil && !isColumnExistsErr(err) { return err } if _, err := tx.Exec(`update signups_inflight set verification_token = lower(hex(randomblob(32))) where verification_token = ''`); err != nil { return err } _, err := tx.Exec(`create unique index if not exists signups_inflight_verification_token_unique on signups_inflight (verification_token)`) return err }); err != nil { return err } if err := orm.RunMigration(conn, logger, "signup-resend-cooldown", func(tx *sql.Tx) error { // "-1 day": an upgraded row must never block a first send if _, err := tx.Exec(`alter table signups_inflight add column last_sent text not null default ''`); err != nil && !isColumnExistsErr(err) { return err } _, err := tx.Exec(`update signups_inflight set last_sent = strftime('%Y-%m-%dT%H:%M:%SZ', 'now', '-1 day') where last_sent = ''`) return err }); err != nil { return err } if err := orm.RunMigration(conn, logger, "notification-knot-did", func(tx *sql.Tx) error { if _, err := tx.Exec(`alter table notifications add column knot_did text not null default ''`); err != nil && !isColumnExistsErr(err) { return err } return nil }); err != nil { return err } if err := orm.RunMigration(conn, logger, "drop-knot-invited-notifications", func(tx *sql.Tx) error { _, err := tx.Exec(`delete from notifications where type = 'knot_invited'`) return err }); err != nil { return err } return nil}
func isColumnExistsErr(err error) bool { return err != nil && strings.Contains(err.Error(), "duplicate column name")}
// reports whether err is a sqlite unique-constraint violation (e.g. the// lower(email) reservation index)func IsUniqueConstraintErr(err error) bool { return err != nil && strings.Contains(err.Error(), "UNIQUE constraint failed")}