diff --git a/app-config.yaml b/app-config.yaml index 81b2a2ed56..de920ee1df 100644 --- a/app-config.yaml +++ b/app-config.yaml @@ -278,7 +278,7 @@ scaffolder: # defaultCommitMessage: 'Initial commit' # Use to customize when will stale tasks be marked as closed # taskTimeout: - # ms: 3600000 # 1hr + # seconds: 120 # taskTimeoutMessage: 'This task was canceled as it timed out' auth: diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 5b2474190a..7793b131b9 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -620,6 +620,8 @@ export interface TaskBroker { // (undocumented) claim(): Promise; // (undocumented) + closeStaleTasks(): Promise; + // (undocumented) dispatch( options: TaskBrokerDispatchOptions, ): Promise; diff --git a/plugins/scaffolder-backend/config.d.ts b/plugins/scaffolder-backend/config.d.ts index 14b2a7e9f0..51add2288d 100644 --- a/plugins/scaffolder-backend/config.d.ts +++ b/plugins/scaffolder-backend/config.d.ts @@ -29,10 +29,10 @@ export interface Config { */ defaultCommitMessage?: string; /** - * To mark stale tasks has closed + * To mark stale tasks as closed */ taskTimeout?: { - ms?: number; + seconds?: number; message?: string; }; }; diff --git a/plugins/scaffolder-backend/package.json b/plugins/scaffolder-backend/package.json index 2812150b39..35f7fe01e2 100644 --- a/plugins/scaffolder-backend/package.json +++ b/plugins/scaffolder-backend/package.json @@ -36,6 +36,7 @@ "dependencies": { "@backstage/backend-common": "workspace:^", "@backstage/backend-plugin-api": "workspace:^", + "@backstage/backend-tasks": "workspace:^", "@backstage/catalog-client": "workspace:^", "@backstage/catalog-model": "workspace:^", "@backstage/config": "workspace:^", diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index 79d6641491..8275f77a17 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -231,39 +231,4 @@ describe('StorageTaskBroker', () => { const promise = broker.list({ createdBy: 'user:default/foo' }); await expect(promise).resolves.toEqual({ tasks: [task] }); }); - - it('should shut down task if task is stale', async () => { - const broker = new StorageTaskBroker( - storage, - logger, - new ConfigReader({ - scaffolder: { - taskTimeout: { - ms: -1, - }, - }, - }), - ); - const { taskId } = await broker.dispatch({ - spec: {} as TaskSpec, - }); - - jest.spyOn(storage, 'getTask').mockResolvedValueOnce({ - status: 'processing', - lastHeartbeatAt: new Date().toISOString(), - } as any); - - jest.spyOn(storage, 'getTask').mockResolvedValueOnce({ - status: 'completed', - } as any); - - jest - .spyOn(storage, 'shutdownTask') - .mockImplementationOnce(() => Promise.resolve()); - - const task = await broker.get(taskId); - - expect(storage.shutdownTask).toHaveBeenCalled(); - expect(task.status).toEqual('completed'); - }); }); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index e7bdcfe2eb..2c0a77173f 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -202,33 +202,11 @@ export class StorageTaskBroker implements TaskBroker { }; } - private isStaleTask(task: SerializedTask) { - const { status, lastHeartbeatAt } = task; - if (status === 'processing' && lastHeartbeatAt) { - const timeDiff = - new Date().getTime() - new Date(lastHeartbeatAt).getTime(); - const timeoutLimit = - this.config.getOptionalNumber('scaffolder.taskTimeout.ms') || 3600000; - return timeDiff >= timeoutLimit; - } - return false; - } - /** * {@inheritdoc TaskBroker.get} */ async get(taskId: string): Promise { - let task = await this.storage.getTask(taskId); - if (this.isStaleTask(task)) { - await this.storage.shutdownTask({ - taskId, - message: this.config.getOptionalString( - 'scaffolder.taskTimeout.message', - ), - }); - task = await this.storage.getTask(taskId); - } - return task; + return this.storage.getTask(taskId); } /** @@ -286,6 +264,20 @@ export class StorageTaskBroker implements TaskBroker { ); } + /** + * {@inheritdoc TaskBroker.closeStaleTasks} + */ + async closeStaleTasks(): Promise { + const timeoutS = + this.config.getOptionalNumber('scaffolder.taskTimeout.seconds') || 300; + const { tasks } = await this.storage.listStaleTasks({ timeoutS }); + + for (const task of tasks) { + this.logger.info(`Successfully closed stale task ${task.taskId}`); + await this.storage.shutdownTask(task); + } + } + private waitForDispatch() { return this.deferredDispatch.promise; } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index d8f95c8423..4ac515055e 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -128,6 +128,7 @@ export interface TaskBroker { options: TaskBrokerDispatchOptions, ): Promise; vacuumTasks(options: { timeoutS: number }): Promise; + closeStaleTasks(): Promise; event$(options: { taskId: string; after: number | undefined; diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index a616892d50..3848e8801d 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -131,16 +131,17 @@ describe('createRouter', () => { }, }; - beforeEach(async () => { - const logger = getVoidLogger(); - const databaseTaskStore = await DatabaseTaskStore.create({ - database: createDatabase(), - }); - taskBroker = new StorageTaskBroker( - databaseTaskStore, - logger, - new ConfigReader({ scaffolder: {} }), - ); + describe('not providing an identity api', () => { + beforeEach(async () => { + const logger = getVoidLogger(); + const databaseTaskStore = await DatabaseTaskStore.create({ + database: createDatabase(), + }); + taskBroker = new StorageTaskBroker( + databaseTaskStore, + logger, + new ConfigReader({ scaffolder: {} }), + ); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); @@ -562,11 +563,9 @@ describe('createRouter', () => { expect(responseDataFn).toHaveBeenCalledTimes(2); expect(responseDataFn).toHaveBeenCalledWith(`event: log data: {"id":0,"taskId":"a-random-id","type":"log","createdAt":"","body":{"message":"My log message"}} - `); expect(responseDataFn).toHaveBeenCalledWith(`event: completion data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{"message":"Finished!"}} - `); expect(taskBroker.event$).toHaveBeenCalledTimes(1); @@ -722,7 +721,11 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ const databaseTaskStore = await DatabaseTaskStore.create({ database: createDatabase(), }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + taskBroker = new StorageTaskBroker( + databaseTaskStore, + logger, + new ConfigReader({}), + ); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); @@ -1142,11 +1145,9 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ expect(responseDataFn).toHaveBeenCalledTimes(2); expect(responseDataFn).toHaveBeenCalledWith(`event: log data: {"id":0,"taskId":"a-random-id","type":"log","createdAt":"","body":{"message":"My log message"}} - `); expect(responseDataFn).toHaveBeenCalledWith(`event: completion data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{"message":"Finished!"}} - `); expect(taskBroker.event$).toHaveBeenCalledTimes(1); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index cdd900d3e7..a1e938a8dd 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -15,6 +15,7 @@ */ import { PluginDatabaseManager, UrlReader } from '@backstage/backend-common'; +import { TaskScheduler } from '@backstage/backend-tasks'; import { CatalogApi } from '@backstage/catalog-client'; import { Entity, @@ -203,6 +204,16 @@ export async function createRouter( actionsToRegister.forEach(action => actionRegistry.register(action)); workers.forEach(worker => worker.start()); + const scheduler = TaskScheduler.fromConfig(config).forPlugin('scaffolder'); + await scheduler.scheduleTask({ + id: 'close_stale_tasks', + frequency: { cron: '*/5 * * * *' }, // every 5 minutes, also supports Duration + timeout: { minutes: 15 }, + fn: async () => { + await taskBroker.closeStaleTasks(); + }, + }); + const dryRunner = createDryRunner({ actionRegistry, integrations, diff --git a/yarn.lock b/yarn.lock index 7e6f645f58..bc0ca4bb6f 100644 --- a/yarn.lock +++ b/yarn.lock @@ -6345,6 +6345,7 @@ __metadata: dependencies: "@backstage/backend-common": "workspace:^" "@backstage/backend-plugin-api": "workspace:^" + "@backstage/backend-tasks": "workspace:^" "@backstage/backend-test-utils": "workspace:^" "@backstage/catalog-client": "workspace:^" "@backstage/catalog-model": "workspace:^"