Checkpoint

Signed-off-by: bnechyporenko <bnechyporenko@bol.com>
This commit is contained in:
bnechyporenko
2024-02-04 21:30:32 +01:00
parent 911557bae8
commit 07c2c2e7a7
7 changed files with 39 additions and 8 deletions
+1
View File
@@ -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:^",
@@ -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<void> {
async listCheckpoints({
taskId,
}: {
taskId: string;
}): Promise<{ state: TaskState }> {
const state = await this.db<RawDbTaskRow>('tasks')
.where({ id: taskId })
.select('state');
return { state: JSON.stringify(state) as unknown as TaskState };
}
async saveCheckpoint(options: TaskStoreStateOptions): Promise<void> {
if (options.state) {
const serializedState = JSON.stringify(options.state);
await this.db<RawDbTaskRow>('tasks')
@@ -332,6 +332,7 @@ export class NunjucksWorkflowRunner implements WorkflowRunner {
}
const tmpDirs = new Array<string>();
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<U>,
) {
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',
@@ -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<void> {
async getCheckpoints?(): Promise<{ state: TaskState } | undefined> {
return this.storage.listCheckpoints?.({ taskId: this.task.taskId });
}
async updateCheckpoint?(options: CheckpointRecord): Promise<void> {
if (this.task.state) {
this.task.state[options.key] = { ...options };
} else {
@@ -199,6 +199,12 @@ export interface TaskStore {
emitLogEvent(options: TaskStoreEmitOptions): Promise<void>;
listCheckpoints?({
taskId,
}: {
taskId: string;
}): Promise<{ state: TaskState }>;
saveCheckpoint?(options: TaskStoreStateOptions): Promise<void>;
listEvents(
+1 -1
View File
@@ -26,5 +26,5 @@ export type {
TaskEventType,
TaskState,
TaskStatus,
UpdateCheckpointOptions,
CheckpointRecord,
} from './types';
+5 -3
View File
@@ -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<void>;
updateCheckpoint?(options: UpdateCheckpointOptions): Promise<void>;
getCheckpoints?(): Promise<{ state: TaskState } | undefined>;
updateCheckpoint?(options: CheckpointRecord): Promise<void>;
getWorkspaceName(): Promise<string>;
}