mailship/src/database.ts
2026-09-23 12:12:51 -04:00

448 lines
14 KiB
TypeScript

/* eslint @typescript-eslint/no-unused-vars: 0 */
import sqlite3 from 'sqlite3'
import crypto from 'crypto'
import { instrument } from 'succinct-async'
import { now } from '@welshman/lib'
import type { Subscription } from './alert.js'
const DATA_DIR = process.env.DATA_DIR || '.'
const db = new sqlite3.Database(DATA_DIR + '/mailship.db')
type Param = number | string | boolean | null | undefined
type Row = Record<string, any>
const run = (query: string, params: Param[] = []) =>
new Promise((resolve, reject) => {
db.run(query, params, function (err) {
return err ? reject(err) : resolve(this.changes > 0)
})
})
// prettier-ignore
const all = <T=Row>(query: string, params: Param[] = []) =>
new Promise<T[]>((resolve, reject) => {
db.all(query, params, (err, rows: T[]) => (err ? reject(err) : resolve(rows)))
})
// prettier-ignore
const get = <T=Row>(query: string, params: Param[] = []) =>
new Promise<T | undefined>((resolve, reject) => {
db.get(query, params, (err, row) => {
if (err) {
reject(err)
} else if (row) {
resolve(row as T)
} else {
resolve(undefined)
}
})
})
async function assertResult<T>(p: T | Promise<T>) {
return (await p)!
}
// Migrations
export const migrate = () =>
new Promise<void>(async (resolve, reject) => {
try {
db.serialize(async () => {
await run(
`
CREATE TABLE IF NOT EXISTS subscriptions (
id TEXT PRIMARY KEY,
key TEXT NOT NULL UNIQUE,
pubkey TEXT NOT NULL,
email TEXT NOT NULL,
frequency TEXT NOT NULL DEFAULT 'daily',
hour INTEGER NOT NULL DEFAULT 17,
minute INTEGER NOT NULL DEFAULT 0,
day_of_week INTEGER,
timezone TEXT NOT NULL DEFAULT 'UTC',
created_at INTEGER NOT NULL,
confirmed_at INTEGER,
unsubscribed_at INTEGER,
last_digest_at INTEGER
)
`
)
await run(
`
CREATE TABLE IF NOT EXISTS events (
id TEXT NOT NULL,
subscription_id TEXT NOT NULL,
event JSON NOT NULL,
relay TEXT NOT NULL,
received_at INTEGER NOT NULL,
PRIMARY KEY (id, subscription_id)
)
`
)
await run(
`CREATE INDEX IF NOT EXISTS idx_events_subscription_received ON events (subscription_id, received_at)`
)
// Tombstone older duplicate active rows, so the unique index below can be created
// even if a previous version of the server let races insert more than one.
await run(
`
UPDATE subscriptions
SET unsubscribed_at = created_at
WHERE unsubscribed_at IS NULL
AND id NOT IN (
SELECT id FROM (
SELECT id, ROW_NUMBER() OVER (
PARTITION BY pubkey ORDER BY created_at DESC, id DESC
) AS rn
FROM subscriptions
WHERE unsubscribed_at IS NULL
)
WHERE rn = 1
)
`
)
// At most one active subscription per pubkey. Enforced in the DB so a
// check-then-insert race can never create duplicate digest rows.
await run(
`
CREATE UNIQUE INDEX IF NOT EXISTS idx_subscriptions_active_pubkey
ON subscriptions (pubkey)
WHERE unsubscribed_at IS NULL
`
)
// Idempotent migration: add schedule columns to existing databases.
// ALTER TABLE ADD COLUMN fails if the column already exists, so we
// run each one and ignore "duplicate column" errors.
const addCol = (colDef: string) =>
run(`ALTER TABLE subscriptions ADD COLUMN ${colDef}`).catch((err: any) => {
if (!err?.message?.includes('duplicate column')) throw err
})
await addCol('hour INTEGER NOT NULL DEFAULT 17')
await addCol('minute INTEGER NOT NULL DEFAULT 0')
await addCol('day_of_week INTEGER')
await addCol("timezone TEXT NOT NULL DEFAULT 'UTC'")
resolve()
})
} catch (err) {
reject(err)
}
})
// Subscriptions
const parseSubscription = (row: any): Subscription | undefined => {
if (row) {
return row as Subscription
}
}
export const updateSubscription = instrument(
'database.updateSubscription',
async (existing: Subscription, email: string, frequency: string, hour?: number, minute?: number, dayOfWeek?: number, timezone?: string) => {
const scheduleChanged =
(hour !== undefined && hour !== existing.hour) ||
(minute !== undefined && minute !== existing.minute) ||
(dayOfWeek !== undefined && dayOfWeek !== existing.day_of_week) ||
(timezone !== undefined && timezone !== existing.timezone)
if (existing.email === email && existing.frequency === frequency && !scheduleChanged) {
return existing
}
const hr = hour ?? existing.hour
const mn = minute ?? existing.minute
const dow = dayOfWeek ?? existing.day_of_week
const tz = timezone ?? existing.timezone
// Update by id, not pubkey, so tombstoned rows for the same account are never
// re-activated alongside this one (the unique index would reject that anyway).
if (existing.email === email) {
return parseSubscription(
await get(
`UPDATE subscriptions SET frequency = ?, hour = ?, minute = ?, day_of_week = ?, timezone = ?, unsubscribed_at = NULL
WHERE id = ? RETURNING *`,
[frequency, hr, mn, dow, tz, existing.id]
)
)
}
return parseSubscription(
await get(
`UPDATE subscriptions SET email = ?, frequency = ?, hour = ?, minute = ?, day_of_week = ?, timezone = ?, confirmed_at = NULL, unsubscribed_at = NULL
WHERE id = ? RETURNING *`,
[email, frequency, hr, mn, dow, tz, existing.id]
)
)
}
)
const getMostRecentTombstonedSubscription = async (pubkey: string) =>
parseSubscription(
await get(
`SELECT * FROM subscriptions
WHERE pubkey = ? AND unsubscribed_at IS NOT NULL
ORDER BY created_at DESC, id DESC LIMIT 1`,
[pubkey]
)
)
const reactivateSubscription = async (tombstoned: Subscription, frequency: string) =>
parseSubscription(
await get(
`UPDATE subscriptions SET key = ?, frequency = ?, unsubscribed_at = NULL
WHERE id = ? RETURNING *`,
[crypto.randomBytes(32).toString('hex'), frequency, tombstoned.id]
)
)
export const insertSubscription = instrument(
'database.insertSubscription',
async (pubkey: string, email: string, frequency: string, hour?: number, minute?: number, dayOfWeek?: number, timezone?: string) => {
const existing = await getSubscriptionByPubkey(pubkey)
if (existing) {
return assertResult(await updateSubscription(existing, email, frequency, hour, minute, dayOfWeek, timezone))
}
// No active row. If this account previously subscribed to the same email
// and was unsubscribed, reactivate that row instead of inserting a fresh
// one. Otherwise every turn-on → turn-off → turn-on cycle would create a
// new row, abandoning the confirmed state (and the already-known key).
const tombstoned = await getMostRecentTombstonedSubscription(pubkey)
if (tombstoned?.email === email) {
try {
return assertResult(await reactivateSubscription(tombstoned, frequency))
} catch (err: any) {
if (err.message?.includes('UNIQUE constraint')) {
const concurrent = await getSubscriptionByPubkey(pubkey)
if (concurrent) {
return assertResult(await updateSubscription(concurrent, email, frequency))
}
}
throw err
}
}
try {
return assertResult(
parseSubscription(
await get(
`INSERT INTO subscriptions (id, key, pubkey, email, frequency, hour, minute, day_of_week, timezone, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING *`,
[
crypto.randomUUID(),
crypto.randomBytes(32).toString('hex'),
pubkey,
email,
frequency,
hour ?? 17,
minute ?? 0,
dayOfWeek ?? null,
timezone ?? 'UTC',
now(),
]
)
)
)
} catch (err: any) {
// A concurrent request inserted the active subscription between our select
// and insert. The partial unique index forces one active row per pubkey,
// so fall back to updating the row that won the race.
if (err.message?.includes('UNIQUE constraint')) {
const concurrent = await getSubscriptionByPubkey(pubkey)
if (concurrent) {
return assertResult(await updateSubscription(concurrent, email, frequency, hour, minute, dayOfWeek, timezone))
}
}
throw err
}
}
)
export type ConfirmResult = {
sub: Subscription
alreadyConfirmed: boolean
}
export const confirmSubscription = instrument(
'database.confirmSubscription',
async (key: string): Promise<ConfirmResult | undefined> => {
// Try to update an unconfirmed, active row
const updated = parseSubscription(
await get(
`UPDATE subscriptions SET confirmed_at = unixepoch()
WHERE key = ? AND confirmed_at IS NULL AND unsubscribed_at IS NULL RETURNING *`,
[key]
)
)
if (updated) {
return { sub: updated, alreadyConfirmed: false }
}
// No unconfirmed row was updated. Check if the key exists and is
// already confirmed — the user is re-clicking a used link. If the row
// is unsubscribed, treat it as invalid (expired).
const existing = await getSubscriptionByKey(key)
if (!existing || existing.unsubscribed_at) {
return undefined
}
return { sub: existing, alreadyConfirmed: true }
}
)
export const unsubscribeSubscription = instrument(
'database.unsubscribeSubscription',
async (key: string) => {
return parseSubscription(
await get(
`UPDATE subscriptions SET unsubscribed_at = unixepoch() WHERE key = ? RETURNING *`,
[key]
)
)
}
)
export const getSubscriptionById = instrument(
'database.getSubscriptionById',
async (id: string) => {
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE id = ?`, [id]))
}
)
export const getSubscriptionByKey = instrument(
'database.getSubscriptionByKey',
async (key: string) => {
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE key = ?`, [key]))
}
)
export const getSubscriptionByPubkey = instrument(
'database.getSubscriptionByPubkey',
async (pubkey: string) => {
return parseSubscription(
await get(`SELECT * FROM subscriptions WHERE pubkey = ? AND unsubscribed_at IS NULL`, [
pubkey,
])
)
}
)
export const getActiveSubscriptions = instrument('database.getActiveSubscriptions', async () => {
const rows = await all(
`SELECT * FROM subscriptions
WHERE confirmed_at IS NOT NULL
AND unsubscribed_at IS NULL`
)
return rows.map(parseSubscription) as Subscription[]
})
// Every row ever created for a pubkey — both active and tombstoned. Exposed
// primarily for tests asserting re-subscribe never grows the table.
export const getAllSubscriptionsByPubkey = instrument(
'database.getAllSubscriptionsByPubkey',
async (pubkey: string) => {
const rows = await all<Row>(
`SELECT * FROM subscriptions WHERE pubkey = ? ORDER BY created_at`,
[pubkey]
)
return rows.map(parseSubscription) as Subscription[]
}
)
export const updateLastDigestAt = instrument(
'database.updateLastDigestAt',
async (id: string, timestamp: number) => {
await run(`UPDATE subscriptions SET last_digest_at = ? WHERE id = ?`, [timestamp, id])
}
)
// Events
export type StoredEvent = {
id: string
subscription_id: string
event: any
relay: string
received_at: number
}
export const insertEvent = instrument(
'database.insertEvent',
async (eventId: string, subscriptionId: string, event: any, relay: string) => {
try {
await run(
`INSERT INTO events (id, subscription_id, event, relay, received_at)
VALUES (?, ?, ?, ?, ?)`,
[eventId, subscriptionId, JSON.stringify(event), relay, now()]
)
return true
} catch (err: any) {
// PRIMARY KEY collision = dedup, not an error
if (err.message?.includes('UNIQUE constraint')) {
return false
}
throw err
}
}
)
export const getEventsForSubscription = instrument(
'database.getEventsForSubscription',
async (subscriptionId: string, since: number) => {
const rows = await all<StoredEvent>(
`SELECT * FROM events
WHERE subscription_id = ? AND received_at > ?
ORDER BY received_at DESC`,
[subscriptionId, since]
)
return rows.map((row) => ({
...row,
event: JSON.parse(row.event as any),
}))
}
)
export const deleteEventsForSubscription = instrument(
'database.deleteEventsForSubscription',
async (subscriptionId: string, since: number) => {
await run(`DELETE FROM events WHERE subscription_id = ? AND received_at > ?`, [
subscriptionId,
since,
])
}
)
export const deleteEventsByIds = instrument(
'database.deleteEventsByIds',
async (subscriptionId: string, eventIds: string[]) => {
if (eventIds.length === 0) return
const placeholders = eventIds.map(() => '?').join(',')
await run(
`DELETE FROM events WHERE subscription_id = ? AND id IN (${placeholders})`,
[subscriptionId, ...eventIds]
)
}
)
export const purgeEventsOlderThan = instrument(
'database.purgeEventsOlderThan',
async (timestamp: number) => {
await run(`DELETE FROM events WHERE received_at < ?`, [timestamp])
}
)