From ce59c158193a4e1b1662667c5d711e58bfc55423 Mon Sep 17 00:00:00 2001 From: mplorentz Date: Thu, 3 Sep 2026 13:42:50 -0400 Subject: [PATCH] dedupe and serialize subscription writes. --- src/database.ts | 187 +++++++++++++++++++++++++++++------------------- 1 file changed, 115 insertions(+), 72 deletions(-) diff --git a/src/database.ts b/src/database.ts index 268675a..336e293 100644 --- a/src/database.ts +++ b/src/database.ts @@ -63,7 +63,7 @@ export const migrate = () => unsubscribed_at INTEGER, last_digest_at INTEGER ) - `, + ` ) await run( ` @@ -75,10 +75,38 @@ export const migrate = () => 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)`, + `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() }) @@ -95,60 +123,76 @@ const parseSubscription = (row: any): Subscription | undefined => { } } +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) { - // If nothing changed, keep confirmation and don't re-validate - if (existing.email === email && existing.frequency === frequency) { - return assertResult(parseSubscription(existing)) - } - - // Update existing. Only a change of email address invalidates the - // existing confirmation (the new address must be verified); changing - // the frequency keeps it confirmed. - if (existing.email === email) { - return assertResult( - parseSubscription( - await get( - `UPDATE subscriptions SET frequency = ?, unsubscribed_at = NULL - WHERE pubkey = ? RETURNING *`, - [frequency, pubkey], - ), - ), - ) - } + return assertResult(await updateSubscription(existing, email, frequency)) + } + try { return assertResult( parseSubscription( await get( - `UPDATE subscriptions SET email = ?, frequency = ?, confirmed_at = NULL, unsubscribed_at = NULL - WHERE pubkey = ? RETURNING *`, - [email, frequency, pubkey], - ), - ), + `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) - 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(), - ], - ), - ), - ) - }, + if (concurrent) { + return assertResult(await updateSubscription(concurrent, email, frequency)) + } + } + + throw err + } + } ) export const confirmSubscription = instrument( @@ -157,11 +201,11 @@ export const confirmSubscription = instrument( return parseSubscription( await get( `UPDATE subscriptions SET confirmed_at = unixepoch() - WHERE key = ? AND confirmed_at IS NULL RETURNING *`, - [key], - ), + WHERE key = ? AND confirmed_at IS NULL AND unsubscribed_at IS NULL RETURNING *`, + [key] + ) ) - }, + } ) export const unsubscribeSubscription = instrument( @@ -170,43 +214,42 @@ export const unsubscribeSubscription = instrument( return parseSubscription( await get( `UPDATE subscriptions SET unsubscribed_at = unixepoch() WHERE key = ? RETURNING *`, - [key], - ), + [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], - ), + 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`, + AND unsubscribed_at IS NULL` ) return rows.map(parseSubscription) as Subscription[] @@ -216,7 +259,7 @@ export const updateLastDigestAt = instrument( 'database.updateLastDigestAt', async (id: string, timestamp: number) => { await run(`UPDATE subscriptions SET last_digest_at = ? WHERE id = ?`, [timestamp, id]) - }, + } ) // Events @@ -236,7 +279,7 @@ export const insertEvent = instrument( await run( `INSERT INTO events (id, subscription_id, event, relay, received_at) VALUES (?, ?, ?, ?, ?)`, - [eventId, subscriptionId, JSON.stringify(event), relay, now()], + [eventId, subscriptionId, JSON.stringify(event), relay, now()] ) return true } catch (err: any) { @@ -246,7 +289,7 @@ export const insertEvent = instrument( } throw err } - }, + } ) export const getEventsForSubscription = instrument( @@ -256,29 +299,29 @@ export const getEventsForSubscription = instrument( `SELECT * FROM events WHERE subscription_id = ? AND received_at > ? ORDER BY received_at DESC`, - [subscriptionId, since], + [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], - ) - }, + 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]) - }, -) \ No newline at end of file + } +)