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,