Merge pull request #13429 from OscarDHdz/shutdown-stale-scaffolder-tasks
Close scaffolder stale task executions
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
'@backstage/plugin-scaffolder-backend': minor
|
||||
---
|
||||
|
||||
Add functionality to shutdown scaffolder tasks if they are stale
|
||||
@@ -33,5 +33,6 @@ export default async function createPlugin(
|
||||
catalogClient: catalogClient,
|
||||
reader: env.reader,
|
||||
identity: env.identity,
|
||||
scheduler: env.scheduler,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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<void>;
|
||||
}
|
||||
|
||||
// @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<void>;
|
||||
}
|
||||
|
||||
// @public
|
||||
@@ -767,6 +774,11 @@ export type TaskStoreListEventsOptions = {
|
||||
after?: number | undefined;
|
||||
};
|
||||
|
||||
// @public
|
||||
export type TaskStoreShutDownTaskOptions = {
|
||||
taskId: string;
|
||||
};
|
||||
|
||||
// @public
|
||||
export class TaskWorker {
|
||||
// (undocumented)
|
||||
|
||||
@@ -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:^",
|
||||
|
||||
@@ -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<void> {
|
||||
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,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -24,6 +24,7 @@ export type {
|
||||
TaskCompletionState,
|
||||
TaskStoreEmitOptions,
|
||||
TaskStoreListEventsOptions,
|
||||
TaskStoreShutDownTaskOptions,
|
||||
SerializedTask,
|
||||
SerializedTaskEvent,
|
||||
TaskStatus,
|
||||
|
||||
@@ -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<void>;
|
||||
}
|
||||
|
||||
export type WorkflowResponse = { output: { [key: string]: JsonValue } };
|
||||
|
||||
@@ -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<any>[];
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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:^"
|
||||
|
||||
Reference in New Issue
Block a user