Close stale scaffolder tasks on schedule basis
Signed-off-by: OscarDHdz <v-ohernandez@expediagroup.com>
This commit is contained in:
@@ -620,6 +620,8 @@ export interface TaskBroker {
|
||||
// (undocumented)
|
||||
claim(): Promise<TaskContext>;
|
||||
// (undocumented)
|
||||
closeStaleTasks(): Promise<void>;
|
||||
// (undocumented)
|
||||
dispatch(
|
||||
options: TaskBrokerDispatchOptions,
|
||||
): Promise<TaskBrokerDispatchResult>;
|
||||
|
||||
+2
-2
@@ -29,10 +29,10 @@ export interface Config {
|
||||
*/
|
||||
defaultCommitMessage?: string;
|
||||
/**
|
||||
* To mark stale tasks has closed
|
||||
* To mark stale tasks as closed
|
||||
*/
|
||||
taskTimeout?: {
|
||||
ms?: number;
|
||||
seconds?: number;
|
||||
message?: string;
|
||||
};
|
||||
};
|
||||
|
||||
@@ -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:^",
|
||||
|
||||
@@ -231,39 +231,4 @@ describe('StorageTaskBroker', () => {
|
||||
const promise = broker.list({ createdBy: 'user:default/foo' });
|
||||
await expect(promise).resolves.toEqual({ tasks: [task] });
|
||||
});
|
||||
|
||||
it('should shut down task if task is stale', async () => {
|
||||
const broker = new StorageTaskBroker(
|
||||
storage,
|
||||
logger,
|
||||
new ConfigReader({
|
||||
scaffolder: {
|
||||
taskTimeout: {
|
||||
ms: -1,
|
||||
},
|
||||
},
|
||||
}),
|
||||
);
|
||||
const { taskId } = await broker.dispatch({
|
||||
spec: {} as TaskSpec,
|
||||
});
|
||||
|
||||
jest.spyOn(storage, 'getTask').mockResolvedValueOnce({
|
||||
status: 'processing',
|
||||
lastHeartbeatAt: new Date().toISOString(),
|
||||
} as any);
|
||||
|
||||
jest.spyOn(storage, 'getTask').mockResolvedValueOnce({
|
||||
status: 'completed',
|
||||
} as any);
|
||||
|
||||
jest
|
||||
.spyOn(storage, 'shutdownTask')
|
||||
.mockImplementationOnce(() => Promise.resolve());
|
||||
|
||||
const task = await broker.get(taskId);
|
||||
|
||||
expect(storage.shutdownTask).toHaveBeenCalled();
|
||||
expect(task.status).toEqual('completed');
|
||||
});
|
||||
});
|
||||
|
||||
@@ -202,33 +202,11 @@ export class StorageTaskBroker implements TaskBroker {
|
||||
};
|
||||
}
|
||||
|
||||
private isStaleTask(task: SerializedTask) {
|
||||
const { status, lastHeartbeatAt } = task;
|
||||
if (status === 'processing' && lastHeartbeatAt) {
|
||||
const timeDiff =
|
||||
new Date().getTime() - new Date(lastHeartbeatAt).getTime();
|
||||
const timeoutLimit =
|
||||
this.config.getOptionalNumber('scaffolder.taskTimeout.ms') || 3600000;
|
||||
return timeDiff >= timeoutLimit;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc TaskBroker.get}
|
||||
*/
|
||||
async get(taskId: string): Promise<SerializedTask> {
|
||||
let task = await this.storage.getTask(taskId);
|
||||
if (this.isStaleTask(task)) {
|
||||
await this.storage.shutdownTask({
|
||||
taskId,
|
||||
message: this.config.getOptionalString(
|
||||
'scaffolder.taskTimeout.message',
|
||||
),
|
||||
});
|
||||
task = await this.storage.getTask(taskId);
|
||||
}
|
||||
return task;
|
||||
return this.storage.getTask(taskId);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -286,6 +264,20 @@ export class StorageTaskBroker implements TaskBroker {
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritdoc TaskBroker.closeStaleTasks}
|
||||
*/
|
||||
async closeStaleTasks(): Promise<void> {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -128,6 +128,7 @@ export interface TaskBroker {
|
||||
options: TaskBrokerDispatchOptions,
|
||||
): Promise<TaskBrokerDispatchResult>;
|
||||
vacuumTasks(options: { timeoutS: number }): Promise<void>;
|
||||
closeStaleTasks(): Promise<void>;
|
||||
event$(options: {
|
||||
taskId: string;
|
||||
after: number | undefined;
|
||||
|
||||
@@ -131,16 +131,17 @@ describe('createRouter', () => {
|
||||
},
|
||||
};
|
||||
|
||||
beforeEach(async () => {
|
||||
const logger = getVoidLogger();
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({
|
||||
database: createDatabase(),
|
||||
});
|
||||
taskBroker = new StorageTaskBroker(
|
||||
databaseTaskStore,
|
||||
logger,
|
||||
new ConfigReader({ scaffolder: {} }),
|
||||
);
|
||||
describe('not providing an identity api', () => {
|
||||
beforeEach(async () => {
|
||||
const logger = getVoidLogger();
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({
|
||||
database: createDatabase(),
|
||||
});
|
||||
taskBroker = new StorageTaskBroker(
|
||||
databaseTaskStore,
|
||||
logger,
|
||||
new ConfigReader({ scaffolder: {} }),
|
||||
);
|
||||
|
||||
jest.spyOn(taskBroker, 'dispatch');
|
||||
jest.spyOn(taskBroker, 'get');
|
||||
@@ -562,11 +563,9 @@ describe('createRouter', () => {
|
||||
expect(responseDataFn).toHaveBeenCalledTimes(2);
|
||||
expect(responseDataFn).toHaveBeenCalledWith(`event: log
|
||||
data: {"id":0,"taskId":"a-random-id","type":"log","createdAt":"","body":{"message":"My log message"}}
|
||||
|
||||
`);
|
||||
expect(responseDataFn).toHaveBeenCalledWith(`event: completion
|
||||
data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{"message":"Finished!"}}
|
||||
|
||||
`);
|
||||
|
||||
expect(taskBroker.event$).toHaveBeenCalledTimes(1);
|
||||
@@ -722,7 +721,11 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({
|
||||
database: createDatabase(),
|
||||
});
|
||||
taskBroker = new StorageTaskBroker(databaseTaskStore, logger);
|
||||
taskBroker = new StorageTaskBroker(
|
||||
databaseTaskStore,
|
||||
logger,
|
||||
new ConfigReader({}),
|
||||
);
|
||||
|
||||
jest.spyOn(taskBroker, 'dispatch');
|
||||
jest.spyOn(taskBroker, 'get');
|
||||
@@ -1142,11 +1145,9 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{
|
||||
expect(responseDataFn).toHaveBeenCalledTimes(2);
|
||||
expect(responseDataFn).toHaveBeenCalledWith(`event: log
|
||||
data: {"id":0,"taskId":"a-random-id","type":"log","createdAt":"","body":{"message":"My log message"}}
|
||||
|
||||
`);
|
||||
expect(responseDataFn).toHaveBeenCalledWith(`event: completion
|
||||
data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{"message":"Finished!"}}
|
||||
|
||||
`);
|
||||
|
||||
expect(taskBroker.event$).toHaveBeenCalledTimes(1);
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
*/
|
||||
|
||||
import { PluginDatabaseManager, UrlReader } from '@backstage/backend-common';
|
||||
import { TaskScheduler } from '@backstage/backend-tasks';
|
||||
import { CatalogApi } from '@backstage/catalog-client';
|
||||
import {
|
||||
Entity,
|
||||
@@ -203,6 +204,16 @@ 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();
|
||||
},
|
||||
});
|
||||
|
||||
const dryRunner = createDryRunner({
|
||||
actionRegistry,
|
||||
integrations,
|
||||
|
||||
Reference in New Issue
Block a user