Merge pull request #7426 from RoadieHQ/scaffolder-improvement

changes to allow the task workers run outside the backend process
This commit is contained in:
Ben Lambert
2021-10-29 16:41:14 +02:00
committed by GitHub
17 changed files with 631 additions and 130 deletions
+5
View File
@@ -0,0 +1,5 @@
---
'@backstage/plugin-scaffolder-backend': patch
---
Expose some classes and interfaces public so TaskWorkers can run externally from the scaffolder API.
+275 -3
View File
@@ -16,6 +16,7 @@ import { Entity } from '@backstage/catalog-model';
import express from 'express';
import { JsonObject } from '@backstage/types';
import { JsonValue } from '@backstage/types';
import { Knex } from 'knex';
import { LocationSpec } from '@backstage/catalog-model';
import { Logger as Logger_2 } from 'winston';
import { Octokit } from '@octokit/rest';
@@ -55,6 +56,9 @@ export class CatalogEntityClient {
): Promise<TemplateEntityV1beta2>;
}
// @public
export type CompletedTaskState = 'failed' | 'completed';
// Warning: (ae-missing-release-tag) "createBuiltinActions" is exported by the package, but it is missing a release tag (@alpha, @beta, @public, or @internal)
//
// @public (undocumented)
@@ -189,6 +193,64 @@ export const createTemplateAction: <
templateAction: TemplateAction<Input>,
) => TemplateAction<any>;
// @public
export type CreateWorkerOptions = {
taskBroker: TaskBroker;
actionRegistry: TemplateActionRegistry;
integrations: ScmIntegrations;
workingDirectory: string;
logger: Logger_2;
};
// @public
export class DatabaseTaskStore implements TaskStore {
constructor(options: DatabaseTaskStoreOptions);
// (undocumented)
claimTask(): Promise<SerializedTask | undefined>;
// (undocumented)
completeTask({
taskId,
status,
eventBody,
}: {
taskId: string;
status: Status;
eventBody: JsonObject;
}): Promise<void>;
// Warning: (ae-forgotten-export) The symbol "DatabaseTaskStoreOptions" needs to be exported by the entry point index.d.ts
//
// (undocumented)
static create(options: DatabaseTaskStoreOptions): Promise<DatabaseTaskStore>;
// (undocumented)
createTask(
spec: TaskSpec,
secrets?: TaskSecrets,
): Promise<{
taskId: string;
}>;
// (undocumented)
emitLogEvent({ taskId, body }: TaskStoreEmitOptions): Promise<void>;
// (undocumented)
getTask(taskId: string): Promise<SerializedTask>;
// (undocumented)
heartbeatTask(taskId: string): Promise<void>;
// (undocumented)
listEvents({ taskId, after }: TaskStoreListEventsOptions): Promise<{
events: SerializedTaskEvent[];
}>;
// (undocumented)
listStaleTasks({ timeoutS }: { timeoutS: number }): Promise<{
tasks: {
taskId: string;
}[];
}>;
}
// @public
export type DispatchResult = {
taskId: string;
};
// Warning: (ae-missing-release-tag) "fetchContents" is exported by the package, but it is missing a release tag (@alpha, @beta, @public, or @internal)
//
// @public (undocumented)
@@ -216,9 +278,7 @@ export class OctokitProvider {
getOctokit(repoUrl: string): Promise<OctokitIntegration>;
}
// Warning: (ae-missing-release-tag) "RouterOptions" is exported by the package, but it is missing a release tag (@alpha, @beta, @public, or @internal)
//
// @public (undocumented)
// @public
export interface RouterOptions {
// (undocumented)
actions?: TemplateAction<any>[];
@@ -235,6 +295,8 @@ export interface RouterOptions {
// (undocumented)
reader: UrlReader;
// (undocumented)
taskBroker?: TaskBroker;
// (undocumented)
taskWorkers?: number;
}
@@ -260,6 +322,216 @@ export class ScaffolderEntitiesProcessor implements CatalogProcessor {
validateEntityKind(entity: Entity): Promise<boolean>;
}
// @public
export type SerializedTask = {
id: string;
spec: TaskSpec;
status: Status;
createdAt: string;
lastHeartbeatAt?: string;
secrets?: TaskSecrets;
};
// @public
export type SerializedTaskEvent = {
id: number;
taskId: string;
body: JsonObject;
type: TaskEventType;
createdAt: string;
};
// @public
export type Status =
| 'open'
| 'processing'
| 'failed'
| 'cancelled'
| 'completed';
// @public
export interface TaskBroker {
// (undocumented)
claim(): Promise<TaskContext>;
// (undocumented)
dispatch(spec: TaskSpec, secrets?: TaskSecrets): Promise<DispatchResult>;
// (undocumented)
get(taskId: string): Promise<SerializedTask>;
// (undocumented)
observe(
options: {
taskId: string;
after: number | undefined;
},
callback: (
error: Error | undefined,
result: {
events: SerializedTaskEvent[];
},
) => void,
): {
unsubscribe: () => void;
};
// (undocumented)
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;
};
// @public
export type TaskSpec = TaskSpecV1beta2 | TaskSpecV1beta3;
// @public
export interface TaskSpecV1beta2 {
// (undocumented)
apiVersion: 'backstage.io/v1beta2';
// (undocumented)
baseUrl?: string;
// (undocumented)
output: {
[name: string]: string;
};
// (undocumented)
steps: Array<{
id: string;
name: string;
action: string;
input?: JsonObject;
if?: string | boolean;
}>;
// (undocumented)
values: JsonObject;
}
// @public
export interface TaskSpecV1beta3 {
// (undocumented)
apiVersion: 'scaffolder.backstage.io/v1beta3';
// (undocumented)
baseUrl?: string;
// (undocumented)
output: {
[name: string]: JsonValue;
};
// (undocumented)
parameters: JsonObject;
// Warning: (ae-forgotten-export) The symbol "TaskStep" needs to be exported by the entry point index.d.ts
//
// (undocumented)
steps: TaskStep[];
}
// @public
export interface TaskState {
// (undocumented)
secrets?: TaskSecrets;
// (undocumented)
spec: TaskSpec;
// (undocumented)
taskId: string;
}
// @public
export interface TaskStore {
// (undocumented)
claimTask(): Promise<SerializedTask | undefined>;
// (undocumented)
completeTask(options: {
taskId: string;
status: Status;
eventBody: JsonObject;
}): Promise<void>;
// (undocumented)
createTask(
task: TaskSpec,
secrets?: TaskSecrets,
): Promise<{
taskId: string;
}>;
// (undocumented)
emitLogEvent({ taskId, body }: TaskStoreEmitOptions): Promise<void>;
// (undocumented)
getTask(taskId: string): Promise<SerializedTask>;
// (undocumented)
heartbeatTask(taskId: string): Promise<void>;
// (undocumented)
listEvents({ taskId, after }: TaskStoreListEventsOptions): Promise<{
events: SerializedTaskEvent[];
}>;
// (undocumented)
listStaleTasks(options: { timeoutS: number }): Promise<{
tasks: {
taskId: string;
}[];
}>;
}
// @public
export type TaskStoreEmitOptions = {
taskId: string;
body: JsonObject;
};
// @public
export type TaskStoreListEventsOptions = {
taskId: string;
after?: number | undefined;
};
// @public
export class TaskWorker {
// (undocumented)
static create(options: CreateWorkerOptions): Promise<TaskWorker>;
// (undocumented)
runOneTask(task: TaskContext): Promise<void>;
// (undocumented)
start(): void;
}
// Warning: (ae-missing-release-tag) "TemplateAction" is exported by the package, but it is missing a release tag (@alpha, @beta, @public, or @internal)
//
// @public (undocumented)
@@ -6,7 +6,6 @@ metadata:
spec:
targets:
- ./remote-templates.yaml
# For local development of a template, you can reference your local templates here.
# Examples:
#
@@ -14,3 +14,4 @@
* limitations under the License.
*/
export * from './actions';
export * from './tasks';
@@ -20,15 +20,15 @@ import { ConflictError, NotFoundError } from '@backstage/errors';
import { Knex } from 'knex';
import { v4 as uuid } from 'uuid';
import {
DbTaskEventRow,
DbTaskRow,
SerializedTaskEvent,
SerializedTask,
Status,
TaskEventType,
TaskSecrets,
TaskSpec,
TaskStore,
TaskStoreEmitOptions,
TaskStoreGetEventsOptions,
TaskStoreListEventsOptions,
} from './types';
import { DateTime } from 'luxon';
@@ -54,17 +54,37 @@ export type RawDbTaskEventRow = {
created_at: string;
};
/**
* DatabaseTaskStore
*
* @public
*/
export type DatabaseTaskStoreOptions = {
database: Knex;
};
/**
* DatabaseTaskStore
*
* @public
*/
export class DatabaseTaskStore implements TaskStore {
static async create(knex: Knex): Promise<DatabaseTaskStore> {
await knex.migrate.latest({
private readonly db: Knex;
static async create(
options: DatabaseTaskStoreOptions,
): Promise<DatabaseTaskStore> {
await options.database.migrate.latest({
directory: migrationsDir,
});
return new DatabaseTaskStore(knex);
return new DatabaseTaskStore(options);
}
constructor(private readonly db: Knex) {}
constructor(options: DatabaseTaskStoreOptions) {
this.db = options.database;
}
async getTask(taskId: string): Promise<DbTaskRow> {
async getTask(taskId: string): Promise<SerializedTask> {
const [result] = await this.db<RawDbTaskRow>('tasks')
.where({ id: taskId })
.select();
@@ -101,7 +121,7 @@ export class DatabaseTaskStore implements TaskStore {
return { taskId };
}
async claimTask(): Promise<DbTaskRow | undefined> {
async claimTask(): Promise<SerializedTask | undefined> {
return this.db.transaction(async tx => {
const [task] = await tx<RawDbTaskRow>('tasks')
.where({
@@ -243,7 +263,7 @@ export class DatabaseTaskStore implements TaskStore {
async listEvents({
taskId,
after,
}: TaskStoreGetEventsOptions): Promise<{ events: DbTaskEventRow[] }> {
}: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }> {
const rawEvents = await this.db<RawDbTaskEventRow>('task_events')
.where({
task_id: taskId,
@@ -18,16 +18,16 @@ import mockFs from 'mock-fs';
import * as winston from 'winston';
import { getVoidLogger } from '@backstage/backend-common';
import { DefaultWorkflowRunner } from './DefaultWorkflowRunner';
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();
let actionRegistry = new TemplateActionRegistry();
let runner: DefaultWorkflowRunner;
let runner: NunjucksWorkflowRunner;
let fakeActionHandler: jest.Mock;
const integrations = ScmIntegrations.fromConfig(
@@ -38,7 +38,7 @@ describe('DefaultWorkflowRunner', () => {
}),
);
const createMockTaskWithSpec = (spec: TaskSpec): Task => ({
const createMockTaskWithSpec = (spec: TaskSpec): TaskContext => ({
spec,
complete: async () => {},
done: false,
@@ -88,7 +88,7 @@ describe('DefaultWorkflowRunner', () => {
},
});
runner = new DefaultWorkflowRunner({
runner = new NunjucksWorkflowRunner({
actionRegistry,
integrations,
workingDirectory: '/tmp',
@@ -14,7 +14,7 @@
* limitations under the License.
*/
import {
Task,
TaskContext,
WorkflowRunner,
WorkflowResponse,
TaskSpecV1beta2,
@@ -47,7 +47,7 @@ const isValidTaskSpec = (taskSpec: TaskSpec): taskSpec is TaskSpecV1beta2 =>
* This is the legacy workflow runner, which supports handlebars. This entire implementation will be replaced
* with the default workflow runner interface in the future so this entire thing can go bye bye.
*/
export class LegacyWorkflowRunner implements WorkflowRunner {
export class HandlebarsWorkflowRunner implements WorkflowRunner {
private readonly handlebars: typeof Handlebars;
constructor(private readonly options: Options) {
@@ -72,7 +72,7 @@ export class LegacyWorkflowRunner 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`);
}
@@ -20,12 +20,12 @@ import { createTemplateAction, TemplateActionRegistry } from '../actions';
import { ScmIntegrations } from '@backstage/integration';
import { ConfigReader } from '@backstage/config';
import { getVoidLogger } from '@backstage/backend-common';
import { LegacyWorkflowRunner } from './LegacyWorkflowRunner';
import { Task, TaskSpec } from './types';
import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner';
import { TaskContext, TaskSpec } from './types';
import { RepoSpec } from '../actions/builtin/publish/util';
describe('LegacyWorkflowRunner', () => {
let runner: LegacyWorkflowRunner;
let runner: HandlebarsWorkflowRunner;
const logger = getVoidLogger();
let actionRegistry = new TemplateActionRegistry();
@@ -37,7 +37,7 @@ describe('LegacyWorkflowRunner', () => {
}),
);
const createMockTaskWithSpec = (spec: TaskSpec): Task => ({
const createMockTaskWithSpec = (spec: TaskSpec): TaskContext => ({
spec,
complete: async () => {},
done: false,
@@ -60,7 +60,7 @@ describe('LegacyWorkflowRunner', () => {
},
});
runner = new LegacyWorkflowRunner({
runner = new HandlebarsWorkflowRunner({
actionRegistry,
integrations,
workingDirectory: '/tmp',
@@ -15,7 +15,7 @@
*/
import { ScmIntegrations } from '@backstage/integration';
import {
Task,
TaskContext,
TaskSpec,
TaskSpecV1beta3,
TaskStep,
@@ -34,7 +34,7 @@ import { validate as validateJsonSchema } from 'jsonschema';
import { parseRepoUrl } from '../actions/builtin/publish/util';
import { TemplateActionRegistry } from '../actions';
type Options = {
type NunjucksWorkflowRunnerOptions = {
workingDirectory: string;
actionRegistry: TemplateActionRegistry;
integrations: ScmIntegrations;
@@ -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',
@@ -77,7 +83,7 @@ const createStepLogger = ({ task, step }: { task: Task; step: TaskStep }) => {
return { taskLogger, streamLogger };
};
export class DefaultWorkflowRunner implements WorkflowRunner {
export class NunjucksWorkflowRunner implements WorkflowRunner {
private readonly nunjucks: nunjucks.Environment;
private readonly nunjucksOptions: nunjucks.ConfigureOptions = {
@@ -88,7 +94,7 @@ export class DefaultWorkflowRunner implements WorkflowRunner {
},
};
constructor(private readonly options: Options) {
constructor(private readonly options: NunjucksWorkflowRunnerOptions) {
this.nunjucks = nunjucks.configure(this.nunjucksOptions);
// TODO(blam): let's work out how we can deprecate these.
@@ -162,7 +168,7 @@ export class DefaultWorkflowRunner 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,8 +17,8 @@
import { getVoidLogger, DatabaseManager } from '@backstage/backend-common';
import { ConfigReader } from '@backstage/config';
import { DatabaseTaskStore } from './DatabaseTaskStore';
import { StorageTaskBroker, TaskAgent } from './StorageTaskBroker';
import { TaskSecrets, TaskSpec, DbTaskEventRow } from './types';
import { StorageTaskBroker, TaskManager } from './StorageTaskBroker';
import { TaskSecrets, TaskSpec, SerializedTaskEvent } from './types';
async function createStore(): Promise<DatabaseTaskStore> {
const manager = DatabaseManager.fromConfig(
@@ -31,7 +31,9 @@ async function createStore(): Promise<DatabaseTaskStore> {
},
}),
).forPlugin('scaffolder');
return await DatabaseTaskStore.create(await manager.getClient());
return await DatabaseTaskStore.create({
database: await manager.getClient(),
});
}
describe('StorageTaskBroker', () => {
@@ -46,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 () => {
@@ -56,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 () => {
@@ -68,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');
@@ -127,8 +129,8 @@ describe('StorageTaskBroker', () => {
const { taskId } = await broker1.dispatch({} as TaskSpec);
const logPromise = new Promise<DbTaskEventRow[]>(resolve => {
const observedEvents = new Array<DbTaskEventRow>();
const logPromise = new Promise<SerializedTaskEvent[]>(resolve => {
const observedEvents = new Array<SerializedTaskEvent>();
broker2.observe({ taskId, after: undefined }, (_err, { events }) => {
observedEvents.push(...events);
@@ -18,23 +18,28 @@ import { assertError } from '@backstage/errors';
import { Logger } from 'winston';
import {
CompletedTaskState,
Task,
TaskContext,
TaskSecrets,
TaskSpec,
TaskStore,
TaskBroker,
DispatchResult,
DbTaskEventRow,
DbTaskRow,
SerializedTaskEvent,
SerializedTask,
} from './types';
export class TaskAgent implements Task {
/**
* TaskManager
*
* @public
*/
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;
}
@@ -104,7 +109,12 @@ export class TaskAgent implements Task {
}
}
interface TaskState {
/**
* TaskState
*
* @public
*/
export interface TaskState {
spec: TaskSpec;
taskId: string;
secrets?: TaskSecrets;
@@ -125,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,
@@ -155,7 +165,7 @@ export class StorageTaskBroker implements TaskBroker {
};
}
async get(taskId: string): Promise<DbTaskRow> {
async get(taskId: string): Promise<SerializedTask> {
return this.storage.getTask(taskId);
}
@@ -166,9 +176,9 @@ export class StorageTaskBroker implements TaskBroker {
},
callback: (
error: Error | undefined,
result: { events: DbTaskEventRow[] },
result: { events: SerializedTaskEvent[] },
) => void,
): () => void {
): { unsubscribe: () => void } {
const { taskId } = options;
let cancelled = false;
@@ -195,7 +205,7 @@ export class StorageTaskBroker implements TaskBroker {
}
})();
return unsubscribe;
return { unsubscribe };
}
async vacuumTasks(timeoutS: { timeoutS: number }): Promise<void> {
@@ -19,8 +19,20 @@ import { ConfigReader } from '@backstage/config';
import { DatabaseTaskStore } from './DatabaseTaskStore';
import { StorageTaskBroker } from './StorageTaskBroker';
import { TaskWorker } from './TaskWorker';
import { WorkflowRunner } from './types';
import { LegacyWorkflowRunner } from './LegacyWorkflowRunner';
import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner';
import { ScmIntegrations } from '@backstage/integration';
import { TemplateActionRegistry } from '../actions';
import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner';
jest.mock('./HandlebarsWorkflowRunner');
const MockedHandlebarsWorkflowRunner =
HandlebarsWorkflowRunner as jest.Mock<HandlebarsWorkflowRunner>;
MockedHandlebarsWorkflowRunner.mockImplementation();
jest.mock('./NunjucksWorkflowRunner');
const MockedNunjucksWorkflowRunner =
NunjucksWorkflowRunner as jest.Mock<NunjucksWorkflowRunner>;
MockedNunjucksWorkflowRunner.mockImplementation();
async function createStore(): Promise<DatabaseTaskStore> {
const manager = DatabaseManager.fromConfig(
@@ -33,18 +45,26 @@ async function createStore(): Promise<DatabaseTaskStore> {
},
}),
).forPlugin('scaffolder');
return await DatabaseTaskStore.create(await manager.getClient());
return await DatabaseTaskStore.create({
database: await manager.getClient(),
});
}
describe('TaskWorker', () => {
let storage: DatabaseTaskStore;
const workflowRunner: WorkflowRunner = {
execute: jest.fn(),
} as unknown as WorkflowRunner;
const legacyWorkflowRunner: LegacyWorkflowRunner = {
const integrations: ScmIntegrations = {} as ScmIntegrations;
const actionRegistry: TemplateActionRegistry = {} as TemplateActionRegistry;
const workingDirectory = '/tmp/scaffolder';
const handlebarsWorkflowRunner: HandlebarsWorkflowRunner = {
execute: jest.fn(),
} as unknown as LegacyWorkflowRunner;
} as unknown as HandlebarsWorkflowRunner;
const workflowRunner: NunjucksWorkflowRunner = {
execute: jest.fn(),
} as unknown as NunjucksWorkflowRunner;
beforeAll(async () => {
storage = await createStore();
@@ -52,18 +72,22 @@ describe('TaskWorker', () => {
beforeEach(() => {
jest.resetAllMocks();
MockedHandlebarsWorkflowRunner.mockImplementation(
() => handlebarsWorkflowRunner,
);
MockedNunjucksWorkflowRunner.mockImplementation(() => workflowRunner);
});
const logger = getVoidLogger();
it('should call the legacy workflow runner when the apiVersion is not beta3', async () => {
const broker = new StorageTaskBroker(storage, logger);
const taskWorker = new TaskWorker({
const taskWorker = await TaskWorker.create({
logger,
workingDirectory,
integrations,
taskBroker: broker,
runners: {
legacyWorkflowRunner,
workflowRunner,
},
actionRegistry,
});
await broker.dispatch({
@@ -78,17 +102,23 @@ describe('TaskWorker', () => {
const task = await broker.claim();
await taskWorker.runOneTask(task);
expect(legacyWorkflowRunner.execute).toHaveBeenCalled();
expect(MockedHandlebarsWorkflowRunner).toBeCalledWith({
actionRegistry,
integrations,
logger,
workingDirectory,
});
expect(handlebarsWorkflowRunner.execute).toHaveBeenCalled();
});
it('should call the default workflow runner when the apiVersion is beta3', async () => {
const broker = new StorageTaskBroker(storage, logger);
const taskWorker = new TaskWorker({
const taskWorker = await TaskWorker.create({
logger,
workingDirectory,
integrations,
taskBroker: broker,
runners: {
legacyWorkflowRunner,
workflowRunner,
},
actionRegistry,
});
await broker.dispatch({
@@ -112,12 +142,12 @@ describe('TaskWorker', () => {
});
const broker = new StorageTaskBroker(storage, logger);
const taskWorker = new TaskWorker({
const taskWorker = await TaskWorker.create({
logger,
workingDirectory,
integrations,
taskBroker: broker,
runners: {
legacyWorkflowRunner,
workflowRunner,
},
actionRegistry,
});
const { taskId } = await broker.dispatch({
@@ -14,20 +14,77 @@
* limitations under the License.
*/
import { Task, TaskBroker, WorkflowRunner } from './types';
import { LegacyWorkflowRunner } from './LegacyWorkflowRunner';
import { TaskContext, TaskBroker, WorkflowRunner } from './types';
import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner';
import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner';
import { Logger } from 'winston';
import { TemplateActionRegistry } from '../actions';
import { ScmIntegrations } from '@backstage/integration';
import { assertError } from '@backstage/errors';
type Options = {
/**
* TaskWorkerOptions
*
* @public
*/
export type TaskWorkerOptions = {
taskBroker: TaskBroker;
runners: {
legacyWorkflowRunner: LegacyWorkflowRunner;
legacyWorkflowRunner: HandlebarsWorkflowRunner;
workflowRunner: WorkflowRunner;
};
};
/**
* CreateWorkerOptions
*
* @public
*/
export type CreateWorkerOptions = {
taskBroker: TaskBroker;
actionRegistry: TemplateActionRegistry;
integrations: ScmIntegrations;
workingDirectory: string;
logger: Logger;
};
/**
* TaskWorker
*
* @public
*/
export class TaskWorker {
constructor(private readonly options: Options) {}
private constructor(private readonly options: TaskWorkerOptions) {}
static async create(options: CreateWorkerOptions): Promise<TaskWorker> {
const {
taskBroker,
logger,
actionRegistry,
integrations,
workingDirectory,
} = options;
const legacyWorkflowRunner = new HandlebarsWorkflowRunner({
logger,
actionRegistry,
integrations,
workingDirectory,
});
const workflowRunner = new NunjucksWorkflowRunner({
actionRegistry,
integrations,
logger,
workingDirectory,
});
return new TaskWorker({
taskBroker: taskBroker,
runners: { legacyWorkflowRunner, workflowRunner },
});
}
start() {
(async () => {
for (;;) {
@@ -37,7 +94,7 @@ export class TaskWorker {
})();
}
async runOneTask(task: Task) {
async runOneTask(task: TaskContext) {
try {
const { output } =
task.spec.apiVersion === 'scaffolder.backstage.io/v1beta3'
@@ -13,7 +13,25 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
export { DatabaseTaskStore } from './DatabaseTaskStore';
export { StorageTaskBroker } from './StorageTaskBroker';
export { TaskManager } from './StorageTaskBroker';
export type { TaskState } from './StorageTaskBroker';
export { TaskWorker } from './TaskWorker';
export type { CreateWorkerOptions } from './TaskWorker';
export type {
TaskSecrets,
TaskSpec,
CompletedTaskState,
TaskStoreEmitOptions,
TaskStoreListEventsOptions,
SerializedTask,
SerializedTaskEvent,
TaskSpecV1beta2,
TaskSpecV1beta3,
Status,
TaskEventType,
TaskBroker,
TaskContext,
TaskStore,
DispatchResult,
} from './types';
@@ -16,6 +16,11 @@
import { JsonValue, JsonObject } from '@backstage/types';
/**
* Status
*
* @public
*/
export type Status =
| 'open'
| 'processing'
@@ -23,9 +28,19 @@ export type Status =
| 'cancelled'
| 'completed';
/**
* CompletedTaskState
*
* @public
*/
export type CompletedTaskState = 'failed' | 'completed';
export type DbTaskRow = {
/**
* SerializedTask
*
* @public
*/
export type SerializedTask = {
id: string;
spec: TaskSpec;
status: Status;
@@ -34,8 +49,19 @@ export type DbTaskRow = {
secrets?: TaskSecrets;
};
/**
* TaskEventType
*
* @public
*/
export type TaskEventType = 'completion' | 'log';
export type DbTaskEventRow = {
/**
* SerializedTaskEvent
*
* @public
*/
export type SerializedTaskEvent = {
id: number;
taskId: string;
body: JsonObject;
@@ -43,6 +69,11 @@ export type DbTaskEventRow = {
createdAt: string;
};
/**
* TaskSpecV1beta2
*
* @public
*/
export interface TaskSpecV1beta2 {
apiVersion: 'backstage.io/v1beta2';
baseUrl?: string;
@@ -64,6 +95,12 @@ export interface TaskStep {
input?: JsonObject;
if?: string | boolean;
}
/**
* TaskSpecV1beta3
*
* @public
*/
export interface TaskSpecV1beta3 {
apiVersion: 'scaffolder.backstage.io/v1beta3';
baseUrl?: string;
@@ -72,17 +109,37 @@ export interface TaskSpecV1beta3 {
output: { [name: string]: JsonValue };
}
/**
* TaskSpec
*
* @public
*/
export type TaskSpec = TaskSpecV1beta2 | TaskSpecV1beta3;
/**
* TaskSecrets
*
* @public
*/
export type TaskSecrets = {
token: string | undefined;
};
/**
* DispatchResult
*
* @public
*/
export type DispatchResult = {
taskId: string;
};
export interface Task {
/**
* Task
*
* @public
*/
export interface TaskContext {
spec: TaskSpec;
secrets?: TaskSecrets;
done: boolean;
@@ -91,8 +148,13 @@ export interface Task {
getWorkspaceName(): Promise<string>;
}
/**
* TaskBroker
*
* @public
*/
export interface TaskBroker {
claim(): Promise<Task>;
claim(): Promise<TaskContext>;
dispatch(spec: TaskSpec, secrets?: TaskSecrets): Promise<DispatchResult>;
vacuumTasks(timeoutS: { timeoutS: number }): Promise<void>;
observe(
@@ -102,28 +164,44 @@ export interface TaskBroker {
},
callback: (
error: Error | undefined,
result: { events: DbTaskEventRow[] },
result: { events: SerializedTaskEvent[] },
) => void,
): () => void;
): { unsubscribe: () => void };
get(taskId: string): Promise<SerializedTask>;
}
/**
* TaskStoreEmitOptions
*
* @public
*/
export type TaskStoreEmitOptions = {
taskId: string;
body: JsonObject;
};
export type TaskStoreGetEventsOptions = {
/**
* TaskStoreListEventsOptions
*
* @public
*/
export type TaskStoreListEventsOptions = {
taskId: string;
after?: number | undefined;
};
/**
* TaskStore
*
* @public
*/
export interface TaskStore {
createTask(
task: TaskSpec,
secrets?: TaskSecrets,
): Promise<{ taskId: string }>;
getTask(taskId: string): Promise<DbTaskRow>;
claimTask(): Promise<DbTaskRow | undefined>;
getTask(taskId: string): Promise<SerializedTask>;
claimTask(): Promise<SerializedTask | undefined>;
completeTask(options: {
taskId: string;
status: Status;
@@ -138,10 +216,10 @@ export interface TaskStore {
listEvents({
taskId,
after,
}: TaskStoreGetEventsOptions): Promise<{ events: DbTaskEventRow[] }>;
}: TaskStoreListEventsOptions): Promise<{ events: SerializedTaskEvent[] }>;
}
export type WorkflowResponse = { output: { [key: string]: JsonValue } };
export interface WorkflowRunner {
execute(task: Task): Promise<WorkflowResponse>;
execute(task: TaskContext): Promise<WorkflowResponse>;
}
@@ -40,7 +40,13 @@ import { ConfigReader } from '@backstage/config';
import express from 'express';
import request from 'supertest';
import { TemplateEntityV1beta2 } from '@backstage/catalog-model';
import { createRouter } from './router';
/**
* TODO: The following should import directly from the router file.
* Due to a circular dependency between this plugin and the
* plugin-scaffolder-backend-module-cookiecutter plugin, it results in an error:
* TypeError: _pluginscaffolderbackend.createTemplateAction is not a function
*/
import { createRouter } from '../index';
const createCatalogClient = (templates: any[] = []) =>
({
@@ -22,10 +22,12 @@ import { CatalogEntityClient } from '../lib/catalog';
import { validate } from 'jsonschema';
import {
DatabaseTaskStore,
StorageTaskBroker,
TemplateActionRegistry,
TaskWorker,
} from '../scaffolder/tasks';
import { TemplateActionRegistry } from '../scaffolder/actions/TemplateActionRegistry';
TemplateAction,
createBuiltinActions,
} from '../scaffolder';
import { StorageTaskBroker } from '../scaffolder/tasks/StorageTaskBroker';
import { getEntityBaseUrl, getWorkingDirectory } from './helpers';
import {
ContainerRunner,
@@ -38,12 +40,13 @@ import { TemplateEntityV1beta2, Entity } from '@backstage/catalog-model';
import { TemplateEntityV1beta3 } from '@backstage/plugin-scaffolder-common';
import { ScmIntegrations } from '@backstage/integration';
import { TemplateAction } from '../scaffolder/actions';
import { createBuiltinActions } from '../scaffolder/actions/builtin/createBuiltinActions';
import { LegacyWorkflowRunner } from '../scaffolder/tasks/LegacyWorkflowRunner';
import { DefaultWorkflowRunner } from '../scaffolder/tasks/DefaultWorkflowRunner';
import { TaskSpec } from '../scaffolder/tasks/types';
import { TaskBroker, TaskSpec } from '../scaffolder/tasks/types';
/**
* RouterOptions
*
* @public
*/
export interface RouterOptions {
logger: Logger;
config: Config;
@@ -53,6 +56,7 @@ export interface RouterOptions {
actions?: TemplateAction<any>[];
taskWorkers?: number;
containerRunner: ContainerRunner;
taskBroker?: TaskBroker;
}
function isSupportedTemplate(
@@ -85,34 +89,27 @@ export async function createRouter(
const workingDirectory = await getWorkingDirectory(config, logger);
const entityClient = new CatalogEntityClient(catalogClient);
const integrations = ScmIntegrations.fromConfig(config);
let taskBroker: TaskBroker;
if (!options.taskBroker) {
const databaseTaskStore = await DatabaseTaskStore.create({
database: await database.getClient(),
});
taskBroker = new StorageTaskBroker(databaseTaskStore, logger);
} else {
taskBroker = options.taskBroker;
}
const databaseTaskStore = await DatabaseTaskStore.create(
await database.getClient(),
);
const taskBroker = new StorageTaskBroker(databaseTaskStore, logger);
const actionRegistry = new TemplateActionRegistry();
const legacyWorkflowRunner = new LegacyWorkflowRunner({
logger,
actionRegistry,
integrations,
workingDirectory,
});
const workflowRunner = new DefaultWorkflowRunner({
actionRegistry,
integrations,
logger,
workingDirectory,
});
const workers = [];
for (let i = 0; i < (taskWorkers || 1); i++) {
const worker = new TaskWorker({
const worker = await TaskWorker.create({
taskBroker,
runners: {
legacyWorkflowRunner,
workflowRunner,
},
actionRegistry,
integrations,
logger,
workingDirectory,
});
workers.push(worker);
}
@@ -261,7 +258,7 @@ export async function createRouter(
});
// After client opens connection send all events as string
const unsubscribe = taskBroker.observe(
const { unsubscribe } = taskBroker.observe(
{ taskId, after },
(error, { events }) => {
if (error) {