From 07f875c18c03a8d1766040bc40034c1023715fcd Mon Sep 17 00:00:00 2001 From: Dries Augustyns Date: Thu, 11 Dec 2025 14:36:26 +0100 Subject: [PATCH] fix: display email progress instead of scheduling progress for campaigns --- apps/api/src/jobs/email-processor.ts | 53 +++++++++++++++++++++++- apps/api/src/services/CampaignService.ts | 31 -------------- apps/web/src/pages/campaigns/[id].tsx | 2 +- 3 files changed, 53 insertions(+), 33 deletions(-) diff --git a/apps/api/src/jobs/email-processor.ts b/apps/api/src/jobs/email-processor.ts index 7076ccd..e929b8c 100644 --- a/apps/api/src/jobs/email-processor.ts +++ b/apps/api/src/jobs/email-processor.ts @@ -3,7 +3,7 @@ * Processes individual emails from the queue (for all sources: transactional, campaign, workflow) */ -import {EmailSourceType, EmailStatus} from '@plunk/db'; +import {CampaignStatus, EmailSourceType, EmailStatus} from '@plunk/db'; import {type Job, Worker} from 'bullmq'; import signale from 'signale'; @@ -170,6 +170,57 @@ export async function createEmailWorker() { sourceType: email.sourceType, sentAt: new Date().toISOString(), }); + + // If this email belongs to a campaign, check if all campaign emails have been sent + if (email.campaignId) { + const campaign = await prisma.campaign.findUnique({ + where: {id: email.campaignId}, + select: { + id: true, + name: true, + status: true, + totalRecipients: true, + projectId: true, + project: { + select: {name: true}, + }, + }, + }); + + // Only check if campaign is still in SENDING status + if (campaign && campaign.status === CampaignStatus.SENDING) { + // Count how many emails have been sent for this campaign + const sentCount = await prisma.email.count({ + where: { + campaignId: email.campaignId, + sentAt: {not: null}, + }, + }); + + // If all emails have been sent, mark campaign as SENT + if (sentCount >= campaign.totalRecipients) { + await prisma.campaign.update({ + where: {id: email.campaignId}, + data: { + status: CampaignStatus.SENT, + }, + }); + + signale.success( + `[EMAIL-PROCESSOR] Campaign ${campaign.name} completed: ${sentCount}/${campaign.totalRecipients} emails sent`, + ); + + // Send notification about campaign send completed + const {NtfyService} = await import('../services/NtfyService.js'); + await NtfyService.notifyCampaignSendCompleted( + campaign.name, + campaign.project.name, + campaign.projectId, + campaign.totalRecipients, + ); + } + } + } } catch (error) { signale.error(`[EMAIL-PROCESSOR] Failed to send email ${emailId}:`, error); diff --git a/apps/api/src/services/CampaignService.ts b/apps/api/src/services/CampaignService.ts index 7babc2d..65b21bf 100644 --- a/apps/api/src/services/CampaignService.ts +++ b/apps/api/src/services/CampaignService.ts @@ -513,16 +513,6 @@ export class CampaignService { } } - // Update sent count - await prisma.campaign.update({ - where: {id: campaignId}, - data: { - sentCount: { - increment: contacts.length, - }, - }, - }); - // Queue next batch if there are more contacts if (hasMore && nextCursor) { await QueueService.queueCampaignBatch({ @@ -532,27 +522,6 @@ export class CampaignService { limit, cursor: nextCursor, }); - } else { - // All batches processed, mark campaign as SENT - const completedCampaign = await prisma.campaign.update({ - where: {id: campaignId}, - data: { - status: CampaignStatus.SENT, - }, - include: { - project: { - select: {name: true}, - }, - }, - }); - - // Send notification about campaign send completed - await NtfyService.notifyCampaignSendCompleted( - completedCampaign.name, - completedCampaign.project.name, - completedCampaign.projectId, - completedCampaign.totalRecipients || 0, - ); } } diff --git a/apps/web/src/pages/campaigns/[id].tsx b/apps/web/src/pages/campaigns/[id].tsx index cb7f53c..3e7982a 100644 --- a/apps/web/src/pages/campaigns/[id].tsx +++ b/apps/web/src/pages/campaigns/[id].tsx @@ -87,7 +87,7 @@ export default function CampaignDetailsPage() { id && campaign?.data.status !== CampaignStatus.DRAFT ? `/campaigns/${id}/stats` : null, { revalidateOnFocus: false, - refreshInterval: campaign?.data.status === CampaignStatus.SENDING ? 5000 : 0, // Refresh every 5s if sending + refreshInterval: campaign?.data.status === CampaignStatus.SENDING ? 15000 : 0, // Refresh every 15s while sending }, );