diff --git a/ghost/core/core/server/data/importer/import-manager.js b/ghost/core/core/server/data/importer/import-manager.js index d8e170a9139..e0ef5bb943b 100644 --- a/ghost/core/core/server/data/importer/import-manager.js +++ b/ghost/core/core/server/data/importer/import-manager.js @@ -9,6 +9,7 @@ const {extract} = require('@tryghost/zip'); const tpl = require('@tryghost/tpl'); const debug = require('@tryghost/debug')('import-manager'); const logging = require('@tryghost/logging'); +const jobLogging = require('../../services/jobs/job-logging'); const errors = require('@tryghost/errors'); const ImageHandler = require('./handlers/image'); const ImporterContentFileHandler = require('./handlers/importer-content-file-handler'); @@ -507,11 +508,29 @@ class ImportManager { const env = config.get('env'); if (!env?.startsWith('testing') && !importOptions.runningInJob) { + jobLogging.info('[Background Job] site-content-import queued'); return jobManager.addJob({ - job: () => this.importFromFile(file, Object.assign({}, importOptions, { - runningInJob: true, - data: importData - })), + job: async () => { + const startedAt = Date.now(); + jobLogging.info('[Background Job] site-content-import started'); + try { + const result = await this.importFromFile(file, Object.assign({}, importOptions, { + runningInJob: true, + data: importData + })); + // importFromFile swallows its own failures and returns undefined, + // so an absent result is the only signal that the import failed. + if (result === undefined) { + jobLogging.info(`[Background Job] site-content-import failed after ${Date.now() - startedAt}ms`); + } else { + jobLogging.info(`[Background Job] site-content-import completed in ${Date.now() - startedAt}ms`); + } + return result; + } catch (err) { + jobLogging.error(err, `[Background Job] site-content-import failed after ${Date.now() - startedAt}ms`); + throw err; + } + }, offloaded: false }); } @@ -530,7 +549,7 @@ class ImportManager { return importResult; } catch (err) { - logging.error(err, 'Content import was unsuccessful'); + jobLogging.error(err, '[Background Job] site-content-import error'); const errorDetails = err.errorDetails || [err]; importResult = {data: {errors: errorDetails}}; } finally { diff --git a/ghost/core/core/server/services/content-import/import/importer.ts b/ghost/core/core/server/services/content-import/import/importer.ts index 72cb4576b6e..4d30eb5f6f4 100644 --- a/ghost/core/core/server/services/content-import/import/importer.ts +++ b/ghost/core/core/server/services/content-import/import/importer.ts @@ -4,6 +4,7 @@ import type {PostImportRow} from './row'; import type {Clock, ImportRunStore, RowOutcome} from './store'; const errors = require('@tryghost/errors'); +const logging = require('@tryghost/logging'); const tpl = require('@tryghost/tpl'); // The CSV is parsed inside the request (the uploaded temp file is deleted when the @@ -41,6 +42,14 @@ const messages = { urlResolutionFailed: 'Content import could not resolve a URL for {count} created {postNoun}.' }; +function logLifecycle(message: string): void { + try { + logging.info(`[Background Job] content-import ${message}`); + } catch { + // Observability must not change whether an import is queued or resolves. + } +} + const MAX_POSTS = 100; interface ImporterDeps { @@ -115,6 +124,7 @@ class ContentCSVImporter { const importTagNames = buildImportTagNames(runId, this._getTimezone(), this._now()); this._store.create(runId, rows.length); + logLifecycle('queued'); this._addJob({ job: () => this.runImportJob(runId, importTagNames, rows), offloaded: false, @@ -127,8 +137,11 @@ class ContentCSVImporter { // Must resolve in every case: the job manager reads a rejected inline job as a // defect in the job itself, and there is no retry behind it. private async runImportJob(runId: string, importTagNames: string[], rows: PostImportRow[]): Promise { + const startedAt = Date.now(); + logLifecycle('started'); let urlFailureCount = 0; let firstUrlFailure: unknown; + let failed = false; try { const htmlToLexical = this._getHtmlToLexical(); @@ -217,9 +230,13 @@ class ContentCSVImporter { this.reportUrlFailures(urlFailureCount, firstUrlFailure); this._store.finish(runId); } catch (error) { + failed = true; this.reportUrlFailures(urlFailureCount, firstUrlFailure); this._report(error); this._store.fail(runId, messageOf(error)); + } finally { + const outcome = failed ? 'failed after' : 'completed in'; + logLifecycle(`${outcome} ${Date.now() - startedAt}ms`); } } diff --git a/ghost/core/core/server/services/content-import/index.ts b/ghost/core/core/server/services/content-import/index.ts index d13721c2fee..436ce9fa763 100644 --- a/ghost/core/core/server/services/content-import/index.ts +++ b/ghost/core/core/server/services/content-import/index.ts @@ -29,7 +29,7 @@ function makeImporter(): ContentCSVImporter { // offloaded worker path only, so a throw here would be seen by nobody. const report: FailureReporter = (error) => { try { - logging.error({event: {name: 'content.import.error'}, err: error}, 'Content import failure'); + logging.error({event: {name: 'content.import.error'}, err: error}, '[Background Job] content-import error'); sentry.captureException(error); } catch { // Callers report from catch blocks, so this must not throw. diff --git a/ghost/core/core/server/services/email-analytics/email-analytics-service-wrapper.ts b/ghost/core/core/server/services/email-analytics/email-analytics-service-wrapper.ts index fb30cb4639a..bebf6bb7417 100644 --- a/ghost/core/core/server/services/email-analytics/email-analytics-service-wrapper.ts +++ b/ghost/core/core/server/services/email-analytics/email-analytics-service-wrapper.ts @@ -9,6 +9,8 @@ import type {BatchEventProcessor} from './batch-event-processor'; import type {Queries} from './lib/queries'; import {fetchMailgunEvents} from './fetch-mailgun-events'; +const jobLogging = require('../jobs/job-logging'); + export class EmailAnalyticsServiceWrapper { #logName: string; #config?: Pick; @@ -22,6 +24,19 @@ export class EmailAnalyticsServiceWrapper { return `[EmailAnalytics:${this.#logName}]`; } + get #backgroundJobName(): string { + switch (this.#logName) { + case 'newsletters': + return 'email-analytics-fetch-latest'; + case 'automations': + return 'email-analytics-automation-fetch-latest'; + case 'gifts': + return 'email-analytics-gift-fetch-latest'; + default: + return `email-analytics-${this.#logName}-fetch-latest`; + } + } + constructor({logName}: {logName: string}) { this.#logName = logName; } @@ -108,14 +123,14 @@ export class EmailAnalyticsServiceWrapper { const batchMode = config.get('emailAnalytics:batchProcessing') ? 'BATCHED' : 'SEQUENTIAL'; const logMessage = [ - `${this.#logPrefix} Job complete: ${jobType}`, + `[Background Job] ${this.#backgroundJobName} processed ${jobType} | ${this.#logPrefix}`, `${eventCount} events in ${(totalDurationMs / 1000).toFixed(1)}s (${throughput.toFixed(2)} events/s)`, `Mode: ${batchMode}`, `Timings: API ${(apiPollingTimeMs / 1000).toFixed(1)}s (${apiPercent}%) / Processing ${(processingTimeMs / 1000).toFixed(1)}s (${processingPercent}%) / Aggregation ${(aggregationTimeMs / 1000).toFixed(1)}s (${aggregationPercent}%) [Email ${(emailAggregationTimeMs / 1000).toFixed(1)}s / Member ${(memberAggregationTimeMs / 1000).toFixed(1)}s]`, `Events: opened=${result.opened} delivered=${result.delivered} failed=${result.permanentFailed + result.temporaryFailed} unprocessable=${result.unprocessable}` ].join(' | '); - logging.info(logMessage); + jobLogging.info(logMessage); // We're only concerned with open throughput as this is displayed to users and is most sensitive to being up to date if (jobType === 'latest-opened') { @@ -199,13 +214,19 @@ export class EmailAnalyticsServiceWrapper { } async startFetch(): Promise { + const startedAt = Date.now(); if (!this.#restoredSchedule) { this.#restoredSchedule = true; - await this.service.restoreScheduled(); + try { + await this.service.restoreScheduled(); + } catch (e) { + jobLogging.error(e, `[Background Job] ${this.#backgroundJobName} failed while restoring scheduled events after ${Date.now() - startedAt}ms`); + throw e; + } } if (this.#fetching) { - logging.info(`Email analytics fetch for ${this.#logName} already running, skipping`); + jobLogging.info(`[Background Job] ${this.#backgroundJobName} skipped because a fetch is already running`); return; } this.#fetching = true; @@ -240,24 +261,21 @@ export class EmailAnalyticsServiceWrapper { return; } - // Log summary if no events were found across all jobs - if (c1 + c2 + c3 + c4 === 0) { - logging.info(`${this.#logPrefix} Job complete - No events`); - } + jobLogging.info(`[Background Job] ${this.#backgroundJobName} completed in ${Date.now() - startedAt}ms with ${c1 + c2 + c3 + c4} events | ${this.#logPrefix}`); this.#fetching = false; } catch (e) { - logging.error(e, `Error while fetching email analytics for ${this.#logName}`); + jobLogging.error(e, `[Background Job] ${this.#backgroundJobName} failed after ${Date.now() - startedAt}ms`); // Log again only the error, otherwise we lose the stack trace - logging.error(e); + jobLogging.error(e); } this.#fetching = false; } _restartFetch(reason: string): void { this.#fetching = false; - logging.info(`${this.#logPrefix} Restarting fetch due to ${reason}`); + jobLogging.info(`[Background Job] ${this.#backgroundJobName} continuing due to ${reason}`); this.startFetch(); } } diff --git a/ghost/core/core/server/services/email-analytics/jobs/automation-fetch-latest/index.js b/ghost/core/core/server/services/email-analytics/jobs/automation-fetch-latest/index.js index d97ec90b356..b22d0ed5780 100644 --- a/ghost/core/core/server/services/email-analytics/jobs/automation-fetch-latest/index.js +++ b/ghost/core/core/server/services/email-analytics/jobs/automation-fetch-latest/index.js @@ -2,6 +2,5 @@ const {run} = require('../fetch-latest-job'); const {StartAutomationEmailAnalyticsJobEvent} = require('../../events/start-automation-email-analytics-job-event'); run({ - event: StartAutomationEmailAnalyticsJobEvent, - logName: 'automations' + event: StartAutomationEmailAnalyticsJobEvent }); 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 edd78199825..25cb0d6efc7 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 @@ -1,6 +1,8 @@ import * as path from 'node:path'; import moment from 'moment'; +const jobLogging = require('../../jobs/job-logging'); + type CountableQuery = { where(column: string, operator: string, value: unknown): CountableQuery; count(): Promise; @@ -86,8 +88,10 @@ export class EmailAnalyticsJobScheduler { .count()); if (emailCount > 0 && !this.#hasScheduledNewslettersJob) { + const at = randomFiveMinuteCron(); + jobLogging.info(`[Background Job] email-analytics-fetch-latest scheduled at ${at}`); this.#jobManager.addJob({ - at: randomFiveMinuteCron(), + at, job: path.resolve(__dirname, 'fetch-latest/index.js'), name: 'email-analytics-fetch-latest' }); @@ -119,8 +123,10 @@ export class EmailAnalyticsJobScheduler { return; } + const at = randomFiveMinuteCron(); + jobLogging.info(`[Background Job] email-analytics-automation-fetch-latest scheduled at ${at}`); this.#jobManager.addJob({ - at: randomFiveMinuteCron(), + at, job: path.resolve(__dirname, 'automation-fetch-latest/index.js'), name: 'email-analytics-automation-fetch-latest' }); @@ -145,8 +151,10 @@ export class EmailAnalyticsJobScheduler { return; } + const at = randomFiveMinuteCron(); + jobLogging.info(`[Background Job] email-analytics-gift-fetch-latest scheduled at ${at}`); this.#jobManager.addJob({ - at: randomFiveMinuteCron(), + at, job: path.resolve(__dirname, 'gift-fetch-latest/index.js'), name: 'email-analytics-gift-fetch-latest' }); diff --git a/ghost/core/core/server/services/email-analytics/jobs/fetch-latest-job.js b/ghost/core/core/server/services/email-analytics/jobs/fetch-latest-job.js index b110951c879..5e54f3df486 100644 --- a/ghost/core/core/server/services/email-analytics/jobs/fetch-latest-job.js +++ b/ghost/core/core/server/services/email-analytics/jobs/fetch-latest-job.js @@ -5,12 +5,11 @@ const {parentPort} = require('worker_threads'); // Exit early when cancelled to prevent stalling shutdown. No cleanup needed when cancelling as everything is idempotent and will pick up // where it left off on next run /** - * @param {string} logName * @returns {void} */ -function cancel(logName) { +function cancel() { if (parentPort) { - parentPort.postMessage(`Email analytics fetch-latest job for ${logName} cancelled before completion`); + parentPort.postMessage('cancelled before completion'); parentPort.postMessage('cancelled'); } else { setTimeout(() => { @@ -22,17 +21,16 @@ function cancel(logName) { /** * @param {object} options * @param {{name: string}} options.event - * @param {string} options.logName * @returns {void} */ exports.run = ({ - event, - logName + event }) => { if (parentPort) { + parentPort.postMessage('execution started'); parentPort.once('message', (message) => { if (message === 'cancel') { - cancel(logName); + cancel(); return; } }); diff --git a/ghost/core/core/server/services/email-analytics/jobs/fetch-latest/index.js b/ghost/core/core/server/services/email-analytics/jobs/fetch-latest/index.js index ceed664397f..50a4f4c6d2a 100644 --- a/ghost/core/core/server/services/email-analytics/jobs/fetch-latest/index.js +++ b/ghost/core/core/server/services/email-analytics/jobs/fetch-latest/index.js @@ -2,6 +2,5 @@ const {run} = require('../fetch-latest-job'); const {StartEmailAnalyticsJobEvent} = require('../../events/start-email-analytics-job-event'); run({ - event: StartEmailAnalyticsJobEvent, - logName: 'newsletters' + event: StartEmailAnalyticsJobEvent }); 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 index bdd506e8d4a..40b46ef5f89 100644 --- 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 @@ -2,6 +2,5 @@ const {run} = require('../fetch-latest-job'); const {StartGiftEmailAnalyticsJobEvent} = require('../../events/start-gift-email-analytics-job-event'); run({ - event: StartGiftEmailAnalyticsJobEvent, - logName: 'gifts' + event: StartGiftEmailAnalyticsJobEvent }); diff --git a/ghost/core/core/server/services/email-service/batch-sending-service.js b/ghost/core/core/server/services/email-service/batch-sending-service.js index cd4c97b855b..668f1d5110b 100644 --- a/ghost/core/core/server/services/email-service/batch-sending-service.js +++ b/ghost/core/core/server/services/email-service/batch-sending-service.js @@ -1,4 +1,5 @@ const logging = require('@tryghost/logging'); +const jobLogging = require('../jobs/job-logging'); const ObjectID = require('bson-objectid').default; const errors = require('@tryghost/errors'); const tpl = require('@tryghost/tpl'); @@ -190,6 +191,7 @@ class BatchSendingService { * @returns {void} */ scheduleEmail(email) { + jobLogging.info(`[Background Job] batch-sending-service-job queued for email ${email.id}`); return this.#jobsService.addJob({ name: 'batch-sending-service-job', job: this.emailJob.bind(this), @@ -203,20 +205,26 @@ class BatchSendingService { * @param {{emailId: string}} data Data passed from the job service. We only need the emailId because we need to refetch the email anyway to make sure the status is right and 'locked'. */ async emailJob({emailId}) { - logging.info(`Starting email job for email ${emailId}`); + jobLogging.info(`[Background Job] batch-sending-service-job started for email ${emailId}`); const startTime = Date.now(); // Check if email is 'pending' only + change status to submitting in one transaction. // This allows us to have a lock around the email job that makes sure an email can only have one active job. - let email = await this.retryDb( - async () => { - return await this.updateStatusLock(this.#models.Email, emailId, 'submitting', ['pending', 'failed']); - }, - {...this.#getBeforeRetryConfig(), description: `updateStatusLock email ${emailId} -> submitting`} - ); + let email; + try { + email = await this.retryDb( + async () => { + return await this.updateStatusLock(this.#models.Email, emailId, 'submitting', ['pending', 'failed']); + }, + {...this.#getBeforeRetryConfig(), description: `updateStatusLock email ${emailId} -> submitting`} + ); + } catch (err) { + jobLogging.error(err, `[Background Job] batch-sending-service-job failed while acquiring the status lock after ${Date.now() - startTime}ms`); + throw err; + } if (!email) { - logging.error(`Tried sending email that is not pending or failed ${emailId}`); + jobLogging.error(`[Background Job] batch-sending-service-job skipped because email ${emailId} is not pending or failed`); return; } @@ -238,12 +246,13 @@ class BatchSendingService { error: null }, {patch: true, autoRefresh: false}); }, {...this.#getAfterRetryConfig(), description: `email ${emailId} -> submitted`}); + jobLogging.info(`[Background Job] batch-sending-service-job completed for email ${emailId} in ${Date.now() - startTime}ms`); } catch (e) { // Any failure while shutting down counts as interrupted, not failed: // collapsed budgets surface transient errors as hard failures, and `failed` // drops the email out of the boot resume scan. if ((e && e.code === SHUTDOWN_CODE) || this.#shuttingDown) { - logging.info(`Email ${email.id} send stopped because the container is shutting down — leaving status=submitting so it can resume on next boot`); + jobLogging.info(`[Background Job] batch-sending-service-job send stopped because the container is shutting down — leaving email ${email.id} status=submitting so it can resume on next boot`); return; } const ghostError = new errors.EmailError({ @@ -252,7 +261,7 @@ class BatchSendingService { message: `Error sending email ${email.id}` }); - logging.error(ghostError); + jobLogging.error(ghostError, `[Background Job] batch-sending-service-job failed for email ${emailId} after ${Date.now() - startTime}ms`); if (this.#sentry) { // Log the original error to Sentry this.#sentry.captureException(e); diff --git a/ghost/core/core/server/services/gifts/index.ts b/ghost/core/core/server/services/gifts/index.ts index 3913b20753a..10894cda195 100644 --- a/ghost/core/core/server/services/gifts/index.ts +++ b/ghost/core/core/server/services/gifts/index.ts @@ -36,6 +36,7 @@ export async function init(options: GiftServiceInitOptions): Promise { const staffService = require('../staff'); const DomainEvents = require('@tryghost/domain-events'); const logging = require('@tryghost/logging'); + const jobLogging = require('../jobs/job-logging'); const {SubscriptionActivatedEvent} = require('../../../shared/events'); const StartGiftReminderFlushEvent = require('./events/start-gift-reminder-flush-event'); const StartGiftCleanupEvent = require('./events/start-gift-cleanup-event'); @@ -122,12 +123,13 @@ export async function init(options: GiftServiceInitOptions): Promise { DomainEvents.subscribe(StartGiftReminderFlushEvent, async () => { const start = Date.now(); + jobLogging.info('[Background Job] send-gift-reminders started'); try { const {remindedCount, skippedCount, failedCount} = await giftService.processReminders(); - logging.info(`Sent ${remindedCount} gift reminders, skipped ${skippedCount}, failed ${failedCount} in ${Date.now() - start}ms`); + jobLogging.info(`[Background Job] send-gift-reminders completed in ${Date.now() - start}ms: ${remindedCount} sent, ${skippedCount} not due, ${failedCount} rejected`); } catch (err) { - logging.error(err, 'Failed to process gift reminders'); + jobLogging.error(err, `[Background Job] send-gift-reminders failed after ${Date.now() - start}ms`); } }); @@ -142,41 +144,46 @@ export async function init(options: GiftServiceInitOptions): Promise { }); DomainEvents.subscribe(StartGiftCleanupEvent, async () => { + const cleanupStart = Date.now(); + jobLogging.info('[Background Job] clean-gifts started'); + const checkoutStart = Date.now(); try { const {deletedCount} = await giftService.processAbandonedCheckouts(); - logging.info(`Deleted ${deletedCount} abandoned gift checkouts in ${Date.now() - checkoutStart}ms`); + jobLogging.info(`[Background Job] clean-gifts processed abandoned checkouts: deleted ${deletedCount} in ${Date.now() - checkoutStart}ms`); } catch (err) { - logging.error(err, 'Failed to clean abandoned gift checkouts'); + jobLogging.error(err, '[Background Job] clean-gifts error processing abandoned checkouts'); } const consumedStart = Date.now(); try { const {consumedCount, updatedMemberCount} = await giftService.processConsumed(); - logging.info(`Consumed ${consumedCount} gifts, updated ${updatedMemberCount} members in ${Date.now() - consumedStart}ms`); + jobLogging.info(`[Background Job] clean-gifts processed consumed gifts: consumed ${consumedCount}, updated ${updatedMemberCount} members in ${Date.now() - consumedStart}ms`); } catch (err) { - logging.error(err, 'Failed to process consumed gifts'); + jobLogging.error(err, '[Background Job] clean-gifts error processing consumed gifts'); } const expiredStart = Date.now(); try { const {expiredCount} = await giftService.processExpired(); - logging.info(`Expired ${expiredCount} gifts in ${Date.now() - expiredStart}ms`); + jobLogging.info(`[Background Job] clean-gifts processed expired gifts: expired ${expiredCount} in ${Date.now() - expiredStart}ms`); } catch (err) { - logging.error(err, 'Failed to process expired gifts'); + jobLogging.error(err, '[Background Job] clean-gifts error processing expired gifts'); } try { const {sentCount, skippedCount, failedCount} = await giftDeliveryService.recoverPending(); if (sentCount + skippedCount + failedCount > 0) { - logging.info(`Gift delivery recovery during cleanup: ${sentCount} sent, ${skippedCount} skipped, ${failedCount} failed`); + jobLogging.info(`[Background Job] clean-gifts processed pending gift deliveries: ${sentCount} sent, ${skippedCount} not due, ${failedCount} rejected`); } } catch (err) { - logging.error(err, 'Failed to recover pending gift deliveries during cleanup'); + jobLogging.error(err, '[Background Job] clean-gifts error processing pending gift deliveries'); } + + jobLogging.info(`[Background Job] clean-gifts completed in ${Date.now() - cleanupStart}ms`); }); jobs.scheduleGiftCleanupJob(); diff --git a/ghost/core/core/server/services/gifts/jobs/clean-gifts-job.js b/ghost/core/core/server/services/gifts/jobs/clean-gifts-job.js index d1521815e4f..60b243a9e42 100644 --- a/ghost/core/core/server/services/gifts/jobs/clean-gifts-job.js +++ b/ghost/core/core/server/services/gifts/jobs/clean-gifts-job.js @@ -10,7 +10,7 @@ const StartGiftCleanupEvent = require('../events/start-gift-cleanup-event'); // off on next run function cancel() { if (parentPort) { - parentPort.postMessage('Gift cleanup job cancelled before completion'); + parentPort.postMessage('cancelled before completion'); parentPort.postMessage('cancelled'); } else { setTimeout(() => { @@ -37,6 +37,7 @@ if (parentPort) { type: StartGiftCleanupEvent.name } }); + parentPort.postMessage('dispatched to main process'); parentPort.postMessage('done'); } else { setTimeout(() => { diff --git a/ghost/core/core/server/services/gifts/jobs/index.js b/ghost/core/core/server/services/gifts/jobs/index.js index e9946b3ee6d..1005d327249 100644 --- a/ghost/core/core/server/services/gifts/jobs/index.js +++ b/ghost/core/core/server/services/gifts/jobs/index.js @@ -1,4 +1,5 @@ const path = require('path'); +const jobLogging = require('../../jobs/job-logging'); const jobsService = require('../../jobs'); let hasScheduled = { @@ -18,8 +19,11 @@ function scheduleJob(key, name, jobFile) { const m = Math.floor(Math.random() * 60); const h = Math.floor(Math.random() * 6); + const at = `${s} ${m} ${h} * * *`; + + jobLogging.info(`[Background Job] ${name} scheduled at ${at}`); jobsService.addJob({ - at: `${s} ${m} ${h} * * *`, + at, job: path.resolve(__dirname, jobFile), name }); diff --git a/ghost/core/core/server/services/gifts/jobs/send-gift-reminders-job.js b/ghost/core/core/server/services/gifts/jobs/send-gift-reminders-job.js index 33f457fad79..40b40e527f3 100644 --- a/ghost/core/core/server/services/gifts/jobs/send-gift-reminders-job.js +++ b/ghost/core/core/server/services/gifts/jobs/send-gift-reminders-job.js @@ -10,7 +10,7 @@ const StartGiftReminderFlushEvent = require('../events/start-gift-reminder-flush // its reminder recorded will be picked up on the next run function cancel() { if (parentPort) { - parentPort.postMessage('Gift reminder job cancelled before completion'); + parentPort.postMessage('cancelled before completion'); parentPort.postMessage('cancelled'); } else { setTimeout(() => { @@ -37,6 +37,7 @@ if (parentPort) { type: StartGiftReminderFlushEvent.name } }); + parentPort.postMessage('dispatched to main process'); parentPort.postMessage('done'); } else { setTimeout(() => { diff --git a/ghost/core/core/server/services/jobs-service/register-job-handlers.ts b/ghost/core/core/server/services/jobs-service/register-job-handlers.ts index f21b09bfec5..53cddc63722 100644 --- a/ghost/core/core/server/services/jobs-service/register-job-handlers.ts +++ b/ghost/core/core/server/services/jobs-service/register-job-handlers.ts @@ -1,21 +1,30 @@ -import logging from '@tryghost/logging'; import {getInstance} from './index'; import CleanTokensJob from '../members/jobs/clean-tokens-job'; import cleanTokens from '../members/jobs/clean-tokens-task'; +const jobLogging = require('../jobs/job-logging'); + export default function registerJobHandlers(): void { const jobsService = getInstance(); const db = require('../../data/db'); jobsService.handle(CleanTokensJob, async () => { const startedAt = Date.now(); - const deletedCount = await cleanTokens({db}); - logging.info({ - system: { - event: 'clean_tokens.completed', - deleted_count: deletedCount, - duration_ms: Date.now() - startedAt - } - }, `Removed ${deletedCount} tokens older than 24 hours`); + jobLogging.info('[Background Job] clean-tokens started'); + + try { + const deletedCount = await cleanTokens({db}); + const durationMs = Date.now() - startedAt; + jobLogging.info({ + system: { + event: 'clean_tokens.completed', + deleted_count: deletedCount, + duration_ms: durationMs + } + }, `[Background Job] clean-tokens completed in ${durationMs}ms: removed ${deletedCount} tokens older than 24 hours`); + } catch (error) { + jobLogging.error(error, `[Background Job] clean-tokens failed after ${Date.now() - startedAt}ms`); + throw error; + } }); } diff --git a/ghost/core/core/server/services/jobs/job-logging.ts b/ghost/core/core/server/services/jobs/job-logging.ts new file mode 100644 index 00000000000..2dbeb9f6fc4 --- /dev/null +++ b/ghost/core/core/server/services/jobs/job-logging.ts @@ -0,0 +1,19 @@ +import logging from '@tryghost/logging'; + +type LogMethod = 'info' | 'error'; + +function bestEffort(method: LogMethod, args: unknown[]): void { + try { + logging[method](...args); + } catch { + // Observability must not control background-job execution. + } +} + +export function info(...args: unknown[]): void { + bestEffort('info', args); +} + +export function error(...args: unknown[]): void { + bestEffort('error', args); +} diff --git a/ghost/core/core/server/services/jobs/job-service.js b/ghost/core/core/server/services/jobs/job-service.js index ce3def9a46e..733718c4245 100644 --- a/ghost/core/core/server/services/jobs/job-service.js +++ b/ghost/core/core/server/services/jobs/job-service.js @@ -5,14 +5,14 @@ const JobManager = require('@tryghost/job-manager'); const logging = require('@tryghost/logging'); +const jobLogging = require('./job-logging'); const models = require('../../models'); const sentry = require('../../../shared/sentry'); const domainEvents = require('@tryghost/domain-events'); const config = require('../../../shared/config'); const WorkerModelEventBridge = require('./worker-model-event-bridge'); const errorHandler = (error, workerMeta) => { - logging.info(`Capturing error for worker during execution of job: ${workerMeta.name}`); - logging.error(error); + jobLogging.error(error, `[Background Job] ${workerMeta.name} failed`); sentry.captureException(error); }; const events = require('../../lib/common/events'); @@ -26,8 +26,8 @@ const workerMessageHandler = ({name, message}) => { return; } - if (typeof message === 'string') { - logging.info(`Worker for job ${name} sent a message: ${message}`); + if (typeof message === 'string' && !['done', 'cancelled'].includes(message)) { + jobLogging.info(`[Background Job] ${name}: ${message}`); } }; diff --git a/ghost/core/core/server/services/media-inliner/service.js b/ghost/core/core/server/services/media-inliner/service.js index 0664ad479c3..988c4a2416a 100644 --- a/ghost/core/core/server/services/media-inliner/service.js +++ b/ghost/core/core/server/services/media-inliner/service.js @@ -2,6 +2,7 @@ module.exports = { async init() { const debug = require('@tryghost/debug')('mediaInliner'); const MediaInliner = require('./external-media-inliner'); + const jobLogging = require('../jobs/job-logging'); const models = require('../../models'); const jobsService = require('../jobs'); const adapterManager = require('../../services/adapter-manager').default; @@ -45,10 +46,20 @@ module.exports = { // @NOTE: the job is "inline" (aka non-offloaded into a thread), because usecases are currently // limited to migrational, so there is no expectations for site's availability etc. + jobLogging.info('[Background Job] external-media-inliner queued'); await jobsService.addJob({ name: 'external-media-inliner', - job: (data) => { - return mediaInliner.inline(data.domains); + job: async (data) => { + const startedAt = Date.now(); + jobLogging.info('[Background Job] external-media-inliner started'); + try { + const result = await mediaInliner.inline(data.domains); + jobLogging.info(`[Background Job] external-media-inliner completed in ${Date.now() - startedAt}ms`); + return result; + } catch (err) { + jobLogging.error(err, `[Background Job] external-media-inliner failed after ${Date.now() - startedAt}ms`); + throw err; + } }, data: {domains}, offloaded: false diff --git a/ghost/core/core/server/services/members/import-export/import/importer.ts b/ghost/core/core/server/services/members/import-export/import/importer.ts index 15b14f56a17..123301aad3c 100644 --- a/ghost/core/core/server/services/members/import-export/import/importer.ts +++ b/ghost/core/core/server/services/members/import-export/import/importer.ts @@ -8,6 +8,7 @@ import type {RowSpool, SpooledRows} from './spool'; const metrics = require('@tryghost/metrics'); const errors = require('@tryghost/errors'); +const jobLogging = require('../../../jobs/job-logging'); const tpl = require('@tryghost/tpl'); // The members CSV importer, sliced into one concern per method. Two entry points by @@ -264,6 +265,7 @@ class MembersCSVImporter { const emailRecipient: string = requestUserEmail ?? await this._email.getDefaultRecipient(); const spooled = await this._spool.write(rows); + jobLogging.info('[Background Job] members-import queued'); this._addJob({ job: () => this.runImportJob(spooled, {labelName, extraLabels, emailRecipient}, verificationTrigger), offloaded: false, @@ -278,6 +280,8 @@ class MembersCSVImporter { {labelName, extraLabels, emailRecipient}: {labelName: string; extraLabels: Label[]; emailRecipient: string}, verificationTrigger: VerificationTrigger ): Promise { + const startedAt = Date.now(); + jobLogging.info('[Background Job] members-import started'); // Null until the import produces one: parsing and mapping already happened inside // the request, so anything failing from here is ours rather than the file's. let result: ImportResult | null = null; @@ -299,6 +303,12 @@ class MembersCSVImporter { labelName, links: this._email.links }))); + + if (result) { + jobLogging.info(`[Background Job] members-import completed in ${Date.now() - startedAt}ms: imported ${result.imported}, ${result.errors.length} row(s) rejected`); + } else { + jobLogging.info(`[Background Job] members-import failed after ${Date.now() - startedAt}ms`); + } } // Only the write itself may throw. Callers rely on that to tell an import that never diff --git a/ghost/core/core/server/services/members/import-export/index.ts b/ghost/core/core/server/services/members/import-export/index.ts index d0f1e8277a2..d207d6ee649 100644 --- a/ghost/core/core/server/services/members/import-export/index.ts +++ b/ghost/core/core/server/services/members/import-export/index.ts @@ -105,7 +105,7 @@ export function makeImporter(deps: ImporterServices) { // offloaded worker path only, so a throw here would be seen by nobody. const report: FailureReporter = (error) => { try { - logging.error({event: {name: 'members.import.error'}, err: error}, 'Members import failure'); + logging.error({event: {name: 'members.import.error'}, err: error}, '[Background Job] members-import error'); sentry.captureException(error); } catch { // Callers report from catch and finally blocks, so this must not throw. diff --git a/ghost/core/core/server/services/members/jobs/clean-expired-comped.js b/ghost/core/core/server/services/members/jobs/clean-expired-comped.js index 838aafc9d65..71ea153ed43 100644 --- a/ghost/core/core/server/services/members/jobs/clean-expired-comped.js +++ b/ghost/core/core/server/services/members/jobs/clean-expired-comped.js @@ -10,7 +10,7 @@ const moment = require('moment'); // where it left off on next run function cancel() { if (parentPort) { - parentPort.postMessage('Expired complimentary subscriptions cleanup cancelled before completion'); + parentPort.postMessage('cancelled before completion'); parentPort.postMessage('cancelled'); } else { setTimeout(() => { @@ -30,6 +30,9 @@ if (parentPort) { (async () => { const cleanupStartDate = new Date(); const db = require('../../../data/db'); + if (parentPort) { + parentPort.postMessage('execution started'); + } debug(`Starting cleanup of expired comp subscriptions`); const expiredCompedRows = await db.knex('members_products') .where('expiry_at', '<', moment.utc().startOf('day').toISOString()) @@ -116,7 +119,7 @@ if (parentPort) { debug(`Removed ${deletedExpiredSubs} expired subscriptions, updated ${updatedMembers} members in ${cleanupEndDate.valueOf() - cleanupStartDate.valueOf()}ms`); if (parentPort) { - parentPort.postMessage(`Removed ${deletedExpiredSubs} expired subscriptions, updated ${updatedMembers} members in ${cleanupEndDate.valueOf() - cleanupStartDate.valueOf()}ms`); + parentPort.postMessage(`completed in ${cleanupEndDate.valueOf() - cleanupStartDate.valueOf()}ms: removed ${deletedExpiredSubs} expired subscriptions and updated ${updatedMembers} members`); parentPort.postMessage('done'); } else { // give the logging pipes time finish writing before exit diff --git a/ghost/core/core/server/services/members/jobs/index.js b/ghost/core/core/server/services/members/jobs/index.js index 3c2182dfca1..14c99aa194e 100644 --- a/ghost/core/core/server/services/members/jobs/index.js +++ b/ghost/core/core/server/services/members/jobs/index.js @@ -1,4 +1,5 @@ const path = require('path'); +const jobLogging = require('../../jobs/job-logging'); const jobsService = require('../../jobs'); const CleanTokensJob = require('./clean-tokens-job').default; @@ -27,8 +28,11 @@ function scheduleJob(key, name, jobFile, maxHour = 6) { const m = Math.floor(Math.random() * 60); const h = Math.floor(Math.random() * maxHour); + const at = `${s} ${m} ${h} * * *`; + + jobLogging.info(`[Background Job] ${name} scheduled at ${at}`); jobsService.addJob({ - at: `${s} ${m} ${h} * * *`, + at, job: path.resolve(__dirname, jobFile), name }); @@ -49,7 +53,9 @@ module.exports = { } const classBasedJobs = require('../../jobs-service').getInstance(); - await classBasedJobs.scheduleRecurring(new CleanTokensJob(), {cron: randomDailyCron()}); + const cron = randomDailyCron(); + jobLogging.info(`[Background Job] clean-tokens scheduled at ${cron}`); + await classBasedJobs.scheduleRecurring(new CleanTokensJob(), {cron}); hasScheduled.tokens = true; } diff --git a/ghost/core/core/server/services/members/service.js b/ghost/core/core/server/services/members/service.js index b22a9ad8ce4..0416203f5f6 100644 --- a/ghost/core/core/server/services/members/service.js +++ b/ghost/core/core/server/services/members/service.js @@ -9,6 +9,7 @@ const {resolveInlineThreshold} = require('./import-export/config'); const MembersStats = require('./stats/members-stats'); const memberJobs = require('./jobs'); const logging = require('@tryghost/logging'); +const jobLogging = require('../jobs/job-logging'); const urlUtils = require('../../../shared/url-utils').default; const settingsCache = require('../../../shared/settings-cache'); const config = require('../../../shared/config'); @@ -175,13 +176,27 @@ module.exports = { if (!env?.startsWith('testing')) { const membersMigrationJobName = 'members-migrations'; if (!(await jobsService.hasExecutedSuccessfully(membersMigrationJobName))) { + jobLogging.info(`[Background Job] ${membersMigrationJobName} queued`); jobsService.addOneOffJob({ name: membersMigrationJobName, offloaded: false, - job: stripeService.migrations.execute.bind(stripeService.migrations) + job: async () => { + const startedAt = Date.now(); + jobLogging.info(`[Background Job] ${membersMigrationJobName} started`); + try { + const result = await stripeService.migrations.execute(); + jobLogging.info(`[Background Job] ${membersMigrationJobName} completed in ${Date.now() - startedAt}ms`); + return result; + } catch (err) { + jobLogging.error(err, `[Background Job] ${membersMigrationJobName} failed after ${Date.now() - startedAt}ms`); + throw err; + } + } }); await jobsService.awaitOneOffCompletion(membersMigrationJobName); + } else { + jobLogging.info(`[Background Job] ${membersMigrationJobName} skipped because it has already run`); } } diff --git a/ghost/core/core/server/services/mentions-jobs/job-service.js b/ghost/core/core/server/services/mentions-jobs/job-service.js index 1edd38ba0d6..80795bd2427 100644 --- a/ghost/core/core/server/services/mentions-jobs/job-service.js +++ b/ghost/core/core/server/services/mentions-jobs/job-service.js @@ -5,19 +5,19 @@ const JobManager = require('@tryghost/job-manager'); const logging = require('@tryghost/logging'); +const jobLogging = require('../jobs/job-logging'); const models = require('../../models'); const sentry = require('../../../shared/sentry'); const domainEvents = require('@tryghost/domain-events'); const errorHandler = (error, workerMeta) => { - logging.info(`Capturing error for worker during execution of job: ${workerMeta.name}`); - logging.error(error); + jobLogging.error(error, `[Background Job] ${workerMeta.name} failed`); sentry.captureException(error); }; const workerMessageHandler = ({name, message}) => { - if (typeof message === 'string') { - logging.info(`Worker for job ${name} sent a message: ${message}`); + if (typeof message === 'string' && !['done', 'cancelled'].includes(message)) { + jobLogging.info(`[Background Job] ${name}: ${message}`); } }; diff --git a/ghost/core/core/server/services/mentions/mention-controller.js b/ghost/core/core/server/services/mentions/mention-controller.js index 9db7712504f..7bf4008a9a9 100644 --- a/ghost/core/core/server/services/mentions/mention-controller.js +++ b/ghost/core/core/server/services/mentions/mention-controller.js @@ -131,7 +131,7 @@ module.exports = class MentionController { payload }); } catch (err) { - logging.error(err); + logging.error(err, '[Webmention] Failed processing webmention'); } }); } diff --git a/ghost/core/core/server/services/mentions/service.js b/ghost/core/core/server/services/mentions/service.js index 499a39880de..f6d15fbc88e 100644 --- a/ghost/core/core/server/services/mentions/service.js +++ b/ghost/core/core/server/services/mentions/service.js @@ -14,6 +14,7 @@ const outputSerializerUrlUtil = require('../../../server/api/endpoints/utils/ser const urlService = require('../url'); const settingsCache = require('../../../shared/settings-cache'); const DomainEvents = require('@tryghost/domain-events'); +const jobLogging = require('../jobs/job-logging'); const jobsService = require('../mentions-jobs'); // Serializes a post model to the data the URL service needs, loading the @@ -37,6 +38,33 @@ function getPostUrl(id, postData) { return jsonModel.url; } +// Reports the same queued/started/finished/failed lifecycle for every mentions +// background job. The wrapped callback's result and errors pass through unchanged, +// so job manager outcomes and retries are unaffected. +function makeLoggingJobService() { + return { + async addJob(name, fn) { + jobLogging.info(`[Background Job] ${name} queued`); + jobsService.addJob({ + name, + job: async () => { + const startedAt = Date.now(); + jobLogging.info(`[Background Job] ${name} started`); + try { + const result = await fn(); + jobLogging.info(`[Background Job] ${name} completed in ${Date.now() - startedAt}ms`); + return result; + } catch (err) { + jobLogging.error(err, `[Background Job] ${name} failed after ${Date.now() - startedAt}ms`); + throw err; + } + }, + offloaded: false + }); + } + }; +} + module.exports = { /** @type {import('./mentions-api')} */ api: null, @@ -82,15 +110,7 @@ module.exports = { this.controller.init({ api, - jobService: { - async addJob(name, fn) { - jobsService.addJob({ - name, - job: fn, - offloaded: false - }); - } - }, + jobService: makeLoggingJobService(), mentionResourceService: { async getByID(id) { if (!id) { @@ -117,15 +137,7 @@ module.exports = { getPostData: post => getPostData(post), getPostUrl: (id, data) => getPostUrl(id, data), isEnabled: () => !settingsCache.get('is_private'), - jobService: { - async addJob(name, fn) { - jobsService.addJob({ - name, - job: fn, - offloaded: false - }); - } - } + jobService: makeLoggingJobService() }); sendingService.listen(events); @@ -136,3 +148,4 @@ module.exports = { // exposed for testing module.exports.getPostData = getPostData; module.exports.getPostUrl = getPostUrl; +module.exports.makeLoggingJobService = makeLoggingJobService; diff --git a/ghost/core/core/server/services/update-check/index.js b/ghost/core/core/server/services/update-check/index.js index 15ec246b961..2930ea69f16 100644 --- a/ghost/core/core/server/services/update-check/index.js +++ b/ghost/core/core/server/services/update-check/index.js @@ -1,6 +1,7 @@ const api = require('../../api').endpoints; const config = require('../../../shared/config'); +const jobLogging = require('../jobs/job-logging'); const urlUtils = require('../../../shared/url-utils').default; const jobsService = require('../jobs'); @@ -72,14 +73,17 @@ module.exports.scheduleRecurringJobs = () => { const m = Math.floor(Math.random() * 60); // 0-59 const h = Math.floor(Math.random() * 24); // 0-23 + const at = `${s} ${m} ${h} * * *`; + jobLogging.info(`[Background Job] update-check scheduled at ${at}`); jobsService.addJob({ - at: `${s} ${m} ${h} * * *`, // Every day + at, // Every day job: require('path').resolve(__dirname, 'run-update-check.js'), name: 'update-check' }); }; module.exports.scheduleBootJob = () => { + jobLogging.info('[Background Job] update-check-boot queued'); jobsService.addJob({ job: require('path').resolve(__dirname, 'run-update-check.js'), name: 'update-check-boot' diff --git a/ghost/core/core/server/services/update-check/run-update-check.js b/ghost/core/core/server/services/update-check/run-update-check.js index 3191ac50391..121429e3694 100644 --- a/ghost/core/core/server/services/update-check/run-update-check.js +++ b/ghost/core/core/server/services/update-check/run-update-check.js @@ -9,7 +9,7 @@ const postParentPortMessage = (message) => { // Exit early when cancelled to prevent stalling shutdown. No cleanup needed when cancelling as everything is idempotent and will pick up // where it left off on next run function cancel() { - postParentPortMessage('Update check job cancelled before completion'); + postParentPortMessage('cancelled before completion'); if (parentPort) { postParentPortMessage('cancelled'); @@ -29,6 +29,8 @@ if (parentPort) { } (async () => { + const startedAt = Date.now(); + postParentPortMessage('execution started'); const updateCheck = require('./'); // INIT required services @@ -48,7 +50,7 @@ if (parentPort) { updateCheckUrl: workerData.updateCheckUrl }); - postParentPortMessage(`Ran update check`); + postParentPortMessage(`completed in ${Date.now() - startedAt}ms`); if (parentPort) { postParentPortMessage('done'); diff --git a/ghost/core/test/integration/services/email-service/batch-sending.test.js b/ghost/core/test/integration/services/email-service/batch-sending.test.js index 0f6019eeb45..9bda959fa66 100644 --- a/ghost/core/test/integration/services/email-service/batch-sending.test.js +++ b/ghost/core/test/integration/services/email-service/batch-sending.test.js @@ -229,7 +229,7 @@ describe('Batch sending tests', function () { // config, an unrelated transient failure elsewhere in the same window // would log a non-string Error object and crash assert.match instead of // failing the assertion cleanly. - const guardLogs = errorLog.getCalls().filter(call => typeof call.args[0] === 'string' && /Tried sending email that is not pending or failed/.test(call.args[0])); + const guardLogs = errorLog.getCalls().filter(call => typeof call.args[0] === 'string' && /\[Background Job\] batch-sending-service-job skipped/.test(call.args[0])); assert.ok(guardLogs.length > 0, 'expected at least one "not pending or failed" guard error log'); }); diff --git a/ghost/core/test/unit/server/services/content-import/import/importer.test.ts b/ghost/core/test/unit/server/services/content-import/import/importer.test.ts index d1ee6e40ffd..5840f59982d 100644 --- a/ghost/core/test/unit/server/services/content-import/import/importer.test.ts +++ b/ghost/core/test/unit/server/services/content-import/import/importer.test.ts @@ -1,4 +1,6 @@ import assert from 'node:assert/strict'; +import sinon from 'sinon'; +import logging from '@tryghost/logging'; import ContentCSVImporter from '../../../../../../core/server/services/content-import/import/importer'; import {ImportRunStore} from '../../../../../../core/server/services/content-import/import/store'; import type {PostImportRow} from '../../../../../../core/server/services/content-import/import/row'; @@ -77,6 +79,16 @@ function harness(rows: PostImportRow[] = [row('First'), row('Second')]) { } describe('ContentCSVImporter', function () { + let infoLog: sinon.SinonStub; + + beforeEach(function () { + infoLog = sinon.stub(logging, 'info'); + }); + + afterEach(function () { + sinon.restore(); + }); + it('accepts the upload with the row count and defers the writes to one inline job', async function () { const h = harness(); @@ -90,6 +102,28 @@ describe('ContentCSVImporter', function () { assert.equal(h.store.get('run_test')?.status, 'running', 'the run is registered before the job starts'); }); + it('logs the searchable lifecycle of the inline job', async function () { + const h = harness(); + + await h.run(); + + sinon.assert.calledWithExactly(infoLog, '[Background Job] content-import queued'); + sinon.assert.calledWithExactly(infoLog, '[Background Job] content-import started'); + sinon.assert.calledWithMatch(infoLog, /^\[Background Job\] content-import completed in \d+ms$/); + }); + + it('keeps lifecycle logging best-effort', async function () { + const h = harness(); + infoLog.throws(new Error('Logging unavailable')); + + const accepted = await h.run(); + + assert.deepEqual(accepted, {importId: 'run_test', total: 2}); + assert.equal(h.jobs.length, 1, 'the job is queued even when the queued log fails'); + assert.equal(h.store.get('run_test')?.status, 'complete', 'the job still resolves and finalizes its run'); + assert.equal(h.created.length, 2); + }); + it('writes one post per row, in order, under the importing options', async function () { const h = harness(); @@ -236,6 +270,7 @@ describe('ContentCSVImporter', function () { assert.equal(h.store.get('run_test')?.status, 'failed'); assert.equal(h.store.get('run_test')?.failureReason, 'converter unavailable'); assert.ok(h.store.get('run_test')?.finishedAt instanceof Date); + sinon.assert.calledWithMatch(infoLog, /^\[Background Job\] content-import failed after \d+ms$/); }); it('keeps a successfully written post created when its URL cannot be resolved', async function () { diff --git a/ghost/core/test/unit/server/services/email-analytics/email-analytics-service-wrapper.test.ts b/ghost/core/test/unit/server/services/email-analytics/email-analytics-service-wrapper.test.ts index 2a14da620b9..b2adc10e26a 100644 --- a/ghost/core/test/unit/server/services/email-analytics/email-analytics-service-wrapper.test.ts +++ b/ghost/core/test/unit/server/services/email-analytics/email-analytics-service-wrapper.test.ts @@ -1,4 +1,6 @@ +import assert from 'node:assert/strict'; import sinon from 'sinon'; +import logging from '@tryghost/logging'; import {EmailAnalyticsServiceWrapper} from '../../../../../core/server/services/email-analytics/email-analytics-service-wrapper'; import {EventProcessingResult} from '../../../../../core/server/services/email-analytics/event-processing-result'; import {Queries} from '../../../../../core/server/services/email-analytics/lib/queries'; @@ -73,6 +75,8 @@ describe('EmailAnalyticsServiceWrapper', function () { memberAggregationTimeMs: 200, result: new EventProcessingResult() }, 2000); + + return wrapper; } it('uses existing open throughput metric name for newsletters', function () { @@ -95,6 +99,72 @@ describe('EmailAnalyticsServiceWrapper', function () { }); }); + it('uses the gift analytics job name in lifecycle logs', function () { + const infoLog = sinon.stub(logging, 'info'); + + logLatestOpenedJob('gifts'); + + sinon.assert.calledWith(infoLog, sinon.match('[Background Job] email-analytics-gift-fetch-latest processed')); + }); + + it('does not let completion logging failures interrupt event processing', function () { + const wrapper = logLatestOpenedJob('newsletters'); + sinon.stub(logging, 'info').throws(new Error('Logger unavailable')); + + assert.doesNotThrow(() => wrapper._logJobCompletion('latest-opened', { + eventCount: 10, + apiPollingTimeMs: 500, + processingTimeMs: 1000, + aggregationTimeMs: 500, + emailAggregationTimeMs: 300, + memberAggregationTimeMs: 200, + result: new EventProcessingResult() + }, 2000)); + + assert.equal(metricStub.callCount, 2); + }); + + it('logs and preserves initial schedule restoration failures', async function () { + const errorLog = sinon.stub(logging, 'error'); + const wrapper = logLatestOpenedJob('newsletters'); + const restoreError = new Error('Restore failed'); + sinon.stub(wrapper.service, 'restoreScheduled').rejects(restoreError); + + await assert.rejects(wrapper.startFetch(), error => error === restoreError); + + sinon.assert.calledOnceWithExactly( + errorLog, + restoreError, + sinon.match('[Background Job] email-analytics-fetch-latest failed while restoring scheduled events') + ); + }); + + it('logs exactly one terminal event with a run duration', async function () { + const infoLog = sinon.stub(logging, 'info'); + const wrapper = logLatestOpenedJob('newsletters'); + infoLog.resetHistory(); + sinon.stub(wrapper.service, 'restoreScheduled').resolves(); + sinon.stub(wrapper, 'fetchLatestOpenedEvents').resolves(1); + sinon.stub(wrapper, 'fetchLatestNonOpenedEvents').resolves(0); + sinon.stub(wrapper, 'fetchMissing').resolves(0); + sinon.stub(wrapper, 'fetchScheduled').resolves(0); + + await wrapper.startFetch(); + + const completions = infoLog.args.filter(([message]) => typeof message === 'string' && message.startsWith('[Background Job] email-analytics-fetch-latest completed')); + assert.equal(completions.length, 1); + assert.match(completions[0][0] as string, /^\[Background Job\] email-analytics-fetch-latest completed in \d+ms with 1 events /); + }); + + it('does not let failure logging escape the fetch error handler', async function () { + const wrapper = logLatestOpenedJob('newsletters'); + sinon.stub(logging, 'error').throws(new Error('Logger unavailable')); + sinon.stub(wrapper.service, 'restoreScheduled').resolves(); + sinon.stub(wrapper, 'fetchLatestOpenedEvents').rejects(new Error('Fetch failed')); + + await assert.doesNotReject(wrapper.startFetch()); + }); + it('skips opened event polling when the cursor seed has no opened column', async function () { const wrapper = new EmailAnalyticsServiceWrapper({logName: 'gifts'}); wrapper.init({ diff --git a/ghost/core/test/unit/server/services/email-service/batch-sending-service.test.js b/ghost/core/test/unit/server/services/email-service/batch-sending-service.test.js index b9ab10a0f65..b8ccc65fa39 100644 --- a/ghost/core/test/unit/server/services/email-service/batch-sending-service.test.js +++ b/ghost/core/test/unit/server/services/email-service/batch-sending-service.test.js @@ -14,10 +14,11 @@ const simulateSleep = async (ms, clock) => { describe('Batch Sending Service', function () { let errorLog; + let infoLog; beforeEach(function () { errorLog = sinon.stub(logging, 'error'); - sinon.stub(logging, 'info'); + infoLog = sinon.stub(logging, 'info'); }); afterEach(function () { @@ -52,6 +53,20 @@ describe('Batch Sending Service', function () { }); describe('emailJob', function () { + it('logs and preserves status lock failures', async function () { + const lockError = new Error('Database unavailable'); + const service = new BatchSendingService({}); + sinon.stub(service, 'retryDb').rejects(lockError); + + await assert.rejects(service.emailJob({emailId: '123'}), error => error === lockError); + + sinon.assert.calledOnceWithExactly( + errorLog, + lockError, + sinon.match('[Background Job] batch-sending-service-job failed while acquiring the status lock') + ); + }); + it('does not send if already submitting', async function () { const Email = createModelClass({ findOne: { @@ -64,7 +79,7 @@ describe('Batch Sending Service', function () { const result = await service.emailJob({emailId: '123'}); assert.equal(result, undefined); sinon.assert.calledOnce(errorLog); - sinon.assert.calledWith(errorLog, 'Tried sending email that is not pending or failed 123'); + sinon.assert.calledWith(errorLog, '[Background Job] batch-sending-service-job skipped because email 123 is not pending or failed'); }); it('does not send if already submitted', async function () { @@ -79,7 +94,7 @@ describe('Batch Sending Service', function () { const result = await service.emailJob({emailId: '123'}); assert.equal(result, undefined); sinon.assert.calledOnce(errorLog); - sinon.assert.calledWith(errorLog, 'Tried sending email that is not pending or failed 123'); + sinon.assert.calledWith(errorLog, '[Background Job] batch-sending-service-job skipped because email 123 is not pending or failed'); }); it('does send email if pending', async function () { @@ -111,6 +126,33 @@ describe('Batch Sending Service', function () { assert.equal(afterEmailModel.get('error'), null); }); + it('keeps the email submitted when completion logging fails', async function () { + const Email = createModelClass({ + findOne: { + status: 'pending' + } + }); + const service = new BatchSendingService({ + models: {Email} + }); + let emailModel; + sinon.stub(service, 'sendEmail').callsFake((email) => { + emailModel = email; + return Promise.resolve(); + }); + infoLog.callsFake((message) => { + if (message.startsWith('[Background Job] batch-sending-service-job completed')) { + throw new Error('Logger unavailable'); + } + }); + + await service.emailJob({emailId: '123'}); + + assert.equal(emailModel.get('status'), 'submitted'); + assert.equal(emailModel.get('error'), null); + sinon.assert.notCalled(errorLog); + }); + it('saves error state if sending fails', async function () { const Email = createModelClass({ findOne: { diff --git a/ghost/core/test/unit/server/services/jobs/job-logging.test.ts b/ghost/core/test/unit/server/services/jobs/job-logging.test.ts new file mode 100644 index 00000000000..0bbbf6b7062 --- /dev/null +++ b/ghost/core/test/unit/server/services/jobs/job-logging.test.ts @@ -0,0 +1,30 @@ +import assert from 'node:assert/strict'; +import sinon from 'sinon'; +import logging from '@tryghost/logging'; +import * as jobLogging from '../../../../../core/server/services/jobs/job-logging'; + +describe('Background job logging', function () { + afterEach(function () { + sinon.restore(); + }); + + it('forwards lifecycle logs to the Ghost logger', function () { + const info = sinon.stub(logging, 'info'); + const error = sinon.stub(logging, 'error'); + const failure = new Error('Job failed'); + + jobLogging.info('[Background Job] test-job started'); + jobLogging.error(failure, '[Background Job] test-job failed'); + + sinon.assert.calledOnceWithExactly(info, '[Background Job] test-job started'); + sinon.assert.calledOnceWithExactly(error, failure, '[Background Job] test-job failed'); + }); + + it('does not let synchronous logger failures escape into job execution', function () { + sinon.stub(logging, 'info').throws(new Error('Info logger unavailable')); + sinon.stub(logging, 'error').throws(new Error('Error logger unavailable')); + + assert.doesNotThrow(() => jobLogging.info('[Background Job] test-job started')); + assert.doesNotThrow(() => jobLogging.error(new Error('Job failed'), '[Background Job] test-job failed')); + }); +}); diff --git a/ghost/core/test/unit/server/services/jobs/job-service.test.js b/ghost/core/test/unit/server/services/jobs/job-service.test.js index 473011806a1..e4804951a87 100644 --- a/ghost/core/test/unit/server/services/jobs/job-service.test.js +++ b/ghost/core/test/unit/server/services/jobs/job-service.test.js @@ -4,19 +4,27 @@ const sinon = require('sinon'); describe('JobService', function () { const jobServicePath = '../../../../../core/server/services/jobs/job-service'; + const mentionsJobServicePath = '../../../../../core/server/services/mentions-jobs/job-service'; + const jobLoggingPath = '../../../../../core/server/services/jobs/job-logging'; let originalLoad; let workerMessageHandler; + let workerErrorHandler; let handleModelEvent; + let info; + let errorLog; beforeEach(function () { originalLoad = Module._load; handleModelEvent = sinon.stub().resolves(true); + info = sinon.stub(); + errorLog = sinon.stub(); Module._load = function (request, parent, isMain) { if (request === '@tryghost/job-manager') { return class JobManager { constructor(options) { workerMessageHandler = options.workerMessageHandler; + workerErrorHandler = options.errorHandler; } }; } @@ -35,9 +43,9 @@ describe('JobService', function () { if (request === '@tryghost/logging') { return { - info: sinon.stub(), + info, warn: sinon.stub(), - error: sinon.stub() + error: errorLog }; } @@ -65,12 +73,15 @@ describe('JobService', function () { }; delete require.cache[require.resolve(jobServicePath)]; + delete require.cache[require.resolve(jobLoggingPath)]; require(jobServicePath); }); afterEach(function () { Module._load = originalLoad; delete require.cache[require.resolve(jobServicePath)]; + delete require.cache[require.resolve(mentionsJobServicePath)]; + delete require.cache[require.resolve(jobLoggingPath)]; sinon.restore(); }); @@ -100,6 +111,43 @@ describe('JobService', function () { assert.equal(message.event, undefined); assert.equal(message.eventName, 'member.edited'); }); + + it('adds the common marker to worker messages', function () { + workerMessageHandler({name: 'clean-tokens', message: 'completed'}); + + sinon.assert.calledOnceWithExactly(info, '[Background Job] clean-tokens: completed'); + }); + + it('does not log worker control messages', function () { + workerMessageHandler({name: 'clean-tokens', message: 'done'}); + workerMessageHandler({name: 'clean-tokens', message: 'cancelled'}); + + sinon.assert.notCalled(info); + }); + + it('does not let worker status logging failures escape', function () { + info.throws(new Error('Logger unavailable')); + + assert.doesNotThrow(() => workerMessageHandler({name: 'clean-tokens', message: 'execution started'})); + }); + + it('adds the common marker to worker failures', function () { + const error = new Error('Job failed'); + + workerErrorHandler(error, {name: 'clean-tokens'}); + + sinon.assert.calledOnceWithExactly(errorLog, error, '[Background Job] clean-tokens failed'); + sinon.assert.notCalled(info); + }); + + it('does not let worker failure logging failures escape either job manager', function () { + errorLog.throws(new Error('Logger unavailable')); + + assert.doesNotThrow(() => workerErrorHandler(new Error('Job failed'), {name: 'clean-tokens'})); + + require(mentionsJobServicePath); + assert.doesNotThrow(() => workerErrorHandler(new Error('Job failed'), {name: 'send-webmentions'})); + }); }); describe('JobService model-event bridge wiring', function () { diff --git a/ghost/core/test/unit/server/services/members/jobs/schedule-token-cleanup.test.ts b/ghost/core/test/unit/server/services/members/jobs/schedule-token-cleanup.test.ts index ac070b7f4f5..9a78d7a6bb2 100644 --- a/ghost/core/test/unit/server/services/members/jobs/schedule-token-cleanup.test.ts +++ b/ghost/core/test/unit/server/services/members/jobs/schedule-token-cleanup.test.ts @@ -1,6 +1,7 @@ import assert from 'node:assert/strict'; import sinon from 'sinon'; import {describe, it, beforeEach, afterEach} from 'vitest'; +import logging from '@tryghost/logging'; // require, not import: these must resolve to the same CommonJS module // instances that core/server/services/members/jobs/index.js loads, so the @@ -29,8 +30,9 @@ describe('member jobs: token cleanup scheduling', function () { assert.ok(scheduleStub.notCalled, 'token cleanup must not be scheduled under NODE_ENV=test*'); }); - it('schedules a daily clean-tokens job outside the test environment', async function () { + it('schedules a daily clean-tokens job outside the test environment even when logging fails', async function () { const originalEnv = process.env.NODE_ENV; + sinon.stub(logging, 'info').throws(new Error('Logger unavailable')); process.env.NODE_ENV = 'production'; try { await memberJobs.scheduleTokenCleanupJob(); diff --git a/ghost/core/test/unit/server/services/mentions/service.test.js b/ghost/core/test/unit/server/services/mentions/service.test.js index dae9f22404b..826a1cb3791 100644 --- a/ghost/core/test/unit/server/services/mentions/service.test.js +++ b/ghost/core/test/unit/server/services/mentions/service.test.js @@ -2,7 +2,8 @@ const assert = require('node:assert/strict'); const sinon = require('sinon'); const urlService = require('../../../../../core/server/services/url'); const outputSerializerUrlUtil = require('../../../../../core/server/api/endpoints/utils/serializers/output/utils/url'); -const {getPostData, getPostUrl} = require('../../../../../core/server/services/mentions/service'); +const jobsService = require('../../../../../core/server/services/mentions-jobs'); +const {getPostData, getPostUrl, makeLoggingJobService} = require('../../../../../core/server/services/mentions/service'); describe('Mentions service post url helpers', function () { afterEach(function () { @@ -77,3 +78,58 @@ describe('Mentions service post url helpers', function () { sinon.assert.notCalled(post.load); }); }); + +// Both mentions jobs are queued through this wrapper, so it has to hand the job +// manager the same job it was given: same name, same inline flag, same result, +// same error. The lifecycle logging it adds is deliberately not asserted here: +// that would mean stubbing the shared logger, which is order-dependent under the +// unit project's `isolate: false`. +describe('Mentions service background job wrapper', function () { + let addJob; + + // Runs the job the wrapper handed to the job service, the way the job + // manager runs an inline job. + function runQueuedJob() { + return addJob.firstCall.args[0].job(); + } + + beforeEach(function () { + addJob = sinon.stub(jobsService, 'addJob'); + }); + + afterEach(function () { + sinon.restore(); + }); + + it('queues the job under its own name without running it', async function () { + const fn = sinon.stub().resolves(); + + await makeLoggingJobService().addJob('processWebmention', fn); + + sinon.assert.calledOnce(addJob); + assert.equal(addJob.firstCall.args[0].name, 'processWebmention'); + assert.equal(addJob.firstCall.args[0].offloaded, false); + assert.notEqual(addJob.firstCall.args[0].job, fn, 'the job is wrapped'); + assert.ok(fn.notCalled, 'the job is not run at queue time'); + }); + + it('runs the job once and returns its result untouched', async function () { + const result = {mentions: 1}; + const fn = sinon.stub().resolves(result); + + await makeLoggingJobService().addJob('sendWebmentions', fn); + const returned = await runQueuedJob(); + + sinon.assert.calledOnce(fn); + assert.equal(returned, result, 'the wrapped result is passed through by reference'); + }); + + it('rethrows the original error', async function () { + const failure = new Error('Job failed'); + const fn = sinon.stub().rejects(failure); + + await makeLoggingJobService().addJob('sendWebmentions', fn); + + await assert.rejects(runQueuedJob(), error => error === failure); + }); +});