dedupe and serialize subscription writes.

This commit is contained in:
mplorentz 2026-09-03 13:42:50 -04:00
parent db490f5c99
commit ce59c15819

View file

@ -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])
},
}
)