From ddd76ac98d50945f30987860d7c89c5be1cb7e04 Mon Sep 17 00:00:00 2001 From: Alec Jacobs Date: Fri, 29 Sep 2023 11:06:21 -0700 Subject: [PATCH 1/2] fix(backend-tasks): prevent immediately triggering tasks out of band Signed-off-by: Alec Jacobs --- .changeset/nine-garlics-attend.md | 5 + .../src/tasks/TaskWorker.test.ts | 108 +++++++++++++++++- .../backend-tasks/src/tasks/TaskWorker.ts | 47 ++++---- 3 files changed, 138 insertions(+), 22 deletions(-) create mode 100644 .changeset/nine-garlics-attend.md diff --git a/.changeset/nine-garlics-attend.md b/.changeset/nine-garlics-attend.md new file mode 100644 index 0000000000..a9cdac4a32 --- /dev/null +++ b/.changeset/nine-garlics-attend.md @@ -0,0 +1,5 @@ +--- +'@backstage/backend-tasks': patch +--- + +Fix bug where backend tasks that are defined with HumanDuration are immediately triggered on application startup diff --git a/packages/backend-tasks/src/tasks/TaskWorker.test.ts b/packages/backend-tasks/src/tasks/TaskWorker.test.ts index 3022f59060..4e1e6dc23b 100644 --- a/packages/backend-tasks/src/tasks/TaskWorker.test.ts +++ b/packages/backend-tasks/src/tasks/TaskWorker.test.ts @@ -16,7 +16,7 @@ import { getVoidLogger } from '@backstage/backend-common'; import { TestDatabases } from '@backstage/backend-test-utils'; -import { Duration } from 'luxon'; +import { Duration, DateTime } from 'luxon'; import waitForExpect from 'wait-for-expect'; import { migrateBackendTasks } from '../database/migrateBackendTasks'; import { DbTasksRow, DB_TASKS_TABLE } from '../database/tables'; @@ -338,7 +338,7 @@ describe('TaskWorker', () => { ); it.each(databases.eachSupportedId())( - 'next_run_start_at is always the min between schedule changes, %p', + 'next_run_start_at is always the min between schedule changes from cron frequency, %p', async databaseId => { const knex = await databases.init(databaseId); await migrateBackendTasks(knex); @@ -376,4 +376,108 @@ describe('TaskWorker', () => { await knex.destroy(); }, ); + + it.each(databases.eachSupportedId())( + 'next_run_start_at is always the min between schedule changes when using human duration frequency, %p', + async databaseId => { + const knex = await databases.init(databaseId); + await migrateBackendTasks(knex); + + const fn = jest.fn( + async () => new Promise(resolve => setTimeout(resolve, 50)), + ); + + const initialSettings: TaskSettingsV2 = { + version: 2, + cadence: 'PT120M', + timeoutAfterDuration: 'PT1M', + }; + + const worker = new TaskWorker('task99', fn, knex, logger); + await worker.persistTask(initialSettings); + // replicate task running, sets next_run_start_at based on cadence + await worker.tryClaimTask('ticket', initialSettings); + await worker.tryReleaseTask('ticket', initialSettings); + + // grab initial row for comparisons later + const initialRow = (await knex(DB_TASKS_TABLE))[0]; + + const settings: TaskSettingsV2 = { + ...initialSettings, + cadence: 'PT60M', + }; + await worker.persistTask(settings); + const row1 = (await knex(DB_TASKS_TABLE))[0]; + + expect( + new Date(row1.next_run_start_at) < + new Date(initialRow.next_run_start_at), + ).toBeTruthy(); // ensure that next start at is sooner than initial + expect(new Date(row1.next_run_start_at) > new Date()).toBeTruthy(); // ensure that next start at is later than now + + const settings2 = { + ...settings, + }; + await worker.persistTask(settings2); + const row2 = (await knex(DB_TASKS_TABLE))[0]; + + expect(row2.next_run_start_at).toStrictEqual(row1.next_run_start_at); + + await knex.destroy(); + }, + ); + + it.each(databases.eachSupportedId())( + 'next_run_start_at is always the min between schedule changes when using human duration frequency with initial start delay, %p', + async databaseId => { + const knex = await databases.init(databaseId); + await migrateBackendTasks(knex); + + const fn = jest.fn( + async () => new Promise(resolve => setTimeout(resolve, 50)), + ); + + const initialSettings: TaskSettingsV2 = { + version: 2, + cadence: 'PT120M', + initialDelayDuration: 'PT2M', + timeoutAfterDuration: 'PT1M', + }; + + const worker = new TaskWorker('task99', fn, knex, logger); + await worker.persistTask(initialSettings); + // replicate task running, sets next_run_start_at based on cadence + await worker.tryClaimTask('ticket', initialSettings); + await worker.tryReleaseTask('ticket', initialSettings); + + // grab initial row for comparisons later + const initialRow = (await knex(DB_TASKS_TABLE))[0]; + + const settings: TaskSettingsV2 = { + ...initialSettings, + cadence: 'PT60M', + }; + await worker.persistTask(settings); + const row1 = (await knex(DB_TASKS_TABLE))[0]; + + expect( + new Date(row1.next_run_start_at) < + new Date(initialRow.next_run_start_at), + ).toBeTruthy(); // ensure that next start at is sooner than initial + expect( + new Date(row1.next_run_start_at) > + DateTime.now().plus({ minutes: 3 }).toJSDate(), + ).toBeTruthy(); // ensure that next start at is later than initial delay start time since next run at should be an hour from now + + const settings2 = { + ...settings, + }; + await worker.persistTask(settings2); + const row2 = (await knex(DB_TASKS_TABLE))[0]; + + expect(row2.next_run_start_at).toStrictEqual(row1.next_run_start_at); + + await knex.destroy(); + }, + ); }); diff --git a/packages/backend-tasks/src/tasks/TaskWorker.ts b/packages/backend-tasks/src/tasks/TaskWorker.ts index c1007b4964..6e3a88410b 100644 --- a/packages/backend-tasks/src/tasks/TaskWorker.ts +++ b/packages/backend-tasks/src/tasks/TaskWorker.ts @@ -165,27 +165,26 @@ export class TaskWorker { const isCron = !settings?.cadence.startsWith('P'); - let startAt: Knex.Raw; + let startAt: Knex.Raw | undefined; + let nextStartAt: Knex.Raw | undefined; if (settings.initialDelayDuration) { startAt = nowPlus( Duration.fromISO(settings.initialDelayDuration), this.knex, ); - } else if (isCron) { + } + + if (isCron) { const time = new CronTime(settings.cadence) .sendAt() .minus({ seconds: 1 }) // immediately, if "* * * * * *" .toUTC(); - if (this.knex.client.config.client.includes('sqlite3')) { - startAt = this.knex.raw('datetime(?)', [time.toISO()]); - } else if (this.knex.client.config.client.includes('mysql')) { - startAt = this.knex.raw(`?`, [time.toSQL({ includeOffset: false })]); - } else { - startAt = this.knex.raw(`?`, [time.toISO()]); - } + nextStartAt = this.nextRunAtRaw(time); + startAt ||= nextStartAt; } else { - startAt = this.knex.fn.now(); + startAt ||= this.knex.fn.now(); + nextStartAt = nowPlus(Duration.fromISO(settings.cadence), this.knex); } this.logger.debug(`task: ${this.taskId} configured to run at: ${startAt}`); @@ -206,7 +205,12 @@ export class TaskWorker { settings_json: settingsJson, next_run_start_at: this.knex.raw( `CASE WHEN ?? < ?? THEN ?? ELSE ?? END`, - [startAt, 'next_run_start_at', startAt, 'next_run_start_at'], + [ + nextStartAt, + 'next_run_start_at', + nextStartAt, + 'next_run_start_at', + ], ), } : { @@ -214,9 +218,9 @@ export class TaskWorker { next_run_start_at: this.knex.raw( `CASE WHEN ?? < ?? THEN ?? ELSE ?? END`, [ - 'excluded.next_run_start_at', + nextStartAt, `${DB_TASKS_TABLE}.next_run_start_at`, - 'excluded.next_run_start_at', + nextStartAt, `${DB_TASKS_TABLE}.next_run_start_at`, ], ), @@ -308,13 +312,7 @@ export class TaskWorker { const time = new CronTime(settings.cadence).sendAt().toUTC(); this.logger.debug(`task: ${this.taskId} will next occur around ${time}`); - if (this.knex.client.config.client.includes('sqlite3')) { - nextRun = this.knex.raw('datetime(?)', [time.toISO()]); - } else if (this.knex.client.config.client.includes('mysql')) { - nextRun = this.knex.raw(`?`, [time.toSQL({ includeOffset: false })]); - } else { - nextRun = this.knex.raw(`?`, [time.toISO()]); - } + nextRun = this.nextRunAtRaw(time); } else { const dt = Duration.fromISO(settings.cadence).as('seconds'); this.logger.debug( @@ -351,4 +349,13 @@ export class TaskWorker { return rows === 1; } + + private nextRunAtRaw(time: DateTime): Knex.Raw { + if (this.knex.client.config.client.includes('sqlite3')) { + return this.knex.raw('datetime(?)', [time.toISO()]); + } else if (this.knex.client.config.client.includes('mysql')) { + return this.knex.raw(`?`, [time.toSQL({ includeOffset: false })]); + } + return this.knex.raw(`?`, [time.toISO()]); + } } From 06451c88d176c1e7e9cc7a7d70937c5a4c55089a Mon Sep 17 00:00:00 2001 From: Alec Jacobs Date: Mon, 2 Oct 2023 12:14:41 -0700 Subject: [PATCH 2/2] test(backend-tasks): refine tests to be more explicit Signed-off-by: Alec Jacobs --- .../src/tasks/TaskWorker.test.ts | 43 +++++++++++++------ 1 file changed, 31 insertions(+), 12 deletions(-) diff --git a/packages/backend-tasks/src/tasks/TaskWorker.test.ts b/packages/backend-tasks/src/tasks/TaskWorker.test.ts index 4e1e6dc23b..efad4e649b 100644 --- a/packages/backend-tasks/src/tasks/TaskWorker.test.ts +++ b/packages/backend-tasks/src/tasks/TaskWorker.test.ts @@ -400,7 +400,9 @@ describe('TaskWorker', () => { await worker.tryReleaseTask('ticket', initialSettings); // grab initial row for comparisons later - const initialRow = (await knex(DB_TASKS_TABLE))[0]; + const rowAfterClaimAndRelease = ( + await knex(DB_TASKS_TABLE) + )[0]; const settings: TaskSettingsV2 = { ...initialSettings, @@ -409,11 +411,20 @@ describe('TaskWorker', () => { await worker.persistTask(settings); const row1 = (await knex(DB_TASKS_TABLE))[0]; + const rowAfterClaimAndReleaseNextStartAt = DateTime.fromJSDate( + new Date(rowAfterClaimAndRelease.next_run_start_at), + ); + const row1NextStartAt = DateTime.fromJSDate( + new Date(row1.next_run_start_at), + ); + const now = DateTime.now(); expect( - new Date(row1.next_run_start_at) < - new Date(initialRow.next_run_start_at), - ).toBeTruthy(); // ensure that next start at is sooner than initial - expect(new Date(row1.next_run_start_at) > new Date()).toBeTruthy(); // ensure that next start at is later than now + rowAfterClaimAndReleaseNextStartAt.diff(row1NextStartAt).as('minutes'), + ).toBeCloseTo(60, 1); // ensure that next start at is sooner than initial by one hour + expect(row1NextStartAt.diff(now).as('minutes')).toBeCloseTo(60, 1); // ensure that next start at is later than now by one hour + expect( + rowAfterClaimAndReleaseNextStartAt.diff(now).as('minutes'), + ).toBeCloseTo(120, 1); const settings2 = { ...settings, @@ -451,7 +462,9 @@ describe('TaskWorker', () => { await worker.tryReleaseTask('ticket', initialSettings); // grab initial row for comparisons later - const initialRow = (await knex(DB_TASKS_TABLE))[0]; + const rowAfterClaimAndRelease = ( + await knex(DB_TASKS_TABLE) + )[0]; const settings: TaskSettingsV2 = { ...initialSettings, @@ -460,14 +473,20 @@ describe('TaskWorker', () => { await worker.persistTask(settings); const row1 = (await knex(DB_TASKS_TABLE))[0]; + const rowAfterClaimAndReleaseNextStartAt = DateTime.fromJSDate( + new Date(rowAfterClaimAndRelease.next_run_start_at), + ); + const row1NextStartAt = DateTime.fromJSDate( + new Date(row1.next_run_start_at), + ); + const now = DateTime.now(); expect( - new Date(row1.next_run_start_at) < - new Date(initialRow.next_run_start_at), - ).toBeTruthy(); // ensure that next start at is sooner than initial + rowAfterClaimAndReleaseNextStartAt.diff(row1NextStartAt).as('minutes'), + ).toBeCloseTo(62, 1); // ensure that next start at is sooner than initial by one hour, plus the 2 minute delay (set my tryReleaseTask) + expect(row1NextStartAt.diff(now).as('minutes')).toBeCloseTo(60, 1); // ensure that next start at is later than now by one hour (2 minute delay doesn't take effect here) expect( - new Date(row1.next_run_start_at) > - DateTime.now().plus({ minutes: 3 }).toJSDate(), - ).toBeTruthy(); // ensure that next start at is later than initial delay start time since next run at should be an hour from now + rowAfterClaimAndReleaseNextStartAt.diff(now).as('minutes'), + ).toBeCloseTo(122, 1); // includes 2 minute start delay (which is persisted from tryReleaseTask) const settings2 = { ...settings,