/* 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 db = new sqlite3.Database('mailship.db') type Param = number | string | boolean type Row = Record 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 = (query: string, params: Param[] = []) => new Promise((resolve, reject) => { db.all(query, params, (err, rows: T[]) => (err ? reject(err) : resolve(rows))) }) // prettier-ignore const get = (query: string, params: Param[] = []) => new Promise((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(p: T | Promise) { return (await p)! } // Migrations export const migrate = () => new Promise(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)`, ) resolve() }) } catch (err) { reject(err) } }) // Subscriptions const parseSubscription = (row: any): Subscription | undefined => { if (row) { return row as Subscription } } export const insertSubscription = instrument( 'database.insertSubscription', async (pubkey: string, email: string, frequency: string) => { const existing = await getSubscriptionByPubkey(pubkey) if (existing) { // Update existing return assertResult( parseSubscription( await get( `UPDATE subscriptions SET email = ?, frequency = ?, confirmed_at = NULL, unsubscribed_at = NULL WHERE pubkey = ? RETURNING *`, [email, frequency, pubkey], ), ), ) } // Create new 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(), ], ), ), ) }, ) 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 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( `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]) }, )