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
91 lines
No EOL
2.7 KiB
TypeScript
91 lines
No EOL
2.7 KiB
TypeScript
import { CronJob } from 'cron'
|
|
import type { Subscription } from '../alert.js'
|
|
import { getCronExpression } from '../alert.js'
|
|
import { Digest } from '../digest.js'
|
|
import * as db from '../database.js'
|
|
|
|
const jobsById = new Map<string, CronJob>()
|
|
|
|
// Test-only accessor to inspect stored jobs
|
|
export const getJobCronSource = (id: string): string | undefined => {
|
|
const source = jobsById.get(id)?.cronTime.source
|
|
return typeof source === 'string' ? source : undefined
|
|
}
|
|
|
|
export const runJob = async (sub: Subscription) => {
|
|
try {
|
|
if (!sub.confirmed_at || sub.unsubscribed_at) {
|
|
console.log('worker: job skipped', sub.id, 'not confirmed or unsubscribed')
|
|
return false
|
|
}
|
|
|
|
console.log('worker: job starting', sub.id)
|
|
|
|
const start = Date.now()
|
|
const since = sub.last_digest_at || sub.created_at
|
|
const events = await db.getEventsForSubscription(sub.id, since)
|
|
|
|
if (events.length === 0) {
|
|
console.log('worker: job skipped', sub.id, 'ok', 'no data received')
|
|
return false
|
|
}
|
|
|
|
const digest = new Digest(sub)
|
|
await digest.sendFromStoredEvents(events)
|
|
|
|
// 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')
|
|
return true
|
|
} catch (e) {
|
|
console.log('worker: job failed', sub.id, e)
|
|
return false
|
|
}
|
|
}
|
|
|
|
const createJob = (sub: Subscription) => {
|
|
const cron = getCronExpression(sub.frequency)
|
|
|
|
const run = async () => {
|
|
// 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({
|
|
cronTime: cron,
|
|
onTick: run,
|
|
start: true,
|
|
timeZone: 'UTC',
|
|
})
|
|
}
|
|
|
|
export const addJob = (sub: Subscription) => {
|
|
jobsById.get(sub.id)?.stop()
|
|
jobsById.set(sub.id, createJob(sub))
|
|
}
|
|
|
|
export const removeJob = (sub: Subscription) => {
|
|
jobsById.get(sub.id)?.stop()
|
|
jobsById.delete(sub.id)
|
|
}
|
|
|
|
// Daily purge of events older than 7 days
|
|
CronJob.from({
|
|
cronTime: '0 0 3 * * *', // 3am UTC daily
|
|
onTick: async () => {
|
|
const weekAgo = Math.floor(Date.now() / 1000) - 7 * 24 * 3600
|
|
await db.purgeEventsOlderThan(weekAgo)
|
|
console.log('worker: purged events older than 7 days')
|
|
},
|
|
start: true,
|
|
timeZone: 'UTC',
|
|
}) |