diff --git a/.changeset/clean-suits-hear.md b/.changeset/clean-suits-hear.md new file mode 100644 index 0000000000..911a46a55a --- /dev/null +++ b/.changeset/clean-suits-hear.md @@ -0,0 +1,6 @@ +--- +'@backstage/plugin-scaffolder-backend': patch +--- + +- **DEPRECATED** - Deprecated the `runCommand` export in favour of `executeShellCommand`. Please migrate to using the new method. +- Added a type parameter to `TaskStoreEmitOptions` to type the `body` property diff --git a/.changeset/six-cobras-allow.md b/.changeset/six-cobras-allow.md new file mode 100644 index 0000000000..db34ecb1f9 --- /dev/null +++ b/.changeset/six-cobras-allow.md @@ -0,0 +1,6 @@ +--- +'@backstage/plugin-scaffolder-backend': minor +--- + +- **BREAKING** - Removed the re-export of types `TaskSpec` `TaskSpecV1Beta2` and `TaskSpecV1Beta3` these should now be import from `@backstage/plugin-scaffolder-common` directly. +- **BREAKING** - Removed the `observe` method from the `TaskBroker` interface, this has now been replaced with an `Observable` implementation under `event$`. diff --git a/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.test.ts b/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.test.ts index 2c1ecd06f9..18a6722663 100644 --- a/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.test.ts +++ b/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.test.ts @@ -13,14 +13,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -const runCommand = jest.fn(); +const executeShellCommand = jest.fn(); const commandExists = jest.fn(); const fetchContents = jest.fn(); jest.mock('@backstage/plugin-scaffolder-backend', () => ({ ...jest.requireActual('@backstage/plugin-scaffolder-backend'), fetchContents, - runCommand, + executeShellCommand, })); jest.mock('command-exists', () => commandExists); @@ -115,8 +115,8 @@ describe('fetch:cookiecutter', () => { }); }); - // Mock when runCommand is called it creats some new files in the mock filesystem - runCommand.mockImplementation(async () => { + // Mock when executeShellCommand is called it creats some new files in the mock filesystem + executeShellCommand.mockImplementation(async () => { mockFs({ [`${join(mockTmpDir, 'intermediate')}`]: { 'testfile.json': '{}', @@ -163,12 +163,12 @@ describe('fetch:cookiecutter', () => { ); }); - it('should call out to cookiecutter using runCommand when cookiecutter is installed', async () => { + it('should call out to cookiecutter using executeShellCommand when cookiecutter is installed', async () => { commandExists.mockResolvedValue(true); await action.handler(mockContext); - expect(runCommand).toHaveBeenCalledWith( + expect(executeShellCommand).toHaveBeenCalledWith( expect.objectContaining({ command: 'cookiecutter', args: [ diff --git a/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.ts b/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.ts index 0455f01989..6b0cb4695d 100644 --- a/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.ts +++ b/plugins/scaffolder-backend-module-cookiecutter/src/actions/fetch/cookiecutter.ts @@ -27,9 +27,9 @@ import fs from 'fs-extra'; import path, { resolve as resolvePath } from 'path'; import { Writable } from 'stream'; import { - runCommand, createTemplateAction, fetchContents, + executeShellCommand, } from '@backstage/plugin-scaffolder-backend'; export class CookiecutterRunner { @@ -89,7 +89,7 @@ export class CookiecutterRunner { () => false, ); if (cookieCutterInstalled) { - await runCommand({ + await executeShellCommand({ command: 'cookiecutter', args: ['--no-input', '-o', intermediateDir, templateDir, '--verbose'], logStream, diff --git a/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.test.ts b/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.test.ts index 03d894f832..0767e63f29 100644 --- a/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.test.ts +++ b/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.test.ts @@ -14,10 +14,12 @@ * limitations under the License. */ -const runCommand = jest.fn(); +const executeShellCommand = jest.fn(); const commandExists = jest.fn(); -jest.mock('@backstage/plugin-scaffolder-backend', () => ({ runCommand })); +jest.mock('@backstage/plugin-scaffolder-backend', () => ({ + executeShellCommand, +})); jest.mock('command-exists', () => commandExists); jest.mock('fs-extra'); @@ -195,7 +197,7 @@ describe('Rails Templater', () => { logStream: stream, }); - expect(runCommand).toHaveBeenCalledWith({ + expect(executeShellCommand).toHaveBeenCalledWith({ command: 'rails', args: expect.arrayContaining([ 'new', @@ -225,7 +227,7 @@ describe('Rails Templater', () => { logStream: stream, }); - expect(runCommand).toHaveBeenCalledWith({ + expect(executeShellCommand).toHaveBeenCalledWith({ command: 'rails', args: expect.arrayContaining([ 'new', diff --git a/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.ts b/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.ts index f7817c7fb5..4a4672795a 100644 --- a/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.ts +++ b/plugins/scaffolder-backend-module-rails/src/actions/fetch/rails/railsNewRunner.ts @@ -17,7 +17,7 @@ import { ContainerRunner } from '@backstage/backend-common'; import fs from 'fs-extra'; import path from 'path'; -import { runCommand } from '@backstage/plugin-scaffolder-backend'; +import { executeShellCommand } from '@backstage/plugin-scaffolder-backend'; import commandExists from 'command-exists'; import { railsArgumentResolver, @@ -64,7 +64,7 @@ export class RailsNewRunner { railsArguments as RailsRunOptions, ); - await runCommand({ + await executeShellCommand({ command: baseCommand, args: [ ...baseArguments, diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 1c627759df..54d5bf0280 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -20,6 +20,7 @@ import { JsonValue } from '@backstage/types'; import { Knex } from 'knex'; import { LocationSpec } from '@backstage/catalog-model'; import { Logger as Logger_2 } from 'winston'; +import { Observable } from '@backstage/types'; import { Octokit } from 'octokit'; import { PluginDatabaseManager } from '@backstage/backend-common'; import { Schema } from 'jsonschema'; @@ -27,8 +28,6 @@ import { ScmIntegrationRegistry } from '@backstage/integration'; import { ScmIntegrations } from '@backstage/integration'; import { SpawnOptionsWithoutStdio } from 'child_process'; import { TaskSpec } from '@backstage/plugin-scaffolder-common'; -import { TaskSpecV1beta2 } from '@backstage/plugin-scaffolder-common'; -import { TaskSpecV1beta3 } from '@backstage/plugin-scaffolder-common'; import { TemplateInfo } from '@backstage/plugin-scaffolder-common'; import { TemplateMetadata } from '@backstage/plugin-scaffolder-common'; import { UrlReader } from '@backstage/backend-common'; @@ -355,7 +354,13 @@ export class DatabaseTaskStore implements TaskStore { options: TaskStoreCreateTaskOptions, ): Promise; // (undocumented) - emitLogEvent({ taskId, body }: TaskStoreEmitOptions): Promise; + emitLogEvent( + options: TaskStoreEmitOptions< + { + message: string; + } & JsonObject + >, + ): Promise; // (undocumented) getTask(taskId: string): Promise; // (undocumented) @@ -375,6 +380,11 @@ export class DatabaseTaskStore implements TaskStore { // @public @deprecated export type DispatchResult = TaskBrokerDispatchResult; +// Warning: (ae-forgotten-export) The symbol "RunCommandOptions" needs to be exported by the entry point index.d.ts +// +// @public +export const executeShellCommand: (options: RunCommandOptions) => Promise; + // @public export function fetchContents({ reader, @@ -433,16 +443,8 @@ export interface RouterOptions { taskWorkers?: number; } -// Warning: (ae-forgotten-export) The symbol "RunCommandOptions" needs to be exported by the entry point index.d.ts -// Warning: (ae-missing-release-tag) "runCommand" is exported by the package, but it is missing a release tag (@alpha, @beta, @public, or @internal) -// -// @public -export const runCommand: ({ - command, - args, - logStream, - options, -}: RunCommandOptions) => Promise; +// @public @deprecated +export const runCommand: (options: RunCommandOptions) => Promise; // @public (undocumented) export class ScaffolderEntitiesProcessor implements CatalogProcessor { @@ -489,22 +491,11 @@ export interface TaskBroker { options: TaskBrokerDispatchOptions, ): Promise; // (undocumented) - get(taskId: string): Promise; + event$(options: { taskId: string; after: number | undefined }): Observable<{ + events: SerializedTaskEvent[]; + }>; // (undocumented) - observe( - options: { - taskId: string; - after: number | undefined; - }, - callback: ( - error: Error | undefined, - result: { - events: SerializedTaskEvent[]; - }, - ) => void, - ): { - unsubscribe: () => void; - }; + get(taskId: string): Promise; // (undocumented) vacuumTasks(options: { timeoutS: number }): Promise; } @@ -569,12 +560,6 @@ export type TaskSecrets = Record & { backstageToken?: string; }; -export { TaskSpec }; - -export { TaskSpecV1beta2 }; - -export { TaskSpecV1beta3 }; - // @public @deprecated export type TaskState = CurrentClaimedTask; @@ -630,9 +615,9 @@ export type TaskStoreCreateTaskResult = { }; // @public -export type TaskStoreEmitOptions = { +export type TaskStoreEmitOptions = { taskId: string; - body: JsonObject; + body: TBody; }; // @public @@ -680,6 +665,4 @@ export class TemplateActionRegistry { // // @public (undocumented) export type TemplateFilter = (...args: JsonValue[]) => JsonValue | undefined; - -export { TemplateMetadata }; ``` diff --git a/plugins/scaffolder-backend/package.json b/plugins/scaffolder-backend/package.json index 8c61ee65c4..5d23ea0aff 100644 --- a/plugins/scaffolder-backend/package.json +++ b/plugins/scaffolder-backend/package.json @@ -72,7 +72,8 @@ "uuid": "^8.2.0", "winston": "^3.2.1", "yaml": "^1.10.0", - "vm2": "^3.9.6" + "vm2": "^3.9.6", + "zen-observable": "^0.8.15" }, "devDependencies": { "@backstage/cli": "^0.14.0", @@ -83,6 +84,7 @@ "@types/mock-fs": "^4.13.0", "@types/nunjucks": "^3.1.4", "@types/supertest": "^2.0.8", + "@types/zen-observable": "^0.8.0", "esbuild": "^0.14.1", "jest-when": "^3.1.0", "mock-fs": "^5.1.0", diff --git a/plugins/scaffolder-backend/src/scaffolder/actions/builtin/helpers.ts b/plugins/scaffolder-backend/src/scaffolder/actions/builtin/helpers.ts index 9303b5e487..7ee334a3f7 100644 --- a/plugins/scaffolder-backend/src/scaffolder/actions/builtin/helpers.ts +++ b/plugins/scaffolder-backend/src/scaffolder/actions/builtin/helpers.ts @@ -34,15 +34,18 @@ export type RunCommandOptions = { /** * Run a command in a sub-process, normally a shell command. + * + * @public */ -export const runCommand = async ({ - command, - args, - logStream = new PassThrough(), - options, -}: RunCommandOptions) => { +export const executeShellCommand = async (options: RunCommandOptions) => { + const { + command, + args, + options: spawnOptions, + logStream = new PassThrough(), + } = options; await new Promise((resolve, reject) => { - const process = spawn(command, args, options); + const process = spawn(command, args, spawnOptions); process.stdout.on('data', stream => { logStream.write(stream); @@ -67,6 +70,13 @@ export const runCommand = async ({ }); }; +/** + * Run a command in a sub-process, normally a shell command. + * @public + * @deprecated use {@link executeShellCommand} instead + */ +export const runCommand = executeShellCommand; + export async function initRepoAndPush({ dir, remoteUrl, diff --git a/plugins/scaffolder-backend/src/scaffolder/actions/builtin/index.ts b/plugins/scaffolder-backend/src/scaffolder/actions/builtin/index.ts index 5f4c04cc7d..8930c23126 100644 --- a/plugins/scaffolder-backend/src/scaffolder/actions/builtin/index.ts +++ b/plugins/scaffolder-backend/src/scaffolder/actions/builtin/index.ts @@ -25,4 +25,4 @@ export * from './github'; /** @deprecated please add this package to your own installation manually */ export { createFetchCookiecutterAction } from '@backstage/plugin-scaffolder-backend-module-cookiecutter'; -export { runCommand } from './helpers'; +export { runCommand, executeShellCommand } from './helpers'; diff --git a/plugins/scaffolder-backend/src/scaffolder/actions/types.ts b/plugins/scaffolder-backend/src/scaffolder/actions/types.ts index b2227426b8..1aa4f81bdd 100644 --- a/plugins/scaffolder-backend/src/scaffolder/actions/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/actions/types.ts @@ -18,8 +18,11 @@ import { Logger } from 'winston'; import { Writable } from 'stream'; import { JsonValue, JsonObject } from '@backstage/types'; import { Schema } from 'jsonschema'; -import { TaskSecrets, TemplateMetadata } from '../tasks/types'; -import { TemplateInfo } from '@backstage/plugin-scaffolder-common'; +import { TaskSecrets } from '../tasks/types'; +import { + TemplateInfo, + TemplateMetadata, +} from '@backstage/plugin-scaffolder-common'; export type ActionContext = { /** diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 389ecae5c6..65b890c265 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -252,12 +252,15 @@ export class DatabaseTaskStore implements TaskStore { }); } - async emitLogEvent({ taskId, body }: TaskStoreEmitOptions): Promise { - const serliazedBody = JSON.stringify(body); + async emitLogEvent( + options: TaskStoreEmitOptions<{ message: string } & JsonObject>, + ): Promise { + const { taskId, body } = options; + const serializedBody = JSON.stringify(body); await this.db('task_events').insert({ task_id: taskId, event_type: 'log', - body: serliazedBody, + body: serializedBody, }); } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.test.ts index 9f4135b3d1..dcad1fc5d0 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.test.ts @@ -21,8 +21,9 @@ import { ScmIntegrations } from '@backstage/integration'; import { ConfigReader } from '@backstage/config'; import { getVoidLogger } from '@backstage/backend-common'; import { HandlebarsWorkflowRunner } from './HandlebarsWorkflowRunner'; -import { TaskContext, TaskSpec } from './types'; +import { TaskContext } from './types'; import { RepoSpec } from '../actions/builtin/publish/util'; +import { TaskSpec } from '@backstage/plugin-scaffolder-common'; describe('LegacyWorkflowRunner', () => { let runner: HandlebarsWorkflowRunner; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.ts index e525e6bd14..c82da51ead 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/HandlebarsWorkflowRunner.ts @@ -13,13 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -import { - TaskContext, - WorkflowRunner, - WorkflowResponse, - TaskSpecV1beta2, - TaskSpec, -} from './types'; +import { TaskContext, WorkflowRunner, WorkflowResponse } from './types'; import * as Handlebars from 'handlebars'; import { TemplateActionRegistry } from '..'; import { ScmIntegrations } from '@backstage/integration'; @@ -33,6 +27,7 @@ import fs from 'fs-extra'; import { validate as validateJsonSchema } from 'jsonschema'; import { JsonObject, JsonValue } from '@backstage/types'; import { InputError } from '@backstage/errors'; +import { TaskSpec, TaskSpecV1beta2 } from '@backstage/plugin-scaffolder-common'; type Options = { workingDirectory: string; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.test.ts index b66c69ee17..1b635b695c 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.test.ts @@ -22,7 +22,8 @@ import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner'; import { TemplateActionRegistry } from '../actions'; import { ScmIntegrations } from '@backstage/integration'; import { ConfigReader } from '@backstage/config'; -import { TaskContext, TaskSpec, TaskSecrets } from './types'; +import { TaskContext, TaskSecrets } from './types'; +import { TaskSpec } from '@backstage/plugin-scaffolder-common'; const realFiles = Object.fromEntries( [ diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts index b1177ca49c..fb9da56c22 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts @@ -15,14 +15,7 @@ */ import { ScmIntegrations } from '@backstage/integration'; -import { - TaskContext, - TaskSpec, - TaskSpecV1beta3, - TaskStep, - WorkflowResponse, - WorkflowRunner, -} from './types'; +import { TaskContext, WorkflowResponse, WorkflowRunner } from './types'; import * as winston from 'winston'; import fs from 'fs-extra'; import path from 'path'; @@ -39,6 +32,11 @@ import { SecureTemplater, SecureTemplateRenderer, } from '../../lib/templating/SecureTemplater'; +import { + TaskSpec, + TaskSpecV1beta3, + TaskStep, +} from '@backstage/plugin-scaffolder-common'; type NunjucksWorkflowRunnerOptions = { workingDirectory: string; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index daf444cdd8..52fd175dc8 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -16,9 +16,10 @@ import { getVoidLogger, DatabaseManager } from '@backstage/backend-common'; import { ConfigReader } from '@backstage/config'; +import { TaskSpec } from '@backstage/plugin-scaffolder-common'; import { DatabaseTaskStore } from './DatabaseTaskStore'; import { StorageTaskBroker, TaskManager } from './StorageTaskBroker'; -import { TaskSecrets, TaskSpec, SerializedTaskEvent } from './types'; +import { TaskSecrets, SerializedTaskEvent } from './types'; async function createStore(): Promise { const manager = DatabaseManager.fromConfig( @@ -124,7 +125,7 @@ describe('StorageTaskBroker', () => { const logPromise = new Promise(resolve => { const observedEvents = new Array(); - broker2.observe({ taskId, after: undefined }, (_err, { events }) => { + broker2.event$({ taskId, after: undefined }).subscribe(({ events }) => { observedEvents.push(...events); if (events.some(e => e.type === 'completion')) { resolve(observedEvents); @@ -146,9 +147,11 @@ describe('StorageTaskBroker', () => { ]); const afterLogs = await new Promise(resolve => { - broker2.observe({ taskId, after: logs[1].id }, (_err, { events }) => - resolve(events.map(e => e.body.message as string)), - ); + broker2 + .event$({ taskId, after: logs[1].id }) + .subscribe(({ events }) => + resolve(events.map(e => e.body.message as string)), + ); }); expect(afterLogs).toEqual([ 'log 3', diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 898516ef3b..0fd593f2f3 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -13,14 +13,14 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -import { JsonObject } from '@backstage/types'; -import { assertError } from '@backstage/errors'; +import { JsonObject, Observable } from '@backstage/types'; +import ObservableImpl from 'zen-observable'; +import { TaskSpec } from '@backstage/plugin-scaffolder-common'; import { Logger } from 'winston'; import { TaskCompletionState, TaskContext, TaskSecrets, - TaskSpec, TaskStore, TaskBroker, SerializedTaskEvent, @@ -176,43 +176,33 @@ export class StorageTaskBroker implements TaskBroker { return this.storage.getTask(taskId); } - observe( - options: { - taskId: string; - after: number | undefined; - }, - callback: ( - error: Error | undefined, - result: { events: SerializedTaskEvent[] }, - ) => void, - ): { unsubscribe: () => void } { - const { taskId } = options; + event$(options: { + taskId: string; + after?: number; + }): Observable<{ events: SerializedTaskEvent[] }> { + return new ObservableImpl(observer => { + const { taskId } = options; - let cancelled = false; - const unsubscribe = () => { - cancelled = true; - }; - - (async () => { let after = options.after; - while (!cancelled) { - const result = await this.storage.listEvents({ taskId, after: after }); - const { events } = result; - if (events.length) { - after = events[events.length - 1].id; - try { - callback(undefined, result); - } catch (error) { - assertError(error); - callback(error, { events: [] }); + let cancelled = false; + + (async () => { + while (!cancelled) { + const result = await this.storage.listEvents({ taskId, after }); + const { events } = result; + if (events.length) { + after = events[events.length - 1].id; + observer.next(result); } + + await new Promise(resolve => setTimeout(resolve, 1000)); } + })(); - await new Promise(resolve => setTimeout(resolve, 1000)); - } - })(); - - return { unsubscribe }; + return () => { + cancelled = true; + }; + }); } async vacuumTasks(options: { timeoutS: number }): Promise { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts index aa1640dae2..a2bf5f0e1c 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts @@ -20,15 +20,12 @@ export { TaskWorker } from './TaskWorker'; export type { CreateWorkerOptions } from './TaskWorker'; export type { TaskSecrets, - TaskSpec, TaskCompletionState, CompletedTaskState, TaskStoreEmitOptions, TaskStoreListEventsOptions, SerializedTask, SerializedTaskEvent, - TaskSpecV1beta2, - TaskSpecV1beta3, Status, TaskStatus, TaskEventType, @@ -40,5 +37,4 @@ export type { TaskBrokerDispatchOptions, TaskStoreCreateTaskOptions, TaskStoreCreateTaskResult, - TemplateMetadata, } from './types'; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index e17f8b2c3b..57b92e0222 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -14,22 +14,8 @@ * limitations under the License. */ -import { JsonValue, JsonObject } from '@backstage/types'; -import { - TaskSpec, - TaskStep, - TemplateMetadata, - TaskSpecV1beta2, - TaskSpecV1beta3, -} from '@backstage/plugin-scaffolder-common'; - -export type { - TaskSpec, - TaskStep, - TemplateMetadata, - TaskSpecV1beta2, - TaskSpecV1beta3, -}; +import { JsonValue, JsonObject, Observable } from '@backstage/types'; +import { TaskSpec } from '@backstage/plugin-scaffolder-common'; /** * The status of each step of the Task @@ -162,16 +148,10 @@ export interface TaskBroker { options: TaskBrokerDispatchOptions, ): Promise; vacuumTasks(options: { timeoutS: number }): Promise; - observe( - options: { - taskId: string; - after: number | undefined; - }, - callback: ( - error: Error | undefined, - result: { events: SerializedTaskEvent[] }, - ) => void, - ): { unsubscribe: () => void }; + event$(options: { + taskId: string; + after: number | undefined; + }): Observable<{ events: SerializedTaskEvent[] }>; get(taskId: string): Promise; } @@ -180,9 +160,9 @@ export interface TaskBroker { * * @public */ -export type TaskStoreEmitOptions = { +export type TaskStoreEmitOptions = { taskId: string; - body: JsonObject; + body: TBody; }; /** diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index b20dfa5d24..f8336179eb 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -38,6 +38,7 @@ import { import { CatalogApi } from '@backstage/catalog-client'; import { TemplateEntityV1beta2 } from '@backstage/plugin-scaffolder-common'; import { ConfigReader } from '@backstage/config'; +import ObservableImpl from 'zen-observable'; import express from 'express'; import request from 'supertest'; /** @@ -113,7 +114,7 @@ describe('createRouter', () => { jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); - jest.spyOn(taskBroker, 'observe'); + jest.spyOn(taskBroker, 'event$'); const router = await createRouter({ logger: getVoidLogger(), @@ -194,37 +195,38 @@ describe('createRouter', () => { describe('GET /v2/tasks/:taskId/eventstream', () => { it('should return log messages', async () => { - const unsubscribe = jest.fn(); + let subscriber: ZenObservable.SubscriptionObserver; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation(({ taskId }, callback) => { - // emit after this function returned - setImmediate(() => { - callback(undefined, { - events: [ - { - id: 0, - taskId, - type: 'log', - createdAt: '', - body: { message: 'My log message' }, - }, - ], - }); - callback(undefined, { - events: [ - { - id: 1, - taskId, - type: 'completion', - createdAt: '', - body: { message: 'Finished!' }, - }, - ], + taskBroker.event$ as jest.Mocked['event$'] + ).mockImplementation(({ taskId }) => { + return new ObservableImpl(observer => { + subscriber = observer; + setImmediate(() => { + observer.next({ + events: [ + { + id: 0, + taskId, + type: 'log', + createdAt: '', + body: { message: 'My log message' }, + }, + ], + }); + observer.next({ + events: [ + { + id: 1, + taskId, + type: 'completion', + createdAt: '', + body: { message: 'Finished!' }, + }, + ], + }); }); }); - - return { unsubscribe }; + // emit after this function returned }); let statusCode: any = undefined; @@ -264,34 +266,32 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ `); - expect(taskBroker.observe).toBeCalledTimes(1); - expect(taskBroker.observe).toBeCalledWith( - { taskId: 'a-random-id' }, - expect.any(Function), - ); - - expect(unsubscribe).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledWith({ taskId: 'a-random-id' }); + expect(subscriber!.closed).toBe(true); }); it('should return log messages with after query', async () => { - const unsubscribe = jest.fn(); + let subscriber: ZenObservable.SubscriptionObserver; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation(({ taskId }, callback) => { - setImmediate(() => { - callback(undefined, { - events: [ - { - id: 1, - taskId, - type: 'completion', - createdAt: '', - body: { message: 'Finished!' }, - }, - ], + taskBroker.event$ as jest.Mocked['event$'] + ).mockImplementation(({ taskId }) => { + return new ObservableImpl(observer => { + subscriber = observer; + setImmediate(() => { + observer.next({ + events: [ + { + id: 1, + taskId, + type: 'completion', + createdAt: '', + body: { message: 'Finished!' }, + }, + ], + }); }); }); - return { unsubscribe }; }); let statusCode: any = undefined; @@ -318,41 +318,43 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ expect(statusCode).toBe(200); expect(headers['content-type']).toBe('text/event-stream'); - expect(taskBroker.observe).toBeCalledTimes(1); - expect(taskBroker.observe).toBeCalledWith( - { taskId: 'a-random-id', after: 10 }, - expect.any(Function), - ); + expect(taskBroker.event$).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledWith({ + taskId: 'a-random-id', + after: 10, + }); - expect(unsubscribe).toBeCalledTimes(1); + expect(subscriber!.closed).toBe(true); }); }); describe('GET /v2/tasks/:taskId/events', () => { it('should return log messages', async () => { - const unsubscribe = jest.fn(); + let subscriber: ZenObservable.SubscriptionObserver; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation(({ taskId }, callback) => { - callback(undefined, { - events: [ - { - id: 0, - taskId, - type: 'log', - createdAt: '', - body: { message: 'My log message' }, - }, - { - id: 1, - taskId, - type: 'completion', - createdAt: '', - body: { message: 'Finished!' }, - }, - ], + taskBroker.event$ as jest.Mocked['event$'] + ).mockImplementation(({ taskId }) => { + return new ObservableImpl(observer => { + subscriber = observer; + observer.next({ + events: [ + { + id: 0, + taskId, + type: 'log', + createdAt: '', + body: { message: 'My log message' }, + }, + { + id: 1, + taskId, + type: 'completion', + createdAt: '', + body: { message: 'Finished!' }, + }, + ], + }); }); - return { unsubscribe }; }); const response = await request(app).get('/v2/tasks/a-random-id/events'); @@ -375,21 +377,20 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ }, ]); - expect(taskBroker.observe).toBeCalledTimes(1); - expect(taskBroker.observe).toBeCalledWith( - { taskId: 'a-random-id' }, - expect.any(Function), - ); - expect(unsubscribe).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledWith({ taskId: 'a-random-id' }); + expect(subscriber!.closed).toBe(true); }); it('should return log messages with after query', async () => { - const unsubscribe = jest.fn(); + let subscriber: ZenObservable.SubscriptionObserver; ( - taskBroker.observe as jest.Mocked['observe'] - ).mockImplementation((_, callback) => { - callback(undefined, { events: [] }); - return { unsubscribe }; + taskBroker.event$ as jest.Mocked['event$'] + ).mockImplementation(() => { + return new ObservableImpl(observer => { + subscriber = observer; + observer.next({ events: [] }); + }); }); const response = await request(app) @@ -399,12 +400,12 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ expect(response.status).toEqual(200); expect(response.body).toEqual([]); - expect(taskBroker.observe).toBeCalledTimes(1); - expect(taskBroker.observe).toBeCalledWith( - { taskId: 'a-random-id', after: 10 }, - expect.any(Function), - ); - expect(unsubscribe).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledTimes(1); + expect(taskBroker.event$).toBeCalledWith({ + taskId: 'a-random-id', + after: 10, + }); + expect(subscriber!.closed).toBe(true); }); }); }); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 371fc47dc2..b36670a340 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -29,6 +29,7 @@ import { TemplateEntityV1beta2, TemplateEntityV1beta3, TaskSpecV1beta3, + TaskSpec, TaskSpecV1beta2, } from '@backstage/plugin-scaffolder-common'; import express from 'express'; @@ -40,7 +41,6 @@ import { createBuiltinActions, DatabaseTaskStore, TaskBroker, - TaskSpec, TaskWorker, TemplateAction, TemplateActionRegistry, @@ -301,15 +301,13 @@ export async function createRouter( }); // After client opens connection send all events as string - const { unsubscribe } = taskBroker.observe( - { taskId, after }, - (error, { events }) => { - if (error) { - logger.error( - `Received error from event stream when observing taskId '${taskId}', ${error}`, - ); - } - + const subscription = taskBroker.event$({ taskId, after }).subscribe({ + error: error => { + logger.error( + `Received error from event stream when observing taskId '${taskId}', ${error}`, + ); + }, + next: ({ events }) => { let shouldUnsubscribe = false; for (const event of events) { res.write( @@ -317,19 +315,18 @@ export async function createRouter( ); if (event.type === 'completion') { shouldUnsubscribe = true; - // Closing the event stream here would cause the frontend - // to automatically reconnect because it lost connection. } } // res.flush() is only available with the compression middleware res.flush?.(); - if (shouldUnsubscribe) unsubscribe(); + if (shouldUnsubscribe) subscription.unsubscribe(); }, - ); + }); + // When client closes connection we update the clients list // avoiding the disconnected one req.on('close', () => { - unsubscribe(); + subscription.unsubscribe(); logger.debug(`Event stream observing taskId '${taskId}' closed`); }); }) @@ -337,36 +334,29 @@ export async function createRouter( const { taskId } = req.params; const after = Number(req.query.after) || undefined; - let unsubscribe = () => {}; - // cancel the request after 30 seconds. this aligns with the recommendations of RFC 6202. const timeout = setTimeout(() => { - unsubscribe(); res.json([]); }, 30_000); // Get all known events after an id (always includes the completion event) and return the first callback - ({ unsubscribe } = taskBroker.observe( - { taskId, after }, - (error, { events }) => { - // stop the timeout + const subscription = taskBroker.event$({ taskId, after }).subscribe({ + error: error => { + logger.error( + `Received error from event stream when observing taskId '${taskId}', ${error}`, + ); + }, + next: ({ events }) => { clearTimeout(timeout); - unsubscribe(); - - if (error) { - logger.error( - `Received error from log when observing taskId '${taskId}', ${error}`, - ); - } - + subscription.unsubscribe(); res.json(events); }, - )); + }); // When client closes connection we update the clients list // avoiding the disconnected one req.on('close', () => { - unsubscribe(); + subscription.unsubscribe(); clearTimeout(timeout); }); });