dedupe and serialize subscription writes.
This commit is contained in:
parent
5abec46cb7
commit
d5f54845b6
1 changed files with 115 additions and 72 deletions
187
src/database.ts
187
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])
|
||||
},
|
||||
)
|
||||
}
|
||||
)
|
||||
|
|
|
|||
Loading…
Reference in a new issue