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/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 0d5fbcb31f..35e3dab0b5 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'; @@ -499,6 +500,8 @@ export class DatabaseTaskStore implements TaskStore { taskId: string; }[]; }>; + // (undocumented) + shutdownTask({ taskId }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -547,6 +550,8 @@ export interface RouterOptions { // (undocumented) reader: UrlReader; // (undocumented) + scheduler?: PluginTaskScheduler; + // (undocumented) taskBroker?: TaskBroker; // (undocumented) taskWorkers?: number; @@ -741,6 +746,8 @@ export interface TaskStore { taskId: string; }[]; }>; + // (undocumented) + shutdownTask?({ taskId }: TaskStoreShutDownTaskOptions): Promise; } // @public @@ -767,6 +774,11 @@ export type TaskStoreListEventsOptions = { after?: number | undefined; }; +// @public +export type TaskStoreShutDownTaskOptions = { + taskId: string; +}; + // @public export class TaskWorker { // (undocumented) 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/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 206dc2ac1d..b61c3d2676 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'; @@ -268,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(), ]), ); @@ -380,4 +380,42 @@ export class DatabaseTaskStore implements TaskStore { }); return { events }; } + + 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, + ); + + 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, + stepId: step, + status: 'failed', + }, + }); + } + + await this.completeTask({ + taskId, + status: 'failed', + eventBody: { + message, + }, + }); + } } 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..5d3f8113ee 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -156,6 +156,15 @@ export type TaskStoreListEventsOptions = { after?: number | undefined; }; +/** + * TaskStoreShutDownTaskOptions + * + * @public + */ +export type TaskStoreShutDownTaskOptions = { + taskId: string; +}; + /** * The options passed to {@link TaskStore.createTask} * @public @@ -201,6 +210,7 @@ export interface TaskStore { taskId, after, }: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }>; + shutdownTask?({ taskId }: TaskStoreShutDownTaskOptions): Promise; } export type WorkflowResponse = { output: { [key: string]: JsonValue } }; diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index f8e24d1acf..28dd1ca413 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 { PluginTaskScheduler } from '@backstage/backend-tasks'; import { CatalogApi } from '@backstage/catalog-client'; import { Entity, @@ -63,6 +64,7 @@ export interface RouterOptions { reader: UrlReader; database: PluginDatabaseManager; catalogClient: CatalogApi; + scheduler?: PluginTaskScheduler; actions?: TemplateAction[]; taskWorkers?: number; @@ -156,6 +158,7 @@ export async function createRouter( catalogClient, actions, taskWorkers, + scheduler, additionalTemplateFilters, } = options; @@ -166,11 +169,29 @@ export async function createRouter( const workingDirectory = await getWorkingDirectory(config, logger); const integrations = ScmIntegrations.fromConfig(config); - let taskBroker: TaskBroker; + 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: 86400, + }); + + for (const task of tasks) { + await databaseTaskStore.shutdownTask(task); + logger.info(`Successfully closed stale task ${task.taskId}`); + } + }, + }); + } } else { taskBroker = options.taskBroker; } 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:^"