rename task agent to task manager
Signed-off-by: Brian Fletcher <brian@roadie.io>
This commit is contained in:
@@ -349,48 +349,10 @@ export type Status =
|
||||
| 'cancelled'
|
||||
| 'completed';
|
||||
|
||||
// @public
|
||||
export interface Task {
|
||||
// (undocumented)
|
||||
complete(result: CompletedTaskState, metadata?: JsonValue): Promise<void>;
|
||||
// (undocumented)
|
||||
done: boolean;
|
||||
// (undocumented)
|
||||
emitLog(message: string, metadata?: JsonValue): Promise<void>;
|
||||
// (undocumented)
|
||||
getWorkspaceName(): Promise<string>;
|
||||
// (undocumented)
|
||||
secrets?: TaskSecrets;
|
||||
// (undocumented)
|
||||
spec: TaskSpec;
|
||||
}
|
||||
|
||||
// @public
|
||||
export class TaskAgent implements Task {
|
||||
// (undocumented)
|
||||
complete(result: CompletedTaskState, metadata?: JsonObject): Promise<void>;
|
||||
// (undocumented)
|
||||
static create(
|
||||
state: TaskState,
|
||||
storage: TaskStore,
|
||||
logger: Logger_2,
|
||||
): TaskAgent;
|
||||
// (undocumented)
|
||||
get done(): boolean;
|
||||
// (undocumented)
|
||||
emitLog(message: string, metadata?: JsonObject): Promise<void>;
|
||||
// (undocumented)
|
||||
getWorkspaceName(): Promise<string>;
|
||||
// (undocumented)
|
||||
get secrets(): TaskSecrets | undefined;
|
||||
// (undocumented)
|
||||
get spec(): TaskSpec;
|
||||
}
|
||||
|
||||
// @public
|
||||
export interface TaskBroker {
|
||||
// (undocumented)
|
||||
claim(): Promise<Task>;
|
||||
claim(): Promise<TaskContext>;
|
||||
// (undocumented)
|
||||
dispatch(spec: TaskSpec, secrets?: TaskSecrets): Promise<DispatchResult>;
|
||||
// (undocumented)
|
||||
@@ -412,9 +374,47 @@ export interface TaskBroker {
|
||||
vacuumTasks(timeoutS: { timeoutS: number }): Promise<void>;
|
||||
}
|
||||
|
||||
// @public
|
||||
export interface TaskContext {
|
||||
// (undocumented)
|
||||
complete(result: CompletedTaskState, metadata?: JsonValue): Promise<void>;
|
||||
// (undocumented)
|
||||
done: boolean;
|
||||
// (undocumented)
|
||||
emitLog(message: string, metadata?: JsonValue): Promise<void>;
|
||||
// (undocumented)
|
||||
getWorkspaceName(): Promise<string>;
|
||||
// (undocumented)
|
||||
secrets?: TaskSecrets;
|
||||
// (undocumented)
|
||||
spec: TaskSpec;
|
||||
}
|
||||
|
||||
// @public
|
||||
export type TaskEventType = 'completion' | 'log';
|
||||
|
||||
// @public
|
||||
export class TaskManager implements TaskContext {
|
||||
// (undocumented)
|
||||
complete(result: CompletedTaskState, metadata?: JsonObject): Promise<void>;
|
||||
// (undocumented)
|
||||
static create(
|
||||
state: TaskState,
|
||||
storage: TaskStore,
|
||||
logger: Logger_2,
|
||||
): TaskManager;
|
||||
// (undocumented)
|
||||
get done(): boolean;
|
||||
// (undocumented)
|
||||
emitLog(message: string, metadata?: JsonObject): Promise<void>;
|
||||
// (undocumented)
|
||||
getWorkspaceName(): Promise<string>;
|
||||
// (undocumented)
|
||||
get secrets(): TaskSecrets | undefined;
|
||||
// (undocumented)
|
||||
get spec(): TaskSpec;
|
||||
}
|
||||
|
||||
// @public
|
||||
export type TaskSecrets = {
|
||||
token: string | undefined;
|
||||
@@ -525,7 +525,7 @@ export class TaskWorker {
|
||||
// (undocumented)
|
||||
static createWorker(options: CreateWorkerOptions): TaskWorker;
|
||||
// (undocumented)
|
||||
runOneTask(task: Task): Promise<void>;
|
||||
runOneTask(task: TaskContext): Promise<void>;
|
||||
// (undocumented)
|
||||
start(): void;
|
||||
}
|
||||
|
||||
@@ -22,7 +22,7 @@ import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner';
|
||||
import { TemplateActionRegistry } from '../actions';
|
||||
import { ScmIntegrations } from '@backstage/integration';
|
||||
import { ConfigReader } from '@backstage/config';
|
||||
import { Task, TaskSpec } from './types';
|
||||
import { TaskContext, TaskSpec } from './types';
|
||||
|
||||
describe('DefaultWorkflowRunner', () => {
|
||||
const logger = getVoidLogger();
|
||||
@@ -38,7 +38,7 @@ describe('DefaultWorkflowRunner', () => {
|
||||
}),
|
||||
);
|
||||
|
||||
const createMockTaskWithSpec = (spec: TaskSpec): Task => ({
|
||||
const createMockTaskWithSpec = (spec: TaskSpec): TaskContext => ({
|
||||
spec,
|
||||
complete: async () => {},
|
||||
done: false,
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
import {
|
||||
Task,
|
||||
TaskContext,
|
||||
WorkflowRunner,
|
||||
WorkflowResponse,
|
||||
TaskSpecV1beta2,
|
||||
@@ -72,7 +72,7 @@ export class HandlebarsWorkflowRunner implements WorkflowRunner {
|
||||
this.handlebars.registerHelper('eq', (a, b) => a === b);
|
||||
}
|
||||
|
||||
async execute(task: Task): Promise<WorkflowResponse> {
|
||||
async execute(task: TaskContext): Promise<WorkflowResponse> {
|
||||
if (!isValidTaskSpec(task.spec)) {
|
||||
throw new InputError(`Task spec is not a valid v1beta2 task spec`);
|
||||
}
|
||||
|
||||
@@ -21,7 +21,7 @@ import { ScmIntegrations } from '@backstage/integration';
|
||||
import { ConfigReader } from '@backstage/config';
|
||||
import { getVoidLogger } from '@backstage/backend-common';
|
||||
import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner';
|
||||
import { Task, TaskSpec } from './types';
|
||||
import { TaskContext, TaskSpec } from './types';
|
||||
import { RepoSpec } from '../actions/builtin/publish/util';
|
||||
|
||||
describe('LegacyWorkflowRunner', () => {
|
||||
@@ -37,7 +37,7 @@ describe('LegacyWorkflowRunner', () => {
|
||||
}),
|
||||
);
|
||||
|
||||
const createMockTaskWithSpec = (spec: TaskSpec): Task => ({
|
||||
const createMockTaskWithSpec = (spec: TaskSpec): TaskContext => ({
|
||||
spec,
|
||||
complete: async () => {},
|
||||
done: false,
|
||||
|
||||
@@ -15,7 +15,7 @@
|
||||
*/
|
||||
import { ScmIntegrations } from '@backstage/integration';
|
||||
import {
|
||||
Task,
|
||||
TaskContext,
|
||||
TaskSpec,
|
||||
TaskSpecV1beta3,
|
||||
TaskStep,
|
||||
@@ -52,7 +52,13 @@ const isValidTaskSpec = (taskSpec: TaskSpec): taskSpec is TaskSpecV1beta3 => {
|
||||
return taskSpec.apiVersion === 'scaffolder.backstage.io/v1beta3';
|
||||
};
|
||||
|
||||
const createStepLogger = ({ task, step }: { task: Task; step: TaskStep }) => {
|
||||
const createStepLogger = ({
|
||||
task,
|
||||
step,
|
||||
}: {
|
||||
task: TaskContext;
|
||||
step: TaskStep;
|
||||
}) => {
|
||||
const metadata = { stepId: step.id };
|
||||
const taskLogger = winston.createLogger({
|
||||
level: process.env.LOG_LEVEL || 'info',
|
||||
@@ -162,7 +168,7 @@ export class NunjucksWorkflowRunner implements WorkflowRunner {
|
||||
});
|
||||
}
|
||||
|
||||
async execute(task: Task): Promise<WorkflowResponse> {
|
||||
async execute(task: TaskContext): Promise<WorkflowResponse> {
|
||||
if (!isValidTaskSpec(task.spec)) {
|
||||
throw new InputError(
|
||||
'Wrong template version executed with the workflow engine',
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
import { getVoidLogger, DatabaseManager } from '@backstage/backend-common';
|
||||
import { ConfigReader } from '@backstage/config';
|
||||
import { DatabaseTaskStore } from './DatabaseTaskStore';
|
||||
import { StorageTaskBroker, TaskAgent } from './StorageTaskBroker';
|
||||
import { StorageTaskBroker, TaskManager } from './StorageTaskBroker';
|
||||
import { TaskSecrets, TaskSpec, SerializedTaskEvent } from './types';
|
||||
|
||||
async function createStore(): Promise<DatabaseTaskStore> {
|
||||
@@ -48,7 +48,7 @@ describe('StorageTaskBroker', () => {
|
||||
it('should claim a dispatched work item', async () => {
|
||||
const broker = new StorageTaskBroker(storage, logger);
|
||||
await broker.dispatch({} as TaskSpec);
|
||||
await expect(broker.claim()).resolves.toEqual(expect.any(TaskAgent));
|
||||
await expect(broker.claim()).resolves.toEqual(expect.any(TaskManager));
|
||||
});
|
||||
|
||||
it('should wait for a dispatched work item', async () => {
|
||||
@@ -58,7 +58,7 @@ describe('StorageTaskBroker', () => {
|
||||
await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting');
|
||||
|
||||
await broker.dispatch({} as TaskSpec);
|
||||
await expect(promise).resolves.toEqual(expect.any(TaskAgent));
|
||||
await expect(promise).resolves.toEqual(expect.any(TaskManager));
|
||||
});
|
||||
|
||||
it('should dispatch multiple items and claim them in order', async () => {
|
||||
@@ -70,9 +70,9 @@ describe('StorageTaskBroker', () => {
|
||||
const taskA = await broker.claim();
|
||||
const taskB = await broker.claim();
|
||||
const taskC = await broker.claim();
|
||||
await expect(taskA).toEqual(expect.any(TaskAgent));
|
||||
await expect(taskB).toEqual(expect.any(TaskAgent));
|
||||
await expect(taskC).toEqual(expect.any(TaskAgent));
|
||||
await expect(taskA).toEqual(expect.any(TaskManager));
|
||||
await expect(taskB).toEqual(expect.any(TaskManager));
|
||||
await expect(taskC).toEqual(expect.any(TaskManager));
|
||||
await expect(taskA.spec.steps[0].id).toBe('a');
|
||||
await expect(taskB.spec.steps[0].id).toBe('b');
|
||||
await expect(taskC.spec.steps[0].id).toBe('c');
|
||||
|
||||
@@ -18,7 +18,7 @@ import { assertError } from '@backstage/errors';
|
||||
import { Logger } from 'winston';
|
||||
import {
|
||||
CompletedTaskState,
|
||||
Task,
|
||||
TaskContext,
|
||||
TaskSecrets,
|
||||
TaskSpec,
|
||||
TaskStore,
|
||||
@@ -29,17 +29,17 @@ import {
|
||||
} from './types';
|
||||
|
||||
/**
|
||||
* TaskAgent
|
||||
* TaskManager
|
||||
*
|
||||
* @public
|
||||
*/
|
||||
export class TaskAgent implements Task {
|
||||
export class TaskManager implements TaskContext {
|
||||
private isDone = false;
|
||||
|
||||
private heartbeatTimeoutId?: ReturnType<typeof setInterval>;
|
||||
|
||||
static create(state: TaskState, storage: TaskStore, logger: Logger) {
|
||||
const agent = new TaskAgent(state, storage, logger);
|
||||
const agent = new TaskManager(state, storage, logger);
|
||||
agent.startTimeout();
|
||||
return agent;
|
||||
}
|
||||
@@ -135,11 +135,11 @@ export class StorageTaskBroker implements TaskBroker {
|
||||
) {}
|
||||
private deferredDispatch = defer();
|
||||
|
||||
async claim(): Promise<Task> {
|
||||
async claim(): Promise<TaskContext> {
|
||||
for (;;) {
|
||||
const pendingTask = await this.storage.claimTask();
|
||||
if (pendingTask) {
|
||||
return TaskAgent.create(
|
||||
return TaskManager.create(
|
||||
{
|
||||
taskId: pendingTask.id,
|
||||
spec: pendingTask.spec,
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
import { Task, TaskBroker, WorkflowRunner } from './types';
|
||||
import { TaskContext, TaskBroker, WorkflowRunner } from './types';
|
||||
import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner';
|
||||
import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner';
|
||||
import { Logger } from 'winston';
|
||||
@@ -94,7 +94,7 @@ export class TaskWorker {
|
||||
})();
|
||||
}
|
||||
|
||||
async runOneTask(task: Task) {
|
||||
async runOneTask(task: TaskContext) {
|
||||
try {
|
||||
const { output } =
|
||||
task.spec.apiVersion === 'scaffolder.backstage.io/v1beta3'
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
export { DatabaseTaskStore } from './DatabaseTaskStore';
|
||||
export { TaskAgent } from './StorageTaskBroker';
|
||||
export { TaskManager } from './StorageTaskBroker';
|
||||
export type { TaskState } from './StorageTaskBroker';
|
||||
export { TaskWorker } from './TaskWorker';
|
||||
export type { CreateWorkerOptions } from './TaskWorker';
|
||||
@@ -31,7 +31,7 @@ export type {
|
||||
Status,
|
||||
TaskEventType,
|
||||
TaskBroker,
|
||||
Task,
|
||||
TaskContext,
|
||||
TaskStore,
|
||||
DispatchResult,
|
||||
} from './types';
|
||||
|
||||
@@ -139,7 +139,7 @@ export type DispatchResult = {
|
||||
*
|
||||
* @public
|
||||
*/
|
||||
export interface Task {
|
||||
export interface TaskContext {
|
||||
spec: TaskSpec;
|
||||
secrets?: TaskSecrets;
|
||||
done: boolean;
|
||||
@@ -154,7 +154,7 @@ export interface Task {
|
||||
* @public
|
||||
*/
|
||||
export interface TaskBroker {
|
||||
claim(): Promise<Task>;
|
||||
claim(): Promise<TaskContext>;
|
||||
dispatch(spec: TaskSpec, secrets?: TaskSecrets): Promise<DispatchResult>;
|
||||
vacuumTasks(timeoutS: { timeoutS: number }): Promise<void>;
|
||||
observe(
|
||||
@@ -221,5 +221,5 @@ export interface TaskStore {
|
||||
|
||||
export type WorkflowResponse = { output: { [key: string]: JsonValue } };
|
||||
export interface WorkflowRunner {
|
||||
execute(task: Task): Promise<WorkflowResponse>;
|
||||
execute(task: TaskContext): Promise<WorkflowResponse>;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user