Make it possible to shutdown unfinished scaffolder ongoing tasks after restart

Signed-off-by: Bogdan Nechyporenko <bnechyporenko@bol.com>
This commit is contained in:
Bogdan Nechyporenko
2023-10-26 19:05:51 +02:00
parent c938f60df6
commit 4b47afef98
4 changed files with 5 additions and 52 deletions
+1 -1
View File
@@ -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.
@@ -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}`,
@@ -149,7 +149,6 @@ export class DatabaseTaskStore implements TaskStore {
async list(options: {
createdBy?: string;
status?: TaskStatus;
}): Promise<{ tasks: SerializedTask[] }> {
const queryBuilder = this.db<RawDbTaskRow>('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 => ({
@@ -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<void> {
if (this.taskQueue.pending < this.options.concurrentTasksLimit) {
return Promise.resolve();