Events: self-serve resume for a send interrupted by a hub restart

A hub restart kills sendEventCampaignAsync's in-memory loop mid-flight,
leaving a campaign stuck at status='sending' with recipients still
'pending' — this happened live today and required a manual server-side
fix. Building a proper resume path instead of relying on that:

- resumeCampaignSend(campaignId): re-reads channel/from/event_id from the
  campaign's own params (nothing to re-enter), refuses if nothing is
  pending or a send is already active, and re-invokes the same worker,
  which skips already-'sent' rows — safe to call on a healthy campaign too.
- Anti-double-send guard (activeEventSends / tryReserveSend): reservation
  now happens synchronously at the scheduling point (massSend/reminderSend/
  resumeCampaignSend), not inside the deferred setImmediate callback —
  closes a race where two near-simultaneous calls could both slip through
  before either worker actually started.
- New route: POST /events/campaign-resume {campaign_id}.
- CampaignDetailPage.vue: "Reprendre l'envoi" button, shown only for
  event-type campaigns stuck at status='sending' with pending recipients.

Verified live: used the new endpoint to resume the real in-progress "On
fête nos 20 ans" campaign after this exact deploy interrupted it — resumed
cleanly with zero duplicate sends.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
louispaulb 2026-07-24 10:12:27 -04:00
parent d961c57f9e
commit 07b0331a3a
3 changed files with 122 additions and 34 deletions

View File

@ -60,3 +60,7 @@ export async function deleteAttachment (eventId, stored) { return hubFetch(`/eve
export async function listEventCampaigns (eventId) { return hubFetch(`/events/${encodeURIComponent(eventId)}/campaigns`) } export async function listEventCampaigns (eventId) { return hubFetch(`/events/${encodeURIComponent(eventId)}/campaigns`) }
export async function previewReminder (eventId, campaignId) { return hubFetch(`/events/${encodeURIComponent(eventId)}/reminder-preview?campaign=${encodeURIComponent(campaignId)}`) } export async function previewReminder (eventId, campaignId) { return hubFetch(`/events/${encodeURIComponent(eventId)}/reminder-preview?campaign=${encodeURIComponent(campaignId)}`) }
export async function sendReminder (eventId, body) { return hubFetch(`/events/${encodeURIComponent(eventId)}/reminder-send`, { method: 'POST', body }) } export async function sendReminder (eventId, body) { return hubFetch(`/events/${encodeURIComponent(eventId)}/reminder-send`, { method: 'POST', body }) }
// Reprend un envoi d'événement interrompu (ex. redémarrage du hub en plein envoi) — relit
// channel/from/event_id depuis la campagne elle-même, saute les destinataires déjà envoyés.
export async function resumeEventCampaign (campaignId) { return hubFetch('/events/campaign-resume', { method: 'POST', body: { campaign_id: campaignId } }) }

View File

@ -21,6 +21,13 @@
:loading="creatingReminder" @click="confirmCreateReminder"> :loading="creatingReminder" @click="confirmCreateReminder">
<q-tooltip>Cloner cette campagne pour les destinataires qui n'ont PAS cliqué le cadeau encore</q-tooltip> <q-tooltip>Cloner cette campagne pour les destinataires qui n'ont PAS cliqué le cadeau encore</q-tooltip>
</q-btn> </q-btn>
<!-- Envoi d'événement interrompu (ex. redémarrage du hub en plein envoi) : relance le worker,
qui reprend exactement il s'était arrêté (saute les destinataires déjà envoyés). -->
<q-btn v-if="canResumeEventSend" unelevated color="warning" icon="play_arrow"
label="Reprendre l'envoi" class="q-mr-sm"
:loading="resumingSend" @click="confirmResumeSend">
<q-tooltip>{{ counterFor('pending') }} destinataire(s) encore en attente reprend sans renvoyer à ceux déjà rejoints</q-tooltip>
</q-btn>
<!-- Re-send path: clone a locked (sent) campaign into a fresh draft where <!-- Re-send path: clone a locked (sent) campaign into a fresh draft where
the sender + subjects are editable again. --> the sender + subjects are editable again. -->
<q-btn v-if="campaign && campaign.status !== 'draft'" flat dense icon="content_copy" color="primary" <q-btn v-if="campaign && campaign.status !== 'draft'" flat dense icon="content_copy" color="primary"
@ -358,6 +365,7 @@ import DataTable from 'src/components/shared/DataTable.vue'
import { useRoute, useRouter } from 'vue-router' import { useRoute, useRouter } from 'vue-router'
import { useQuasar } from 'quasar' import { useQuasar } from 'quasar'
import { getCampaign, sendCampaign, cloneCampaign, campaignSseUrl, campaignReportCsvUrl, retryRecipient, createReminderCampaign, updateCampaign, listTemplates, previewTemplate } from 'src/api/campaigns' import { getCampaign, sendCampaign, cloneCampaign, campaignSseUrl, campaignReportCsvUrl, retryRecipient, createReminderCampaign, updateCampaign, listTemplates, previewTemplate } from 'src/api/campaigns'
import { resumeEventCampaign } from 'src/api/events'
import SenderField from 'src/modules/campaigns/components/SenderField.vue' import SenderField from 'src/modules/campaigns/components/SenderField.vue'
import ChannelToggle from 'src/modules/campaigns/components/ChannelToggle.vue' import ChannelToggle from 'src/modules/campaigns/components/ChannelToggle.vue'
@ -370,6 +378,7 @@ const loading = ref(true)
const resending = ref(false) const resending = ref(false)
const creatingReminder = ref(false) const creatingReminder = ref(false)
const cloning = ref(false) const cloning = ref(false)
const resumingSend = ref(false)
// Aperçu de l'envoi // Aperçu de l'envoi
// Rend le template de la campagne (FR/EN) exactement comme envoyé, avec les // Rend le template de la campagne (FR/EN) exactement comme envoyé, avec les
@ -569,6 +578,15 @@ const canCreateReminder = computed(() =>
&& ['sending', 'completed'].includes(campaign.value.status) && ['sending', 'completed'].includes(campaign.value.status)
&& nonClickedCount.value > 0, && nonClickedCount.value > 0,
) )
// Envoi d'événement resté bloqué à « sending » (ex. redémarrage du hub en plein envoi) avec des
// destinataires encore 'pending' uniquement pour les campagnes d'événement (params.type==='event'),
// ce n'est pas pertinent pour les campagnes cadeau génériques.
const canResumeEventSend = computed(() =>
campaign.value
&& campaign.value.params?.type === 'event'
&& campaign.value.status === 'sending'
&& counterFor('pending') > 0,
)
// Pick a representative expiry to show in the confirmation dialog // Pick a representative expiry to show in the confirmation dialog
// almost always identical across recipients since they were sent in the // almost always identical across recipients since they were sent in the
@ -797,6 +815,31 @@ function confirmCreateReminder () {
}) })
} }
// Reprend un envoi d'événement interrompu le worker relit channel/from depuis la campagne elle-même
// et saute automatiquement les destinataires déjà 'sent' (aucun risque de doublon).
function confirmResumeSend () {
const n = counterFor('pending')
$q.dialog({
title: "Reprendre l'envoi",
message: `Reprendre l'envoi pour les <b>${n}</b> destinataire(s) encore en attente ?<br>Ceux déjà rejoints ne recevront PAS de second courriel.`,
html: true,
cancel: { label: 'Annuler', flat: true, noCaps: true },
ok: { label: `Reprendre (${n})`, color: 'warning', unelevated: true, noCaps: true },
persistent: true,
}).onOk(async () => {
resumingSend.value = true
try {
await resumeEventCampaign(id)
$q.notify({ type: 'positive', message: `Reprise lancée — ${n} destinataire(s) en attente`, icon: 'play_arrow' })
await load()
} catch (e) {
$q.notify({ type: 'negative', message: 'Erreur : ' + e.message })
} finally {
resumingSend.value = false
}
})
}
onMounted(async () => { onMounted(async () => {
await load() await load()
// Auto-subscribe to SSE if still running (or about to run) // Auto-subscribe to SSE if still running (or about to run)

View File

@ -49,6 +49,14 @@ const LOGO = 'https://msg.gigafibre.ca/campaigns/assets/6a2bcb6057d9881c08304a5d
const esc = (s) => String(s == null ? '' : s).replace(/[&<>"]/g, c => ({ '&': '&amp;', '<': '&lt;', '>': '&gt;', '"': '&quot;' }[c])) const esc = (s) => String(s == null ? '' : s).replace(/[&<>"]/g, c => ({ '&': '&amp;', '<': '&lt;', '>': '&gt;', '"': '&quot;' }[c]))
const normLang = (l) => /^en/i.test(String(l || '')) ? 'en' : 'fr' const normLang = (l) => /^en/i.test(String(l || '')) ? 'en' : 'fr'
const firstName = (n) => String(n || '').trim().split(/\s+/)[0] || '' const firstName = (n) => String(n || '').trim().split(/\s+/)[0] || ''
// Garde anti-double-exécution : un redémarrage du hub en plein envoi interrompt sendEventCampaignAsync
// (boucle en mémoire, perdue au restart) — la reprise (resumeCampaignSend) relance le même worker, qui
// doit refuser un 2e lancement concurrent sur le MÊME id (sinon deux boucles pourraient lire/écrire le
// même destinataire 'pending' en même temps → doublon d'envoi).
// RÉSERVATION SYNCHRONE (avant tout setImmediate) : ferme la fenêtre de course entre deux appels rapprochés
// (ex. double-clic) — sendEventCampaignAsync ne fait que LIBÉRER (finally), jamais re-vérifier/ajouter.
const activeEventSends = new Set()
function tryReserveSend (id) { if (activeEventSends.has(id)) return false; activeEventSends.add(id); return true }
// ── Libellés d'interface COMMUNS (indépendants de l'événement) ───────────── // ── Libellés d'interface COMMUNS (indépendants de l'événement) ─────────────
// Séparés du contenu éditable : mêmes champs/erreurs/remerciements pour tous les événements. // Séparés du contenu éditable : mêmes champs/erreurs/remerciements pour tous les événements.
@ -490,45 +498,68 @@ function massSend (eventId, { channel = 'mailjet', from = '' } = {}) {
recipients: recipients.map(r => ({ email: r.email, firstname: r.firstname || '', lastname: r.lastname || '', language: normLang(r.language), customer_id: r.customer_id || '', source: r.source || '', status: 'pending' })), recipients: recipients.map(r => ({ email: r.email, firstname: r.firstname || '', lastname: r.lastname || '', language: normLang(r.language), customer_id: r.customer_id || '', source: r.source || '', status: 'pending' })),
} }
campaigns.saveCampaign(campaign) campaigns.saveCampaign(campaign)
tryReserveSend(id) // id tout neuf → toujours dispo, mais on passe par le même chemin pour rester cohérent
setImmediate(() => sendEventCampaignAsync(id, eventId, ch, from).catch(e => log('event blast async: ' + e.message))) setImmediate(() => sendEventCampaignAsync(id, eventId, ch, from).catch(e => log('event blast async: ' + e.message)))
return { campaign_id: id, count: campaign.recipients.length } return { campaign_id: id, count: campaign.recipients.length }
} }
// Le SLOT anti-doublon doit déjà être réservé par l'appelant (tryReserveSend), AVANT le setImmediate —
// cette fonction se contente de le LIBÉRER en sortie (finally). Ne PAS re-vérifier/ajouter ici (fenêtre
// de course : deux appels synchrones rapprochés passeraient tous les deux avant que le 1er setImmediate ne s'exécute).
async function sendEventCampaignAsync (id, eventId, channel, from, opts = {}) { async function sendEventCampaignAsync (id, eventId, channel, from, opts = {}) {
const isReminder = !!opts.isReminder try {
const campaigns = require('./campaigns') const isReminder = !!opts.isReminder
const c = campaigns.loadCampaign(id); if (!c) return const campaigns = require('./campaigns')
const event = getEvent(eventId); if (!event) return const c = campaigns.loadCampaign(id); if (!c) return
let sse = null; try { sse = require('./sse') } catch { /* SSE optionnel */ } const event = getEvent(eventId); if (!event) return
const topic = 'campaign:' + id let sse = null; try { sse = require('./sse') } catch { /* SSE optionnel */ }
const bcast = (ev, data) => { try { sse && sse.broadcast(topic, ev, data) } catch { /* */ } } const topic = 'campaign:' + id
const payloadByLang = { fr: attachmentPayload(channel, loadSendAttachments(event, 'fr')), en: attachmentPayload(channel, loadSendAttachments(event, 'en')) } const bcast = (ev, data) => { try { sse && sse.broadcast(topic, ev, data) } catch { /* */ } }
const sleep = (ms) => new Promise(r => setTimeout(r, ms)) const payloadByLang = { fr: attachmentPayload(channel, loadSendAttachments(event, 'fr')), en: attachmentPayload(channel, loadSendAttachments(event, 'en')) }
bcast('campaign-status', { id, status: 'sending' }) const sleep = (ms) => new Promise(r => setTimeout(r, ms))
for (let i = 0; i < c.recipients.length; i++) { bcast('campaign-status', { id, status: 'sending' })
const r = c.recipients[i] for (let i = 0; i < c.recipients.length; i++) {
if (r.status !== 'pending') continue const r = c.recipients[i]
const lang = normLang(r.language) if (r.status !== 'pending') continue
const name = ((r.firstname || '') + ' ' + (r.lastname || '')).trim() const lang = normLang(r.language)
const rsvpUrl = rsvpLink(eventId, r.customer_id || ('x-' + r.email), name, r.email, 60 * 24, lang) const name = ((r.firstname || '') + ' ' + (r.lastname || '')).trim()
r.rsvp_url = rsvpUrl const rsvpUrl = rsvpLink(eventId, r.customer_id || ('x-' + r.email), name, r.email, 60 * 24, lang)
const html = inviteEmail(eventId, lang, name, rsvpUrl, isReminder) r.rsvp_url = rsvpUrl
const subject = inviteSubject(eventId, lang, isReminder) const html = inviteEmail(eventId, lang, name, rsvpUrl, isReminder)
const customId = id + ':' + i const subject = inviteSubject(eventId, lang, isReminder)
try { const customId = id + ':' + i
const res = await sendInviteMessage({ channel, to: r.email, subject, html, from, attachments: payloadByLang[lang] || [], customId }) try {
if (res.ok) { r.status = 'sent'; r.sent_at = new Date().toISOString(); if (channel !== 'gmail') { r.mailjet_custom_id = customId; if (res.id) r.mailjet_uuid = res.id } } const res = await sendInviteMessage({ channel, to: r.email, subject, html, from, attachments: payloadByLang[lang] || [], customId })
else { r.status = 'failed'; r.error = res.error || 'send_failed' } if (res.ok) { r.status = 'sent'; r.sent_at = new Date().toISOString(); if (channel !== 'gmail') { r.mailjet_custom_id = customId; if (res.id) r.mailjet_uuid = res.id } }
} catch (e) { r.status = 'failed'; r.error = e.message } else { r.status = 'failed'; r.error = res.error || 'send_failed' }
} catch (e) { r.status = 'failed'; r.error = e.message }
campaigns.saveCampaign(c)
bcast('recipient-update', { i, recipient: r })
if (i < c.recipients.length - 1) await sleep(600) // throttle (comme le worker campagne)
}
c.status = 'completed'; c.send_completed_at = new Date().toISOString()
campaigns.saveCampaign(c) campaigns.saveCampaign(c)
bcast('recipient-update', { i, recipient: r }) const sent = c.recipients.filter(x => x.status === 'sent').length
if (i < c.recipients.length - 1) await sleep(600) // throttle (comme le worker campagne) log(`Event blast ${eventId} → campaign ${id} : ${sent}/${c.recipients.length} envoyés (${channel})`)
} bcast('campaign-done', { id, counters: c.counters })
c.status = 'completed'; c.send_completed_at = new Date().toISOString() } finally { activeEventSends.delete(id) }
campaigns.saveCampaign(c) }
const sent = c.recipients.filter(x => x.status === 'sent').length // Reprend un envoi interrompu (ex. redémarrage du hub en plein envoi) : relit channel/from/event_id
log(`Event blast ${eventId} → campaign ${id} : ${sent}/${c.recipients.length} envoyés (${channel})`) // DEPUIS la campagne elle-même (rien à ressaisir), refuse si déjà en cours ou s'il ne reste rien en
bcast('campaign-done', { id, counters: c.counters }) // attente. Le worker saute automatiquement les destinataires déjà 'sent' → reprise sûre, sans doublon.
function resumeCampaignSend (campaignId) {
const campaigns = require('./campaigns')
const c = campaigns.loadCampaign(campaignId)
if (!c) return { error: 'campaign_not_found' }
const p = c.params || {}
if (p.type !== 'event' || !p.event_id) return { error: 'not_an_event_campaign' }
if (!getEvent(p.event_id)) return { error: 'event_not_found' }
const pending = (c.recipients || []).filter(r => r.status === 'pending').length
if (!pending) return { error: 'nothing_pending' }
if (!tryReserveSend(campaignId)) return { error: 'already_sending' } // réservation SYNCHRONE : ferme la fenêtre de course
const ch = p.channel === 'gmail' ? 'gmail' : 'mailjet'
setImmediate(() => sendEventCampaignAsync(campaignId, p.event_id, ch, p.from || '', { isReminder: !!c.reminder_of }).catch(e => log('event resume async: ' + e.message)))
return { ok: true, pending }
} }
// ── Relance (« Petit rappel ») aux invités PAS ENCORE inscrits ───────────── // ── Relance (« Petit rappel ») aux invités PAS ENCORE inscrits ─────────────
@ -589,6 +620,7 @@ function reminderSend (eventId, sourceCampaignId, { channel = 'mailjet', from =
recipients: r.candidates.map(x => ({ email: x.email, firstname: x.firstname || '', lastname: x.lastname || '', language: normLang(x.language), customer_id: x.customer_id || '', source: 'Rappel : ' + r.source.name, status: 'pending' })), recipients: r.candidates.map(x => ({ email: x.email, firstname: x.firstname || '', lastname: x.lastname || '', language: normLang(x.language), customer_id: x.customer_id || '', source: 'Rappel : ' + r.source.name, status: 'pending' })),
} }
campaigns.saveCampaign(campaign) campaigns.saveCampaign(campaign)
tryReserveSend(id)
setImmediate(() => sendEventCampaignAsync(id, eventId, ch, from, { isReminder: true }).catch(e => log('event reminder async: ' + e.message))) setImmediate(() => sendEventCampaignAsync(id, eventId, ch, from, { isReminder: true }).catch(e => log('event reminder async: ' + e.message)))
return { campaign_id: id, count: campaign.recipients.length } return { campaign_id: id, count: campaign.recipients.length }
} }
@ -1181,6 +1213,15 @@ async function handle (req, res, method, path, url) {
if (r.error) return json(res, 400, { error: r.error }) if (r.error) return json(res, 400, { error: r.error })
return json(res, 200, { source: r.source, count: r.count, breakdown: r.breakdown, sample: r.candidates.slice(0, 50).map(x => ({ email: x.email, name: ((x.firstname || '') + ' ' + (x.lastname || '')).trim() || x.email, language: x.language, status: x.status, opened_at: x.opened_at || null, clicked_at: x.clicked_at || null })) }) return json(res, 200, { source: r.source, count: r.count, breakdown: r.breakdown, sample: r.candidates.slice(0, 50).map(x => ({ email: x.email, name: ((x.firstname || '') + ' ' + (x.lastname || '')).trim() || x.email, language: x.language, status: x.status, opened_at: x.opened_at || null, clicked_at: x.clicked_at || null })) })
} }
// Staff : POST /events/campaign-resume {campaign_id} → reprend un envoi événement interrompu (relit channel/from/event_id depuis la campagne)
if (path === '/events/campaign-resume' && method === 'POST') {
const b = await parseBody(req)
if (!b.campaign_id) return json(res, 400, { error: 'campaign_id_required' })
const r = resumeCampaignSend(b.campaign_id)
if (r.error) return json(res, 400, { error: r.error })
log(`Event campaign RESUME — ${b.campaign_id} · ${r.pending} en attente relancé(s)`)
return json(res, 202, { ok: true, pending: r.pending })
}
// Staff : POST /events/<id>/reminder-send {campaign_id, channel, from} → lance la relance réelle (recompte les candidats à l'envoi) // Staff : POST /events/<id>/reminder-send {campaign_id, channel, from} → lance la relance réelle (recompte les candidats à l'envoi)
mm = path.match(/^\/events\/([A-Za-z0-9_-]+)\/reminder-send$/) mm = path.match(/^\/events\/([A-Za-z0-9_-]+)\/reminder-send$/)
if (mm && method === 'POST') { if (mm && method === 'POST') {
@ -1236,4 +1277,4 @@ async function handle (req, res, method, path, url) {
} catch (e) { log('events handle: ' + e.message); return json(res, 500, { error: e.message }) } } catch (e) { log('events handle: ' + e.message); return json(res, 500, { error: e.message }) }
} }
module.exports = { handle, getEvent, listEvents, listRsvps, deleteRsvp, rsvpLink, inviteEmail, inviteSubject, sendTest, massSend, sendEventCampaignAsync, matchAndValidate, resolveAudience, parseAudienceCsv, page, normLang, addAttachment, removeAttachment, loadSendAttachments, headcountOf, audienceListAdd, audienceListRemove, audienceListClear, loadAudienceList, audienceSummary, listEventCampaigns, reminderCandidates, reminderSend } module.exports = { handle, getEvent, listEvents, listRsvps, deleteRsvp, rsvpLink, inviteEmail, inviteSubject, sendTest, massSend, sendEventCampaignAsync, matchAndValidate, resolveAudience, parseAudienceCsv, page, normLang, addAttachment, removeAttachment, loadSendAttachments, headcountOf, audienceListAdd, audienceListRemove, audienceListClear, loadAudienceList, audienceSummary, listEventCampaigns, reminderCandidates, reminderSend, resumeCampaignSend }