From 11b9a08e92af9639a79d3ec20a236c55434efa68 Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 14 Jan 2024 16:18:42 +0100 Subject: [PATCH 01/10] Recoverable tasks [version 1] Signed-off-by: Bogdan Nechyporenko Signed-off-by: bnechyporenko --- .changeset/serious-carpets-learn.md | 9 ++ .../techdocs/creating-and-publishing.md | 17 +++ packages/backend/src/plugins/scaffolder.ts | 1 + plugins/scaffolder-backend/api-report.md | 19 ++++ plugins/scaffolder-backend/config.d.ts | 16 +++ plugins/scaffolder-backend/package.json | 1 + .../src/ScaffolderPlugin.ts | 3 + .../src/scaffolder/tasks/DatabaseTaskStore.ts | 107 ++++++++++++++---- .../tasks/NunjucksWorkflowRunner.ts | 4 + .../tasks/StorageTaskBroker.test.ts | 31 ++--- .../src/scaffolder/tasks/StorageTaskBroker.ts | 29 +++++ .../src/scaffolder/tasks/TaskWorker.ts | 15 +++ .../src/scaffolder/tasks/helper.ts | 14 +++ .../src/scaffolder/tasks/index.ts | 1 + .../tasks/taskRecoveryHelper.test.ts | 47 ++++++++ .../scaffolder/tasks/taskRecoveryHelper.ts | 41 +++++++ .../src/scaffolder/tasks/types.ts | 12 +- .../src/service/router.test.ts | 6 +- .../scaffolder-backend/src/service/router.ts | 18 ++- plugins/scaffolder-common/api-report.md | 15 +++ plugins/scaffolder-common/src/TaskSpec.ts | 28 +++++ .../src/Template.v1beta3.schema.json | 10 ++ .../src/TemplateEntityV1beta3.ts | 21 ++++ plugins/scaffolder-common/src/index.ts | 1 + plugins/scaffolder-node/api-report.md | 4 +- plugins/scaffolder-node/src/tasks/types.ts | 4 +- plugins/scaffolder-react/api-report.md | 2 +- plugins/scaffolder-react/src/api/types.ts | 2 +- .../src/hooks/useEventStream.ts | 55 ++++++--- plugins/scaffolder/src/api.ts | 1 + yarn.lock | 1 + 31 files changed, 473 insertions(+), 62 deletions(-) create mode 100644 .changeset/serious-carpets-learn.md create mode 100644 plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts create mode 100644 plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts diff --git a/.changeset/serious-carpets-learn.md b/.changeset/serious-carpets-learn.md new file mode 100644 index 0000000000..1f13420976 --- /dev/null +++ b/.changeset/serious-carpets-learn.md @@ -0,0 +1,9 @@ +--- +'@backstage/plugin-scaffolder-backend': minor +'@backstage/plugin-scaffolder-common': minor +'@backstage/plugin-scaffolder-react': minor +'@backstage/plugin-scaffolder-node': minor +'@backstage/plugin-scaffolder': patch +--- + +Introduced the first version of recoverable tasks. diff --git a/docs/features/techdocs/creating-and-publishing.md b/docs/features/techdocs/creating-and-publishing.md index 1251f67276..3fca490b6d 100644 --- a/docs/features/techdocs/creating-and-publishing.md +++ b/docs/features/techdocs/creating-and-publishing.md @@ -41,6 +41,23 @@ default, we highly recommend you to set that up. Follow our how-to guide [How to add documentation setup to your software templates](./how-to-guides.md#how-to-add-the-documentation-setup-to-your-software-templates) to get started. +### Use the documentation template + +There could be _some_ situations where you don't want to keep your docs close to +your code, but still want to publish documentation - for example, an onboarding +tutorial. For this use case, we have put together a documentation template. Your +Backstage instance should by default have a documentation template added. If +not, copy the catalog locations from the +[create-app template](https://github.com/backstage/backstage/blob/master/packages/create-app/templates/default-app/app-config.yaml.hbs) +to add the documentation template. The template creates a component with +**only** TechDocs configuration and default markdown files, and is otherwise +empty. + +![Documentation Template](../../assets/techdocs/documentation-template.png) + +Create an entity from the documentation template and you will get the needed +setup for free. + ### Enable documentation for an already existing entity Prerequisites: diff --git a/packages/backend/src/plugins/scaffolder.ts b/packages/backend/src/plugins/scaffolder.ts index 505a514344..6a801f0392 100644 --- a/packages/backend/src/plugins/scaffolder.ts +++ b/packages/backend/src/plugins/scaffolder.ts @@ -54,6 +54,7 @@ export default async function createPlugin( config: env.config, database: env.database, catalogClient: catalogClient, + eventBroker: env.eventBroker, reader: env.reader, identity: env.identity, scheduler: env.scheduler, diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index c928e5341e..c11a75e555 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -9,6 +9,7 @@ import * as bitbucket from '@backstage/plugin-scaffolder-backend-module-bitbucke import { CatalogApi } from '@backstage/catalog-client'; import { Config } from '@backstage/config'; import { Duration } from 'luxon'; +import { EventBroker } from '@backstage/plugin-events-node'; import { executeShellCommand as executeShellCommand_2 } from '@backstage/plugin-scaffolder-node'; import { ExecuteShellCommandOptions } from '@backstage/plugin-scaffolder-node'; import express from 'express'; @@ -20,6 +21,7 @@ import { HumanDuration } from '@backstage/types'; import { IdentityApi } from '@backstage/plugin-auth-node'; import { JsonObject } from '@backstage/types'; import { Knex } from 'knex'; +import { LifecycleService } from '@backstage/backend-plugin-api'; import { Logger } from 'winston'; import { PermissionEvaluator } from '@backstage/plugin-permission-common'; import { PermissionRule } from '@backstage/plugin-permission-node'; @@ -40,6 +42,7 @@ import { TaskBrokerDispatchResult as TaskBrokerDispatchResult_2 } from '@backsta import { TaskCompletionState as TaskCompletionState_2 } from '@backstage/plugin-scaffolder-node'; import { TaskContext as TaskContext_2 } from '@backstage/plugin-scaffolder-node'; import { TaskEventType as TaskEventType_2 } from '@backstage/plugin-scaffolder-node'; +import { TaskRecovery } from '@backstage/plugin-scaffolder-common'; import { TaskSecrets as TaskSecrets_2 } from '@backstage/plugin-scaffolder-node'; import { TaskSpec } from '@backstage/plugin-scaffolder-common'; import { TaskSpecV1beta3 } from '@backstage/plugin-scaffolder-common'; @@ -400,9 +403,12 @@ export class DatabaseTaskStore implements TaskStore { listStaleTasks(options: { timeoutS: number }): Promise<{ tasks: { taskId: string; + recovery?: TaskRecovery; }[]; }>; // (undocumented) + recoverTasks(options: TaskStoreRecoverTaskOptions): Promise; + // (undocumented) shutdownTask(options: TaskStoreShutDownTaskOptions): Promise; } @@ -433,8 +439,12 @@ export interface RouterOptions { // (undocumented) database: PluginDatabaseManager; // (undocumented) + eventBroker?: EventBroker; + // (undocumented) identity?: IdentityApi; // (undocumented) + lifecycle?: LifecycleService; + // (undocumented) logger: Logger; // (undocumented) permissionRules?: Array< @@ -552,6 +562,8 @@ export interface TaskStore { }[]; }>; // (undocumented) + recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise; + // (undocumented) shutdownTask?(options: TaskStoreShutDownTaskOptions): Promise; } @@ -579,6 +591,11 @@ export type TaskStoreListEventsOptions = { after?: number | undefined; }; +// @public +export type TaskStoreRecoverTaskOptions = { + timeoutS: HumanDuration; +}; + // @public export type TaskStoreShutDownTaskOptions = { taskId: string; @@ -591,6 +608,8 @@ export class TaskWorker { // (undocumented) protected onReadyToClaimTask(): Promise; // (undocumented) + recoverTasks(): Promise; + // (undocumented) runOneTask(task: TaskContext): Promise; // (undocumented) start(): void; diff --git a/plugins/scaffolder-backend/config.d.ts b/plugins/scaffolder-backend/config.d.ts index 5914cddce7..477defe026 100644 --- a/plugins/scaffolder-backend/config.d.ts +++ b/plugins/scaffolder-backend/config.d.ts @@ -40,6 +40,22 @@ export interface Config { */ concurrentTasksLimit?: number; + /** + * Sets the tasks recoverability on system start up. + * + * If not specified, the default value is false. + */ + EXPERIMENTAL_recoverTasks?: boolean; + + /** + * Every task which is in progress state and having a last heartbeat longer than a specified timeout is going to + * be attempted to recover. + * + * If not specified, the default value is 5 seconds. + * + */ + EXPERIMENTAL_recoverTasksTimeout?: HumanDuration; + /** * Makes sure to auto-expire and clean up things that time out or for other reasons should not be left lingering. * diff --git a/plugins/scaffolder-backend/package.json b/plugins/scaffolder-backend/package.json index 55640ceff4..a54f203dc3 100644 --- a/plugins/scaffolder-backend/package.json +++ b/plugins/scaffolder-backend/package.json @@ -57,6 +57,7 @@ "@backstage/plugin-auth-node": "workspace:^", "@backstage/plugin-catalog-backend-module-scaffolder-entity-model": "workspace:^", "@backstage/plugin-catalog-node": "workspace:^", + "@backstage/plugin-events-node": "workspace:^", "@backstage/plugin-permission-common": "workspace:^", "@backstage/plugin-permission-node": "workspace:^", "@backstage/plugin-scaffolder-backend-module-azure": "workspace:^", diff --git a/plugins/scaffolder-backend/src/ScaffolderPlugin.ts b/plugins/scaffolder-backend/src/ScaffolderPlugin.ts index beabaa50aa..e0303b0a30 100644 --- a/plugins/scaffolder-backend/src/ScaffolderPlugin.ts +++ b/plugins/scaffolder-backend/src/ScaffolderPlugin.ts @@ -75,6 +75,7 @@ export const scaffolderPlugin = createBackendPlugin({ deps: { logger: coreServices.logger, config: coreServices.rootConfig, + lifecycle: coreServices.rootLifecycle, reader: coreServices.urlReader, permissions: coreServices.permissions, database: coreServices.database, @@ -84,6 +85,7 @@ export const scaffolderPlugin = createBackendPlugin({ async init({ logger, config, + lifecycle, reader, database, httpRouter, @@ -115,6 +117,7 @@ export const scaffolderPlugin = createBackendPlugin({ database, catalogClient, reader, + lifecycle, actions, taskBroker, additionalTemplateFilters, diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 6beb01e401..7f7dd01a6d 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -29,6 +29,7 @@ import { TaskStoreCreateTaskOptions, TaskStoreCreateTaskResult, TaskStoreShutDownTaskOptions, + TaskStoreRecoverTaskOptions, } from './types'; import { SerializedTaskEvent, @@ -36,7 +37,9 @@ import { TaskStatus, TaskEventType, } from '@backstage/plugin-scaffolder-node'; -import { DateTime } from 'luxon'; +import { DateTime, Duration } from 'luxon'; +import { TaskRecovery, TaskSpec } from '@backstage/plugin-scaffolder-common'; +import { compactEvents } from './taskRecoveryHelper'; const migrationsDir = resolvePackagePath( '@backstage/plugin-scaffolder-backend', @@ -114,6 +117,20 @@ export class DatabaseTaskStore implements TaskStore { return new DatabaseTaskStore(client); } + private isRecoverableTask(spec: TaskSpec): boolean { + return ['startOver'].includes( + spec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy ?? 'none', + ); + } + + private parseSpec({ spec, id }: { spec: string; id: string }): TaskSpec { + try { + return JSON.parse(spec); + } catch (error) { + throw new Error(`Failed to parse spec of task '${id}', ${error}`); + } + } + private static async getClient( database: PluginDatabaseManager | Knex, ): Promise { @@ -223,34 +240,31 @@ export class DatabaseTaskStore implements TaskStore { return undefined; } + const spec = this.parseSpec(task); + const updateCount = await tx('tasks') .where({ id: task.id, status: 'open' }) .update({ status: 'processing', last_heartbeat_at: this.db.fn.now(), - // remove the secrets when moving to processing state. - secrets: null, + // remove the secrets for non-recoverable tasks when moving to processing state. + secrets: this.isRecoverableTask(spec) ? task.secrets : null, }); if (updateCount < 1) { return undefined; } - try { - const spec = JSON.parse(task.spec); - const secrets = task.secrets ? JSON.parse(task.secrets) : undefined; - return { - id: task.id, - spec, - status: 'processing', - lastHeartbeatAt: task.last_heartbeat_at, - createdAt: task.created_at, - createdBy: task.created_by ?? undefined, - secrets, - }; - } catch (error) { - throw new Error(`Failed to parse spec of task '${task.id}', ${error}`); - } + const secrets = task.secrets ? JSON.parse(task.secrets) : undefined; + return { + id: task.id, + spec, + status: 'processing', + lastHeartbeatAt: task.last_heartbeat_at, + createdAt: task.created_at, + createdBy: task.created_by ?? undefined, + secrets, + } as SerializedTask; }); } @@ -266,7 +280,7 @@ export class DatabaseTaskStore implements TaskStore { } async listStaleTasks(options: { timeoutS: number }): Promise<{ - tasks: { taskId: string }[]; + tasks: { taskId: string; recovery?: TaskRecovery }[]; }> { const { timeoutS } = options; let heartbeatInterval = this.db.raw(`? - interval '${timeoutS} seconds'`, [ @@ -285,6 +299,7 @@ export class DatabaseTaskStore implements TaskStore { .where('status', 'processing') .andWhere('last_heartbeat_at', '<=', heartbeatInterval); const tasks = rawRows.map(row => ({ + recovery: (JSON.parse(row.spec) as TaskSpec).EXPERIMENTAL_recovery, taskId: row.id, })); return { tasks }; @@ -297,7 +312,7 @@ export class DatabaseTaskStore implements TaskStore { }): Promise { const { taskId, status, eventBody } = options; - let oldStatus: string; + let oldStatus: TaskStatus; if (['failed', 'completed', 'cancelled'].includes(status)) { oldStatus = 'processing'; } else { @@ -322,6 +337,7 @@ export class DatabaseTaskStore implements TaskStore { .where(criteria) .update({ status, + secrets: null, }); if (updateCount !== 1) { @@ -409,7 +425,8 @@ export class DatabaseTaskStore implements TaskStore { ); } }); - return { events }; + + return compactEvents(events); } async shutdownTask(options: TaskStoreShutDownTaskOptions): Promise { @@ -462,4 +479,52 @@ export class DatabaseTaskStore implements TaskStore { body: serializedBody, }); } + + async recoverTasks(options: TaskStoreRecoverTaskOptions): Promise { + const taskIdsToRecover: string[] = []; + const timeoutS = Duration.fromObject(options.timeoutS).as('seconds'); + + await this.db.transaction(async tx => { + let heartbeatInterval = this.db.raw( + `? - interval '${timeoutS} seconds'`, + [this.db.fn.now()], + ); + if (this.db.client.config.client.includes('mysql')) { + heartbeatInterval = this.db.raw( + `date_sub(now(), interval ${timeoutS} second)`, + ); + } else if (this.db.client.config.client.includes('sqlite3')) { + heartbeatInterval = this.db.raw(`datetime('now', ?)`, [ + `-${timeoutS} seconds`, + ]); + } + + const result = await tx('tasks') + .where('status', 'processing') + .andWhere('last_heartbeat_at', '<=', heartbeatInterval) + .update( + { + status: 'open', + last_heartbeat_at: this.db.fn.now(), + }, + ['id', 'spec'], + ); + + taskIdsToRecover.push(...result.map(i => i.id)); + + for (const { id, spec } of result) { + const taskSpec = JSON.parse(spec as string) as TaskSpec; + await this.db('task_events').insert({ + task_id: id, + event_type: 'recovered', + body: JSON.stringify({ + recoverStrategy: + taskSpec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy ?? 'none', + }), + }); + } + }); + + return taskIdsToRecover; + } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts index 139d91b7a2..e9c7740357 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts @@ -55,6 +55,7 @@ import { } from '@backstage/plugin-permission-common'; import { scaffolderActionRules } from '../../service/rules'; import { actionExecutePermission } from '@backstage/plugin-scaffolder-common/alpha'; +import { TaskRecovery } from '@backstage/plugin-scaffolder-common'; type NunjucksWorkflowRunnerOptions = { workingDirectory: string; @@ -68,6 +69,7 @@ type NunjucksWorkflowRunnerOptions = { type TemplateContext = { parameters: JsonObject; + EXPERIMENTAL_recovery?: TaskRecovery; steps: { [stepName: string]: { output: { [outputName: string]: JsonValue } }; }; @@ -119,6 +121,7 @@ const isActionAuthorized = createConditionAuthorizer( export class NunjucksWorkflowRunner implements WorkflowRunner { private readonly defaultTemplateFilters: Record; + constructor(private readonly options: NunjucksWorkflowRunnerOptions) { this.defaultTemplateFilters = createDefaultFilters({ integrations: this.options.integrations, @@ -415,6 +418,7 @@ export class NunjucksWorkflowRunner implements WorkflowRunner { const context: TemplateContext = { parameters: task.spec.parameters, + EXPERIMENTAL_recovery: task.spec.EXPERIMENTAL_recovery, steps: {}, user: task.spec.user, }; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts index 2255798d5c..a79624456a 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.test.ts @@ -47,10 +47,16 @@ describe('StorageTaskBroker', () => { storage = await createStore(); }); + const emptyTaskSpec = { spec: { steps: [] } as unknown as TaskSpec }; + const emptyTaskWithFakeSecretsSpec = { + spec: { steps: [] } as unknown as TaskSpec, + secrets: fakeSecrets, + }; + const logger = getVoidLogger(); it('should claim a dispatched work item', async () => { const broker = new StorageTaskBroker(storage, logger); - await broker.dispatch({ spec: {} as TaskSpec }); + await broker.dispatch(emptyTaskSpec); await expect(broker.claim()).resolves.toEqual( expect.any(TaskManager as any), ); @@ -62,7 +68,7 @@ describe('StorageTaskBroker', () => { await expect(Promise.race([promise, 'waiting'])).resolves.toBe('waiting'); - await broker.dispatch({ spec: {} as TaskSpec }); + await broker.dispatch(emptyTaskSpec); await expect(promise).resolves.toEqual(expect.any(TaskManager as any)); }); @@ -85,14 +91,14 @@ describe('StorageTaskBroker', () => { it('should store secrets', async () => { const broker = new StorageTaskBroker(storage, logger); - await broker.dispatch({ spec: {} as TaskSpec, secrets: fakeSecrets }); + await broker.dispatch(emptyTaskWithFakeSecretsSpec); const task = await broker.claim(); expect(task.secrets).toEqual(fakeSecrets); }, 10000); it('should complete a task', async () => { const broker = new StorageTaskBroker(storage, logger); - const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); + const dispatchResult = await broker.dispatch(emptyTaskSpec); const task = await broker.claim(); await task.complete('completed'); const taskRow = await storage.getTask(dispatchResult.taskId); @@ -101,10 +107,7 @@ describe('StorageTaskBroker', () => { it('should remove secrets after picking up a task', async () => { const broker = new StorageTaskBroker(storage, logger); - const dispatchResult = await broker.dispatch({ - spec: {} as TaskSpec, - secrets: fakeSecrets, - }); + const dispatchResult = await broker.dispatch(emptyTaskWithFakeSecretsSpec); await broker.claim(); const taskRow = await storage.getTask(dispatchResult.taskId); @@ -113,7 +116,7 @@ describe('StorageTaskBroker', () => { it('should fail a task', async () => { const broker = new StorageTaskBroker(storage, logger); - const dispatchResult = await broker.dispatch({ spec: {} as TaskSpec }); + const dispatchResult = await broker.dispatch(emptyTaskSpec); const task = await broker.claim(); await task.complete('failed'); const taskRow = await storage.getTask(dispatchResult.taskId); @@ -124,7 +127,7 @@ describe('StorageTaskBroker', () => { const broker1 = new StorageTaskBroker(storage, logger); const broker2 = new StorageTaskBroker(storage, logger); - const { taskId } = await broker1.dispatch({ spec: {} as TaskSpec }); + const { taskId } = await broker1.dispatch(emptyTaskSpec); const logPromise = new Promise(resolve => { const observedEvents = new Array(); @@ -169,7 +172,7 @@ describe('StorageTaskBroker', () => { it('should heartbeat', async () => { const broker = new StorageTaskBroker(storage, logger); - const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); + const { taskId } = await broker.dispatch(emptyTaskSpec); const task = await broker.claim(); const initialTask = await storage.getTask(taskId); @@ -187,7 +190,7 @@ describe('StorageTaskBroker', () => { it('should be update the status to failed if heartbeat fails', async () => { const broker = new StorageTaskBroker(storage, logger); - const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); + const { taskId } = await broker.dispatch(emptyTaskSpec); const task = await broker.claim(); jest @@ -213,7 +216,7 @@ describe('StorageTaskBroker', () => { it('should list all tasks', async () => { const broker = new StorageTaskBroker(storage, logger); - const { taskId } = await broker.dispatch({ spec: {} as TaskSpec }); + const { taskId } = await broker.dispatch(emptyTaskSpec); const promise = broker.list(); await expect(promise).resolves.toEqual({ @@ -228,7 +231,7 @@ describe('StorageTaskBroker', () => { it('should list only tasks createdBy a specific user', async () => { const broker = new StorageTaskBroker(storage, logger); const { taskId } = await broker.dispatch({ - spec: {} as TaskSpec, + spec: { steps: [] } as unknown as TaskSpec, createdBy: 'user:default/foo', }); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index fc5a724e47..467da0f9de 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -14,6 +14,7 @@ * limitations under the License. */ +import { Config } from '@backstage/config'; import { TaskSpec } from '@backstage/plugin-scaffolder-common'; import { TaskSecrets } from '@backstage/plugin-scaffolder-node'; import { JsonObject, Observable } from '@backstage/types'; @@ -28,6 +29,7 @@ import { TaskContext, TaskStore, } from './types'; +import { readDuration } from './helper'; /** * TaskManager @@ -160,6 +162,7 @@ export class StorageTaskBroker implements TaskBroker { constructor( private readonly storage: TaskStore, private readonly logger: Logger, + private readonly config?: Config, ) {} async list(options?: { @@ -202,6 +205,32 @@ export class StorageTaskBroker implements TaskBroker { }); } + public async recoverTasks(): Promise { + const enabled = + (this.config && + this.config.getOptionalBoolean( + 'scaffolder.EXPERIMENTAL_recoverTasks', + )) ?? + false; + + if (enabled) { + const recoveredTaskIds = + (await this.storage.recoverTasks?.({ + timeoutS: readDuration( + this.config, + 'scaffolder.EXPERIMENTAL_recoverTasksTimeout', + { + seconds: 30, + }, + ), + })) ?? []; + recoveredTaskIds.forEach(() => { + this.signalDispatch(); + }); + } + return enabled; + } + /** * {@inheritdoc TaskBroker.claim} */ diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 9f9b86b723..8111a28793 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -23,6 +23,7 @@ import { ScmIntegrations } from '@backstage/integration'; import { assertError } from '@backstage/errors'; import { TemplateFilter, TemplateGlobal } from '../../lib'; import { PermissionEvaluator } from '@backstage/plugin-permission-common'; + /** * TaskWorkerOptions * @@ -111,7 +112,21 @@ export class TaskWorker { }); } + async recoverTasks() { + try { + await this.options.taskBroker.recoverTasks?.(); + } catch (_err) { + // ignore + } + } + start() { + (async () => { + for (;;) { + await new Promise(resolve => setTimeout(resolve, 10000)); + await this.recoverTasks(); + } + })(); (async () => { for (;;) { await this.onReadyToClaimTask(); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/helper.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/helper.ts index 76951afde4..41425affc9 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/helper.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/helper.ts @@ -14,6 +14,9 @@ * limitations under the License. */ +import { Config, readDurationFromConfig } from '@backstage/config'; +import { HumanDuration } from '@backstage/types'; + import { isArray } from 'lodash'; import { Schema } from 'jsonschema'; @@ -53,3 +56,14 @@ export function generateExampleOutput(schema: Schema): unknown { } return ''; } + +export const readDuration = ( + config: Config | undefined, + key: string, + defaultValue: HumanDuration, +) => { + if (config?.has(key)) { + return readDurationFromConfig(config, { key }); + } + return defaultValue; +}; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts index f1c1a7bdac..9d231d7e7c 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/index.ts @@ -35,5 +35,6 @@ export type { TaskBrokerDispatchResult, TaskBrokerDispatchOptions, TaskStoreCreateTaskOptions, + TaskStoreRecoverTaskOptions, TaskStoreCreateTaskResult, } from './types'; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts new file mode 100644 index 0000000000..9a1c012c70 --- /dev/null +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts @@ -0,0 +1,47 @@ +/* + * Copyright 2021 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { compactEvents } from './taskRecoveryHelper'; +import { SerializedTaskEvent } from './types'; + +const toLogEvent = (stepId: string) => + ({ + type: 'log', + body: { stepId }, + } as unknown as SerializedTaskEvent); + +const toRecoveredEvent = (recoverStrategy: string) => + ({ + type: 'recovered', + body: { recoverStrategy }, + } as unknown as SerializedTaskEvent); + +describe('taskRecoveryHelper', () => { + describe('compactEvents', () => { + it('should return only events related to a restarted task. Recover strategy: "startOver"', () => { + const logEvents = [ + 'fetch', + 'mock-step-1', + 'mock-step-2', + 'mock-step-3', + ].map(toLogEvent); + + const events = [...logEvents, toRecoveredEvent('startOver')]; + + expect(compactEvents(events)).toEqual({ events: [] }); + }); + }); +}); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts new file mode 100644 index 0000000000..31b05c6924 --- /dev/null +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts @@ -0,0 +1,41 @@ +/* + * Copyright 2021 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { SerializedTaskEvent } from '@backstage/plugin-scaffolder-node'; +import { TaskRecoverStrategy } from '@backstage/plugin-scaffolder-common'; + +export const compactEvents = ( + events: SerializedTaskEvent[], +): { events: SerializedTaskEvent[] } => { + const recoveredEventInd = events + .slice() + .reverse() + .findIndex(event => event.type === 'recovered'); + + if (recoveredEventInd >= 0) { + const ind = events.length - recoveredEventInd - 1; + const { recoverStrategy } = events[ind].body as { + recoverStrategy: TaskRecoverStrategy; + }; + if (recoverStrategy === 'startOver') { + return { + events: recoveredEventInd === 0 ? [] : events.slice(ind), + }; + } + } + + return { events }; +}; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index 4560f27c0b..c505f0f978 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -14,7 +14,7 @@ * limitations under the License. */ -import { JsonValue, JsonObject } from '@backstage/types'; +import { JsonValue, JsonObject, HumanDuration } from '@backstage/types'; import { TaskSpec, TaskStep } from '@backstage/plugin-scaffolder-common'; import { TaskSecrets } from '@backstage/plugin-scaffolder-node'; import { @@ -142,6 +142,14 @@ export type TaskStoreCreateTaskOptions = { secrets?: TaskSecrets; }; +/** + * The options passed to {@link TaskStore.recoverTasks} + * @public + */ +export type TaskStoreRecoverTaskOptions = { + timeoutS: HumanDuration; +}; + /** * The response from {@link TaskStore.createTask} * @public @@ -162,6 +170,8 @@ export interface TaskStore { options: TaskStoreCreateTaskOptions, ): Promise; + recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise; + getTask(taskId: string): Promise; claimTask(): Promise; diff --git a/plugins/scaffolder-backend/src/service/router.test.ts b/plugins/scaffolder-backend/src/service/router.test.ts index 7a0bcc6877..8bfbd6b915 100644 --- a/plugins/scaffolder-backend/src/service/router.test.ts +++ b/plugins/scaffolder-backend/src/service/router.test.ts @@ -85,6 +85,8 @@ const mockUrlReader = UrlReaders.default({ const getIdentity = jest.fn(); +const config = new ConfigReader({}); + describe('createRouter', () => { let app: express.Express; let loggerSpy: jest.SpyInstance; @@ -181,7 +183,7 @@ describe('createRouter', () => { const databaseTaskStore = await DatabaseTaskStore.create({ database: createDatabase(), }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); @@ -787,7 +789,7 @@ data: {"id":1,"taskId":"a-random-id","type":"completion","createdAt":"","body":{ const databaseTaskStore = await DatabaseTaskStore.create({ database: createDatabase(), }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config); jest.spyOn(taskBroker, 'dispatch'); jest.spyOn(taskBroker, 'get'); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index de0087e7de..a2608f3752 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -77,7 +77,9 @@ import { PermissionRule, } from '@backstage/plugin-permission-node'; import { scaffolderActionRules, scaffolderTemplateRules } from './rules'; +import { EventBroker } from '@backstage/plugin-events-node'; import { Duration } from 'luxon'; +import { LifecycleService } from '@backstage/backend-plugin-api'; /** * @@ -124,9 +126,11 @@ export interface RouterOptions { logger: Logger; config: Config; reader: UrlReader; + lifecycle?: LifecycleService; database: PluginDatabaseManager; catalogClient: CatalogApi; scheduler?: PluginTaskScheduler; + eventBroker?: EventBroker; actions?: TemplateAction[]; /** * @deprecated taskWorkers is deprecated in favor of concurrentTasksLimit option with a single TaskWorker @@ -266,7 +270,7 @@ export async function createRouter( let taskBroker: TaskBroker; if (!options.taskBroker) { const databaseTaskStore = await DatabaseTaskStore.create({ database }); - taskBroker = new StorageTaskBroker(databaseTaskStore, logger); + taskBroker = new StorageTaskBroker(databaseTaskStore, logger, config); if (scheduler && databaseTaskStore.listStaleTasks) { await scheduler.scheduleTask({ @@ -301,7 +305,7 @@ export async function createRouter( const actionRegistry = new TemplateActionRegistry(); - const workers = []; + const workers: TaskWorker[] = []; if (concurrentTasksLimit !== 0) { for (let i = 0; i < (taskWorkers || 1); i++) { const worker = await TaskWorker.create({ @@ -331,7 +335,14 @@ export async function createRouter( }); actionsToRegister.forEach(action => actionRegistry.register(action)); - workers.forEach(worker => worker.start()); + + const launchWorkers = () => workers.forEach(worker => worker.start()); + + if (options.lifecycle) { + options.lifecycle.addStartupHook(launchWorkers); + } else { + launchWorkers(); + } const dryRunner = createDryRunner({ actionRegistry, @@ -462,6 +473,7 @@ export async function createRouter( id: step.id ?? `step-${index + 1}`, name: step.name ?? step.action, })), + EXPERIMENTAL_recovery: template.spec.EXPERIMENTAL_recovery, output: template.spec.output ?? {}, parameters: values, user: { diff --git a/plugins/scaffolder-common/api-report.md b/plugins/scaffolder-common/api-report.md index 8cfbd2aaac..8d7a61910d 100644 --- a/plugins/scaffolder-common/api-report.md +++ b/plugins/scaffolder-common/api-report.md @@ -16,12 +16,21 @@ export const isTemplateEntityV1beta3: ( entity: Entity, ) => entity is TemplateEntityV1beta3; +// @public +export type TaskRecoverStrategy = 'none' | 'startOver'; + +// @public +export interface TaskRecovery { + EXPERIMENTAL_strategy?: TaskRecoverStrategy; +} + // @public export type TaskSpec = TaskSpecV1beta3; // @public export interface TaskSpecV1beta3 { apiVersion: 'scaffolder.backstage.io/v1beta3'; + EXPERIMENTAL_recovery?: TaskRecovery; output: { [name: string]: JsonValue; }; @@ -67,6 +76,7 @@ export interface TemplateEntityV1beta3 extends Entity { spec: { type: string; presentation?: TemplatePresentationV1beta3; + EXPERIMENTAL_recovery?: TemplateRecoveryV1beta3; parameters?: TemplateParametersV1beta3 | TemplateParametersV1beta3[]; steps: Array; output?: { @@ -108,4 +118,9 @@ export interface TemplatePresentationV1beta3 extends JsonObject { reviewButtonText?: string; }; } + +// @public +export interface TemplateRecoveryV1beta3 extends JsonObject { + EXPERIMENTAL_strategy?: 'none' | 'startOver'; +} ``` diff --git a/plugins/scaffolder-common/src/TaskSpec.ts b/plugins/scaffolder-common/src/TaskSpec.ts index 6b144b4f78..ce93c05ef0 100644 --- a/plugins/scaffolder-common/src/TaskSpec.ts +++ b/plugins/scaffolder-common/src/TaskSpec.ts @@ -44,6 +44,30 @@ export type TemplateInfo = { }; }; +/** + * + * none - not recover, let the task be marked as failed + * startOver - do recover, start the execution of the task from the first step. + * + * @public + */ +export type TaskRecoverStrategy = 'none' | 'startOver'; + +/** + * When task didn't have a chance to complete due to system restart you can define the strategy what to do with such tasks, + * by defining a strategy. + * + * By default, it is none, what means to not recover but updating the status from 'processing' to 'failed'. + * + * @public + */ +export interface TaskRecovery { + /** + * Depends on how you designed your task you might tailor the behaviour for each of them. + */ + EXPERIMENTAL_strategy?: TaskRecoverStrategy; +} + /** * An individual step of a scaffolder task, as stored in the database. * @@ -119,6 +143,10 @@ export interface TaskSpecV1beta3 { */ ref?: string; }; + /** + * How to recover the task after system restart or system crash. + */ + EXPERIMENTAL_recovery?: TaskRecovery; } /** diff --git a/plugins/scaffolder-common/src/Template.v1beta3.schema.json b/plugins/scaffolder-common/src/Template.v1beta3.schema.json index d4be9b7538..26f6a6a755 100644 --- a/plugins/scaffolder-common/src/Template.v1beta3.schema.json +++ b/plugins/scaffolder-common/src/Template.v1beta3.schema.json @@ -139,6 +139,16 @@ } ] }, + "EXPERIMENTAL_recovery": { + "type": "object", + "description": "A task recovery section.", + "properties": { + "EXPERIMENTAL_strategy": { + "type": "string", + "description": "Recovery strategy for your task (none or startOver). By default none" + } + } + }, "steps": { "type": "array", "description": "A list of steps to execute.", diff --git a/plugins/scaffolder-common/src/TemplateEntityV1beta3.ts b/plugins/scaffolder-common/src/TemplateEntityV1beta3.ts index 23cc3ca3d8..0e2fec4264 100644 --- a/plugins/scaffolder-common/src/TemplateEntityV1beta3.ts +++ b/plugins/scaffolder-common/src/TemplateEntityV1beta3.ts @@ -51,6 +51,11 @@ export interface TemplateEntityV1beta3 extends Entity { */ presentation?: TemplatePresentationV1beta3; + /** + * Recovery strategy for the template + */ + EXPERIMENTAL_recovery?: TemplateRecoveryV1beta3; + /** * This is a JSONSchema or an array of JSONSchema's which is used to render a form in the frontend * to collect user input and validate it against that schema. This can then be used in the `steps` part below to template @@ -73,6 +78,22 @@ export interface TemplateEntityV1beta3 extends Entity { }; } +/** + * Depends on how you designed your task you might tailor the behaviour for each of them. + * + * @public + */ +export interface TemplateRecoveryV1beta3 extends JsonObject { + /** + * + * none - not recover, let the task be marked as failed + * startOver - do recover, start the execution of the task from the first step. + * + * @public + */ + EXPERIMENTAL_strategy?: 'none' | 'startOver'; +} + /** * The presentation of the template. * diff --git a/plugins/scaffolder-common/src/index.ts b/plugins/scaffolder-common/src/index.ts index 3300aebf46..4d9b5e39c6 100644 --- a/plugins/scaffolder-common/src/index.ts +++ b/plugins/scaffolder-common/src/index.ts @@ -32,4 +32,5 @@ export type { TemplateEntityStepV1beta3, TemplateParametersV1beta3, TemplatePermissionsV1beta3, + TemplateRecoveryV1beta3, } from './TemplateEntityV1beta3'; diff --git a/plugins/scaffolder-node/api-report.md b/plugins/scaffolder-node/api-report.md index 6bb1534679..9a71d41969 100644 --- a/plugins/scaffolder-node/api-report.md +++ b/plugins/scaffolder-node/api-report.md @@ -238,6 +238,8 @@ export interface TaskBroker { tasks: SerializedTask[]; }>; // (undocumented) + recoverTasks?(): Promise; + // (undocumented) vacuumTasks(options: { timeoutS: number }): Promise; } @@ -279,7 +281,7 @@ export interface TaskContext { } // @public -export type TaskEventType = 'completion' | 'log' | 'cancelled'; +export type TaskEventType = 'completion' | 'log' | 'cancelled' | 'recovered'; // @public export type TaskSecrets = Record & { diff --git a/plugins/scaffolder-node/src/tasks/types.ts b/plugins/scaffolder-node/src/tasks/types.ts index 8c9cc29e31..12e6416749 100644 --- a/plugins/scaffolder-node/src/tasks/types.ts +++ b/plugins/scaffolder-node/src/tasks/types.ts @@ -65,7 +65,7 @@ export type SerializedTask = { * * @public */ -export type TaskEventType = 'completion' | 'log' | 'cancelled'; +export type TaskEventType = 'completion' | 'log' | 'cancelled' | 'recovered'; /** * SerializedTaskEvent @@ -131,6 +131,8 @@ export interface TaskBroker { claim(): Promise; + recoverTasks?(): Promise; + dispatch( options: TaskBrokerDispatchOptions, ): Promise; diff --git a/plugins/scaffolder-react/api-report.md b/plugins/scaffolder-react/api-report.md index 8589078f1e..b4edf79cd0 100644 --- a/plugins/scaffolder-react/api-report.md +++ b/plugins/scaffolder-react/api-report.md @@ -152,7 +152,7 @@ export type ListActionsResponse = Array; // @public export type LogEvent = { - type: 'log' | 'completion' | 'cancelled'; + type: 'log' | 'completion' | 'cancelled' | 'recovered'; body: { message: string; stepId?: string; diff --git a/plugins/scaffolder-react/src/api/types.ts b/plugins/scaffolder-react/src/api/types.ts index 2adecdd0d2..3f80382f77 100644 --- a/plugins/scaffolder-react/src/api/types.ts +++ b/plugins/scaffolder-react/src/api/types.ts @@ -105,7 +105,7 @@ export type ScaffolderTaskOutput = { * @public */ export type LogEvent = { - type: 'log' | 'completion' | 'cancelled'; + type: 'log' | 'completion' | 'cancelled' | 'recovered'; body: { message: string; stepId?: string; diff --git a/plugins/scaffolder-react/src/hooks/useEventStream.ts b/plugins/scaffolder-react/src/hooks/useEventStream.ts index ad802dba99..f63f6024e6 100644 --- a/plugins/scaffolder-react/src/hooks/useEventStream.ts +++ b/plugins/scaffolder-react/src/hooks/useEventStream.ts @@ -62,12 +62,14 @@ type ReducerLogEntry = { message: string; output?: ScaffolderTaskOutput; error?: Error; + recoverStrategy?: 'none' | 'startOver'; }; }; type ReducerAction = | { type: 'INIT'; data: ScaffolderTask } | { type: 'CANCELLED' } + | { type: 'RECOVERED'; data: ReducerLogEntry } | { type: 'LOGS'; data: ReducerLogEntry[] } | { type: 'COMPLETED'; data: ReducerLogEntry } | { type: 'ERROR'; data: Error }; @@ -105,17 +107,19 @@ function reducer(draft: TaskStream, action: ReducerAction) { const currentStepLog = draft.stepLogs?.[entry.body.stepId]; const currentStep = draft.steps?.[entry.body.stepId]; - if (entry.body.status && entry.body.status !== currentStep.status) { - currentStep.status = entry.body.status; + if (currentStep) { + if (entry.body.status && entry.body.status !== currentStep.status) { + currentStep.status = entry.body.status; - if (currentStep.status === 'processing') { - currentStep.startedAt = entry.createdAt; - } + if (currentStep.status === 'processing') { + currentStep.startedAt = entry.createdAt; + } - if ( - ['cancelled', 'completed', 'failed'].includes(currentStep.status) - ) { - currentStep.endedAt = entry.createdAt; + if ( + ['cancelled', 'completed', 'failed'].includes(currentStep.status) + ) { + currentStep.endedAt = entry.createdAt; + } } } @@ -138,6 +142,17 @@ function reducer(draft: TaskStream, action: ReducerAction) { return; } + case 'RECOVERED': { + for (const stepId in draft.steps) { + if (draft.steps.hasOwnProperty(stepId)) { + draft.steps[stepId].startedAt = undefined; + draft.steps[stepId].endedAt = undefined; + draft.steps[stepId].status = 'open'; + } + } + return; + } + case 'ERROR': { draft.error = action.data; draft.loading = false; @@ -202,6 +217,7 @@ export const useTaskEventStream = (taskId: string): TaskStream => { subscription = observable.subscribe({ next: event => { + retryCount = 1; switch (event.type) { case 'log': return collectedLogEvents.push(event); @@ -212,6 +228,9 @@ export const useTaskEventStream = (taskId: string): TaskStream => { emitLogs(); dispatch({ type: 'COMPLETED', data: event }); return undefined; + case 'recovered': + dispatch({ type: 'RECOVERED', data: event }); + return undefined; default: throw new Error( `Unhandled event type ${event.type} in observer`, @@ -226,16 +245,18 @@ export const useTaskEventStream = (taskId: string): TaskStream => { // just to restart the fetch process // details here https://github.com/backstage/backstage/issues/15002 + const maxRetries = 3; + if (!error.message) { - error.message = `We cannot connect at the moment, trying again in some seconds... Retrying (${retryCount}/3 retries)`; + error.message = `We cannot connect at the moment, trying again in some seconds... Retrying (${ + retryCount > maxRetries ? maxRetries : retryCount + }/${maxRetries} retries)`; } - if (retryCount <= 3) { - setTimeout(() => { - retryCount += 1; - startStreamLogProcess(); - }, 15000); - } + setTimeout(() => { + retryCount += 1; + void startStreamLogProcess(); + }, 15000); dispatch({ type: 'ERROR', data: error }); }, @@ -247,7 +268,7 @@ export const useTaskEventStream = (taskId: string): TaskStream => { } }, ); - startStreamLogProcess(); + void startStreamLogProcess(); return () => { didCancel = true; if (subscription) { diff --git a/plugins/scaffolder/src/api.ts b/plugins/scaffolder/src/api.ts index a66ac522c2..fe0d7c11d0 100644 --- a/plugins/scaffolder/src/api.ts +++ b/plugins/scaffolder/src/api.ts @@ -252,6 +252,7 @@ export class ScaffolderClient implements ScaffolderApi { : {}, }); eventSource.addEventListener('log', processEvent); + eventSource.addEventListener('recovered', processEvent); eventSource.addEventListener('cancelled', processEvent); eventSource.addEventListener('completion', (event: any) => { processEvent(event); diff --git a/yarn.lock b/yarn.lock index 230f6d3e8d..8c0ef18c0d 100644 --- a/yarn.lock +++ b/yarn.lock @@ -8383,6 +8383,7 @@ __metadata: "@backstage/plugin-auth-node": "workspace:^" "@backstage/plugin-catalog-backend-module-scaffolder-entity-model": "workspace:^" "@backstage/plugin-catalog-node": "workspace:^" + "@backstage/plugin-events-node": "workspace:^" "@backstage/plugin-permission-common": "workspace:^" "@backstage/plugin-permission-node": "workspace:^" "@backstage/plugin-scaffolder-backend-module-azure": "workspace:^" From 0e0bb6d09eca6e1ec89fd861f4fcda9e0fd74ac7 Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Mon, 15 Jan 2024 10:43:37 +0100 Subject: [PATCH 02/10] wip Signed-off-by: bnechyporenko --- packages/backend/src/plugins/scaffolder.ts | 1 - plugins/scaffolder-backend/api-report.md | 3 --- plugins/scaffolder-backend/src/service/router.ts | 2 -- 3 files changed, 6 deletions(-) diff --git a/packages/backend/src/plugins/scaffolder.ts b/packages/backend/src/plugins/scaffolder.ts index 6a801f0392..505a514344 100644 --- a/packages/backend/src/plugins/scaffolder.ts +++ b/packages/backend/src/plugins/scaffolder.ts @@ -54,7 +54,6 @@ export default async function createPlugin( config: env.config, database: env.database, catalogClient: catalogClient, - eventBroker: env.eventBroker, reader: env.reader, identity: env.identity, scheduler: env.scheduler, diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index c11a75e555..a3a3531bb3 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -9,7 +9,6 @@ import * as bitbucket from '@backstage/plugin-scaffolder-backend-module-bitbucke import { CatalogApi } from '@backstage/catalog-client'; import { Config } from '@backstage/config'; import { Duration } from 'luxon'; -import { EventBroker } from '@backstage/plugin-events-node'; import { executeShellCommand as executeShellCommand_2 } from '@backstage/plugin-scaffolder-node'; import { ExecuteShellCommandOptions } from '@backstage/plugin-scaffolder-node'; import express from 'express'; @@ -439,8 +438,6 @@ export interface RouterOptions { // (undocumented) database: PluginDatabaseManager; // (undocumented) - eventBroker?: EventBroker; - // (undocumented) identity?: IdentityApi; // (undocumented) lifecycle?: LifecycleService; diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index a2608f3752..365ec8748d 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -77,7 +77,6 @@ import { PermissionRule, } from '@backstage/plugin-permission-node'; import { scaffolderActionRules, scaffolderTemplateRules } from './rules'; -import { EventBroker } from '@backstage/plugin-events-node'; import { Duration } from 'luxon'; import { LifecycleService } from '@backstage/backend-plugin-api'; @@ -130,7 +129,6 @@ export interface RouterOptions { database: PluginDatabaseManager; catalogClient: CatalogApi; scheduler?: PluginTaskScheduler; - eventBroker?: EventBroker; actions?: TemplateAction[]; /** * @deprecated taskWorkers is deprecated in favor of concurrentTasksLimit option with a single TaskWorker From d9a0a4b4674d3d16772cd0cbabf4dd4fee34d423 Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Tue, 16 Jan 2024 13:31:06 +0100 Subject: [PATCH 03/10] wip Signed-off-by: bnechyporenko --- plugins/scaffolder-backend/api-report.md | 2 +- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 2 +- .../src/scaffolder/tasks/StorageTaskBroker.ts | 2 +- plugins/scaffolder-backend/src/scaffolder/tasks/types.ts | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index a3a3531bb3..fd2d368f15 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -590,7 +590,7 @@ export type TaskStoreListEventsOptions = { // @public export type TaskStoreRecoverTaskOptions = { - timeoutS: HumanDuration; + timeout: HumanDuration; }; // @public diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 7f7dd01a6d..b7d7517c10 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -482,7 +482,7 @@ export class DatabaseTaskStore implements TaskStore { async recoverTasks(options: TaskStoreRecoverTaskOptions): Promise { const taskIdsToRecover: string[] = []; - const timeoutS = Duration.fromObject(options.timeoutS).as('seconds'); + const timeoutS = Duration.fromObject(options.timeout).as('seconds'); await this.db.transaction(async tx => { let heartbeatInterval = this.db.raw( diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 467da0f9de..5bd72d05a2 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -216,7 +216,7 @@ export class StorageTaskBroker implements TaskBroker { if (enabled) { const recoveredTaskIds = (await this.storage.recoverTasks?.({ - timeoutS: readDuration( + timeout: readDuration( this.config, 'scaffolder.EXPERIMENTAL_recoverTasksTimeout', { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index c505f0f978..d113e23c5b 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -147,7 +147,7 @@ export type TaskStoreCreateTaskOptions = { * @public */ export type TaskStoreRecoverTaskOptions = { - timeoutS: HumanDuration; + timeout: HumanDuration; }; /** From e5a3599819cb90c9969ad74e054658374c1690af Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sat, 20 Jan 2024 15:49:32 +0100 Subject: [PATCH 04/10] incorporated some of feedback Signed-off-by: bnechyporenko --- .../techdocs/creating-and-publishing.md | 17 ------------- plugins/scaffolder-backend/package.json | 1 - .../src/scaffolder/tasks/DatabaseTaskStore.ts | 25 ++++++++++++++----- .../tasks/NunjucksWorkflowRunner.ts | 1 - .../src/scaffolder/tasks/StorageTaskBroker.ts | 21 ++++++++-------- .../src/scaffolder/tasks/TaskWorker.ts | 8 ++++-- .../tasks/taskRecoveryHelper.test.ts | 4 +-- .../scaffolder/tasks/taskRecoveryHelper.ts | 2 +- .../src/scaffolder/tasks/types.ts | 4 ++- plugins/scaffolder-node/src/tasks/types.ts | 2 +- yarn.lock | 1 - 11 files changed, 42 insertions(+), 44 deletions(-) diff --git a/docs/features/techdocs/creating-and-publishing.md b/docs/features/techdocs/creating-and-publishing.md index 3fca490b6d..1251f67276 100644 --- a/docs/features/techdocs/creating-and-publishing.md +++ b/docs/features/techdocs/creating-and-publishing.md @@ -41,23 +41,6 @@ default, we highly recommend you to set that up. Follow our how-to guide [How to add documentation setup to your software templates](./how-to-guides.md#how-to-add-the-documentation-setup-to-your-software-templates) to get started. -### Use the documentation template - -There could be _some_ situations where you don't want to keep your docs close to -your code, but still want to publish documentation - for example, an onboarding -tutorial. For this use case, we have put together a documentation template. Your -Backstage instance should by default have a documentation template added. If -not, copy the catalog locations from the -[create-app template](https://github.com/backstage/backstage/blob/master/packages/create-app/templates/default-app/app-config.yaml.hbs) -to add the documentation template. The template creates a component with -**only** TechDocs configuration and default markdown files, and is otherwise -empty. - -![Documentation Template](../../assets/techdocs/documentation-template.png) - -Create an entity from the documentation template and you will get the needed -setup for free. - ### Enable documentation for an already existing entity Prerequisites: diff --git a/plugins/scaffolder-backend/package.json b/plugins/scaffolder-backend/package.json index 3b2029d997..24325449b3 100644 --- a/plugins/scaffolder-backend/package.json +++ b/plugins/scaffolder-backend/package.json @@ -57,7 +57,6 @@ "@backstage/plugin-auth-node": "workspace:^", "@backstage/plugin-catalog-backend-module-scaffolder-entity-model": "workspace:^", "@backstage/plugin-catalog-node": "workspace:^", - "@backstage/plugin-events-node": "workspace:^", "@backstage/plugin-permission-common": "workspace:^", "@backstage/plugin-permission-node": "workspace:^", "@backstage/plugin-scaffolder-backend-module-azure": "workspace:^", diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index b7d7517c10..0ce4e38659 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -36,10 +36,11 @@ import { SerializedTask, TaskStatus, TaskEventType, + TaskSecrets, } from '@backstage/plugin-scaffolder-node'; import { DateTime, Duration } from 'luxon'; import { TaskRecovery, TaskSpec } from '@backstage/plugin-scaffolder-common'; -import { compactEvents } from './taskRecoveryHelper'; +import { trimEventsTillLastRecovery } from './taskRecoveryHelper'; const migrationsDir = resolvePackagePath( '@backstage/plugin-scaffolder-backend', @@ -131,6 +132,16 @@ export class DatabaseTaskStore implements TaskStore { } } + private parseTaskSecrets(taskRow: RawDbTaskRow): TaskSecrets | undefined { + try { + return taskRow.secrets ? JSON.parse(taskRow.secrets) : undefined; + } catch (error) { + throw new Error( + `Failed to parse secrets of task '${taskRow.id}', ${error}`, + ); + } + } + private static async getClient( database: PluginDatabaseManager | Knex, ): Promise { @@ -255,7 +266,7 @@ export class DatabaseTaskStore implements TaskStore { return undefined; } - const secrets = task.secrets ? JSON.parse(task.secrets) : undefined; + const secrets = this.parseTaskSecrets(task); return { id: task.id, spec, @@ -264,7 +275,7 @@ export class DatabaseTaskStore implements TaskStore { createdAt: task.created_at, createdBy: task.created_by ?? undefined, secrets, - } as SerializedTask; + }; }); } @@ -426,7 +437,7 @@ export class DatabaseTaskStore implements TaskStore { } }); - return compactEvents(events); + return trimEventsTillLastRecovery(events); } async shutdownTask(options: TaskStoreShutDownTaskOptions): Promise { @@ -480,7 +491,9 @@ export class DatabaseTaskStore implements TaskStore { }); } - async recoverTasks(options: TaskStoreRecoverTaskOptions): Promise { + async recoverTasks( + options: TaskStoreRecoverTaskOptions, + ): Promise<{ id: string }[]> { const taskIdsToRecover: string[] = []; const timeoutS = Duration.fromObject(options.timeout).as('seconds'); @@ -525,6 +538,6 @@ export class DatabaseTaskStore implements TaskStore { } }); - return taskIdsToRecover; + return taskIdsToRecover.map(id => ({ id })); } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts index e9c7740357..ce1ad71c33 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts @@ -418,7 +418,6 @@ export class NunjucksWorkflowRunner implements WorkflowRunner { const context: TemplateContext = { parameters: task.spec.parameters, - EXPERIMENTAL_recovery: task.spec.EXPERIMENTAL_recovery, steps: {}, user: task.spec.user, }; diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 5bd72d05a2..96300bec58 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -205,7 +205,7 @@ export class StorageTaskBroker implements TaskBroker { }); } - public async recoverTasks(): Promise { + public async recoverTasks(): Promise { const enabled = (this.config && this.config.getOptionalBoolean( @@ -214,21 +214,20 @@ export class StorageTaskBroker implements TaskBroker { false; if (enabled) { + const defaultTimeout = { seconds: 30 }; + const timeout = readDuration( + this.config, + 'scaffolder.EXPERIMENTAL_recoverTasksTimeout', + defaultTimeout, + ); const recoveredTaskIds = (await this.storage.recoverTasks?.({ - timeout: readDuration( - this.config, - 'scaffolder.EXPERIMENTAL_recoverTasksTimeout', - { - seconds: 30, - }, - ), + timeout, })) ?? []; - recoveredTaskIds.forEach(() => { + if (recoveredTaskIds.length > 0) { this.signalDispatch(); - }); + } } - return enabled; } /** diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 8111a28793..5e116922de 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -20,7 +20,7 @@ import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner'; import { Logger } from 'winston'; import { TemplateActionRegistry } from '../actions'; import { ScmIntegrations } from '@backstage/integration'; -import { assertError } from '@backstage/errors'; +import { assertError, stringifyError } from '@backstage/errors'; import { TemplateFilter, TemplateGlobal } from '../../lib'; import { PermissionEvaluator } from '@backstage/plugin-permission-common'; @@ -36,6 +36,7 @@ export type TaskWorkerOptions = { }; concurrentTasksLimit: number; permissions?: PermissionEvaluator; + logger?: Logger; }; /** @@ -74,8 +75,10 @@ export type CreateWorkerOptions = { */ export class TaskWorker { private taskQueue: PQueue; + private logger: Logger | undefined; private constructor(private readonly options: TaskWorkerOptions) { + this.logger = options.logger; this.taskQueue = new PQueue({ concurrency: options.concurrentTasksLimit, }); @@ -115,7 +118,8 @@ export class TaskWorker { async recoverTasks() { try { await this.options.taskBroker.recoverTasks?.(); - } catch (_err) { + } catch (err) { + this.logger?.error(stringifyError(err)); // ignore } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts index 9a1c012c70..8224feab8e 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.test.ts @@ -14,7 +14,7 @@ * limitations under the License. */ -import { compactEvents } from './taskRecoveryHelper'; +import { trimEventsTillLastRecovery } from './taskRecoveryHelper'; import { SerializedTaskEvent } from './types'; const toLogEvent = (stepId: string) => @@ -41,7 +41,7 @@ describe('taskRecoveryHelper', () => { const events = [...logEvents, toRecoveredEvent('startOver')]; - expect(compactEvents(events)).toEqual({ events: [] }); + expect(trimEventsTillLastRecovery(events)).toEqual({ events: [] }); }); }); }); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts index 31b05c6924..f3b8b321fe 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/taskRecoveryHelper.ts @@ -17,7 +17,7 @@ import { SerializedTaskEvent } from '@backstage/plugin-scaffolder-node'; import { TaskRecoverStrategy } from '@backstage/plugin-scaffolder-common'; -export const compactEvents = ( +export const trimEventsTillLastRecovery = ( events: SerializedTaskEvent[], ): { events: SerializedTaskEvent[] } => { const recoveredEventInd = events diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index d113e23c5b..85c5b6f69e 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -170,7 +170,9 @@ export interface TaskStore { options: TaskStoreCreateTaskOptions, ): Promise; - recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise; + recoverTasks?( + options: TaskStoreRecoverTaskOptions, + ): Promise<{ id: string }[]>; getTask(taskId: string): Promise; diff --git a/plugins/scaffolder-node/src/tasks/types.ts b/plugins/scaffolder-node/src/tasks/types.ts index 12e6416749..7cdc044bdf 100644 --- a/plugins/scaffolder-node/src/tasks/types.ts +++ b/plugins/scaffolder-node/src/tasks/types.ts @@ -131,7 +131,7 @@ export interface TaskBroker { claim(): Promise; - recoverTasks?(): Promise; + recoverTasks?(): Promise; dispatch( options: TaskBrokerDispatchOptions, diff --git a/yarn.lock b/yarn.lock index 2fee1c1d2c..4992eb3807 100644 --- a/yarn.lock +++ b/yarn.lock @@ -8190,7 +8190,6 @@ __metadata: "@backstage/plugin-auth-node": "workspace:^" "@backstage/plugin-catalog-backend-module-scaffolder-entity-model": "workspace:^" "@backstage/plugin-catalog-node": "workspace:^" - "@backstage/plugin-events-node": "workspace:^" "@backstage/plugin-permission-common": "workspace:^" "@backstage/plugin-permission-node": "workspace:^" "@backstage/plugin-scaffolder-backend-module-azure": "workspace:^" From 03cb5dc5d06eed70dd87d51cb1b157af9599df14 Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sat, 20 Jan 2024 15:52:39 +0100 Subject: [PATCH 05/10] incorporated some of feedback Signed-off-by: bnechyporenko --- plugins/scaffolder-backend/api-report.md | 12 ++++++++++-- plugins/scaffolder-node/api-report.md | 2 +- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index fd2d368f15..f6ca4045d3 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -406,7 +406,11 @@ export class DatabaseTaskStore implements TaskStore { }[]; }>; // (undocumented) - recoverTasks(options: TaskStoreRecoverTaskOptions): Promise; + recoverTasks(options: TaskStoreRecoverTaskOptions): Promise< + { + id: string; + }[] + >; // (undocumented) shutdownTask(options: TaskStoreShutDownTaskOptions): Promise; } @@ -559,7 +563,11 @@ export interface TaskStore { }[]; }>; // (undocumented) - recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise; + recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise< + { + id: string; + }[] + >; // (undocumented) shutdownTask?(options: TaskStoreShutDownTaskOptions): Promise; } diff --git a/plugins/scaffolder-node/api-report.md b/plugins/scaffolder-node/api-report.md index 9a71d41969..356fc37cc6 100644 --- a/plugins/scaffolder-node/api-report.md +++ b/plugins/scaffolder-node/api-report.md @@ -238,7 +238,7 @@ export interface TaskBroker { tasks: SerializedTask[]; }>; // (undocumented) - recoverTasks?(): Promise; + recoverTasks?(): Promise; // (undocumented) vacuumTasks(options: { timeoutS: number }): Promise; } From 04bab8a81a11c41bf840618b576005fafa7ab78f Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 21 Jan 2024 12:37:35 +0100 Subject: [PATCH 06/10] stop running new tasks after shutdown Signed-off-by: bnechyporenko --- .../src/scaffolder/tasks/TaskWorker.ts | 17 ++++++++++++----- .../scaffolder-backend/src/service/router.ts | 5 +++++ 2 files changed, 17 insertions(+), 5 deletions(-) diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 5e116922de..cd7dc06a91 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -76,8 +76,10 @@ export type CreateWorkerOptions = { export class TaskWorker { private taskQueue: PQueue; private logger: Logger | undefined; + private stopWorkers: boolean; private constructor(private readonly options: TaskWorkerOptions) { + this.stopWorkers = false; this.logger = options.logger; this.taskQueue = new PQueue({ concurrency: options.concurrentTasksLimit, @@ -120,26 +122,31 @@ export class TaskWorker { await this.options.taskBroker.recoverTasks?.(); } catch (err) { this.logger?.error(stringifyError(err)); - // ignore } } start() { (async () => { - for (;;) { + while (!this.stopWorkers) { await new Promise(resolve => setTimeout(resolve, 10000)); await this.recoverTasks(); } })(); (async () => { - for (;;) { + while (!this.stopWorkers) { await this.onReadyToClaimTask(); - const task = await this.options.taskBroker.claim(); - this.taskQueue.add(() => this.runOneTask(task)); + if (!this.stopWorkers) { + const task = await this.options.taskBroker.claim(); + void this.taskQueue.add(() => this.runOneTask(task)); + } } })(); } + stop() { + this.stopWorkers = true; + } + protected onReadyToClaimTask(): Promise { if (this.taskQueue.pending < this.options.concurrentTasksLimit) { return Promise.resolve(); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 365ec8748d..c0322af283 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -336,8 +336,13 @@ export async function createRouter( const launchWorkers = () => workers.forEach(worker => worker.start()); + const shutdownWorkers = () => { + workers.forEach(worker => worker.stop()); + }; + if (options.lifecycle) { options.lifecycle.addStartupHook(launchWorkers); + options.lifecycle.addShutdownHook(shutdownWorkers); } else { launchWorkers(); } From 898f513bace1398944ff5a62ffc99c62ce6f0f8c Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 21 Jan 2024 12:53:14 +0100 Subject: [PATCH 07/10] feedback Signed-off-by: bnechyporenko --- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 28 ++----------- .../src/scaffolder/tasks/dbUtil.test.ts | 40 +++++++++++++++++++ .../src/scaffolder/tasks/dbUtil.ts | 32 +++++++++++++++ 3 files changed, 75 insertions(+), 25 deletions(-) create mode 100644 plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.test.ts create mode 100644 plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.ts diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 0ce4e38659..35119368fa 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -41,6 +41,7 @@ import { import { DateTime, Duration } from 'luxon'; import { TaskRecovery, TaskSpec } from '@backstage/plugin-scaffolder-common'; import { trimEventsTillLastRecovery } from './taskRecoveryHelper'; +import { intervalFromNowTill } from './dbUtil'; const migrationsDir = resolvePackagePath( '@backstage/plugin-scaffolder-backend', @@ -294,18 +295,7 @@ export class DatabaseTaskStore implements TaskStore { tasks: { taskId: string; recovery?: TaskRecovery }[]; }> { const { timeoutS } = options; - let heartbeatInterval = this.db.raw(`? - interval '${timeoutS} seconds'`, [ - this.db.fn.now(), - ]); - if (this.db.client.config.client.includes('mysql')) { - heartbeatInterval = this.db.raw( - `date_sub(now(), interval ${timeoutS} second)`, - ); - } else if (this.db.client.config.client.includes('sqlite3')) { - heartbeatInterval = this.db.raw(`datetime('now', ?)`, [ - `-${timeoutS} seconds`, - ]); - } + const heartbeatInterval = intervalFromNowTill(timeoutS, this.db); const rawRows = await this.db('tasks') .where('status', 'processing') .andWhere('last_heartbeat_at', '<=', heartbeatInterval); @@ -498,19 +488,7 @@ export class DatabaseTaskStore implements TaskStore { const timeoutS = Duration.fromObject(options.timeout).as('seconds'); await this.db.transaction(async tx => { - let heartbeatInterval = this.db.raw( - `? - interval '${timeoutS} seconds'`, - [this.db.fn.now()], - ); - if (this.db.client.config.client.includes('mysql')) { - heartbeatInterval = this.db.raw( - `date_sub(now(), interval ${timeoutS} second)`, - ); - } else if (this.db.client.config.client.includes('sqlite3')) { - heartbeatInterval = this.db.raw(`datetime('now', ?)`, [ - `-${timeoutS} seconds`, - ]); - } + const heartbeatInterval = intervalFromNowTill(timeoutS, this.db); const result = await tx('tasks') .where('status', 'processing') diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.test.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.test.ts new file mode 100644 index 0000000000..9ccc17b8a4 --- /dev/null +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.test.ts @@ -0,0 +1,40 @@ +/* + * Copyright 2024 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +import knexFactory, { Knex } from 'knex'; +import { intervalFromNowTill } from './dbUtil'; + +class KnexBuilder { + public build(client: string): Knex { + return knexFactory({ client, useNullAsDefault: true }); + } +} + +describe('util', () => { + describe('intervalFromNowTill', () => { + const databases = [ + { client: 'sqlite3', expected: "datetime('now', '-5 seconds')" }, + { client: 'mysql', expected: 'date_sub(now(), interval 5 second)' }, + { client: 'pg', expected: "CURRENT_TIMESTAMP - interval '5 seconds'" }, + ]; + + it.each(databases)('for client $client', ({ client, expected }) => { + const knex = new KnexBuilder().build(client); + const result = intervalFromNowTill(5, knex); + + expect(result.toString()).toBe(expected); + }); + }); +}); diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.ts new file mode 100644 index 0000000000..afc4e8c0d6 --- /dev/null +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/dbUtil.ts @@ -0,0 +1,32 @@ +/* + * Copyright 2024 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +import { Knex } from 'knex'; + +export const intervalFromNowTill = (timeoutS: number, knex: Knex) => { + let heartbeatInterval = knex.raw(`? - interval '${timeoutS} seconds'`, [ + knex.fn.now(), + ]); + if (knex.client.config.client.includes('mysql')) { + heartbeatInterval = knex.raw( + `date_sub(now(), interval ${timeoutS} second)`, + ); + } else if (knex.client.config.client.includes('sqlite3')) { + heartbeatInterval = knex.raw(`datetime('now', ?)`, [ + `-${timeoutS} seconds`, + ]); + } + return heartbeatInterval; +}; From 6ac91e818035958b403c43a3b78aacaf097ba8ea Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 21 Jan 2024 13:23:35 +0100 Subject: [PATCH 08/10] api-report Signed-off-by: bnechyporenko --- plugins/scaffolder-backend/api-report.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index f6ca4045d3..98f0af527c 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -618,6 +618,8 @@ export class TaskWorker { runOneTask(task: TaskContext): Promise; // (undocumented) start(): void; + // (undocumented) + stop(): void; } // @public @deprecated (undocumented) From 90fbb8e95d304dc114a9954ef6368b4269670bba Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 21 Jan 2024 20:31:16 +0100 Subject: [PATCH 09/10] wip Signed-off-by: bnechyporenko --- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 35119368fa..901bf012ac 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -484,7 +484,7 @@ export class DatabaseTaskStore implements TaskStore { async recoverTasks( options: TaskStoreRecoverTaskOptions, ): Promise<{ id: string }[]> { - const taskIdsToRecover: string[] = []; + const taskIdsToRecover: { id: string }[] = []; const timeoutS = Duration.fromObject(options.timeout).as('seconds'); await this.db.transaction(async tx => { @@ -501,7 +501,7 @@ export class DatabaseTaskStore implements TaskStore { ['id', 'spec'], ); - taskIdsToRecover.push(...result.map(i => i.id)); + taskIdsToRecover.push(...result.map(i => ({ id: i.id }))); for (const { id, spec } of result) { const taskSpec = JSON.parse(spec as string) as TaskSpec; @@ -516,6 +516,6 @@ export class DatabaseTaskStore implements TaskStore { } }); - return taskIdsToRecover.map(id => ({ id })); + return taskIdsToRecover; } } From 49062022a56e9669c43e0eb2695441f1edbfd61c Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Tue, 23 Jan 2024 13:36:17 +0100 Subject: [PATCH 10/10] feedback Signed-off-by: bnechyporenko --- plugins/scaffolder-backend/api-report.md | 16 ++++++---------- .../src/scaffolder/tasks/DatabaseTaskStore.ts | 8 ++++---- .../src/scaffolder/tasks/StorageTaskBroker.ts | 7 +++---- .../src/scaffolder/tasks/types.ts | 2 +- 4 files changed, 14 insertions(+), 19 deletions(-) diff --git a/plugins/scaffolder-backend/api-report.md b/plugins/scaffolder-backend/api-report.md index 98f0af527c..cc16791390 100644 --- a/plugins/scaffolder-backend/api-report.md +++ b/plugins/scaffolder-backend/api-report.md @@ -406,11 +406,9 @@ export class DatabaseTaskStore implements TaskStore { }[]; }>; // (undocumented) - recoverTasks(options: TaskStoreRecoverTaskOptions): Promise< - { - id: string; - }[] - >; + recoverTasks(options: TaskStoreRecoverTaskOptions): Promise<{ + ids: string[]; + }>; // (undocumented) shutdownTask(options: TaskStoreShutDownTaskOptions): Promise; } @@ -563,11 +561,9 @@ export interface TaskStore { }[]; }>; // (undocumented) - recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise< - { - id: string; - }[] - >; + recoverTasks?(options: TaskStoreRecoverTaskOptions): Promise<{ + ids: string[]; + }>; // (undocumented) shutdownTask?(options: TaskStoreShutDownTaskOptions): Promise; } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 901bf012ac..766db4646e 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -483,8 +483,8 @@ export class DatabaseTaskStore implements TaskStore { async recoverTasks( options: TaskStoreRecoverTaskOptions, - ): Promise<{ id: string }[]> { - const taskIdsToRecover: { id: string }[] = []; + ): Promise<{ ids: string[] }> { + const taskIdsToRecover: string[] = []; const timeoutS = Duration.fromObject(options.timeout).as('seconds'); await this.db.transaction(async tx => { @@ -501,7 +501,7 @@ export class DatabaseTaskStore implements TaskStore { ['id', 'spec'], ); - taskIdsToRecover.push(...result.map(i => ({ id: i.id }))); + taskIdsToRecover.push(...result.map(i => i.id)); for (const { id, spec } of result) { const taskSpec = JSON.parse(spec as string) as TaskSpec; @@ -516,6 +516,6 @@ export class DatabaseTaskStore implements TaskStore { } }); - return taskIdsToRecover; + return { ids: taskIdsToRecover }; } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 96300bec58..8b49f492e8 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -220,10 +220,9 @@ export class StorageTaskBroker implements TaskBroker { 'scaffolder.EXPERIMENTAL_recoverTasksTimeout', defaultTimeout, ); - const recoveredTaskIds = - (await this.storage.recoverTasks?.({ - timeout, - })) ?? []; + const { ids: recoveredTaskIds } = (await this.storage.recoverTasks?.({ + timeout, + })) ?? { ids: [] }; if (recoveredTaskIds.length > 0) { this.signalDispatch(); } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index 85c5b6f69e..c5783ccb49 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -172,7 +172,7 @@ export interface TaskStore { recoverTasks?( options: TaskStoreRecoverTaskOptions, - ): Promise<{ id: string }[]>; + ): Promise<{ ids: string[] }>; getTask(taskId: string): Promise;