From 07c2c2e7a793576318aec0c25609c13c765540e4 Mon Sep 17 00:00:00 2001 From: bnechyporenko Date: Sun, 4 Feb 2024 21:30:32 +0100 Subject: [PATCH] Checkpoint Signed-off-by: bnechyporenko --- packages/backend-next/package.json | 1 + .../src/scaffolder/tasks/DatabaseTaskStore.ts | 14 +++++++++++++- .../src/scaffolder/tasks/NunjucksWorkflowRunner.ts | 8 +++++++- .../src/scaffolder/tasks/StorageTaskBroker.ts | 8 ++++++-- .../src/scaffolder/tasks/types.ts | 6 ++++++ plugins/scaffolder-node/src/tasks/index.ts | 2 +- plugins/scaffolder-node/src/tasks/types.ts | 8 +++++--- 7 files changed, 39 insertions(+), 8 deletions(-) diff --git a/packages/backend-next/package.json b/packages/backend-next/package.json index eca36208a9..aca0acbdd3 100644 --- a/packages/backend-next/package.json +++ b/packages/backend-next/package.json @@ -52,6 +52,7 @@ "@backstage/plugin-playlist-backend": "workspace:^", "@backstage/plugin-proxy-backend": "workspace:^", "@backstage/plugin-scaffolder-backend": "workspace:^", + "@backstage/plugin-scaffolder-backend-module-github": "workspace:^", "@backstage/plugin-search-backend": "workspace:^", "@backstage/plugin-search-backend-module-catalog": "workspace:^", "@backstage/plugin-search-backend-module-explore": "workspace:^", diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index e6083a469f..14535d034c 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts @@ -38,6 +38,7 @@ import { TaskStatus, TaskEventType, TaskSecrets, + TaskState, } from '@backstage/plugin-scaffolder-node'; import { DateTime, Duration } from 'luxon'; import { TaskRecovery, TaskSpec } from '@backstage/plugin-scaffolder-common'; @@ -396,7 +397,18 @@ export class DatabaseTaskStore implements TaskStore { }); } - async saveCheckpoint?(options: TaskStoreStateOptions): Promise { + async listCheckpoints({ + taskId, + }: { + taskId: string; + }): Promise<{ state: TaskState }> { + const state = await this.db('tasks') + .where({ id: taskId }) + .select('state'); + return { state: JSON.stringify(state) as unknown as TaskState }; + } + + async saveCheckpoint(options: TaskStoreStateOptions): Promise { if (options.state) { const serializedState = JSON.stringify(options.state); await this.db('tasks') diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts index d746cf859f..407c623793 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/NunjucksWorkflowRunner.ts @@ -332,6 +332,7 @@ export class NunjucksWorkflowRunner implements WorkflowRunner { } const tmpDirs = new Array(); const stepOutput: { [outputName: string]: JsonValue } = {}; + const prevTaskState = await task.getCheckpoints?.(); for (const iteration of iterations) { if (iteration.each) { @@ -354,7 +355,12 @@ export class NunjucksWorkflowRunner implements WorkflowRunner { fn: () => Promise, ) { try { - const value = await fn(); + let prevValue: U | undefined; + if (prevTaskState) { + prevValue = prevTaskState.state[key] as unknown as U; + } + + const value = prevValue ? prevValue : await fn(); task.updateCheckpoint?.({ key, status: 'success', diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 01722edba9..b74c64d7a3 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -19,7 +19,7 @@ import { TaskSpec } from '@backstage/plugin-scaffolder-common'; import { TaskSecrets, TaskState, - UpdateCheckpointOptions, + CheckpointRecord, } from '@backstage/plugin-scaffolder-node'; import { JsonObject, Observable } from '@backstage/types'; import { Logger } from 'winston'; @@ -95,7 +95,11 @@ export class TaskManager implements TaskContext { }); } - async updateCheckpoint?(options: UpdateCheckpointOptions): Promise { + async getCheckpoints?(): Promise<{ state: TaskState } | undefined> { + return this.storage.listCheckpoints?.({ taskId: this.task.taskId }); + } + + async updateCheckpoint?(options: CheckpointRecord): Promise { if (this.task.state) { this.task.state[options.key] = { ...options }; } else { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index b6225d46a9..54c0f19782 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -199,6 +199,12 @@ export interface TaskStore { emitLogEvent(options: TaskStoreEmitOptions): Promise; + listCheckpoints?({ + taskId, + }: { + taskId: string; + }): Promise<{ state: TaskState }>; + saveCheckpoint?(options: TaskStoreStateOptions): Promise; listEvents( diff --git a/plugins/scaffolder-node/src/tasks/index.ts b/plugins/scaffolder-node/src/tasks/index.ts index 4e66f1c0e3..60023bf48e 100644 --- a/plugins/scaffolder-node/src/tasks/index.ts +++ b/plugins/scaffolder-node/src/tasks/index.ts @@ -26,5 +26,5 @@ export type { TaskEventType, TaskState, TaskStatus, - UpdateCheckpointOptions, + CheckpointRecord, } from './types'; diff --git a/plugins/scaffolder-node/src/tasks/types.ts b/plugins/scaffolder-node/src/tasks/types.ts index f57353c917..c6654ea642 100644 --- a/plugins/scaffolder-node/src/tasks/types.ts +++ b/plugins/scaffolder-node/src/tasks/types.ts @@ -116,12 +116,12 @@ export type TaskBrokerDispatchOptions = { }; /** - * The options passed to {@link TaskBroker.updateCheckpoint} + * The record passed to {@link TaskBroker.updateCheckpoint?} * Parameters to store the result of the executed checkpoint * * @public */ -export type UpdateCheckpointOptions = +export type CheckpointRecord = | { key: string; status: 'success'; @@ -151,7 +151,9 @@ export interface TaskContext { emitLog(message: string, logMetadata?: JsonObject): Promise; - updateCheckpoint?(options: UpdateCheckpointOptions): Promise; + getCheckpoints?(): Promise<{ state: TaskState } | undefined>; + + updateCheckpoint?(options: CheckpointRecord): Promise; getWorkspaceName(): Promise; }