mailship/src/database.ts
2026-09-03 13:44:02 -04:00

327 lines
9 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
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',
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
`
)
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) => {
if (existing.email === email && existing.frequency === frequency) {
return existing
}
// 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 = ?, unsubscribed_at = NULL
WHERE id = ? RETURNING *`,
[frequency, existing.id]
)
)
}
return parseSubscription(
await get(
`UPDATE subscriptions SET email = ?, frequency = ?, confirmed_at = NULL, unsubscribed_at = NULL
WHERE id = ? RETURNING *`,
[email, frequency, existing.id]
)
)
}
)
export const insertSubscription = instrument(
'database.insertSubscription',
async (pubkey: string, email: string, frequency: string) => {
const existing = await getSubscriptionByPubkey(pubkey)
if (existing) {
return assertResult(await updateSubscription(existing, email, frequency))
}
try {
return assertResult(
parseSubscription(
await get(
`INSERT INTO subscriptions (id, key, pubkey, email, frequency, created_at)
VALUES (?, ?, ?, ?, ?, ?) RETURNING *`,
[
crypto.randomUUID(),
crypto.randomBytes(32).toString('hex'),
pubkey,
email,
frequency,
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))
}
}
throw err
}
}
)
export const confirmSubscription = instrument(
'database.confirmSubscription',
async (key: string) => {
return parseSubscription(
await get(
`UPDATE subscriptions SET confirmed_at = unixepoch()
WHERE key = ? AND confirmed_at IS NULL AND unsubscribed_at IS NULL RETURNING *`,
[key]
)
)
}
)
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[]
})
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 purgeEventsOlderThan = instrument(
'database.purgeEventsOlderThan',
async (timestamp: number) => {
await run(`DELETE FROM events WHERE received_at < ?`, [timestamp])
}
)