package 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") }