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 }