diff --git a/ghost/core/core/boot.js b/ghost/core/core/boot.js index 7a54eaa613d..7dd9b6a005b 100644 --- a/ghost/core/core/boot.js +++ b/ghost/core/core/boot.js @@ -450,10 +450,17 @@ async function initBackgroundServices({config}) { // Load email analytics recurring jobs if (config.get('backgroundJobs:emailAnalytics')) { const emailAnalyticsJobs = require('./server/services/email-analytics/jobs'); - await Promise.all([ + const analyticsJobs = [ emailAnalyticsJobs.scheduleRecurringNewslettersJob(), emailAnalyticsJobs.scheduleRecurringAutomationsJob() - ]); + ]; + + const labs = require('./shared/labs'); + if (labs.isSet('giftSubCustomization')) { + analyticsJobs.push(emailAnalyticsJobs.scheduleRecurringGiftDeliveriesJob()); + } + + await Promise.all(analyticsJobs); } const updateCheck = require('./server/services/update-check'); diff --git a/ghost/core/core/server/services/email-analytics/events/start-gift-email-analytics-job-event.ts b/ghost/core/core/server/services/email-analytics/events/start-gift-email-analytics-job-event.ts new file mode 100644 index 00000000000..2270cc06bbc --- /dev/null +++ b/ghost/core/core/server/services/email-analytics/events/start-gift-email-analytics-job-event.ts @@ -0,0 +1,11 @@ +export class StartGiftEmailAnalyticsJobEvent { + readonly timestamp: Date; + + constructor(timestamp: Date) { + this.timestamp = timestamp; + } + + static create(timestamp = new Date()) { + return new StartGiftEmailAnalyticsJobEvent(timestamp); + } +} diff --git a/ghost/core/core/server/services/email-analytics/gift-email-analytics-batch-processor.ts b/ghost/core/core/server/services/email-analytics/gift-email-analytics-batch-processor.ts new file mode 100644 index 00000000000..0973fa1dd0b --- /dev/null +++ b/ghost/core/core/server/services/email-analytics/gift-email-analytics-batch-processor.ts @@ -0,0 +1,66 @@ +import type {BatchEventProcessor} from './batch-event-processor'; +import {EventProcessingResult} from './event-processing-result'; + +type GiftService = { + recordDeliveryOutcome(data: { + providerMessageId: string; + outcome: 'delivered' | 'temporary_failed' | 'permanent_failed'; + timestamp: Date; + error: string | null; + }): Promise; +}; + +type EmailAnalyticsEvent = { + type: string; + severity?: string; + providerId: string; + timestamp: Date; + error?: {code?: unknown; message?: unknown; enhancedCode?: unknown} | null; +}; + +const normalizeMessageId = (value: string): string => value.trim().replace(/^<|>$/g, ''); + +export class GiftEmailAnalyticsBatchProcessor implements BatchEventProcessor { + readonly #giftService: GiftService; + + constructor({giftService}: {giftService: GiftService}) { + this.#giftService = giftService; + } + + async processBatch(events: ReadonlyArray, result: EventProcessingResult, fetchData: {lastEventTimestamp?: Date}): Promise { + for (const event of events) { + if (!fetchData.lastEventTimestamp || event.timestamp > fetchData.lastEventTimestamp) { + fetchData.lastEventTimestamp = event.timestamp; + } + + let outcome: 'delivered' | 'temporary_failed' | 'permanent_failed' | null = null; + if (event.type === 'delivered') { + outcome = 'delivered'; + } else if (event.type === 'failed') { + outcome = event.severity === 'temporary' ? 'temporary_failed' : 'permanent_failed'; + } + + if (!outcome) { + result.merge(new EventProcessingResult({unhandled: 1})); + continue; + } + + const updated = await this.#giftService.recordDeliveryOutcome({ + providerMessageId: normalizeMessageId(event.providerId), + outcome, + timestamp: event.timestamp, + error: event.error ? JSON.stringify(event.error) : null + }); + + if (!updated) { + result.merge(new EventProcessingResult({unprocessable: 1})); + } else if (outcome === 'delivered') { + result.merge(new EventProcessingResult({delivered: 1})); + } else if (outcome === 'temporary_failed') { + result.merge(new EventProcessingResult({temporaryFailed: 1})); + } else { + result.merge(new EventProcessingResult({permanentFailed: 1})); + } + } + } +} diff --git a/ghost/core/core/server/services/email-analytics/index.ts b/ghost/core/core/server/services/email-analytics/index.ts index 34da53958d3..0e408fc86d7 100644 --- a/ghost/core/core/server/services/email-analytics/index.ts +++ b/ghost/core/core/server/services/email-analytics/index.ts @@ -24,6 +24,12 @@ import {StartAutomationEmailAnalyticsJobEvent} from './events/start-automation-e import {AUTOMATION_EMAIL_TAG} from '../member-welcome-emails/constants'; import * as automationsApi from '../automations/automations-api'; import {AutomationEmailAnalyticsBatchProcessor} from './automation-email-analytics-batch-processor'; +import {GiftEmailAnalyticsBatchProcessor} from './gift-email-analytics-batch-processor'; +import {StartGiftEmailAnalyticsJobEvent} from './events/start-gift-email-analytics-job-event'; +// @ts-expect-error This CommonJS service wrapper lacks type declarations. +import giftsService from '../gifts'; +// @ts-expect-error This CommonJS service lacks type declarations. +import labs from '../../../shared/labs'; export const newsletters = new EmailAnalyticsServiceWrapper({ logName: 'newsletters' @@ -33,6 +39,10 @@ export const automations = new EmailAnalyticsServiceWrapper({ logName: 'automations', }); +export const gifts = new EmailAnalyticsServiceWrapper({ + logName: 'gifts' +}); + export const init = () => { const newsletterEmailEventProcessor = new EmailEventProcessor({ domainEvents, @@ -106,4 +116,31 @@ export const init = () => { }) ) }); + + if (labs.isSet('giftSubCustomization')) { + gifts.init({ + event: StartGiftEmailAnalyticsJobEvent, + mailgunTags: ['gift-delivery'], + jobNames: { + latestNonOpened: 'email-analytics-gifts-latest-others', + missing: 'email-analytics-gifts-missing', + latestOpened: 'email-analytics-gifts-latest-opened', + scheduled: 'email-analytics-gifts-scheduled' + }, + cursorSeed: { + tableName: 'gift_deliveries', + eventColumns: { + delivered: 'outcome_at', + failed: 'outcome_at' + } + }, + createEventProcessor: () => ( + new GiftEmailAnalyticsBatchProcessor({ + giftService: { + recordDeliveryOutcome: data => giftsService.service.recordDeliveryOutcome(data) + } + }) + ) + }); + } }; diff --git a/ghost/core/core/server/services/email-analytics/jobs/email-analytics-job-scheduler.ts b/ghost/core/core/server/services/email-analytics/jobs/email-analytics-job-scheduler.ts index b0b56cd9650..edd78199825 100644 --- a/ghost/core/core/server/services/email-analytics/jobs/email-analytics-job-scheduler.ts +++ b/ghost/core/core/server/services/email-analytics/jobs/email-analytics-job-scheduler.ts @@ -17,6 +17,9 @@ type Models = { AutomatedEmailRecipient: { query(): ExistingRecipientQuery; }; + GiftDelivery?: { + query(): ExistingRecipientQuery; + }; }; type Config = {get(key: string): unknown}; type JobManager = { @@ -39,6 +42,7 @@ function randomFiveMinuteCron(): string { export class EmailAnalyticsJobScheduler { #hasScheduledNewslettersJob = false; #hasScheduledAutomationsJob = false; + #hasScheduledGiftDeliveriesJob = false; readonly #models: Models; readonly #config: Config; readonly #jobManager: JobManager; @@ -123,4 +127,29 @@ export class EmailAnalyticsJobScheduler { this.#hasScheduledAutomationsJob = true; } + + async scheduleRecurringGiftDeliveriesJob(skipGiftDeliveryCheck: boolean = false): Promise { + if (this.#hasScheduledGiftDeliveriesJob || !this.#isConfigured()) { + return; + } + + const hasGiftDelivery = skipGiftDeliveryCheck || Boolean( + this.#models.GiftDelivery && await this.#models.GiftDelivery + .query() + .where('email_sent_at', '>', moment.utc().subtract(30, 'days').toDate()) + .whereNotNull('email_provider_message_id') + .first('id') + ); + + if (!hasGiftDelivery || this.#hasScheduledGiftDeliveriesJob) { + return; + } + + this.#jobManager.addJob({ + at: randomFiveMinuteCron(), + job: path.resolve(__dirname, 'gift-fetch-latest/index.js'), + name: 'email-analytics-gift-fetch-latest' + }); + this.#hasScheduledGiftDeliveriesJob = true; + } } diff --git a/ghost/core/core/server/services/email-analytics/jobs/gift-fetch-latest/index.js b/ghost/core/core/server/services/email-analytics/jobs/gift-fetch-latest/index.js new file mode 100644 index 00000000000..bdd506e8d4a --- /dev/null +++ b/ghost/core/core/server/services/email-analytics/jobs/gift-fetch-latest/index.js @@ -0,0 +1,7 @@ +const {run} = require('../fetch-latest-job'); +const {StartGiftEmailAnalyticsJobEvent} = require('../../events/start-gift-email-analytics-job-event'); + +run({ + event: StartGiftEmailAnalyticsJobEvent, + logName: 'gifts' +}); diff --git a/ghost/core/core/server/services/email-analytics/jobs/index.js b/ghost/core/core/server/services/email-analytics/jobs/index.js index 485ae9411d7..a8ab9713e70 100644 --- a/ghost/core/core/server/services/email-analytics/jobs/index.js +++ b/ghost/core/core/server/services/email-analytics/jobs/index.js @@ -28,3 +28,13 @@ exports.scheduleRecurringAutomationsJob = async (...args) => { await emailAnalyticsJobScheduler.scheduleRecurringAutomationsJob(...args); } }; + +/** + * @param {Parameters} args + * @returns {Promise} + */ +exports.scheduleRecurringGiftDeliveriesJob = async (...args) => { + if (!process.env.NODE_ENV.startsWith('test')) { + await emailAnalyticsJobScheduler.scheduleRecurringGiftDeliveriesJob(...args); + } +}; diff --git a/ghost/core/core/server/services/gifts/gift-service-wrapper.js b/ghost/core/core/server/services/gifts/gift-service-wrapper.js index 97f73c08d9b..1b368932b11 100644 --- a/ghost/core/core/server/services/gifts/gift-service-wrapper.js +++ b/ghost/core/core/server/services/gifts/gift-service-wrapper.js @@ -42,6 +42,7 @@ class GiftServiceWrapper { const StartGiftDeliveryFlushEvent = require('./events/start-gift-delivery-flush-event'); const StartGiftCleanupEvent = require('./events/start-gift-cleanup-event'); const jobs = require('./jobs'); + const emailAnalyticsJobs = require('../email-analytics/jobs'); const {GhostMailer} = require('../mail'); const settingsCache = require('../../../shared/settings-cache'); @@ -101,7 +102,7 @@ class GiftServiceWrapper { giftReminderScheduler, giftDeliveryScheduler, giftEmailAnalytics: { - schedule: () => Promise.resolve() + schedule: () => emailAnalyticsJobs.scheduleRecurringGiftDeliveriesJob(true) }, checkoutAdapter, labsService, diff --git a/ghost/core/test/unit/server/services/email-analytics/email-analytics-job-scheduler.test.ts b/ghost/core/test/unit/server/services/email-analytics/email-analytics-job-scheduler.test.ts index 69ce2356d72..b2298728d86 100644 --- a/ghost/core/test/unit/server/services/email-analytics/email-analytics-job-scheduler.test.ts +++ b/ghost/core/test/unit/server/services/email-analytics/email-analytics-job-scheduler.test.ts @@ -28,21 +28,27 @@ function buildScheduler({ emailAnalyticsEnabled = true, backgroundJobEnabled = true, emailCount = 1, - automatedEmailRecipient = null + automatedEmailRecipient = null, + giftDelivery = null }: { emailAnalyticsEnabled?: boolean; backgroundJobEnabled?: boolean; emailCount?: string | number; automatedEmailRecipient?: unknown; + giftDelivery?: unknown; } = {}) { const newsletterQuery = buildNewsletterQuery(emailCount); const automationsQuery = buildAutomationsQuery(automatedEmailRecipient); + const giftQuery = buildAutomationsQuery(giftDelivery); const models = { Email: { where: newsletterQuery.where }, AutomatedEmailRecipient: { query: sinon.stub().returns(automationsQuery) + }, + GiftDelivery: { + query: sinon.stub().returns(giftQuery) } }; const config = { @@ -65,6 +71,7 @@ function buildScheduler({ jobManager, newsletterQuery, automationsQuery, + giftQuery, models }; } @@ -281,4 +288,24 @@ describe('EmailAnalyticsJobScheduler', function () { sinon.assert.notCalled(jobManager.addJob); sinon.assert.calledOnce(automationsQuery.first); }); + + it('adds the existing recurring analytics collector for accepted gift email telemetry', async function () { + const {scheduler, jobManager, giftQuery, models} = buildScheduler({ + emailCount: 0, + giftDelivery: {id: 'gift-id'} + }); + + await scheduler.scheduleRecurringGiftDeliveriesJob(); + + sinon.assert.calledOnceWithMatch(jobManager.addJob, { + job: sinon.match((value: unknown) => ( + typeof value === 'string' && value.endsWith('gift-fetch-latest/index.js') + )), + name: 'email-analytics-gift-fetch-latest' + }); + sinon.assert.calledOnce(models.GiftDelivery.query); + sinon.assert.calledOnceWithExactly(giftQuery.where, 'email_sent_at', '>', sinon.match.date); + sinon.assert.calledOnceWithExactly(giftQuery.whereNotNull, 'email_provider_message_id'); + sinon.assert.calledOnceWithExactly(giftQuery.first, 'id'); + }); }); diff --git a/ghost/core/test/unit/server/services/email-analytics/gift-email-analytics-batch-processor.test.ts b/ghost/core/test/unit/server/services/email-analytics/gift-email-analytics-batch-processor.test.ts new file mode 100644 index 00000000000..7d8d154ec04 --- /dev/null +++ b/ghost/core/test/unit/server/services/email-analytics/gift-email-analytics-batch-processor.test.ts @@ -0,0 +1,50 @@ +import assert from 'node:assert/strict'; +import sinon from 'sinon'; +import {GiftEmailAnalyticsBatchProcessor} from '../../../../../core/server/services/email-analytics/gift-email-analytics-batch-processor'; +import {EventProcessingResult} from '../../../../../core/server/services/email-analytics/event-processing-result'; + +describe('GiftEmailAnalyticsBatchProcessor', function () { + it('maps Mailgun delivery and failure events to latest gift outcomes without opens', async function () { + const giftService = {recordDeliveryOutcome: sinon.stub().resolves(true)}; + const processor = new GiftEmailAnalyticsBatchProcessor({giftService}); + const result = new EventProcessingResult(); + const fetchData: {lastEventTimestamp?: Date} = {}; + const deliveredAt = new Date('2026-08-05T12:00:00.000Z'); + const failedAt = new Date('2026-08-05T12:10:00.000Z'); + + await processor.processBatch([ + {type: 'delivered', providerId: '', timestamp: deliveredAt}, + {type: 'failed', severity: 'temporary', providerId: 'provider-123', timestamp: failedAt, error: {code: 421, message: 'try later'}}, + {type: 'opened', providerId: 'provider-123', timestamp: new Date('2026-08-05T12:20:00.000Z')} + ], result, fetchData); + + sinon.assert.calledWithExactly(giftService.recordDeliveryOutcome, sinon.match({ + providerMessageId: 'provider-123', + outcome: 'delivered', + timestamp: deliveredAt, + error: null + })); + assert.deepEqual(giftService.recordDeliveryOutcome.secondCall.firstArg, { + providerMessageId: 'provider-123', + outcome: 'temporary_failed', + timestamp: failedAt, + error: JSON.stringify({code: 421, message: 'try later'}) + }); + assert.equal(result.delivered, 1); + assert.equal(result.temporaryFailed, 1); + assert.equal(result.unhandled, 1); + }); + + it('marks events for unknown message IDs unprocessable', async function () { + const giftService = {recordDeliveryOutcome: sinon.stub().resolves(false)}; + const processor = new GiftEmailAnalyticsBatchProcessor({giftService}); + const result = new EventProcessingResult(); + + await processor.processBatch([ + {type: 'failed', severity: 'permanent', providerId: 'unknown', timestamp: new Date()} + ], result, {}); + + assert.equal(result.unprocessable, 1); + assert.equal(result.permanentFailed, 0); + }); +});