From 07b0331a3ac88716c5fa025bf92ad419a886eb55 Mon Sep 17 00:00:00 2001 From: louispaulb Date: Fri, 24 Jul 2026 10:12:27 -0400 Subject: [PATCH] Events: self-serve resume for a send interrupted by a hub restart MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- apps/ops/src/api/events.js | 4 + .../campaigns/pages/CampaignDetailPage.vue | 43 +++++++ services/targo-hub/lib/events.js | 109 ++++++++++++------ 3 files changed, 122 insertions(+), 34 deletions(-) diff --git a/apps/ops/src/api/events.js b/apps/ops/src/api/events.js index 205031f..d2c6f55 100644 --- a/apps/ops/src/api/events.js +++ b/apps/ops/src/api/events.js @@ -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 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 }) } + +// 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 } }) } diff --git a/apps/ops/src/modules/campaigns/pages/CampaignDetailPage.vue b/apps/ops/src/modules/campaigns/pages/CampaignDetailPage.vue index e14ff58..11dea58 100644 --- a/apps/ops/src/modules/campaigns/pages/CampaignDetailPage.vue +++ b/apps/ops/src/modules/campaigns/pages/CampaignDetailPage.vue @@ -21,6 +21,13 @@ :loading="creatingReminder" @click="confirmCreateReminder"> Cloner cette campagne pour les destinataires qui n'ont PAS cliqué le cadeau encore + + + {{ counterFor('pending') }} destinataire(s) encore en attente — reprend sans renvoyer à ceux déjà rejoints + && ['sending', 'completed'].includes(campaign.value.status) && 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 — // 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 ${n} destinataire(s) encore en attente ?
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 () => { await load() // Auto-subscribe to SSE if still running (or about to run) diff --git a/services/targo-hub/lib/events.js b/services/targo-hub/lib/events.js index 81adb3b..9d1f6e1 100644 --- a/services/targo-hub/lib/events.js +++ b/services/targo-hub/lib/events.js @@ -49,6 +49,14 @@ const LOGO = 'https://msg.gigafibre.ca/campaigns/assets/6a2bcb6057d9881c08304a5d const esc = (s) => String(s == null ? '' : s).replace(/[&<>"]/g, c => ({ '&': '&', '<': '<', '>': '>', '"': '"' }[c])) const normLang = (l) => /^en/i.test(String(l || '')) ? 'en' : 'fr' 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) ───────────── // 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' })), } 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))) 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 = {}) { - const isReminder = !!opts.isReminder - const campaigns = require('./campaigns') - const c = campaigns.loadCampaign(id); if (!c) return - const event = getEvent(eventId); if (!event) return - let sse = null; try { sse = require('./sse') } catch { /* SSE optionnel */ } - const topic = 'campaign:' + id - const bcast = (ev, data) => { try { sse && sse.broadcast(topic, ev, data) } catch { /* */ } } - const payloadByLang = { fr: attachmentPayload(channel, loadSendAttachments(event, 'fr')), en: attachmentPayload(channel, loadSendAttachments(event, 'en')) } - const sleep = (ms) => new Promise(r => setTimeout(r, ms)) - bcast('campaign-status', { id, status: 'sending' }) - for (let i = 0; i < c.recipients.length; i++) { - const r = c.recipients[i] - if (r.status !== 'pending') continue - const lang = normLang(r.language) - const name = ((r.firstname || '') + ' ' + (r.lastname || '')).trim() - const rsvpUrl = rsvpLink(eventId, r.customer_id || ('x-' + r.email), name, r.email, 60 * 24, lang) - r.rsvp_url = rsvpUrl - const html = inviteEmail(eventId, lang, name, rsvpUrl, isReminder) - const subject = inviteSubject(eventId, lang, isReminder) - const customId = id + ':' + i - try { - const res = await sendInviteMessage({ channel, to: r.email, subject, html, from, attachments: payloadByLang[lang] || [], customId }) - 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 } } - else { r.status = 'failed'; r.error = res.error || 'send_failed' } - } catch (e) { r.status = 'failed'; r.error = e.message } + try { + const isReminder = !!opts.isReminder + const campaigns = require('./campaigns') + const c = campaigns.loadCampaign(id); if (!c) return + const event = getEvent(eventId); if (!event) return + let sse = null; try { sse = require('./sse') } catch { /* SSE optionnel */ } + const topic = 'campaign:' + id + const bcast = (ev, data) => { try { sse && sse.broadcast(topic, ev, data) } catch { /* */ } } + const payloadByLang = { fr: attachmentPayload(channel, loadSendAttachments(event, 'fr')), en: attachmentPayload(channel, loadSendAttachments(event, 'en')) } + const sleep = (ms) => new Promise(r => setTimeout(r, ms)) + bcast('campaign-status', { id, status: 'sending' }) + for (let i = 0; i < c.recipients.length; i++) { + const r = c.recipients[i] + if (r.status !== 'pending') continue + const lang = normLang(r.language) + const name = ((r.firstname || '') + ' ' + (r.lastname || '')).trim() + const rsvpUrl = rsvpLink(eventId, r.customer_id || ('x-' + r.email), name, r.email, 60 * 24, lang) + r.rsvp_url = rsvpUrl + const html = inviteEmail(eventId, lang, name, rsvpUrl, isReminder) + const subject = inviteSubject(eventId, lang, isReminder) + const customId = id + ':' + i + try { + const res = await sendInviteMessage({ channel, to: r.email, subject, html, from, attachments: payloadByLang[lang] || [], customId }) + 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 } } + 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) - 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) - const sent = c.recipients.filter(x => x.status === 'sent').length - log(`Event blast ${eventId} → campaign ${id} : ${sent}/${c.recipients.length} envoyés (${channel})`) - bcast('campaign-done', { id, counters: c.counters }) + const sent = c.recipients.filter(x => x.status === 'sent').length + log(`Event blast ${eventId} → campaign ${id} : ${sent}/${c.recipients.length} envoyés (${channel})`) + bcast('campaign-done', { id, counters: c.counters }) + } finally { activeEventSends.delete(id) } +} +// Reprend un envoi interrompu (ex. redémarrage du hub en plein envoi) : relit channel/from/event_id +// DEPUIS la campagne elle-même (rien à ressaisir), refuse si déjà en cours ou s'il ne reste rien en +// 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 ───────────── @@ -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' })), } campaigns.saveCampaign(campaign) + tryReserveSend(id) setImmediate(() => sendEventCampaignAsync(id, eventId, ch, from, { isReminder: true }).catch(e => log('event reminder async: ' + e.message))) 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 }) 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//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$/) 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 }) } } -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 }