dedupe and serialize subscription writes.
This commit is contained in:
parent
db490f5c99
commit
ce59c15819
1 changed files with 115 additions and 72 deletions
167
src/database.ts
167
src/database.ts
|
|
@ -63,7 +63,7 @@ export const migrate = () =>
|
||||||
unsubscribed_at INTEGER,
|
unsubscribed_at INTEGER,
|
||||||
last_digest_at INTEGER
|
last_digest_at INTEGER
|
||||||
)
|
)
|
||||||
`,
|
`
|
||||||
)
|
)
|
||||||
await run(
|
await run(
|
||||||
`
|
`
|
||||||
|
|
@ -75,10 +75,38 @@ export const migrate = () =>
|
||||||
received_at INTEGER NOT NULL,
|
received_at INTEGER NOT NULL,
|
||||||
PRIMARY KEY (id, subscription_id)
|
PRIMARY KEY (id, subscription_id)
|
||||||
)
|
)
|
||||||
`,
|
`
|
||||||
)
|
)
|
||||||
await run(
|
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()
|
resolve()
|
||||||
})
|
})
|
||||||
|
|
@ -95,43 +123,45 @@ 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(
|
export const insertSubscription = instrument(
|
||||||
'database.insertSubscription',
|
'database.insertSubscription',
|
||||||
async (pubkey: string, email: string, frequency: string) => {
|
async (pubkey: string, email: string, frequency: string) => {
|
||||||
const existing = await getSubscriptionByPubkey(pubkey)
|
const existing = await getSubscriptionByPubkey(pubkey)
|
||||||
|
|
||||||
if (existing) {
|
if (existing) {
|
||||||
// If nothing changed, keep confirmation and don't re-validate
|
return assertResult(await updateSubscription(existing, email, frequency))
|
||||||
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(
|
|
||||||
parseSubscription(
|
|
||||||
await get(
|
|
||||||
`UPDATE subscriptions SET email = ?, frequency = ?, confirmed_at = NULL, unsubscribed_at = NULL
|
|
||||||
WHERE pubkey = ? RETURNING *`,
|
|
||||||
[email, frequency, pubkey],
|
|
||||||
),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
return assertResult(
|
return assertResult(
|
||||||
parseSubscription(
|
parseSubscription(
|
||||||
await get(
|
await get(
|
||||||
|
|
@ -144,11 +174,25 @@ export const insertSubscription = instrument(
|
||||||
email,
|
email,
|
||||||
frequency,
|
frequency,
|
||||||
now(),
|
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(
|
export const confirmSubscription = instrument(
|
||||||
|
|
@ -157,11 +201,11 @@ export const confirmSubscription = instrument(
|
||||||
return parseSubscription(
|
return parseSubscription(
|
||||||
await get(
|
await get(
|
||||||
`UPDATE subscriptions SET confirmed_at = unixepoch()
|
`UPDATE subscriptions SET confirmed_at = unixepoch()
|
||||||
WHERE key = ? AND confirmed_at IS NULL RETURNING *`,
|
WHERE key = ? AND confirmed_at IS NULL AND unsubscribed_at IS NULL RETURNING *`,
|
||||||
[key],
|
[key]
|
||||||
),
|
|
||||||
)
|
)
|
||||||
},
|
)
|
||||||
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const unsubscribeSubscription = instrument(
|
export const unsubscribeSubscription = instrument(
|
||||||
|
|
@ -170,43 +214,42 @@ export const unsubscribeSubscription = instrument(
|
||||||
return parseSubscription(
|
return parseSubscription(
|
||||||
await get(
|
await get(
|
||||||
`UPDATE subscriptions SET unsubscribed_at = unixepoch() WHERE key = ? RETURNING *`,
|
`UPDATE subscriptions SET unsubscribed_at = unixepoch() WHERE key = ? RETURNING *`,
|
||||||
[key],
|
[key]
|
||||||
),
|
|
||||||
)
|
)
|
||||||
},
|
)
|
||||||
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const getSubscriptionById = instrument(
|
export const getSubscriptionById = instrument(
|
||||||
'database.getSubscriptionById',
|
'database.getSubscriptionById',
|
||||||
async (id: string) => {
|
async (id: string) => {
|
||||||
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE id = ?`, [id]))
|
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE id = ?`, [id]))
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const getSubscriptionByKey = instrument(
|
export const getSubscriptionByKey = instrument(
|
||||||
'database.getSubscriptionByKey',
|
'database.getSubscriptionByKey',
|
||||||
async (key: string) => {
|
async (key: string) => {
|
||||||
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE key = ?`, [key]))
|
return parseSubscription(await get(`SELECT * FROM subscriptions WHERE key = ?`, [key]))
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const getSubscriptionByPubkey = instrument(
|
export const getSubscriptionByPubkey = instrument(
|
||||||
'database.getSubscriptionByPubkey',
|
'database.getSubscriptionByPubkey',
|
||||||
async (pubkey: string) => {
|
async (pubkey: string) => {
|
||||||
return parseSubscription(
|
return parseSubscription(
|
||||||
await get(
|
await get(`SELECT * FROM subscriptions WHERE pubkey = ? AND unsubscribed_at IS NULL`, [
|
||||||
`SELECT * FROM subscriptions WHERE pubkey = ? AND unsubscribed_at IS NULL`,
|
pubkey,
|
||||||
[pubkey],
|
])
|
||||||
),
|
|
||||||
)
|
)
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const getActiveSubscriptions = instrument('database.getActiveSubscriptions', async () => {
|
export const getActiveSubscriptions = instrument('database.getActiveSubscriptions', async () => {
|
||||||
const rows = await all(
|
const rows = await all(
|
||||||
`SELECT * FROM subscriptions
|
`SELECT * FROM subscriptions
|
||||||
WHERE confirmed_at IS NOT NULL
|
WHERE confirmed_at IS NOT NULL
|
||||||
AND unsubscribed_at IS NULL`,
|
AND unsubscribed_at IS NULL`
|
||||||
)
|
)
|
||||||
|
|
||||||
return rows.map(parseSubscription) as Subscription[]
|
return rows.map(parseSubscription) as Subscription[]
|
||||||
|
|
@ -216,7 +259,7 @@ export const updateLastDigestAt = instrument(
|
||||||
'database.updateLastDigestAt',
|
'database.updateLastDigestAt',
|
||||||
async (id: string, timestamp: number) => {
|
async (id: string, timestamp: number) => {
|
||||||
await run(`UPDATE subscriptions SET last_digest_at = ? WHERE id = ?`, [timestamp, id])
|
await run(`UPDATE subscriptions SET last_digest_at = ? WHERE id = ?`, [timestamp, id])
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
// Events
|
// Events
|
||||||
|
|
@ -236,7 +279,7 @@ export const insertEvent = instrument(
|
||||||
await run(
|
await run(
|
||||||
`INSERT INTO events (id, subscription_id, event, relay, received_at)
|
`INSERT INTO events (id, subscription_id, event, relay, received_at)
|
||||||
VALUES (?, ?, ?, ?, ?)`,
|
VALUES (?, ?, ?, ?, ?)`,
|
||||||
[eventId, subscriptionId, JSON.stringify(event), relay, now()],
|
[eventId, subscriptionId, JSON.stringify(event), relay, now()]
|
||||||
)
|
)
|
||||||
return true
|
return true
|
||||||
} catch (err: any) {
|
} catch (err: any) {
|
||||||
|
|
@ -246,7 +289,7 @@ export const insertEvent = instrument(
|
||||||
}
|
}
|
||||||
throw err
|
throw err
|
||||||
}
|
}
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const getEventsForSubscription = instrument(
|
export const getEventsForSubscription = instrument(
|
||||||
|
|
@ -256,29 +299,29 @@ export const getEventsForSubscription = instrument(
|
||||||
`SELECT * FROM events
|
`SELECT * FROM events
|
||||||
WHERE subscription_id = ? AND received_at > ?
|
WHERE subscription_id = ? AND received_at > ?
|
||||||
ORDER BY received_at DESC`,
|
ORDER BY received_at DESC`,
|
||||||
[subscriptionId, since],
|
[subscriptionId, since]
|
||||||
)
|
)
|
||||||
|
|
||||||
return rows.map((row) => ({
|
return rows.map((row) => ({
|
||||||
...row,
|
...row,
|
||||||
event: JSON.parse(row.event as any),
|
event: JSON.parse(row.event as any),
|
||||||
}))
|
}))
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const deleteEventsForSubscription = instrument(
|
export const deleteEventsForSubscription = instrument(
|
||||||
'database.deleteEventsForSubscription',
|
'database.deleteEventsForSubscription',
|
||||||
async (subscriptionId: string, since: number) => {
|
async (subscriptionId: string, since: number) => {
|
||||||
await run(
|
await run(`DELETE FROM events WHERE subscription_id = ? AND received_at > ?`, [
|
||||||
`DELETE FROM events WHERE subscription_id = ? AND received_at > ?`,
|
subscriptionId,
|
||||||
[subscriptionId, since],
|
since,
|
||||||
)
|
])
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
export const purgeEventsOlderThan = instrument(
|
export const purgeEventsOlderThan = instrument(
|
||||||
'database.purgeEventsOlderThan',
|
'database.purgeEventsOlderThan',
|
||||||
async (timestamp: number) => {
|
async (timestamp: number) => {
|
||||||
await run(`DELETE FROM events WHERE received_at < ?`, [timestamp])
|
await run(`DELETE FROM events WHERE received_at < ?`, [timestamp])
|
||||||
},
|
}
|
||||||
)
|
)
|
||||||
Loading…
Reference in a new issue