Close scaffolder stale task executions
Signed-off-by: OscarDHdz <v-ohernandez@expediagroup.com>
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
'@backstage/plugin-scaffolder-backend': minor
|
||||
---
|
||||
|
||||
Add functionality to shutdown scaffolder tasks if they are stale
|
||||
@@ -276,6 +276,10 @@ scaffolder:
|
||||
# email: scaffolder@backstage.io
|
||||
# Use to customize the default commit message when new components are created
|
||||
# defaultCommitMessage: 'Initial commit'
|
||||
# Use to customize when will stale tasks be marked as closed
|
||||
# taskTimeout:
|
||||
# ms: 3600000 # 1hr
|
||||
# taskTimeoutMessage: 'This task was canceled as it timed out'
|
||||
|
||||
auth:
|
||||
### Add auth.keyStore.provider to more granularly control how to store JWK data when running
|
||||
|
||||
@@ -499,6 +499,11 @@ export class DatabaseTaskStore implements TaskStore {
|
||||
taskId: string;
|
||||
}[];
|
||||
}>;
|
||||
// (undocumented)
|
||||
shutdownTask({
|
||||
taskId,
|
||||
message,
|
||||
}: TaskStoreShutDownTaskOptions): Promise<void>;
|
||||
}
|
||||
|
||||
// @public
|
||||
@@ -741,6 +746,11 @@ export interface TaskStore {
|
||||
taskId: string;
|
||||
}[];
|
||||
}>;
|
||||
// (undocumented)
|
||||
shutdownTask({
|
||||
taskId,
|
||||
message,
|
||||
}: TaskStoreShutDownTaskOptions): Promise<void>;
|
||||
}
|
||||
|
||||
// @public
|
||||
@@ -767,6 +777,12 @@ export type TaskStoreListEventsOptions = {
|
||||
after?: number | undefined;
|
||||
};
|
||||
|
||||
// @public
|
||||
export type TaskStoreShutDownTaskOptions = {
|
||||
taskId: string;
|
||||
message?: string | undefined;
|
||||
};
|
||||
|
||||
// @public
|
||||
export class TaskWorker {
|
||||
// (undocumented)
|
||||
|
||||
+7
@@ -28,5 +28,12 @@ export interface Config {
|
||||
* The commit message used when new components are created.
|
||||
*/
|
||||
defaultCommitMessage?: string;
|
||||
/**
|
||||
* To mark stale tasks has closed
|
||||
*/
|
||||
taskTimeout?: {
|
||||
ms?: number;
|
||||
message?: string;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ import {
|
||||
TaskStoreListEventsOptions,
|
||||
TaskStoreCreateTaskOptions,
|
||||
TaskStoreCreateTaskResult,
|
||||
TaskStoreShutDownTaskOptions,
|
||||
} from './types';
|
||||
import { DateTime } from 'luxon';
|
||||
|
||||
@@ -380,4 +381,46 @@ export class DatabaseTaskStore implements TaskStore {
|
||||
});
|
||||
return { events };
|
||||
}
|
||||
|
||||
async shutdownTask({
|
||||
taskId,
|
||||
message,
|
||||
}: TaskStoreShutDownTaskOptions): Promise<void> {
|
||||
const errorMessage =
|
||||
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: errorMessage,
|
||||
stepId: step,
|
||||
status: 'failed',
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
await this.completeTask({
|
||||
taskId,
|
||||
status: 'failed',
|
||||
eventBody: {
|
||||
message: errorMessage,
|
||||
},
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,6 +41,7 @@ async function createStore(): Promise<DatabaseTaskStore> {
|
||||
describe('StorageTaskBroker', () => {
|
||||
let storage: DatabaseTaskStore;
|
||||
const fakeSecrets = { backstageToken: 'secret' } as TaskSecrets;
|
||||
const config = new ConfigReader({ scaffolder: { title: 'Blah' } });
|
||||
|
||||
beforeAll(async () => {
|
||||
storage = await createStore();
|
||||
@@ -48,13 +49,13 @@ describe('StorageTaskBroker', () => {
|
||||
|
||||
const logger = getVoidLogger();
|
||||
it('should claim a dispatched work item', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
await broker.dispatch({ spec: {} as TaskSpec });
|
||||
await expect(broker.claim()).resolves.toEqual(expect.any(TaskManager));
|
||||
});
|
||||
|
||||
it('should wait for a dispatched work item', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const promise = broker.claim();
|
||||
|
||||
await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting');
|
||||
@@ -64,7 +65,7 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should dispatch multiple items and claim them in order', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
await broker.dispatch({ spec: { steps: [{ id: 'a' }] } as TaskSpec });
|
||||
await broker.dispatch({ spec: { steps: [{ id: 'b' }] } as TaskSpec });
|
||||
await broker.dispatch({ spec: { steps: [{ id: 'c' }] } as TaskSpec });
|
||||
@@ -81,14 +82,14 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should store secrets', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets });
|
||||
const task = await broker.claim();
|
||||
expect(task.secrets).toEqual(fakeSecrets);
|
||||
}, 10000);
|
||||
|
||||
it('should complete a task', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec });
|
||||
const task = await broker.claim();
|
||||
await task.complete('completed');
|
||||
@@ -97,7 +98,7 @@ describe('StorageTaskBroker', () => {
|
||||
}, 10000);
|
||||
|
||||
it('should remove secrets after picking up a task', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const dispatchResult = await broker.dispatch({
|
||||
spec: {} as TaskSpec,
|
||||
secrets: fakeSecrets,
|
||||
@@ -109,7 +110,7 @@ describe('StorageTaskBroker', () => {
|
||||
}, 10000);
|
||||
|
||||
it('should fail a task', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec });
|
||||
const task = await broker.claim();
|
||||
await task.complete('failed');
|
||||
@@ -118,8 +119,8 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('multiple brokers should be able to observe a single task', async () => {
|
||||
const broker1 = new StorageTaskBroker(storage, logger);
|
||||
const broker2 = new StorageTaskBroker(storage, logger);
|
||||
const broker1 = new StorageTaskBroker(storage, logger, config);
|
||||
const broker2 = new StorageTaskBroker(storage, logger, config);
|
||||
|
||||
const { taskId } = await broker1.dispatch({ spec: {} as TaskSpec });
|
||||
|
||||
@@ -161,7 +162,7 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should heartbeat', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const { taskId } = await broker.dispatch({ spec: {} as TaskSpec });
|
||||
const task = await broker.claim();
|
||||
|
||||
@@ -179,7 +180,7 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should be update the status to failed if heartbeat fails', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const { taskId } = await broker.dispatch({ spec: {} as TaskSpec });
|
||||
const task = await broker.claim();
|
||||
|
||||
@@ -205,7 +206,7 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should list all tasks', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const { taskId } = await broker.dispatch({ spec: {} as TaskSpec });
|
||||
|
||||
const promise = broker.list();
|
||||
@@ -219,7 +220,7 @@ describe('StorageTaskBroker', () => {
|
||||
});
|
||||
|
||||
it('should list only tasks createdBy a specific user', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const { taskId } = await broker.dispatch({
|
||||
spec: {} as TaskSpec,
|
||||
createdBy: 'user:default/foo',
|
||||
@@ -230,4 +231,39 @@ 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');
|
||||
});
|
||||
});
|
||||
|
||||
@@ -27,6 +27,7 @@ import {
|
||||
SerializedTask,
|
||||
} from './types';
|
||||
import { TaskBrokerDispatchOptions } from '.';
|
||||
import { Config } from '@backstage/config';
|
||||
|
||||
/**
|
||||
* TaskManager
|
||||
@@ -149,6 +150,7 @@ export class StorageTaskBroker implements TaskBroker {
|
||||
constructor(
|
||||
private readonly storage: TaskStore,
|
||||
private readonly logger: Logger,
|
||||
private readonly config: Config,
|
||||
) {}
|
||||
|
||||
async list(options?: {
|
||||
@@ -200,11 +202,33 @@ 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> {
|
||||
return this.storage.getTask(taskId);
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -52,6 +52,8 @@ describe('TaskWorker', () => {
|
||||
const actionRegistry: TemplateActionRegistry = {} as TemplateActionRegistry;
|
||||
const workingDirectory = '/tmp/scaffolder';
|
||||
|
||||
const config = new ConfigReader({ scaffolder: {} });
|
||||
|
||||
const workflowRunner: NunjucksWorkflowRunner = {
|
||||
execute: jest.fn(),
|
||||
} as unknown as NunjucksWorkflowRunner;
|
||||
@@ -68,7 +70,7 @@ describe('TaskWorker', () => {
|
||||
const logger = getVoidLogger();
|
||||
|
||||
it('should call the default workflow runner when the apiVersion is beta3', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const taskWorker = await TaskWorker.create({
|
||||
logger,
|
||||
workingDirectory,
|
||||
@@ -99,7 +101,7 @@ describe('TaskWorker', () => {
|
||||
output: { testOutput: 'testmockoutput' },
|
||||
});
|
||||
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
const broker = new StorageTaskBroker(storage, logger, config);
|
||||
const taskWorker = await TaskWorker.create({
|
||||
logger,
|
||||
workingDirectory,
|
||||
|
||||
@@ -24,6 +24,7 @@ export type {
|
||||
TaskCompletionState,
|
||||
TaskStoreEmitOptions,
|
||||
TaskStoreListEventsOptions,
|
||||
TaskStoreShutDownTaskOptions,
|
||||
SerializedTask,
|
||||
SerializedTaskEvent,
|
||||
TaskStatus,
|
||||
|
||||
@@ -156,6 +156,16 @@ export type TaskStoreListEventsOptions = {
|
||||
after?: number | undefined;
|
||||
};
|
||||
|
||||
/**
|
||||
* TaskStoreShutDownTaskOptions
|
||||
*
|
||||
* @public
|
||||
*/
|
||||
export type TaskStoreShutDownTaskOptions = {
|
||||
taskId: string;
|
||||
message?: string | undefined;
|
||||
};
|
||||
|
||||
/**
|
||||
* The options passed to {@link TaskStore.createTask}
|
||||
* @public
|
||||
@@ -201,6 +211,10 @@ export interface TaskStore {
|
||||
taskId,
|
||||
after,
|
||||
}: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }>;
|
||||
shutdownTask({
|
||||
taskId,
|
||||
message,
|
||||
}: TaskStoreShutDownTaskOptions): Promise<void>;
|
||||
}
|
||||
|
||||
export type WorkflowResponse = { output: { [key: string]: JsonValue } };
|
||||
|
||||
@@ -131,13 +131,16 @@ describe('createRouter', () => {
|
||||
},
|
||||
};
|
||||
|
||||
describe('not providing an identity api', () => {
|
||||
beforeEach(async () => {
|
||||
const logger = getVoidLogger();
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({
|
||||
database: createDatabase(),
|
||||
});
|
||||
taskBroker = new StorageTaskBroker(databaseTaskStore, logger);
|
||||
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');
|
||||
|
||||
@@ -170,7 +170,7 @@ export async function createRouter(
|
||||
|
||||
if (!options.taskBroker) {
|
||||
const databaseTaskStore = await DatabaseTaskStore.create({ database });
|
||||
taskBroker = new StorageTaskBroker(databaseTaskStore, logger);
|
||||
taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config);
|
||||
} else {
|
||||
taskBroker = options.taskBroker;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user