From 737f949f9ba594fdc480a5c547cc787593cd70ab Mon Sep 17 00:00:00 2001 From: Agent Date: Mon, 14 Sep 2026 13:07:34 -0400 Subject: [PATCH] fix digest job race: delete sent events by ID, not timestamp Bug: runJob computed since = sub.last_digest_at, fetched events with received_at > since, sent the digest (slow), then deleted ALL events with received_at > since. Any /notify event that arrived between the fetch and the delete was also received_at > since, so it was deleted without ever being sent in a digest. Two changes: 1. Delete by exact event IDs (src/database.ts, src/worker/email.ts): Added deleteEventsByIds(subscriptionId, eventIds) which deletes only the events that were actually fetched + sent. The old timestamp-based delete is retained but no longer called from runJob. 2. Re-fetch subscription on each cron tick (src/worker/email.ts): createJob's closure captured the original sub, so sub.last_digest_at stayed stale in memory. Every subsequent tick recomputed since from the old value, re-fetching and re-sending duplicate events. Now each tick re-fetches the subscription from the DB via getSubscriptionById before calling runJob. Fixes bead mailship-091 --- src/database.ts | 12 +++ src/worker/email.ts | 14 +++- test/event-arrival-race.test.js | 139 ++++++++++++++++++++++++++++++++ 3 files changed, 162 insertions(+), 3 deletions(-) create mode 100644 test/event-arrival-race.test.js diff --git a/src/database.ts b/src/database.ts index 336e293..718e225 100644 --- a/src/database.ts +++ b/src/database.ts @@ -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( 'database.purgeEventsOlderThan', async (timestamp: number) => { diff --git a/src/worker/email.ts b/src/worker/email.ts index 7ef4076..06a546e 100644 --- a/src/worker/email.ts +++ b/src/worker/email.ts @@ -33,8 +33,11 @@ export const runJob = async (sub: Subscription) => { const digest = new Digest(sub) await digest.sendFromStoredEvents(events) - // Clean up processed events - await db.deleteEventsForSubscription(sub.id, since) + // Collect the exact IDs that were fetched + sent, then delete ONLY those. + // 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)) 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 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({ diff --git a/test/event-arrival-race.test.js b/test/event-arrival-race.test.js new file mode 100644 index 0000000..8a938c7 --- /dev/null +++ b/test/event-arrival-race.test.js @@ -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) +}) \ No newline at end of file