From 4b47afef98de8774fe1eabe8526af174c3381c54 Mon Sep 17 00:00:00 2001 From: Bogdan Nechyporenko Date: Thu, 26 Oct 2023 19:05:51 +0200 Subject: [PATCH] Make it possible to shutdown unfinished scaffolder ongoing tasks after restart Signed-off-by: Bogdan Nechyporenko --- .changeset/pink-eyes-reflect.md | 2 +- .../src/ScaffolderPlugin.ts | 15 +-------- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 7 ---- .../src/scaffolder/tasks/TaskWorker.ts | 33 ++----------------- 4 files changed, 5 insertions(+), 52 deletions(-) diff --git a/.changeset/pink-eyes-reflect.md b/.changeset/pink-eyes-reflect.md index b6e9f7f2b5..506724421f 100644 --- a/.changeset/pink-eyes-reflect.md +++ b/.changeset/pink-eyes-reflect.md @@ -7,4 +7,4 @@ Made shut down stale tasks configurable. There are two properties exposed: - `scaffolder.processingInterval` - sets the processing interval for staled tasks. -- `scaffolder.taskTimeout` - sets the task's heartbeat timeout, when to consider a task to be staled. +- `scaffolder.taskTimeoutReaperFrequency` - sets the task's heartbeat timeout, when to consider a task to be staled. diff --git a/plugins/scaffolder-backend/src/ScaffolderPlugin.ts b/plugins/scaffolder-backend/src/ScaffolderPlugin.ts index 41c62faf92..beabaa50aa 100644 --- a/plugins/scaffolder-backend/src/ScaffolderPlugin.ts +++ b/plugins/scaffolder-backend/src/ScaffolderPlugin.ts @@ -14,7 +14,6 @@ * limitations under the License. */ -import fetch from 'node-fetch'; import { createBackendPlugin, coreServices, @@ -33,7 +32,7 @@ import { scaffolderTaskBrokerExtensionPoint, scaffolderTemplatingExtensionPoint, } from '@backstage/plugin-scaffolder-node/alpha'; -import { createBuiltinActions, DatabaseTaskStore } from './scaffolder'; +import { createBuiltinActions } from './scaffolder'; import { createRouter } from './service/router'; /** @@ -76,21 +75,17 @@ export const scaffolderPlugin = createBackendPlugin({ deps: { logger: coreServices.logger, config: coreServices.rootConfig, - lifecycle: coreServices.lifecycle, reader: coreServices.urlReader, permissions: coreServices.permissions, database: coreServices.database, - discovery: coreServices.discovery, httpRouter: coreServices.httpRouter, catalogClient: catalogServiceRef, }, async init({ logger, config, - lifecycle, reader, database, - discovery, httpRouter, catalogClient, permissions, @@ -109,14 +104,6 @@ export const scaffolderPlugin = createBackendPlugin({ }), ]; - lifecycle.addShutdownHook(async () => { - const baseUrl = await discovery.getBaseUrl('scaffolder'); - const url = `${baseUrl}/v2/tasks/cancel`; - await fetch(url, { - method: 'POST', - }); - }); - const actionIds = actions.map(action => action.id).join(', '); log.info( `Starting scaffolder with the following actions enabled ${actionIds}`, diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 8cb45240cc..6beb01e401 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -149,7 +149,6 @@ export class DatabaseTaskStore implements TaskStore { async list(options: { createdBy?: string; - status?: TaskStatus; }): Promise<{ tasks: SerializedTask[] }> { const queryBuilder = this.db('tasks'); @@ -159,12 +158,6 @@ export class DatabaseTaskStore implements TaskStore { }); } - if (options.status) { - queryBuilder.where({ - status: options.status, - }); - } - const results = await queryBuilder.orderBy('created_at', 'desc').select(); const tasks = results.map(result => ({ diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 01936a8409..9f9b86b723 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -14,19 +14,14 @@ * limitations under the License. */ -import { WorkflowRunner } from './types'; -import { - TaskContext, - TaskBroker, - TemplateFilter, - TemplateGlobal, -} from '@backstage/plugin-scaffolder-node'; +import { TaskContext, TaskBroker, WorkflowRunner } from './types'; import PQueue from 'p-queue'; import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner'; import { Logger } from 'winston'; import { TemplateActionRegistry } from '../actions'; import { ScmIntegrations } from '@backstage/integration'; import { assertError } from '@backstage/errors'; +import { TemplateFilter, TemplateGlobal } from '../../lib'; import { PermissionEvaluator } from '@backstage/plugin-permission-common'; /** * TaskWorkerOptions @@ -78,10 +73,8 @@ export type CreateWorkerOptions = { */ export class TaskWorker { private taskQueue: PQueue; - private taskQueueAbortController: AbortController; private constructor(private readonly options: TaskWorkerOptions) { - this.taskQueueAbortController = new AbortController(); this.taskQueue = new PQueue({ concurrency: options.concurrentTasksLimit, }); @@ -123,31 +116,11 @@ export class TaskWorker { for (;;) { await this.onReadyToClaimTask(); const task = await this.options.taskBroker.claim(); - const taskId = await task.getWorkspaceName(); - - try { - this.taskQueue.add( - ({ signal }) => { - this.runOneTask(task); - signal?.addEventListener('abort', () => { - this.options.taskBroker.cancel?.(taskId); - }); - }, - { signal: this.taskQueueAbortController.signal }, - ); - } catch (error) { - if (!(error instanceof AbortError)) { - throw error; - } - } + this.taskQueue.add(() => this.runOneTask(task)); } })(); } - cancelAllRunningTasks() { - this.taskQueueAbortController.abort(); - } - protected onReadyToClaimTask(): Promise { if (this.taskQueue.pending < this.options.concurrentTasksLimit) { return Promise.resolve();