Fix digest job race: events received between fetch and delete are dropped unsent #8
3 changed files with 162 additions and 3 deletions
|
|
@ -319,6 +319,18 @@ export const deleteEventsForSubscription = instrument(
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
export const deleteEventsByIds = instrument(
|
||||||
|
'database.deleteEventsByIds',
|
||||||
|
async (subscriptionId: string, eventIds: string[]) => {
|
||||||
|
if (eventIds.length === 0) return
|
||||||
|
const placeholders = eventIds.map(() => '?').join(',')
|
||||||
|
await run(
|
||||||
|
`DELETE FROM events WHERE subscription_id = ? AND id IN (${placeholders})`,
|
||||||
|
[subscriptionId, ...eventIds]
|
||||||
|
)
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
export const purgeEventsOlderThan = instrument(
|
export const purgeEventsOlderThan = instrument(
|
||||||
'database.purgeEventsOlderThan',
|
'database.purgeEventsOlderThan',
|
||||||
async (timestamp: number) => {
|
async (timestamp: number) => {
|
||||||
|
|
|
||||||
|
|
@ -33,8 +33,11 @@ export const runJob = async (sub: Subscription) => {
|
||||||
const digest = new Digest(sub)
|
const digest = new Digest(sub)
|
||||||
await digest.sendFromStoredEvents(events)
|
await digest.sendFromStoredEvents(events)
|
||||||
|
|
||||||
// Clean up processed events
|
// Collect the exact IDs that were fetched + sent, then delete ONLY those.
|
||||||
await db.deleteEventsForSubscription(sub.id, since)
|
// Deleting by timestamp (received_at > since) would also remove any event
|
||||||
|
// that arrived between the fetch and the delete — the race condition.
|
||||||
|
const sentIds = events.map(e => e.id)
|
||||||
|
await db.deleteEventsByIds(sub.id, sentIds)
|
||||||
await db.updateLastDigestAt(sub.id, Math.floor(Date.now() / 1000))
|
await db.updateLastDigestAt(sub.id, Math.floor(Date.now() / 1000))
|
||||||
|
|
||||||
console.log('worker: job completed', sub.id, 'in', Date.now() - start, 'ms')
|
console.log('worker: job completed', sub.id, 'in', Date.now() - start, 'ms')
|
||||||
|
|
@ -49,7 +52,12 @@ const createJob = (sub: Subscription) => {
|
||||||
const cron = getCronExpression(sub.frequency)
|
const cron = getCronExpression(sub.frequency)
|
||||||
|
|
||||||
const run = async () => {
|
const run = async () => {
|
||||||
await runJob(sub)
|
// Re-fetch the subscription to pick up the latest last_digest_at.
|
||||||
|
// The closure-captured `sub` is stale — its last_digest_at never
|
||||||
|
// advances, so every tick would re-fetch and re-send old events.
|
||||||
|
const fresh = await db.getSubscriptionById(sub.id)
|
||||||
|
if (!fresh) return
|
||||||
|
await runJob(fresh)
|
||||||
}
|
}
|
||||||
|
|
||||||
return CronJob.from({
|
return CronJob.from({
|
||||||
|
|
|
||||||
139
test/event-arrival-race.test.js
Normal file
139
test/event-arrival-race.test.js
Normal file
|
|
@ -0,0 +1,139 @@
|
||||||
|
#!/usr/bin/env node
|
||||||
|
// Test: event-arrival race in digest job — fix verification
|
||||||
|
//
|
||||||
|
// The original bug: runJob fetched events via getEventsForSubscription(sub.id, since),
|
||||||
|
// sent the digest (slow), then deleted events via deleteEventsForSubscription(sub.id, since)
|
||||||
|
// which deletes EVERY row with received_at > since. Any event that arrived between
|
||||||
|
// the fetch and the delete was also received_at > since, so it was deleted without
|
||||||
|
// ever being emailed.
|
||||||
|
//
|
||||||
|
// The fix: deleteEventsByIds(sub.id, sentEventIds) deletes only the exact event IDs
|
||||||
|
// that were fetched and sent. A late-arriving event (inserted after the fetch) has
|
||||||
|
// a different ID and is not touched.
|
||||||
|
//
|
||||||
|
// This test simulates the sequence with the fixed approach:
|
||||||
|
// 1. Create a confirmed subscription with a known last_digest_at
|
||||||
|
// 2. Insert event A
|
||||||
|
// 3. Fetch events for the subscription (simulating runJob's fetch)
|
||||||
|
// 4. Insert event B after the fetch (simulating a /notify arriving during digest send)
|
||||||
|
// 5. Delete ONLY the event A IDs (the fix — instead of timestamp-based delete)
|
||||||
|
// 6. Verify event B survives (it arrived after fetch and was never sent)
|
||||||
|
|
||||||
|
import * as db from '../dist/database.js'
|
||||||
|
|
||||||
|
let passed = 0
|
||||||
|
let failed = 0
|
||||||
|
|
||||||
|
function assert(label, ok, detail) {
|
||||||
|
if (ok) {
|
||||||
|
console.log(` ? ${label}`)
|
||||||
|
passed++
|
||||||
|
} else {
|
||||||
|
console.log(` ? ${label} -- ${detail || ''}`)
|
||||||
|
failed++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function sleep(ms) {
|
||||||
|
return new Promise(resolve => setTimeout(resolve, ms))
|
||||||
|
}
|
||||||
|
|
||||||
|
async function main() {
|
||||||
|
await db.migrate()
|
||||||
|
|
||||||
|
const pubkey = 'race-test-' + Date.now()
|
||||||
|
const email = 'race-test-' + Date.now() + '@example.com'
|
||||||
|
|
||||||
|
// Step 1: Create and confirm subscription
|
||||||
|
console.log('1. Create confirmed subscription')
|
||||||
|
const sub = await db.insertSubscription(pubkey, email, 'daily')
|
||||||
|
assert('subscription created', !!sub, 'insert returned null')
|
||||||
|
if (!sub) { process.exit(1) }
|
||||||
|
|
||||||
|
const confirmed = await db.confirmSubscription(sub.key)
|
||||||
|
assert('subscription confirmed', !!confirmed, 'confirm returned null')
|
||||||
|
if (!confirmed) { process.exit(1) }
|
||||||
|
|
||||||
|
// Set last_digest_at to a known time well in the past.
|
||||||
|
// This is the "since" value runJob would use.
|
||||||
|
const since = Math.floor(Date.now() / 1000) - 60 // 60 seconds ago
|
||||||
|
await db.updateLastDigestAt(confirmed.id, since)
|
||||||
|
|
||||||
|
// Wait 1s so event received_at timestamps (whole seconds) are strictly > since.
|
||||||
|
await sleep(1100)
|
||||||
|
|
||||||
|
// Step 2: Insert Event A — simulates events that trigger a digest run
|
||||||
|
console.log('\n2. Insert Event A (triggers digest)')
|
||||||
|
const eventA_id = 'race-event-a-' + Date.now()
|
||||||
|
const eventA = {
|
||||||
|
id: eventA_id,
|
||||||
|
kind: 1,
|
||||||
|
pubkey: 'abc',
|
||||||
|
content: 'event A content',
|
||||||
|
created_at: Math.floor(Date.now() / 1000),
|
||||||
|
tags: []
|
||||||
|
}
|
||||||
|
const storedA = await db.insertEvent(eventA_id, confirmed.id, eventA, 'wss://relay.damus.io')
|
||||||
|
assert('Event A stored', storedA === true, `got ${storedA}`)
|
||||||
|
|
||||||
|
// Step 3: Fetch events — simulating what runJob does with since=last_digest_at
|
||||||
|
console.log('\n3. Fetch events for subscription (simulating runJob fetch)')
|
||||||
|
const fetched = await db.getEventsForSubscription(confirmed.id, since)
|
||||||
|
assert('Event A was fetched', fetched.some(e => e.id === eventA_id),
|
||||||
|
`fetched ids: [${fetched.map(e => e.id).join(', ')}]`)
|
||||||
|
|
||||||
|
// Step 4: Insert Event B after the fetch — simulating a /notify arriving
|
||||||
|
// during the slow digest send (profile loads, MJML render, SMTP).
|
||||||
|
console.log('\n4. Insert Event B AFTER fetch (simulating event arriving during digest send)')
|
||||||
|
await sleep(100) // ensure distinct received_at
|
||||||
|
const eventB_id = 'race-event-b-' + Date.now()
|
||||||
|
const eventB = {
|
||||||
|
id: eventB_id,
|
||||||
|
kind: 1,
|
||||||
|
pubkey: 'def',
|
||||||
|
content: 'event B content — arrived during send',
|
||||||
|
created_at: Math.floor(Date.now() / 1000),
|
||||||
|
tags: []
|
||||||
|
}
|
||||||
|
const storedB = await db.insertEvent(eventB_id, confirmed.id, eventB, 'wss://relay.damus.io')
|
||||||
|
assert('Event B stored', storedB === true, `got ${storedB}`)
|
||||||
|
|
||||||
|
// Step 5: Delete ONLY the exact event IDs that were fetched + sent — THE FIX.
|
||||||
|
// Unlike the old timestamp-based delete, this does NOT touch Event B
|
||||||
|
// because Event B has a different ID.
|
||||||
|
console.log('\n5. Delete events by exact IDs (the fix — deleteEventsByIds)')
|
||||||
|
const sentIds = fetched.map(e => e.id)
|
||||||
|
await db.deleteEventsByIds(confirmed.id, sentIds)
|
||||||
|
console.log(' deleted ids:', JSON.stringify(sentIds))
|
||||||
|
|
||||||
|
// Step 6: Check which events remain.
|
||||||
|
// CORRECT BEHAVIOR: Only Event A (which was sent) is deleted.
|
||||||
|
// Event B (which arrived after fetch) must survive.
|
||||||
|
console.log('\n6. Check remaining events after delete')
|
||||||
|
const remaining = await db.getEventsForSubscription(confirmed.id, since - 10)
|
||||||
|
|
||||||
|
const eventB_survived = remaining.some(e => e.id === eventB_id)
|
||||||
|
assert(
|
||||||
|
'Event B survives the delete (it arrived after fetch and was never sent)',
|
||||||
|
eventB_survived,
|
||||||
|
`Event B was deleted despite never being sent. ` +
|
||||||
|
`remaining events: [${remaining.map(e => e.id).join(', ')}]`
|
||||||
|
)
|
||||||
|
|
||||||
|
// Verify Event A is gone (it was sent, so deletion is correct for A)
|
||||||
|
const eventA_survived = remaining.some(e => e.id === eventA_id)
|
||||||
|
assert(
|
||||||
|
'Event A is deleted (it was fetched and sent)',
|
||||||
|
!eventA_survived,
|
||||||
|
`Event A should have been deleted but is still present`
|
||||||
|
)
|
||||||
|
|
||||||
|
console.log('')
|
||||||
|
console.log(`Results: ${passed} passed, ${failed} failed`)
|
||||||
|
process.exit(failed > 0 ? 1 : 0)
|
||||||
|
}
|
||||||
|
|
||||||
|
main().catch(err => {
|
||||||
|
console.error('Unhandled error in test:', err)
|
||||||
|
process.exit(1)
|
||||||
|
})
|
||||||
Loading…
Reference in a new issue