Close scaffolder stale task executions

Signed-off-by: OscarDHdz <v-ohernandez@expediagroup.com>
This commit is contained in:
OscarDHdz
2022-08-30 16:44:12 -05:00
parent 0788cdb1dd
commit 694bfe2d61
12 changed files with 179 additions and 24 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-scaffolder-backend': minor
---
Add functionality to shutdown scaffolder tasks if they are stale
+4
View File
@@ -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
+16
View File
@@ -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
View File
@@ -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;
}