Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 24 additions & 5 deletions ghost/core/core/server/data/importer/import-manager.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down Expand Up @@ -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
});
}
Expand All @@ -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 {
Expand Down
17 changes: 17 additions & 0 deletions ghost/core/core/server/services/content-import/import/importer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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,
Expand All @@ -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<void> {
const startedAt = Date.now();
logLifecycle('started');
let urlFailureCount = 0;
let firstUrlFailure: unknown;
let failed = false;

try {
const htmlToLexical = this._getHtmlToLexical();
Expand Down Expand Up @@ -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`);
}
}

Expand Down
2 changes: 1 addition & 1 deletion ghost/core/core/server/services/content-import/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<ConfigInstance, 'get'>;
Expand All @@ -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;
}
Expand Down Expand Up @@ -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') {
Expand Down Expand Up @@ -199,13 +214,19 @@ export class EmailAnalyticsServiceWrapper {
}

async startFetch(): Promise<void> {
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;
Expand Down Expand Up @@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
});
Original file line number Diff line number Diff line change
@@ -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<string | number>;
Expand Down Expand Up @@ -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'
});
Expand Down Expand Up @@ -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'
});
Expand All @@ -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'
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(() => {
Expand All @@ -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;
}
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
});
Original file line number Diff line number Diff line change
Expand Up @@ -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
});
Original file line number Diff line number Diff line change
@@ -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');
Expand Down Expand Up @@ -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),
Expand All @@ -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;
}

Expand All @@ -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({
Expand All @@ -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);
Expand Down
Loading
Loading