diff --git a/apps/ember-admin/app/components/posts/debug.hbs b/apps/ember-admin/app/components/posts/debug.hbs
index 6d4661a78f6..dc3506cee73 100644
--- a/apps/ember-admin/app/components/posts/debug.hbs
+++ b/apps/ember-admin/app/components/posts/debug.hbs
@@ -401,13 +401,11 @@
Last event time: |
{{ this.analyticsStatus.scheduled.lastEventTimestamp }} |
- {{#unless this.analyticsStatus.scheduled.canceled}}
-
- |
-
- |
-
- {{/unless}}
+
+ |
+
+ |
+
{{else}}
{{#if this.showCustomSchedule}}
diff --git a/ghost/core/core/server/services/email-analytics/email-analytics-service.ts b/ghost/core/core/server/services/email-analytics/email-analytics-service.ts
index c9b779b2bdf..a981a4394de 100644
--- a/ghost/core/core/server/services/email-analytics/email-analytics-service.ts
+++ b/ghost/core/core/server/services/email-analytics/email-analytics-service.ts
@@ -17,8 +17,6 @@ export type FetchData = {
/** The begin time used during the last fetch */
lastBegin?: Date;
lastEventTimestamp?: Date;
- /** Set to quit the job early */
- canceled?: boolean;
};
type FetchDataScheduled = FetchData & {schedule?: {begin: Date; end: Date}};
@@ -143,14 +141,6 @@ export class EmailAnalyticsService {
}
}
- #clearScheduledData() {
- this.#fetchScheduledData = {
- running: false,
- jobName: this.#jobNames.scheduled
- };
- this.queries.setJobMetadata(this.#jobNames.scheduled, null);
- }
-
getStatus() {
return {
latest: this.#fetchLatestNonOpenedData,
@@ -264,19 +254,14 @@ export class EmailAnalyticsService {
/**
* Cancels the scheduled fetch of email analytics events.
- * If a fetch is currently running, it marks it for cancellation.
- * If no fetch is running, it clears the scheduled fetch data.
+ * An in-progress fetch completes its current pass after its schedule is cleared.
*/
cancelScheduled() {
- if (this.#fetchScheduledData) {
- if (this.#fetchScheduledData.running) {
- this.#fetchScheduledData.canceled = true;
- // Clear metadata eagerly; fetchScheduled() will clear in-memory state next cycle
- this.queries.setJobMetadata(this.#jobNames.scheduled, null);
- } else {
- this.#clearScheduledData();
- }
- }
+ this.#fetchScheduledData = {
+ running: false,
+ jobName: this.#jobNames.scheduled
+ };
+ this.queries.setJobMetadata(this.#jobNames.scheduled, null);
}
/**
@@ -322,11 +307,6 @@ export class EmailAnalyticsService {
return createEmptyResult();
}
- if (this.#fetchScheduledData.canceled) {
- this.#clearScheduledData();
- return createEmptyResult();
- }
-
let begin = this.#fetchScheduledData.schedule.begin;
const end = this.#fetchScheduledData.schedule.end;
@@ -337,16 +317,17 @@ export class EmailAnalyticsService {
if (end <= begin) {
logging.info('[EmailAnalytics] Ending fetchScheduled because end is before begin');
- this.#clearScheduledData();
+ this.cancelScheduled();
return createEmptyResult();
}
- const fetchResult = await this.#fetchEventsForJob(this.#fetchScheduledData, {begin, end, maxEvents});
- if (fetchResult.eventCount === 0 || this.#fetchScheduledData.canceled) {
- this.#clearScheduledData();
+ const fetchData = this.#fetchScheduledData;
+ const fetchResult = await this.#fetchEventsForJob(fetchData, {begin, end, maxEvents});
+ if (fetchResult.eventCount === 0 && this.#fetchScheduledData === fetchData) {
+ this.cancelScheduled();
}
- this.queries.setJobTimestamp(this.#fetchScheduledData.jobName, 'finished', this.#fetchScheduledData.lastEventTimestamp!);
+ this.queries.setJobTimestamp(fetchData.jobName, 'finished', fetchData.lastEventTimestamp!);
return fetchResult;
}
/**
@@ -444,24 +425,14 @@ export class EmailAnalyticsService {
logging.error('[EmailAnalytics] Error while aggregating stats');
logging.error(err);
}
-
- if (fetchData.canceled) {
- throw new errors.InternalServerError({
- message: 'Fetching canceled'
- });
- }
};
try {
await this.#fetchEvents({batchHandler: processBatch, begin, end, maxEvents, events: eventTypes});
} catch (err) {
- if (!(err instanceof Error) || err.message !== 'Fetching canceled') {
- logging.error('[EmailAnalytics] Error while fetching');
- logging.error(err);
- error = err;
- } else {
- logging.error('[EmailAnalytics] Canceled fetching');
- }
+ logging.error('[EmailAnalytics] Error while fetching');
+ logging.error(err);
+ error = err;
}
// Final aggregation.
diff --git a/ghost/core/test/unit/server/services/email-analytics/email-analytics-service.test.js b/ghost/core/test/unit/server/services/email-analytics/email-analytics-service.test.js
index 3ae86d7b0d9..fffb865ee55 100644
--- a/ghost/core/test/unit/server/services/email-analytics/email-analytics-service.test.js
+++ b/ghost/core/test/unit/server/services/email-analytics/email-analytics-service.test.js
@@ -385,7 +385,7 @@ describe('EmailAnalyticsService', function () {
sinon.assert.notCalled(eventProcessor.aggregate);
});
- it('returns 0 when fetch is canceled', async function () {
+ it('returns 0 when scheduled fetch is canceled before it starts', async function () {
await service.schedule({
begin: new Date(2023, 0, 1),
end: new Date(2023, 0, 2)
@@ -411,6 +411,51 @@ describe('EmailAnalyticsService', function () {
assert.deepEqual(eventProcessor.processBatch.getCall(0).args[0], [1,2,3,4,5,6,7,8,9,10]);
});
+ it('finishes an in-progress fetch when its schedule is canceled', async function () {
+ let startFetch;
+ const fetchStarted = new Promise((resolve) => {
+ startFetch = resolve;
+ });
+ let continueFetch;
+ const fetchCanContinue = new Promise((resolve) => {
+ continueFetch = resolve;
+ });
+
+ service = createService({
+ queries: {
+ setJobTimestamp: setJobTimestampStub,
+ setJobStatus: setJobStatusStub,
+ setJobMetadata: setJobMetadataStub
+ },
+ fetchEvents: async ({batchHandler}) => {
+ startFetch();
+ await fetchCanContinue;
+ await batchHandler([1,2,3,4,5,6,7,8,9,10]);
+ },
+ createEventProcessor: () => eventProcessor
+ });
+ await service.schedule({
+ begin: new Date(2023, 0, 1),
+ end: new Date(2023, 0, 2)
+ });
+ setJobMetadataStub.resetHistory();
+
+ const fetch = service.fetchScheduled({maxEvents: 100});
+ await fetchStarted;
+ service.cancelScheduled();
+
+ sinon.assert.calledOnceWithExactly(setJobMetadataStub, 'email-analytics-scheduled', null);
+ assert.equal(service.getStatus().scheduled.running, false);
+
+ continueFetch();
+ const result = await fetch;
+
+ assert.equal(result.eventCount, 10);
+ sinon.assert.calledOnce(eventProcessor.processBatch);
+ assert.equal(service.getStatus().scheduled.running, false);
+ assert.equal(service.getStatus().scheduled.schedule, undefined);
+ });
+
it('bails when end date is before begin date', async function () {
await service.schedule({
begin: new Date(2023, 0, 2),