From fcbfc56be8fcc5417219960a5a64e854151093d4 Mon Sep 17 00:00:00 2001 From: Johan Haals Date: Fri, 22 Jan 2021 11:50:09 +0100 Subject: [PATCH] Add dispatchResult MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Fredrik Adelöw Co-authored-by: blam Co-authored-by: Patrik Oldsberg --- .../src/scaffolder/tasks/TaskBroker.test.ts | 24 ++++++++----- .../src/scaffolder/tasks/TaskBroker.ts | 35 ++++++++++++------- .../src/scaffolder/tasks/types.ts | 6 +++- 3 files changed, 42 insertions(+), 23 deletions(-) diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.test.ts index ef4b60cfd6..388e2fdec8 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.test.ts @@ -14,12 +14,14 @@ * limitations under the License. */ +import { InMemoryDatabase } from './Database'; import { MemoryTaskBroker, TaskAgent } from './TaskBroker'; describe('MemoryTaskBroker', () => { - it('should claim a dispatched work item', async () => { - const broker = new MemoryTaskBroker(); + const storage = new InMemoryDatabase(); + const broker = new MemoryTaskBroker(storage); + it('should claim a dispatched work item', async () => { await broker.dispatch({ metadata: '', }); @@ -27,8 +29,6 @@ describe('MemoryTaskBroker', () => { }); it('should wait for a dispatched work item', async () => { - const broker = new MemoryTaskBroker(); - const promise = broker.claim(); await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting'); @@ -38,8 +38,6 @@ describe('MemoryTaskBroker', () => { }); it('should dispatch multiple items and claim them in order', async () => { - const broker = new MemoryTaskBroker(); - await broker.dispatch({ metadata: 'a' }); await broker.dispatch({ metadata: 'b' }); await broker.dispatch({ metadata: 'c' }); @@ -56,10 +54,18 @@ describe('MemoryTaskBroker', () => { }); it('should complete a task', async () => { - const broker = new MemoryTaskBroker(); - - await broker.dispatch({}); + const dispatchResult = await broker.dispatch({ metadata: 'foo' }); const task = await broker.claim(); await task.complete('COMPLETED'); + const taskRow = await storage.get(dispatchResult.taskId); + expect(taskRow.status).toBe('COMPLETED'); + }); + + it('should fail a task', async () => { + const dispatchResult = await broker.dispatch({ metadata: 'foo' }); + const task = await broker.claim(); + await task.complete('FAILED'); + const taskRow = await storage.get(dispatchResult.taskId); + expect(taskRow.status).toBe('FAILED'); }); }); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.ts index c6fb56b42a..6d6fb0a72f 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskBroker.ts @@ -14,14 +14,20 @@ * limitations under the License. */ -import { CompletedTaskState, Task, TaskSpec, TaskBroker } from './types'; -import { InMemoryDatabase } from './database'; +import { + CompletedTaskState, + Task, + TaskSpec, + TaskBroker, + DispatchResult, +} from './types'; +import { InMemoryDatabase } from './Database'; export class TaskAgent implements Task { private heartbeartInterval?: ReturnType; - static create(state: TaskState, db: InMemoryDatabase) { - const agent = new TaskAgent(state, db); + static create(state: TaskState, storage: InMemoryDatabase) { + const agent = new TaskAgent(state, storage); agent.start(); return agent; } @@ -29,7 +35,7 @@ export class TaskAgent implements Task { // Runs heartbeat internally private constructor( private readonly state: TaskState, - private readonly db: InMemoryDatabase, + private readonly storage: InMemoryDatabase, ) {} get spec() { @@ -41,9 +47,9 @@ export class TaskAgent implements Task { } async complete(result: CompletedTaskState): Promise { - this.db.setStatus( + this.storage.setStatus( this.state.taskId, - result === 'FAILED' ? 'COMPLETED' : 'FAILED', + result === 'FAILED' ? 'FAILED' : 'COMPLETED', ); } @@ -52,7 +58,7 @@ export class TaskAgent implements Task { if (!this.state.runId) { throw new Error('no run id provided'); } - this.db.heartBeat(this.state.runId); + this.storage.heartBeat(this.state.runId); }, 1000); } } @@ -72,12 +78,12 @@ function defer() { } export class MemoryTaskBroker implements TaskBroker { - private readonly db = new InMemoryDatabase(); + constructor(private readonly storage: InMemoryDatabase) {} private deferredDispatch = defer(); async claim(): Promise { for (;;) { - const pendingTask = await this.db.claimTask(); + const pendingTask = await this.storage.claimTask(); if (pendingTask) { return TaskAgent.create( { @@ -85,7 +91,7 @@ export class MemoryTaskBroker implements TaskBroker { taskId: pendingTask.taskId, spec: pendingTask.spec, }, - this.db, + this.storage, ); } @@ -93,9 +99,12 @@ export class MemoryTaskBroker implements TaskBroker { } } - async dispatch(spec: TaskSpec): Promise { - await this.db.createTask(spec); + async dispatch(spec: TaskSpec): Promise { + const taskRow = await this.storage.createTask(spec); this.signalDispatch(); + return { + taskId: taskRow.taskId, + }; } private waitForDispatch() { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index bd8bb0ade5..e8cea06005 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -44,6 +44,10 @@ export type TaskSpec = { metadata: string; }; +export type DispatchResult = { + taskId: string; +}; + export interface Task { spec: TaskSpec; emitLog(message: string): Promise; @@ -52,5 +56,5 @@ export interface Task { export interface TaskBroker { claim(): Promise; - dispatch(spec: TaskSpec): Promise; + dispatch(spec: TaskSpec): Promise; }