From 015668c5d2964c20258ae478b0b5546883d6dc90 Mon Sep 17 00:00:00 2001 From: Fredrik Date: Mon, 9 Mar 2026 21:54:22 +0100 Subject: [PATCH 1/7] Add cancelTask to SchedulerService for cancelling running tasks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds the ability to cancel currently running scheduled tasks via a new cancelTask method on the SchedulerService interface. For global (distributed) tasks, the database lock is released and a periodic liveness check detects the lost ticket and aborts the task function's AbortSignal. For local tasks, the abort signal is triggered directly. Also adds a REST endpoint at POST /.backstage/scheduler/v1/tasks/:id/cancel. Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .changeset/add-scheduler-cancel-task.md | 6 ++ .../scheduler/lib/LocalTaskWorker.test.ts | 49 ++++++++++ .../scheduler/lib/LocalTaskWorker.ts | 17 +++- .../lib/PluginTaskSchedulerImpl.test.ts | 96 +++++++++++++++++++ .../scheduler/lib/PluginTaskSchedulerImpl.ts | 20 ++++ .../scheduler/lib/TaskWorker.test.ts | 76 +++++++++++++++ .../entrypoints/scheduler/lib/TaskWorker.ts | 60 +++++++++++- .../services/definitions/SchedulerService.ts | 10 ++ 8 files changed, 329 insertions(+), 5 deletions(-) create mode 100644 .changeset/add-scheduler-cancel-task.md diff --git a/.changeset/add-scheduler-cancel-task.md b/.changeset/add-scheduler-cancel-task.md new file mode 100644 index 0000000000..1f98bccb1c --- /dev/null +++ b/.changeset/add-scheduler-cancel-task.md @@ -0,0 +1,6 @@ +--- +'@backstage/backend-plugin-api': minor +'@backstage/backend-defaults': patch +--- + +Added `cancelTask` method to the `SchedulerService` interface and implementation, allowing cancellation of currently running scheduled tasks. For global tasks, the database lock is released and a periodic liveness check aborts the running task function. For local tasks, the task's abort signal is triggered directly. A new `POST /.backstage/scheduler/v1/tasks/:id/cancel` endpoint is also available. diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.test.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.test.ts index 14fedfb3c0..04c2b17ad3 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.test.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.test.ts @@ -16,6 +16,7 @@ import { LocalTaskWorker } from './LocalTaskWorker'; import { mockServices } from '@backstage/backend-test-utils'; +import { ConflictError } from '@backstage/errors'; import waitFor from 'wait-for-expect'; jest.setTimeout(10_000); @@ -110,6 +111,54 @@ describe('LocalTaskWorker', () => { controller.abort(); }); + it('can cancel a running task', async () => { + let receivedSignal: AbortSignal | undefined; + const fn = jest.fn(async (signal: AbortSignal) => { + receivedSignal = signal; + await new Promise(r => setTimeout(r, 5000)); + }); + const controller = new AbortController(); + + const worker = new LocalTaskWorker('a', fn, logger); + worker.start( + { + version: 2, + cadence: 'PT10S', + timeoutAfterDuration: 'PT10S', + }, + { signal: controller.signal }, + ); + + await waitFor(() => { + expect(fn).toHaveBeenCalledTimes(1); + }); + + expect(receivedSignal?.aborted).toBe(false); + worker.cancel(); + expect(receivedSignal?.aborted).toBe(true); + + controller.abort(); + }); + + it('cannot cancel a task that is not running', async () => { + const fn = jest.fn(); + const controller = new AbortController(); + + const worker = new LocalTaskWorker('a', fn, logger); + worker.start( + { + version: 2, + initialDelayDuration: 'PT1000S', + cadence: 'PT10S', + timeoutAfterDuration: 'PT10S', + }, + { signal: controller.signal }, + ); + + expect(() => worker.cancel()).toThrow(ConflictError); + controller.abort(); + }); + it('goes through the expected states', async () => { const fn = jest .fn() diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts index efecb1c0ed..168649436c 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/LocalTaskWorker.ts @@ -29,6 +29,7 @@ import { delegateAbortController, serializeError, sleep } from './util'; */ export class LocalTaskWorker { private abortWait: AbortController | undefined; + private taskAbortController: AbortController | undefined; #taskState: Exclude = { status: 'idle', }; @@ -93,6 +94,13 @@ export class LocalTaskWorker { this.abortWait.abort(); } + cancel(): void { + if (!this.taskAbortController) { + throw new ConflictError(`Task ${this.taskId} is not running`); + } + this.taskAbortController.abort(); + } + taskState(): TaskApiTasksResponse['taskState'] { return this.#taskState; } @@ -134,10 +142,10 @@ export class LocalTaskWorker { ): Promise { // Abort the task execution either if the worker is stopped, or if the // task timeout is hit - const taskAbortController = delegateAbortController(signal); + this.taskAbortController = delegateAbortController(signal); const timeoutDuration = Duration.fromISO(settings.timeoutAfterDuration); const timeoutHandle = setTimeout(() => { - taskAbortController.abort(); + this.taskAbortController?.abort(); }, timeoutDuration.as('milliseconds')); this.#taskState = { @@ -152,7 +160,7 @@ export class LocalTaskWorker { }; try { - await this.fn(taskAbortController.signal); + await this.fn(this.taskAbortController.signal); this.#taskState.lastRunEndedAt = DateTime.utc().toISO()!; this.#taskState.lastRunError = undefined; } catch (e) { @@ -162,7 +170,8 @@ export class LocalTaskWorker { // release resources clearTimeout(timeoutHandle); - taskAbortController.abort(); + this.taskAbortController.abort(); + this.taskAbortController = undefined; } /** diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts index 2148148c9e..7921e2863b 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts @@ -399,6 +399,102 @@ describe('PluginTaskManagerImpl', () => { ); }); + describe('cancelTask with local scope', () => { + it('can cancel a running task', async () => { + const { manager } = await init('SQLITE_3'); + + const promise = createDeferred(); + + await manager.scheduleTask({ + id: 'task1', + timeout: Duration.fromMillis(5000), + frequency: Duration.fromObject({ years: 1 }), + fn: async () => { + promise.resolve(); + await new Promise(r => setTimeout(r, 20000)); + }, + scope: 'local', + }); + + await promise; + await expect(manager.cancelTask('task1')).resolves.toBeUndefined(); + }, 60_000); + + it('cannot cancel a task that is not running', async () => { + const { manager } = await init('SQLITE_3'); + + const fn = jest.fn(); + await manager.scheduleTask({ + id: 'task1', + timeout: Duration.fromMillis(5000), + frequency: Duration.fromObject({ years: 1 }), + initialDelay: Duration.fromObject({ years: 1 }), + fn, + scope: 'local', + }); + + await expect(manager.cancelTask('task1')).rejects.toThrow(ConflictError); + }, 60_000); + }); + + describe('cancelTask with global scope', () => { + it.each(databases.eachSupportedId())( + 'can cancel a running task, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const promise = createDeferred(); + + await manager.scheduleTask({ + id: 'task1', + timeout: Duration.fromMillis(5000), + frequency: Duration.fromObject({ years: 1 }), + fn: async () => { + promise.resolve(); + await new Promise(r => setTimeout(r, 20000)); + }, + scope: 'global', + }); + + await promise; + await expect(manager.cancelTask('task1')).resolves.toBeUndefined(); + }, + ); + + it.each(databases.eachSupportedId())( + 'cannot cancel a non-existent task, %p', + async databaseId => { + const { manager } = await init(databaseId); + + await expect(manager.cancelTask('nonexistent')).rejects.toThrow( + NotFoundError, + ); + }, + ); + + it.each(databases.eachSupportedId())( + 'cannot cancel a task that is not running, %p', + async databaseId => { + const { manager } = await init(databaseId); + + const fn = jest.fn(); + const promise = new Promise(resolve => fn.mockImplementation(resolve)); + await manager.scheduleTask({ + id: 'task1', + timeout: Duration.fromMillis(5000), + frequency: Duration.fromObject({ years: 1 }), + initialDelay: Duration.fromObject({ years: 1 }), + fn, + scope: 'global', + }); + + await expect(manager.cancelTask('task1')).rejects.toThrow( + ConflictError, + ); + }, + ); + }); + describe('parseDuration', () => { it('should parse durations', () => { expect(parseDuration({ milliseconds: 5000 })).toEqual('PT5S'); diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts index e711744123..2aad8647fe 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.ts @@ -107,6 +107,17 @@ export class PluginTaskSchedulerImpl implements SchedulerService { await TaskWorker.trigger(knex, id); } + async cancelTask(id: string): Promise { + const localTask = this.localWorkersById.get(id); + if (localTask) { + localTask.cancel(); + return; + } + + const knex = await this.databaseFactory(); + await TaskWorker.cancel(knex, id); + } + async scheduleTask( task: SchedulerServiceTaskScheduleDefinition & SchedulerServiceTaskInvocationDefinition, @@ -206,6 +217,15 @@ export class PluginTaskSchedulerImpl implements SchedulerService { }, ); + router.post( + '/.backstage/scheduler/v1/tasks/:id/cancel', + async (req, res) => { + const { id } = req.params; + await this.cancelTask(id); + res.status(200).end(); + }, + ); + return router; } diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.test.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.test.ts index 5bd3fb38cb..3446d635d1 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.test.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.test.ts @@ -15,6 +15,7 @@ */ import { TestDatabases, mockServices } from '@backstage/backend-test-utils'; +import { ConflictError, NotFoundError } from '@backstage/errors'; import { DateTime, Duration } from 'luxon'; import waitForExpect from 'wait-for-expect'; import { migrateBackendTasks } from '../database/migrateBackendTasks'; @@ -584,4 +585,79 @@ describe('TaskWorker', () => { await knex.destroy(); }, ); + + it.each(databases.eachSupportedId())( + 'can cancel a running task, %p', + async databaseId => { + const knex = await databases.init(databaseId); + await migrateBackendTasks(knex); + + const fn = jest.fn(async () => {}); + const settings: TaskSettingsV2 = { + version: 2, + cadence: '* * * * * *', + initialDelayDuration: undefined, + timeoutAfterDuration: Duration.fromObject({ minutes: 1 }).toISO()!, + }; + + const worker = new TaskWorker('task1', fn, knex, logger); + await worker.persistTask(settings); + await worker.tryClaimTask('ticket', settings); + + // Verify the task is running + let row = (await knex(DB_TASKS_TABLE))[0]; + expect(row.current_run_ticket).toBe('ticket'); + + await TaskWorker.cancel(knex, 'task1'); + + // Verify the task is now idle with a cancellation error recorded + row = (await knex(DB_TASKS_TABLE))[0]; + expect(row.current_run_ticket).toBeNull(); + expect(row.current_run_started_at).toBeNull(); + expect(row.current_run_expires_at).toBeNull(); + expect(row.last_run_ended_at).not.toBeNull(); + expect(row.last_run_error_json).toContain('Task was cancelled'); + + await knex.destroy(); + }, + ); + + it.each(databases.eachSupportedId())( + 'cannot cancel a non-existent task, %p', + async databaseId => { + const knex = await databases.init(databaseId); + await migrateBackendTasks(knex); + + await expect(TaskWorker.cancel(knex, 'nonexistent')).rejects.toThrow( + NotFoundError, + ); + + await knex.destroy(); + }, + ); + + it.each(databases.eachSupportedId())( + 'cannot cancel a task that is not running, %p', + async databaseId => { + const knex = await databases.init(databaseId); + await migrateBackendTasks(knex); + + const fn = jest.fn(async () => {}); + const settings: TaskSettingsV2 = { + version: 2, + cadence: '* * * * * *', + initialDelayDuration: undefined, + timeoutAfterDuration: Duration.fromObject({ minutes: 1 }).toISO()!, + }; + + const worker = new TaskWorker('task1', fn, knex, logger); + await worker.persistTask(settings); + + await expect(TaskWorker.cancel(knex, 'task1')).rejects.toThrow( + ConflictError, + ); + + await knex.destroy(); + }, + ); }); diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts index e47967de90..121f844351 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts @@ -152,6 +152,31 @@ export class TaskWorker { } } + static async cancel(knex: Knex, taskId: string): Promise { + // check if task exists + const rows = await knex(DB_TASKS_TABLE) + .select(knex.raw(1)) + .where('id', '=', taskId); + if (rows.length !== 1) { + throw new NotFoundError(`Task ${taskId} does not exist`); + } + + const dbNull = knex.raw('null'); + const updatedRows = await knex(DB_TASKS_TABLE) + .where('id', '=', taskId) + .whereNotNull('current_run_ticket') + .update({ + current_run_ticket: dbNull, + current_run_started_at: dbNull, + current_run_expires_at: dbNull, + last_run_ended_at: knex.fn.now(), + last_run_error_json: serializeError(new Error('Task was cancelled')), + }); + if (updatedRows < 1) { + throw new ConflictError(`Task ${taskId} is not running`); + } + } + static async taskStates( knex: Knex, ): Promise> { @@ -227,11 +252,16 @@ export class TaskWorker { } // Abort the task execution either if the worker is stopped, or if the - // task timeout is hit + // task timeout is hit, or if the task ticket was lost (e.g. due to + // cancellation from another host) const taskAbortController = delegateAbortController(signal); const timeoutHandle = setTimeout(() => { taskAbortController.abort(); }, Duration.fromISO(taskSettings.timeoutAfterDuration).as('milliseconds')); + const livenessHandle = setInterval( + () => this.checkLiveness(ticket, taskAbortController), + this.workCheckFrequency.as('milliseconds'), + ); try { this.#workerState = { @@ -248,6 +278,7 @@ export class TaskWorker { status: 'idle', }; clearTimeout(timeoutHandle); + clearInterval(livenessHandle); } await this.tryReleaseTask(ticket, taskSettings); @@ -334,6 +365,33 @@ export class TaskWorker { ); } + /** + * Checks whether the current task ticket is still valid in the database. + * If the ticket has been cleared (e.g. by cancellation or janitor cleanup), + * aborts the task execution. + */ + private async checkLiveness( + ticket: string, + taskAbortController: AbortController, + ): Promise { + try { + const [row] = await this.knex(DB_TASKS_TABLE) + .where('id', '=', this.taskId) + .select('current_run_ticket'); + + if (!row || row.current_run_ticket !== ticket) { + this.logger.info( + `Task ticket for "${this.taskId}" is no longer valid; aborting execution`, + ); + taskAbortController.abort(); + } + } catch (e) { + this.logger.warn( + `Failed to check liveness for task "${this.taskId}", ${e}`, + ); + } + } + /** * Check if the task is ready to run */ diff --git a/packages/backend-plugin-api/src/services/definitions/SchedulerService.ts b/packages/backend-plugin-api/src/services/definitions/SchedulerService.ts index b1a0e17df8..2748b8fac7 100644 --- a/packages/backend-plugin-api/src/services/definitions/SchedulerService.ts +++ b/packages/backend-plugin-api/src/services/definitions/SchedulerService.ts @@ -304,6 +304,16 @@ export interface SchedulerService { */ triggerTask(id: string): Promise; + /** + * Cancels a currently running task by ID, marking it as idle. + * + * If the task doesn't exist, a NotFoundError is thrown. If the task is + * not currently running, a ConflictError is thrown. + * + * @param id - The task ID + */ + cancelTask(id: string): Promise; + /** * Schedules a task function for recurring runs. * From 60eec596016a177edf50e3eb6301783a832631c8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Mon, 9 Mar 2026 22:16:33 +0100 Subject: [PATCH 2/7] Fix missing cancelTask in mock SchedulerService implementations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .github/vale/config/vocabularies/Backstage/accept.txt | 1 + .../entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts | 4 +--- .../backend-test-utils/src/services/MockSchedulerService.ts | 4 ++++ packages/backend-test-utils/src/services/mockServices.ts | 1 + .../src/providers/GiteaEntityProvider.test.ts | 1 + 5 files changed, 8 insertions(+), 3 deletions(-) diff --git a/.github/vale/config/vocabularies/Backstage/accept.txt b/.github/vale/config/vocabularies/Backstage/accept.txt index 11df7fc1a0..fe8306afa8 100644 --- a/.github/vale/config/vocabularies/Backstage/accept.txt +++ b/.github/vale/config/vocabularies/Backstage/accept.txt @@ -246,6 +246,7 @@ Levenshtein lightbox Lightsail limitranges +liveness LocalStack lockdown lockfile diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts index 7921e2863b..cc8aac5f51 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/PluginTaskSchedulerImpl.test.ts @@ -477,14 +477,12 @@ describe('PluginTaskManagerImpl', () => { async databaseId => { const { manager } = await init(databaseId); - const fn = jest.fn(); - const promise = new Promise(resolve => fn.mockImplementation(resolve)); await manager.scheduleTask({ id: 'task1', timeout: Duration.fromMillis(5000), frequency: Duration.fromObject({ years: 1 }), initialDelay: Duration.fromObject({ years: 1 }), - fn, + fn: jest.fn(), scope: 'global', }); diff --git a/packages/backend-test-utils/src/services/MockSchedulerService.ts b/packages/backend-test-utils/src/services/MockSchedulerService.ts index aafd6d98ec..e56ef11e20 100644 --- a/packages/backend-test-utils/src/services/MockSchedulerService.ts +++ b/packages/backend-test-utils/src/services/MockSchedulerService.ts @@ -95,6 +95,10 @@ export class MockSchedulerService implements SchedulerService { }); } + async cancelTask(_id: string): Promise { + // No-op in mock + } + async triggerTask(id: string): Promise { const task = this.#tasks.get(id); if (!task) { diff --git a/packages/backend-test-utils/src/services/mockServices.ts b/packages/backend-test-utils/src/services/mockServices.ts index 84783f04d7..4d54269b98 100644 --- a/packages/backend-test-utils/src/services/mockServices.ts +++ b/packages/backend-test-utils/src/services/mockServices.ts @@ -526,6 +526,7 @@ export namespace mockServices { getScheduledTasks: jest.fn(), scheduleTask: jest.fn(), triggerTask: jest.fn(), + cancelTask: jest.fn(), })); } diff --git a/plugins/catalog-backend-module-gitea/src/providers/GiteaEntityProvider.test.ts b/plugins/catalog-backend-module-gitea/src/providers/GiteaEntityProvider.test.ts index d6a31899f1..539d432793 100644 --- a/plugins/catalog-backend-module-gitea/src/providers/GiteaEntityProvider.test.ts +++ b/plugins/catalog-backend-module-gitea/src/providers/GiteaEntityProvider.test.ts @@ -34,6 +34,7 @@ describe('GiteaEntityProvider', () => { triggerTask: jest.fn(), scheduleTask: jest.fn(), getScheduledTasks: jest.fn(), + cancelTask: jest.fn(), }; const mockTaskRunner = { run: jest.fn(), From 164711a61fff9382a75dad7d4d5cbb77e783ad45 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Mon, 9 Mar 2026 22:18:23 +0100 Subject: [PATCH 3/7] Add changeset for backend-test-utils cancelTask addition MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .changeset/add-scheduler-cancel-task-test-utils.md | 5 +++++ 1 file changed, 5 insertions(+) create mode 100644 .changeset/add-scheduler-cancel-task-test-utils.md diff --git a/.changeset/add-scheduler-cancel-task-test-utils.md b/.changeset/add-scheduler-cancel-task-test-utils.md new file mode 100644 index 0000000000..242b103eef --- /dev/null +++ b/.changeset/add-scheduler-cancel-task-test-utils.md @@ -0,0 +1,5 @@ +--- +'@backstage/backend-test-utils': patch +--- + +Added `cancelTask` to `MockSchedulerService` and mock scheduler service factory. From 3e4b37295595a830672b5117a67436e3eef2ad8d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Mon, 9 Mar 2026 22:37:08 +0100 Subject: [PATCH 4/7] Implement cancelTask in MockSchedulerService with proper error types MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- packages/backend-plugin-api/report.api.md | 1 + .../src/services/MockSchedulerService.ts | 14 +++++++++++--- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/packages/backend-plugin-api/report.api.md b/packages/backend-plugin-api/report.api.md index 504be97ec6..96d1a02c0a 100644 --- a/packages/backend-plugin-api/report.api.md +++ b/packages/backend-plugin-api/report.api.md @@ -643,6 +643,7 @@ export interface RootServiceFactoryOptions< // @public export interface SchedulerService { + cancelTask(id: string): Promise; createScheduledTaskRunner( schedule: SchedulerServiceTaskScheduleDefinition, ): SchedulerServiceTaskRunner; diff --git a/packages/backend-test-utils/src/services/MockSchedulerService.ts b/packages/backend-test-utils/src/services/MockSchedulerService.ts index e56ef11e20..a47859a96f 100644 --- a/packages/backend-test-utils/src/services/MockSchedulerService.ts +++ b/packages/backend-test-utils/src/services/MockSchedulerService.ts @@ -23,6 +23,7 @@ import { SchedulerServiceTaskRunner, SchedulerServiceTaskScheduleDefinition, } from '@backstage/backend-plugin-api'; +import { ConflictError, NotFoundError } from '@backstage/errors'; import { createDeferred, DeferredPromise } from '@backstage/types'; export class MockSchedulerService implements SchedulerService { @@ -95,14 +96,21 @@ export class MockSchedulerService implements SchedulerService { }); } - async cancelTask(_id: string): Promise { - // No-op in mock + async cancelTask(id: string): Promise { + const task = this.#tasks.get(id); + if (!task) { + throw new NotFoundError(`Task ${id} not found`); + } + if (!this.#runningTasks.has(id)) { + throw new ConflictError(`Task ${id} is not running`); + } + task.abortControllers.abort(); } async triggerTask(id: string): Promise { const task = this.#tasks.get(id); if (!task) { - throw new Error(`Task ${id} not found`); + throw new NotFoundError(`Task ${id} not found`); } if (this.#runningTasks.has(id)) { return; From 0899cb2cee96c27c4a0b21f3abf427bc7eaba485 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Tue, 10 Mar 2026 11:47:51 +0100 Subject: [PATCH 5/7] Fix cancelTask to reschedule next run instead of permanently stopping task MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The static cancel() now reads task settings from the DB and computes the next_run_start_at, so cancelled tasks get picked up again on their normal schedule. Also stops the liveness check polling immediately on first detection of a cancelled task. Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .../entrypoints/scheduler/lib/TaskWorker.ts | 124 ++++++++++-------- 1 file changed, 67 insertions(+), 57 deletions(-) diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts index 121f844351..e33e112f43 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts @@ -153,22 +153,27 @@ export class TaskWorker { } static async cancel(knex: Knex, taskId: string): Promise { - // check if task exists - const rows = await knex(DB_TASKS_TABLE) - .select(knex.raw(1)) - .where('id', '=', taskId); - if (rows.length !== 1) { + const [row] = await knex(DB_TASKS_TABLE) + .where('id', '=', taskId) + .select('settings_json', 'current_run_ticket'); + if (!row) { throw new NotFoundError(`Task ${taskId} does not exist`); } + if (!row.current_run_ticket) { + throw new ConflictError(`Task ${taskId} is not running`); + } + + const settings = taskSettingsV2Schema.parse(JSON.parse(row.settings_json)); + const nextRun = TaskWorker.computeNextRunStartAt(knex, settings); - const dbNull = knex.raw('null'); const updatedRows = await knex(DB_TASKS_TABLE) .where('id', '=', taskId) - .whereNotNull('current_run_ticket') + .where('current_run_ticket', '=', row.current_run_ticket) .update({ - current_run_ticket: dbNull, - current_run_started_at: dbNull, - current_run_expires_at: dbNull, + next_run_start_at: nextRun, + current_run_ticket: knex.raw('null'), + current_run_started_at: knex.raw('null'), + current_run_expires_at: knex.raw('null'), last_run_ended_at: knex.fn.now(), last_run_error_json: serializeError(new Error('Task was cancelled')), }); @@ -258,8 +263,9 @@ export class TaskWorker { const timeoutHandle = setTimeout(() => { taskAbortController.abort(); }, Duration.fromISO(taskSettings.timeoutAfterDuration).as('milliseconds')); - const livenessHandle = setInterval( - () => this.checkLiveness(ticket, taskAbortController), + const livenessHandle: { ref?: ReturnType } = {}; + livenessHandle.ref = setInterval( + () => this.checkLiveness(ticket, taskAbortController, livenessHandle), this.workCheckFrequency.as('milliseconds'), ); @@ -278,7 +284,7 @@ export class TaskWorker { status: 'idle', }; clearTimeout(timeoutHandle); - clearInterval(livenessHandle); + clearInterval(livenessHandle.ref); } await this.tryReleaseTask(ticket, taskSettings); @@ -314,7 +320,7 @@ export class TaskWorker { // We make a conversion here to make typescript happy, because the luxon versions of the cron library and here may not be the same const timeConverted = DateTime.fromJSDate(time.toJSDate()); - nextStartAt = this.nextRunAtRaw(timeConverted); + nextStartAt = TaskWorker.nextRunAtRaw(this.knex, timeConverted); startAt ||= nextStartAt; } else if (isManual) { nextStartAt = this.knex.raw('null'); @@ -373,6 +379,7 @@ export class TaskWorker { private async checkLiveness( ticket: string, taskAbortController: AbortController, + livenessHandle: { ref?: ReturnType }, ): Promise { try { const [row] = await this.knex(DB_TASKS_TABLE) @@ -383,6 +390,7 @@ export class TaskWorker { this.logger.info( `Task ticket for "${this.taskId}" is no longer valid; aborting execution`, ); + clearInterval(livenessHandle.ref); taskAbortController.abort(); } } catch (e) { @@ -465,48 +473,49 @@ export class TaskWorker { return rows === 1; } + private static computeNextRunStartAt( + knex: Knex, + settings: TaskSettingsV2, + ): Knex.Raw { + const isManual = settings?.cadence === 'manual'; + const isDuration = settings?.cadence.startsWith('P'); + const isCron = !isManual && !isDuration; + + if (isCron) { + const time = new CronTime(settings.cadence).sendAt().toUTC(); + const timeConverted = DateTime.fromJSDate(time.toJSDate()); + return TaskWorker.nextRunAtRaw(knex, timeConverted); + } + + if (isManual) { + return knex.raw('null'); + } + + const dt = Duration.fromISO(settings.cadence).as('seconds'); + + if (knex.client.config.client.includes('sqlite3')) { + return knex.raw(`max(datetime(next_run_start_at, ?), datetime('now'))`, [ + `+${dt} seconds`, + ]); + } + + if (knex.client.config.client.includes('mysql')) { + return knex.raw( + `greatest(next_run_start_at + interval ${dt} second, now())`, + ); + } + + return knex.raw( + `greatest(next_run_start_at + interval '${dt} seconds', now())`, + ); + } + async tryReleaseTask( ticket: string, settings: TaskSettingsV2, error?: Error, ): Promise { - const isManual = settings?.cadence === 'manual'; - const isDuration = settings?.cadence.startsWith('P'); - const isCron = !isManual && !isDuration; - - let nextRun: Knex.Raw; - if (isCron) { - const time = new CronTime(settings.cadence).sendAt().toUTC(); - this.logger.debug(`task: ${this.taskId} will next occur around ${time}`); - // We make a conversion here to make typescript happy, because the luxon versions of the cron library and here may not be the same - const timeConverted = DateTime.fromJSDate(time.toJSDate()); - - nextRun = this.nextRunAtRaw(timeConverted); - } else if (isManual) { - nextRun = this.knex.raw('null'); - } else { - const dt = Duration.fromISO(settings.cadence).as('seconds'); - this.logger.debug( - `task: ${this.taskId} will next occur around ${DateTime.now().plus({ - seconds: dt, - })}`, - ); - - if (this.knex.client.config.client.includes('sqlite3')) { - nextRun = this.knex.raw( - `max(datetime(next_run_start_at, ?), datetime('now'))`, - [`+${dt} seconds`], - ); - } else if (this.knex.client.config.client.includes('mysql')) { - nextRun = this.knex.raw( - `greatest(next_run_start_at + interval ${dt} second, now())`, - ); - } else { - nextRun = this.knex.raw( - `greatest(next_run_start_at + interval '${dt} seconds', now())`, - ); - } - } + const nextRun = TaskWorker.computeNextRunStartAt(this.knex, settings); const rows = await this.knex(DB_TASKS_TABLE) .where('id', '=', this.taskId) @@ -525,12 +534,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 })]); + private static nextRunAtRaw(knex: Knex, time: DateTime): Knex.Raw { + if (knex.client.config.client.includes('sqlite3')) { + return knex.raw('datetime(?)', [time.toISO()]); } - return this.knex.raw(`?`, [time.toISO()]); + if (knex.client.config.client.includes('mysql')) { + return knex.raw(`?`, [time.toSQL({ includeOffset: false })]); + } + return knex.raw(`?`, [time.toISO()]); } } From 25157bec8d4d4fdc81ad8ff923f4b8dfdc727318 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Tue, 10 Mar 2026 13:49:38 +0100 Subject: [PATCH 6/7] Fix AbortController reuse after cancel and prevent overlapping liveness checks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reset AbortController in MockSchedulerService after cancelTask so subsequent triggerTask calls receive a fresh signal. Replace setInterval with a self-scheduling setTimeout loop for liveness checks in TaskWorker to prevent overlapping DB queries when a check takes longer than the polling interval. Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .../entrypoints/scheduler/lib/TaskWorker.ts | 19 +++--- .../src/services/MockSchedulerService.test.ts | 62 +++++++++++++++++++ .../src/services/MockSchedulerService.ts | 1 + 3 files changed, 74 insertions(+), 8 deletions(-) diff --git a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts index e33e112f43..4f32dc6863 100644 --- a/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts +++ b/packages/backend-defaults/src/entrypoints/scheduler/lib/TaskWorker.ts @@ -263,11 +263,16 @@ export class TaskWorker { const timeoutHandle = setTimeout(() => { taskAbortController.abort(); }, Duration.fromISO(taskSettings.timeoutAfterDuration).as('milliseconds')); - const livenessHandle: { ref?: ReturnType } = {}; - livenessHandle.ref = setInterval( - () => this.checkLiveness(ticket, taskAbortController, livenessHandle), - this.workCheckFrequency.as('milliseconds'), - ); + let livenessHandle: ReturnType | undefined; + const scheduleLivenessCheck = () => { + livenessHandle = setTimeout(async () => { + await this.checkLiveness(ticket, taskAbortController); + if (!taskAbortController.signal.aborted) { + scheduleLivenessCheck(); + } + }, this.workCheckFrequency.as('milliseconds')); + }; + scheduleLivenessCheck(); try { this.#workerState = { @@ -284,7 +289,7 @@ export class TaskWorker { status: 'idle', }; clearTimeout(timeoutHandle); - clearInterval(livenessHandle.ref); + clearTimeout(livenessHandle); } await this.tryReleaseTask(ticket, taskSettings); @@ -379,7 +384,6 @@ export class TaskWorker { private async checkLiveness( ticket: string, taskAbortController: AbortController, - livenessHandle: { ref?: ReturnType }, ): Promise { try { const [row] = await this.knex(DB_TASKS_TABLE) @@ -390,7 +394,6 @@ export class TaskWorker { this.logger.info( `Task ticket for "${this.taskId}" is no longer valid; aborting execution`, ); - clearInterval(livenessHandle.ref); taskAbortController.abort(); } } catch (e) { diff --git a/packages/backend-test-utils/src/services/MockSchedulerService.test.ts b/packages/backend-test-utils/src/services/MockSchedulerService.test.ts index faf4e3e26d..1c451da749 100644 --- a/packages/backend-test-utils/src/services/MockSchedulerService.test.ts +++ b/packages/backend-test-utils/src/services/MockSchedulerService.test.ts @@ -206,6 +206,68 @@ describe('MockSchedulerService', () => { await expect(isDone()).resolves.toBe(true); }); + it('should cancel a running task and allow re-triggering with a fresh signal', async () => { + const scheduler = new MockSchedulerService(); + const signals: AbortSignal[] = []; + + scheduler.scheduleTask({ + ...baseOpts, + id: 'test', + fn: async signal => { + signals.push(signal); + // Simulate long-running work that respects cancellation + await new Promise((resolve, reject) => { + if (signal.aborted) { + reject(new Error('aborted')); + return; + } + signal.addEventListener('abort', () => reject(new Error('aborted'))); + setTimeout(1).then(resolve); + }); + }, + }); + + // First run completes normally + await scheduler.triggerTask('test'); + expect(signals).toHaveLength(1); + expect(signals[0].aborted).toBe(false); + + // Start a task that will block until cancelled + const blockingScheduler = new MockSchedulerService(); + let resolveBlock: (() => void) | undefined; + blockingScheduler.scheduleTask({ + ...baseOpts, + id: 'blocking', + fn: async signal => { + signals.push(signal); + await new Promise((resolve, reject) => { + signal.addEventListener('abort', () => reject(new Error('aborted'))); + resolveBlock = resolve; + }); + }, + }); + + const triggerPromise = blockingScheduler.triggerTask('blocking'); + // Give the task fn time to start + await setTimeout(1); + + await blockingScheduler.cancelTask('blocking'); + await triggerPromise.catch(() => {}); + + expect(signals).toHaveLength(2); + expect(signals[1].aborted).toBe(true); + + // Re-trigger should get a fresh non-aborted signal + resolveBlock = undefined; + const triggerPromise2 = blockingScheduler.triggerTask('blocking'); + await setTimeout(1); + resolveBlock?.(); + await triggerPromise2; + + expect(signals).toHaveLength(3); + expect(signals[2].aborted).toBe(false); + }); + it('should abort tasks when shutting down', async () => { let taskSignal: AbortSignal | undefined; diff --git a/packages/backend-test-utils/src/services/MockSchedulerService.ts b/packages/backend-test-utils/src/services/MockSchedulerService.ts index a47859a96f..3309ebc954 100644 --- a/packages/backend-test-utils/src/services/MockSchedulerService.ts +++ b/packages/backend-test-utils/src/services/MockSchedulerService.ts @@ -105,6 +105,7 @@ export class MockSchedulerService implements SchedulerService { throw new ConflictError(`Task ${id} is not running`); } task.abortControllers.abort(); + task.abortControllers = new AbortController(); } async triggerTask(id: string): Promise { From c0b9c63be45988525294de6b5f8ae427c9e51ff5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fredrik=20Adel=C3=B6w?= Date: Tue, 10 Mar 2026 13:59:40 +0100 Subject: [PATCH 7/7] Fix TypeScript error in MockSchedulerService test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Use non-null assertion for resolveBlock since TypeScript cannot track the async reassignment across the setTimeout boundary. Co-Authored-By: Claude Opus 4.6 Signed-off-by: Fredrik Adelöw --- .../src/services/MockSchedulerService.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/backend-test-utils/src/services/MockSchedulerService.test.ts b/packages/backend-test-utils/src/services/MockSchedulerService.test.ts index 1c451da749..2668fa7205 100644 --- a/packages/backend-test-utils/src/services/MockSchedulerService.test.ts +++ b/packages/backend-test-utils/src/services/MockSchedulerService.test.ts @@ -261,7 +261,7 @@ describe('MockSchedulerService', () => { resolveBlock = undefined; const triggerPromise2 = blockingScheduler.triggerTask('blocking'); await setTimeout(1); - resolveBlock?.(); + resolveBlock!(); await triggerPromise2; expect(signals).toHaveLength(3);