From 694bfe2d61ae39279d036ccf67166fdbf1ad27f1 Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Tue, 30 Aug 2022 16:44:12 -0500 Subject: [PATCH 1/7] Close scaffolder stale task executions Signed-off-by: OscarDHdz --- .changeset/tender-drinks-drive.md | 5 ++ app-config.yaml | 4 ++ plugins/scaffolder-backend/api-report.md | 16 +++++ plugins/scaffolder-backend/config.d.ts | 7 +++ .../src/scaffolder/tasks/DatabaseTaskStore.ts | 43 +++++++++++++ .../tasks/StorageTaskBroker.test.ts | 62 +++++++++++++++---- .../src/scaffolder/tasks/StorageTaskBroker.ts | 26 +++++++- .../src/scaffolder/tasks/TaskWorker.test.ts | 6 +- .../src/scaffolder/tasks/index.ts | 1 + .../src/scaffolder/tasks/types.ts | 14 +++++ .../src/service/router.test.ts | 17 ++--- .../scaffolder-backend/src/service/router.ts | 2 +- 12 files changed, 179 insertions(+), 24 deletions(-) create mode 100644 .changeset/tender-drinks-drive.md diff --git a/.changeset/tender-drinks-drive.md b/.changeset/tender-drinks-drive.md new file mode 100644 index 0000000000..0ea69832af --- /dev/null +++ b/.changeset/tender-drinks-drive.md @@ -0,0 +1,5 @@ +--- +'@backstage/plugin-scaffolder-backend': minor +--- + +Add functionality to shutdown scaffolder tasks if they are stale diff --git a/app-config.yaml b/app-config.yaml index 22ca1fcaff..81b2a2ed56 100644 --- a/app-config.yaml +++ b/app-config.yaml @@ -276,6 +276,10 @@ scaffolder: # email: scaffolder@backstage.io # Use to customize the default commit message when new components are created # defaultCommitMessage: 'Initial commit' + # Use to customize when will stale tasks be marked as closed + # taskTimeout: + # ms: 3600000 # 1hr + # taskTimeoutMessage: 'This task was canceled as it timed out' auth: ### Add auth.keyStore.provider to more granularly control how to store JWK data when running diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 0d5fbcb31f..5b2474190a 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -499,6 +499,11 @@ export class DatabaseTaskStore implements TaskStore { taskId: string; }[]; }>; + // (undocumented) + shutdownTask({ + taskId, + message, + }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -741,6 +746,11 @@ export interface TaskStore { taskId: string; }[]; }>; + // (undocumented) + shutdownTask({ + taskId, + message, + }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -767,6 +777,12 @@ export type TaskStoreListEventsOptions = { after?: number | undefined; }; +// @public +export type TaskStoreShutDownTaskOptions = { + taskId: string; + message?: string | undefined; +}; + // @public export class TaskWorker { // (undocumented) diff --git a/plugins/scaffolder-backend/config.d.ts b/plugins/scaffolder-backend/config.d.ts index 7883284165..14b2a7e9f0 100644 --- a/plugins/scaffolder-backend/config.d.ts +++ b/plugins/scaffolder-backend/config.d.ts @@ -28,5 +28,12 @@ export interface Config { * The commit message used when new components are created. */ defaultCommitMessage?: string; + /** + * To mark stale tasks has closed + */ + taskTimeout?: { + ms?: number; + message?: string; + }; }; } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 206dc2ac1d..b7cc4e03ef 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -32,6 +32,7 @@ import { TaskStoreListEventsOptions, TaskStoreCreateTaskOptions, TaskStoreCreateTaskResult, + TaskStoreShutDownTaskOptions, } from './types'; import { DateTime } from 'luxon'; @@ -380,4 +381,46 @@ export class DatabaseTaskStore implements TaskStore { }); return { events }; } + + async shutdownTask({ + taskId, + message, + }: TaskStoreShutDownTaskOptions): Promise { + const errorMessage = + message || `This task was marked as stale as it exceeded its timeout`; + + const statusStepEvents = (await this.listEvents({ taskId })).events.filter( + ({ body }) => body?.stepId, + ); + + const completedSteps = statusStepEvents + .filter( + ({ body: { status } }) => status === 'failed' || status === 'completed', + ) + .map(step => step.body.stepId); + + const hungProcessingSteps = statusStepEvents + .filter(({ body: { status } }) => status === 'processing') + .map(event => event.body.stepId) + .filter(step => !completedSteps.includes(step)); + + for (const step of hungProcessingSteps) { + await this.emitLogEvent({ + taskId, + body: { + message: errorMessage, + stepId: step, + status: 'failed', + }, + }); + } + + await this.completeTask({ + taskId, + status: 'failed', + eventBody: { + message: errorMessage, + }, + }); + } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index 4fc142a991..79d6641491 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -41,6 +41,7 @@ async function createStore(): Promise { describe('StorageTaskBroker', () => { let storage: DatabaseTaskStore; const fakeSecrets = { backstageToken: 'secret' } as TaskSecrets; + const config = new ConfigReader({ scaffolder: { title: 'Blah' } }); beforeAll(async () => { storage = await createStore(); @@ -48,13 +49,13 @@ describe('StorageTaskBroker', () => { const logger = getVoidLogger(); it('should claim a dispatched work item', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); await broker.dispatch({ spec: {} as TaskSpec }); await expect(broker.claim()).resolves.toEqual(expect.any(TaskManager)); }); it('should wait for a dispatched work item', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const promise = broker.claim(); await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting'); @@ -64,7 +65,7 @@ describe('StorageTaskBroker', () => { }); it('should dispatch multiple items and claim them in order', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); await broker.dispatch({ spec: { steps: [{ id: 'a' }] } as TaskSpec }); await broker.dispatch({ spec: { steps: [{ id: 'b' }] } as TaskSpec }); await broker.dispatch({ spec: { steps: [{ id: 'c' }] } as TaskSpec }); @@ -81,14 +82,14 @@ describe('StorageTaskBroker', () => { }); it('should store secrets', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets }); const task = await broker.claim(); expect(task.secrets).toEqual(fakeSecrets); }, 10000); it('should complete a task', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); await task.complete('completed'); @@ -97,7 +98,7 @@ describe('StorageTaskBroker', () => { }, 10000); it('should remove secrets after picking up a task', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets, @@ -109,7 +110,7 @@ describe('StorageTaskBroker', () => { }, 10000); it('should fail a task', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); await task.complete('failed'); @@ -118,8 +119,8 @@ describe('StorageTaskBroker', () => { }); it('multiple brokers should be able to observe a single task', async () => { - const broker1 = new StorageTaskBroker(storage, logger); - const broker2 = new StorageTaskBroker(storage, logger); + const broker1 = new StorageTaskBroker(storage, logger, config); + const broker2 = new StorageTaskBroker(storage, logger, config); const { taskId } = await broker1.dispatch({ spec: {} as TaskSpec }); @@ -161,7 +162,7 @@ describe('StorageTaskBroker', () => { }); it('should heartbeat', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); @@ -179,7 +180,7 @@ describe('StorageTaskBroker', () => { }); it('should be update the status to failed if heartbeat fails', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); @@ -205,7 +206,7 @@ describe('StorageTaskBroker', () => { }); it('should list all tasks', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const promise = broker.list(); @@ -219,7 +220,7 @@ describe('StorageTaskBroker', () => { }); it('should list only tasks createdBy a specific user', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec, createdBy: 'user:default/foo', @@ -230,4 +231,39 @@ 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 9e60bc7727..e7bdcfe2eb 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -27,6 +27,7 @@ import { SerializedTask, } from './types'; import { TaskBrokerDispatchOptions } from '.'; +import { Config } from '@backstage/config'; /** * TaskManager @@ -149,6 +150,7 @@ export class StorageTaskBroker implements TaskBroker { constructor( private readonly storage: TaskStore, private readonly logger: Logger, + private readonly config: Config, ) {} async list(options?: { @@ -200,11 +202,33 @@ 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 { - return this.storage.getTask(taskId); + 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; } /** diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts index ee8ae400c4..d856052b59 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts @@ -52,6 +52,8 @@ describe('TaskWorker', () => { const actionRegistry: TemplateActionRegistry = {} as TemplateActionRegistry; const workingDirectory = '/tmp/scaffolder'; + const config = new ConfigReader({ scaffolder: {} }); + const workflowRunner: NunjucksWorkflowRunner = { execute: jest.fn(), } as unknown as NunjucksWorkflowRunner; @@ -68,7 +70,7 @@ describe('TaskWorker', () => { const logger = getVoidLogger(); it('should call the default workflow runner when the apiVersion is beta3', async () => { - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const taskWorker = await TaskWorker.create({ logger, workingDirectory, @@ -99,7 +101,7 @@ describe('TaskWorker', () => { output: { testOutput: 'testmockoutput' }, }); - const broker = new StorageTaskBroker(storage, logger); + const broker = new StorageTaskBroker(storage, logger, config); const taskWorker = await TaskWorker.create({ logger, workingDirectory, diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts index 457c32d9b5..2af7d6a340 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts @@ -24,6 +24,7 @@ export type { TaskCompletionState, TaskStoreEmitOptions, TaskStoreListEventsOptions, + TaskStoreShutDownTaskOptions, SerializedTask, SerializedTaskEvent, TaskStatus, diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index 4b63a0cdd3..d8f95c8423 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -156,6 +156,16 @@ export type TaskStoreListEventsOptions = { after?: number | undefined; }; +/** + * TaskStoreShutDownTaskOptions + * + * @public + */ +export type TaskStoreShutDownTaskOptions = { + taskId: string; + message?: string | undefined; +}; + /** * The options passed to {@link TaskStore.createTask} * @public @@ -201,6 +211,10 @@ export interface TaskStore { taskId, after, }: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }>; + shutdownTask({ + taskId, + message, + }: TaskStoreShutDownTaskOptions): Promise; } export type WorkflowResponse = { output: { [key: string]: JsonValue } }; diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index 5e3c1c7278..a616892d50 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -131,13 +131,16 @@ describe('createRouter', () => { }, }; - describe('not providing an identity api', () => { - beforeEach(async () => { - const logger = getVoidLogger(); - const databaseTaskStore = await DatabaseTaskStore.create({ - database: createDatabase(), - }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + 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'); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index f8e24d1acf..cdd900d3e7 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -170,7 +170,7 @@ export async function createRouter( if (!options.taskBroker) { const databaseTaskStore = await DatabaseTaskStore.create({ database }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config); } else { taskBroker = options.taskBroker; } From 5ba17223eb6061e68d7d93073adc85bafd8df6f2 Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Fri, 2 Sep 2022 10:22:11 -0500 Subject: [PATCH 2/7] Fix DatabaseTaskStore.listStaleTask PG query Signed-off-by: OscarDHdz --- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index b7cc4e03ef..27e399b762 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -269,8 +269,7 @@ export class DatabaseTaskStore implements TaskStore { '<=', this.db.client.config.client.includes('sqlite3') ? this.db.raw(`datetime('now', ?)`, [`-${timeoutS} seconds`]) - : this.db.raw(`dateadd('second', ?, ?)`, [ - `-${timeoutS}`, + : this.db.raw(`? - interval '${timeoutS} seconds'`, [ this.db.fn.now(), ]), ); From 580dc724b4294ba16e1fea84647646575cf51fbe Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Fri, 2 Sep 2022 11:59:38 -0500 Subject: [PATCH 3/7] Close stale scaffolder tasks on schedule basis Signed-off-by: OscarDHdz --- app-config.yaml | 2 +- plugins/scaffolder-backend/api-report.md | 2 + plugins/scaffolder-backend/config.d.ts | 4 +- plugins/scaffolder-backend/package.json | 1 + .../tasks/StorageTaskBroker.test.ts | 35 ----------------- .../src/scaffolder/tasks/StorageTaskBroker.ts | 38 ++++++++----------- .../src/scaffolder/tasks/types.ts | 1 + .../src/service/router.test.ts | 31 +++++++-------- .../scaffolder-backend/src/service/router.ts | 11 ++++++ yarn.lock | 1 + 10 files changed, 50 insertions(+), 76 deletions(-) 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:^" From 038ff2714a7f966c34b838535118b213da17a49b Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Wed, 7 Sep 2022 15:54:05 -0500 Subject: [PATCH 4/7] Cleanup Scaffolder: Close Stale Tasks Signed-off-by: OscarDHdz --- app-config.yaml | 4 -- plugins/scaffolder-backend/api-report.md | 15 ++----- plugins/scaffolder-backend/config.d.ts | 7 ---- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 12 ++---- .../tasks/StorageTaskBroker.test.ts | 27 ++++++------- .../src/scaffolder/tasks/StorageTaskBroker.ts | 16 -------- .../src/scaffolder/tasks/TaskWorker.test.ts | 6 +-- .../src/scaffolder/tasks/types.ts | 7 +--- .../src/service/router.test.ts | 24 ++++++----- .../scaffolder-backend/src/service/router.ts | 40 +++++++++++++------ 10 files changed, 66 insertions(+), 92 deletions(-) diff --git a/app-config.yaml b/app-config.yaml index de920ee1df..22ca1fcaff 100644 --- a/app-config.yaml +++ b/app-config.yaml @@ -276,10 +276,6 @@ scaffolder: # email: scaffolder@backstage.io # Use to customize the default commit message when new components are created # defaultCommitMessage: 'Initial commit' - # Use to customize when will stale tasks be marked as closed - # taskTimeout: - # seconds: 120 - # taskTimeoutMessage: 'This task was canceled as it timed out' auth: ### Add auth.keyStore.provider to more granularly control how to store JWK data when running diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 7793b131b9..02a6f0d664 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -500,10 +500,7 @@ export class DatabaseTaskStore implements TaskStore { }[]; }>; // (undocumented) - shutdownTask({ - taskId, - message, - }: TaskStoreShutDownTaskOptions): Promise; + shutdownTask({ taskId }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -546,6 +543,8 @@ export interface RouterOptions { // (undocumented) database: PluginDatabaseManager; // (undocumented) + databaseTaskStore?: DatabaseTaskStore; + // (undocumented) identity?: IdentityApi; // (undocumented) logger: Logger; @@ -620,8 +619,6 @@ export interface TaskBroker { // (undocumented) claim(): Promise; // (undocumented) - closeStaleTasks(): Promise; - // (undocumented) dispatch( options: TaskBrokerDispatchOptions, ): Promise; @@ -749,10 +746,7 @@ export interface TaskStore { }[]; }>; // (undocumented) - shutdownTask({ - taskId, - message, - }: TaskStoreShutDownTaskOptions): Promise; + shutdownTask?({ taskId }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -782,7 +776,6 @@ export type TaskStoreListEventsOptions = { // @public export type TaskStoreShutDownTaskOptions = { taskId: string; - message?: string | undefined; }; // @public diff --git a/plugins/scaffolder-backend/config.d.ts b/plugins/scaffolder-backend/config.d.ts index 51add2288d..7883284165 100644 --- a/plugins/scaffolder-backend/config.d.ts +++ b/plugins/scaffolder-backend/config.d.ts @@ -28,12 +28,5 @@ export interface Config { * The commit message used when new components are created. */ defaultCommitMessage?: string; - /** - * To mark stale tasks as closed - */ - taskTimeout?: { - seconds?: number; - message?: string; - }; }; } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 27e399b762..b61c3d2676 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -381,12 +381,8 @@ export class DatabaseTaskStore implements TaskStore { return { events }; } - async shutdownTask({ - taskId, - message, - }: TaskStoreShutDownTaskOptions): Promise { - const errorMessage = - message || `This task was marked as stale as it exceeded its timeout`; + async shutdownTask({ taskId }: TaskStoreShutDownTaskOptions): Promise { + const message = `This task was marked as stale as it exceeded its timeout`; const statusStepEvents = (await this.listEvents({ taskId })).events.filter( ({ body }) => body?.stepId, @@ -407,7 +403,7 @@ export class DatabaseTaskStore implements TaskStore { await this.emitLogEvent({ taskId, body: { - message: errorMessage, + message, stepId: step, status: 'failed', }, @@ -418,7 +414,7 @@ export class DatabaseTaskStore implements TaskStore { taskId, status: 'failed', eventBody: { - message: errorMessage, + message, }, }); } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index 8275f77a17..4fc142a991 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -41,7 +41,6 @@ async function createStore(): Promise { describe('StorageTaskBroker', () => { let storage: DatabaseTaskStore; const fakeSecrets = { backstageToken: 'secret' } as TaskSecrets; - const config = new ConfigReader({ scaffolder: { title: 'Blah' } }); beforeAll(async () => { storage = await createStore(); @@ -49,13 +48,13 @@ describe('StorageTaskBroker', () => { const logger = getVoidLogger(); it('should claim a dispatched work item', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); await broker.dispatch({ spec: {} as TaskSpec }); await expect(broker.claim()).resolves.toEqual(expect.any(TaskManager)); }); it('should wait for a dispatched work item', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const promise = broker.claim(); await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting'); @@ -65,7 +64,7 @@ describe('StorageTaskBroker', () => { }); it('should dispatch multiple items and claim them in order', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); await broker.dispatch({ spec: { steps: [{ id: 'a' }] } as TaskSpec }); await broker.dispatch({ spec: { steps: [{ id: 'b' }] } as TaskSpec }); await broker.dispatch({ spec: { steps: [{ id: 'c' }] } as TaskSpec }); @@ -82,14 +81,14 @@ describe('StorageTaskBroker', () => { }); it('should store secrets', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets }); const task = await broker.claim(); expect(task.secrets).toEqual(fakeSecrets); }, 10000); it('should complete a task', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); await task.complete('completed'); @@ -98,7 +97,7 @@ describe('StorageTaskBroker', () => { }, 10000); it('should remove secrets after picking up a task', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets, @@ -110,7 +109,7 @@ describe('StorageTaskBroker', () => { }, 10000); it('should fail a task', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); await task.complete('failed'); @@ -119,8 +118,8 @@ describe('StorageTaskBroker', () => { }); it('multiple brokers should be able to observe a single task', async () => { - const broker1 = new StorageTaskBroker(storage, logger, config); - const broker2 = new StorageTaskBroker(storage, logger, config); + const broker1 = new StorageTaskBroker(storage, logger); + const broker2 = new StorageTaskBroker(storage, logger); const { taskId } = await broker1.dispatch({ spec: {} as TaskSpec }); @@ -162,7 +161,7 @@ describe('StorageTaskBroker', () => { }); it('should heartbeat', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); @@ -180,7 +179,7 @@ describe('StorageTaskBroker', () => { }); it('should be update the status to failed if heartbeat fails', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const task = await broker.claim(); @@ -206,7 +205,7 @@ describe('StorageTaskBroker', () => { }); it('should list all tasks', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); const promise = broker.list(); @@ -220,7 +219,7 @@ describe('StorageTaskBroker', () => { }); it('should list only tasks createdBy a specific user', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const { taskId } = await broker.dispatch({ spec: {} as TaskSpec, createdBy: 'user:default/foo', diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 2c0a77173f..9e60bc7727 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -27,7 +27,6 @@ import { SerializedTask, } from './types'; import { TaskBrokerDispatchOptions } from '.'; -import { Config } from '@backstage/config'; /** * TaskManager @@ -150,7 +149,6 @@ export class StorageTaskBroker implements TaskBroker { constructor( private readonly storage: TaskStore, private readonly logger: Logger, - private readonly config: Config, ) {} async list(options?: { @@ -264,20 +262,6 @@ 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/TaskWorker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts index d856052b59..ee8ae400c4 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.test.ts @@ -52,8 +52,6 @@ describe('TaskWorker', () => { const actionRegistry: TemplateActionRegistry = {} as TemplateActionRegistry; const workingDirectory = '/tmp/scaffolder'; - const config = new ConfigReader({ scaffolder: {} }); - const workflowRunner: NunjucksWorkflowRunner = { execute: jest.fn(), } as unknown as NunjucksWorkflowRunner; @@ -70,7 +68,7 @@ describe('TaskWorker', () => { const logger = getVoidLogger(); it('should call the default workflow runner when the apiVersion is beta3', async () => { - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const taskWorker = await TaskWorker.create({ logger, workingDirectory, @@ -101,7 +99,7 @@ describe('TaskWorker', () => { output: { testOutput: 'testmockoutput' }, }); - const broker = new StorageTaskBroker(storage, logger, config); + const broker = new StorageTaskBroker(storage, logger); const taskWorker = await TaskWorker.create({ logger, workingDirectory, diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index 4ac515055e..5d3f8113ee 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -128,7 +128,6 @@ export interface TaskBroker { options: TaskBrokerDispatchOptions, ): Promise; vacuumTasks(options: { timeoutS: number }): Promise; - closeStaleTasks(): Promise; event$(options: { taskId: string; after: number | undefined; @@ -164,7 +163,6 @@ export type TaskStoreListEventsOptions = { */ export type TaskStoreShutDownTaskOptions = { taskId: string; - message?: string | undefined; }; /** @@ -212,10 +210,7 @@ export interface TaskStore { taskId, after, }: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }>; - shutdownTask({ - taskId, - message, - }: TaskStoreShutDownTaskOptions): Promise; + shutdownTask?({ taskId }: TaskStoreShutDownTaskOptions): Promise; } export type WorkflowResponse = { output: { [key: string]: JsonValue } }; diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index 3848e8801d..b13e955c15 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -27,6 +27,16 @@ import express from 'express'; import request from 'supertest'; import ObservableImpl from 'zen-observable'; +jest.mock('@backstage/backend-tasks', () => ({ + TaskScheduler: { + fromConfig: () => ({ + forPlugin: () => ({ + scheduleTask: jest.fn(), + }), + }), + }, +})); + /** * TODO: The following should import directly from the router file. * Due to a circular dependency between this plugin and the @@ -137,11 +147,7 @@ describe('createRouter', () => { const databaseTaskStore = await DatabaseTaskStore.create({ database: createDatabase(), }); - taskBroker = new StorageTaskBroker( - databaseTaskStore, - logger, - new ConfigReader({ scaffolder: {} }), - ); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); @@ -155,6 +161,7 @@ describe('createRouter', () => { database: createDatabase(), catalogClient, reader: mockUrlReader, + databaseTaskStore, taskBroker, }); app = express().use(router); @@ -721,11 +728,7 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ const databaseTaskStore = await DatabaseTaskStore.create({ database: createDatabase(), }); - taskBroker = new StorageTaskBroker( - databaseTaskStore, - logger, - new ConfigReader({}), - ); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); @@ -756,6 +759,7 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ database: createDatabase(), catalogClient, reader: mockUrlReader, + databaseTaskStore, taskBroker, identity: { getIdentity }, }); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index a1e938a8dd..4abea52027 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -68,6 +68,7 @@ export interface RouterOptions { actions?: TemplateAction[]; taskWorkers?: number; taskBroker?: TaskBroker; + databaseTaskStore?: DatabaseTaskStore; additionalTemplateFilters?: Record; identity?: IdentityApi; } @@ -167,11 +168,17 @@ export async function createRouter( const workingDirectory = await getWorkingDirectory(config, logger); const integrations = ScmIntegrations.fromConfig(config); - let taskBroker: TaskBroker; + let databaseTaskStore: DatabaseTaskStore; + if (!options.databaseTaskStore) { + databaseTaskStore = await DatabaseTaskStore.create({ database }); + } else { + databaseTaskStore = options.databaseTaskStore; + } + + let taskBroker: TaskBroker; if (!options.taskBroker) { - const databaseTaskStore = await DatabaseTaskStore.create({ database }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger); } else { taskBroker = options.taskBroker; } @@ -204,15 +211,24 @@ 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(); - }, - }); + if (databaseTaskStore.shutdownTask) { + 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 () => { + const { tasks } = await databaseTaskStore.listStaleTasks({ + timeoutS: 3600, + }); + + for (const task of tasks) { + logger.info(`Successfully closed stale task ${task.taskId}`); + await databaseTaskStore.shutdownTask(task); + } + }, + }); + } const dryRunner = createDryRunner({ actionRegistry, From cb69395b43aeaae8a854ef656c4aed254bcad6ce Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Thu, 8 Sep 2022 16:33:24 -0500 Subject: [PATCH 5/7] Scaffolder: Set router scheduler param Signed-off-by: OscarDHdz --- packages/backend/src/plugins/scaffolder.ts | 1 + plugins/scaffolder-backend/api-report.md | 7 +++++-- .../src/service/router.test.ts | 18 ++++++------------ .../scaffolder-backend/src/service/router.ts | 15 ++++++++------- 4 files changed, 20 insertions(+), 21 deletions(-) diff --git a/packages/backend/src/plugins/scaffolder.ts b/packages/backend/src/plugins/scaffolder.ts index 19da3b3ae0..d079b64c28 100644 --- a/packages/backend/src/plugins/scaffolder.ts +++ b/packages/backend/src/plugins/scaffolder.ts @@ -33,5 +33,6 @@ export default async function createPlugin( catalogClient: catalogClient, reader: env.reader, identity: env.identity, + scheduler: env.scheduler, }); } diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 02a6f0d664..3fb0b2631f 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -23,6 +23,7 @@ import { Logger } from 'winston'; import { Observable } from '@backstage/types'; import { Octokit } from 'octokit'; import { PluginDatabaseManager } from '@backstage/backend-common'; +import { PluginTaskScheduler } from '@backstage/backend-tasks'; import { Schema } from 'jsonschema'; import { ScmIntegrationRegistry } from '@backstage/integration'; import { ScmIntegrations } from '@backstage/integration'; @@ -543,16 +544,18 @@ export interface RouterOptions { // (undocumented) database: PluginDatabaseManager; // (undocumented) - databaseTaskStore?: DatabaseTaskStore; - // (undocumented) identity?: IdentityApi; // (undocumented) logger: Logger; // (undocumented) reader: UrlReader; // (undocumented) + scheduler?: PluginTaskScheduler; + // (undocumented) taskBroker?: TaskBroker; // (undocumented) + taskStore?: DatabaseTaskStore; + // (undocumented) taskWorkers?: number; } diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index b13e955c15..1df419fbb3 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -27,16 +27,6 @@ import express from 'express'; import request from 'supertest'; import ObservableImpl from 'zen-observable'; -jest.mock('@backstage/backend-tasks', () => ({ - TaskScheduler: { - fromConfig: () => ({ - forPlugin: () => ({ - scheduleTask: jest.fn(), - }), - }), - }, -})); - /** * TODO: The following should import directly from the router file. * Due to a circular dependency between this plugin and the @@ -161,7 +151,7 @@ describe('createRouter', () => { database: createDatabase(), catalogClient, reader: mockUrlReader, - databaseTaskStore, + taskStore: databaseTaskStore, taskBroker, }); app = express().use(router); @@ -570,9 +560,11 @@ 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); @@ -759,7 +751,7 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ database: createDatabase(), catalogClient, reader: mockUrlReader, - databaseTaskStore, + taskStore: databaseTaskStore, taskBroker, identity: { getIdentity }, }); @@ -1149,9 +1141,11 @@ 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 4abea52027..2a48b0f5d1 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -15,7 +15,7 @@ */ import { PluginDatabaseManager, UrlReader } from '@backstage/backend-common'; -import { TaskScheduler } from '@backstage/backend-tasks'; +import { PluginTaskScheduler } from '@backstage/backend-tasks'; import { CatalogApi } from '@backstage/catalog-client'; import { Entity, @@ -64,11 +64,12 @@ export interface RouterOptions { reader: UrlReader; database: PluginDatabaseManager; catalogClient: CatalogApi; + scheduler?: PluginTaskScheduler; actions?: TemplateAction[]; taskWorkers?: number; taskBroker?: TaskBroker; - databaseTaskStore?: DatabaseTaskStore; + taskStore?: DatabaseTaskStore; additionalTemplateFilters?: Record; identity?: IdentityApi; } @@ -158,6 +159,7 @@ export async function createRouter( catalogClient, actions, taskWorkers, + scheduler, additionalTemplateFilters, } = options; @@ -170,10 +172,10 @@ export async function createRouter( const integrations = ScmIntegrations.fromConfig(config); let databaseTaskStore: DatabaseTaskStore; - if (!options.databaseTaskStore) { + if (!options.taskStore) { databaseTaskStore = await DatabaseTaskStore.create({ database }); } else { - databaseTaskStore = options.databaseTaskStore; + databaseTaskStore = options.taskStore; } let taskBroker: TaskBroker; @@ -211,8 +213,7 @@ export async function createRouter( actionsToRegister.forEach(action => actionRegistry.register(action)); workers.forEach(worker => worker.start()); - if (databaseTaskStore.shutdownTask) { - const scheduler = TaskScheduler.fromConfig(config).forPlugin('scaffolder'); + if (scheduler && databaseTaskStore.listStaleTasks) { await scheduler.scheduleTask({ id: 'close_stale_tasks', frequency: { cron: '*/5 * * * *' }, // every 5 minutes, also supports Duration @@ -223,8 +224,8 @@ export async function createRouter( }); for (const task of tasks) { - logger.info(`Successfully closed stale task ${task.taskId}`); await databaseTaskStore.shutdownTask(task); + logger.info(`Successfully closed stale task ${task.taskId}`); } }, }); From ae60130b78f43a03be39c0f4b8c918ff066c388f Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Mon, 12 Sep 2022 10:03:18 -0500 Subject: [PATCH 6/7] Scaffolder: Make Close stale tasks default Signed-off-by: OscarDHdz --- plugins/scaffolder-backend/api-report.md | 2 - .../src/service/router.test.ts | 2 - .../scaffolder-backend/src/service/router.ts | 45 ++++++++----------- 3 files changed, 19 insertions(+), 30 deletions(-) diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 3fb0b2631f..35e3dab0b5 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -554,8 +554,6 @@ export interface RouterOptions { // (undocumented) taskBroker?: TaskBroker; // (undocumented) - taskStore?: DatabaseTaskStore; - // (undocumented) taskWorkers?: number; } diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index 1df419fbb3..5e3c1c7278 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -151,7 +151,6 @@ describe('createRouter', () => { database: createDatabase(), catalogClient, reader: mockUrlReader, - taskStore: databaseTaskStore, taskBroker, }); app = express().use(router); @@ -751,7 +750,6 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ database: createDatabase(), catalogClient, reader: mockUrlReader, - taskStore: databaseTaskStore, taskBroker, identity: { getIdentity }, }); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 2a48b0f5d1..a423e8afb6 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -69,7 +69,6 @@ export interface RouterOptions { actions?: TemplateAction[]; taskWorkers?: number; taskBroker?: TaskBroker; - taskStore?: DatabaseTaskStore; additionalTemplateFilters?: Record; identity?: IdentityApi; } @@ -171,16 +170,28 @@ export async function createRouter( const workingDirectory = await getWorkingDirectory(config, logger); const integrations = ScmIntegrations.fromConfig(config); - let databaseTaskStore: DatabaseTaskStore; - if (!options.taskStore) { - databaseTaskStore = await DatabaseTaskStore.create({ database }); - } else { - databaseTaskStore = options.taskStore; - } - let taskBroker: TaskBroker; if (!options.taskBroker) { + const databaseTaskStore = await DatabaseTaskStore.create({ database }); taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + + if (scheduler && databaseTaskStore.listStaleTasks) { + await scheduler.scheduleTask({ + id: 'close_stale_tasks', + frequency: { cron: '*/5 * * * *' }, // every 5 minutes, also supports Duration + timeout: { minutes: 15 }, + fn: async () => { + const { tasks } = await databaseTaskStore.listStaleTasks({ + timeoutS: 3600, + }); + + for (const task of tasks) { + await databaseTaskStore.shutdownTask(task); + logger.info(`Successfully closed stale task ${task.taskId}`); + } + }, + }); + } } else { taskBroker = options.taskBroker; } @@ -213,24 +224,6 @@ export async function createRouter( actionsToRegister.forEach(action => actionRegistry.register(action)); workers.forEach(worker => worker.start()); - if (scheduler && databaseTaskStore.listStaleTasks) { - await scheduler.scheduleTask({ - id: 'close_stale_tasks', - frequency: { cron: '*/5 * * * *' }, // every 5 minutes, also supports Duration - timeout: { minutes: 15 }, - fn: async () => { - const { tasks } = await databaseTaskStore.listStaleTasks({ - timeoutS: 3600, - }); - - for (const task of tasks) { - await databaseTaskStore.shutdownTask(task); - logger.info(`Successfully closed stale task ${task.taskId}`); - } - }, - }); - } - const dryRunner = createDryRunner({ actionRegistry, integrations, From b3c9b051891ded01edc8e17e260dcdd3af5797bb Mon Sep 17 00:00:00 2001 From: OscarDHdz Date: Mon, 19 Sep 2022 13:12:32 -0500 Subject: [PATCH 7/7] Scaffolder: Bump stale task timeout from 1hr to 24hrs Signed-off-by: OscarDHdz --- plugins/scaffolder-backend/src/service/router.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index a423e8afb6..28dd1ca413 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -182,7 +182,7 @@ export async function createRouter( timeout: { minutes: 15 }, fn: async () => { const { tasks } = await databaseTaskStore.listStaleTasks({ - timeoutS: 3600, + timeoutS: 86400, }); for (const task of tasks) {