diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/DatabaseTaskStore.ts index 20e9672f04..abf8ae7f41 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 listEvents( options: TaskStoreListEventsOptions, ): Promise<{ events: SerializedTaskEvent[] }> { - const { taskId, after } = options; + const { isTaskRecoverable, taskId, after } = options; const rawEvents = await this.db('task_events') .where({ task_id: taskId, @@ -502,6 +502,7 @@ export class DatabaseTaskStore implements TaskStore { const body = JSON.parse(event.body) as JsonObject; return { id: Number(event.id), + isTaskRecoverable, taskId, body, type: event.event_type, @@ -602,6 +603,38 @@ export class DatabaseTaskStore implements TaskStore { }); } + async retryTask?(options: { taskId: string }): Promise { + await this.db.transaction(async tx => { + const result = await tx('tasks') + .where('id', options.taskId) + .update( + { + status: 'open', + last_heartbeat_at: this.db.fn.now(), + }, + ['id', 'spec'], + ); + + for (const { id, spec } of result) { + const taskSpec = JSON.parse(spec as string) as TaskSpec; + + await tx('task_events') + .where('task_id', id) + .andWhere(q => q.whereIn('event_type', ['cancelled', 'completion'])) + .del(); + + await tx('task_events').insert({ + task_id: id, + event_type: 'recovered', + body: JSON.stringify({ + recoverStrategy: + taskSpec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy ?? 'none', + }), + }); + } + }); + } + async recoverTasks( options: TaskStoreRecoverTaskOptions, ): Promise<{ ids: string[] }> { diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts index 3960228529..1da2b4b349 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/StorageTaskBroker.ts @@ -254,7 +254,9 @@ export interface CurrentClaimedTask { * The creator of the task. */ createdBy?: string; - + /** + * The workspace of the task. + */ workspace?: Promise; } @@ -317,7 +319,7 @@ export class StorageTaskBroker implements TaskBroker { shouldUnsubscribe = true; } - if (event.type === 'completion') { + if (event.type === 'completion' && !event.isTaskRecoverable) { shouldUnsubscribe = true; } } @@ -416,8 +418,17 @@ export class StorageTaskBroker implements TaskBroker { let cancelled = false; (async () => { + const task = await this.storage.getTask(taskId); + const isTaskRecoverable = + task.spec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy === + 'startOver'; + while (!cancelled) { - const result = await this.storage.listEvents({ taskId, after }); + const result = await this.storage.listEvents({ + isTaskRecoverable, + taskId, + after, + }); const { events } = result; if (events.length) { after = events[events.length - 1].id; @@ -485,4 +496,9 @@ export class StorageTaskBroker implements TaskBroker { }, }); } + + async retry?(taskId: string): Promise { + await this.storage.retryTask?.({ taskId }); + this.signalDispatch(); + } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/WorkspaceService.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/WorkspaceService.ts index ab424806bf..bb3ce448fc 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/WorkspaceService.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/WorkspaceService.ts @@ -19,6 +19,7 @@ import { CurrentClaimedTask } from './StorageTaskBroker'; import { WorkspaceProvider } from '@backstage/plugin-scaffolder-node/alpha'; import { DatabaseWorkspaceProvider } from './DatabaseWorkspaceProvider'; import { TaskStore } from './types'; +import fs from 'fs-extra'; export interface WorkspaceService { serializeWorkspace(options: { path: string }): Promise; @@ -74,6 +75,7 @@ export class DefaultWorkspaceService implements WorkspaceService { targetPath: string; }): Promise { if (this.isWorkspaceSerializationEnabled()) { + await fs.mkdirp(options.targetPath); await this.workspaceProvider.rehydrateWorkspace(options); } } diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts index b87e0c4b9b..874793922f 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/types.ts @@ -119,6 +119,7 @@ export type TaskStoreEmitOptions = { * @public */ export type TaskStoreListEventsOptions = { + isTaskRecoverable?: boolean; taskId: string; after?: number | undefined; }; @@ -170,6 +171,8 @@ export interface TaskStore { options: TaskStoreCreateTaskOptions, ): Promise; + retryTask?(options: { taskId: string }): Promise; + recoverTasks?( options: TaskStoreRecoverTaskOptions, ): Promise<{ ids: string[] }>; diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 4584d64257..80406496d2 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -652,6 +652,19 @@ export async function createRouter( await taskBroker.cancel?.(taskId); res.status(200).json({ status: 'cancelled' }); }) + .post('/v2/tasks/:taskId/retry', async (req, res) => { + const credentials = await httpAuth.credentials(req); + // Requires both read and cancel permissions + await checkPermission({ + credentials, + permissions: [taskCreatePermission, taskReadPermission], + permissionService: permissions, + }); + + const { taskId } = req.params; + await taskBroker.retry?.(taskId); + res.status(201).json({ id: taskId }); + }) .get('/v2/tasks/:taskId/eventstream', async (req, res) => { const credentials = await httpAuth.credentials(req); await checkPermission({ @@ -687,7 +700,7 @@ export async function createRouter( res.write( `event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`, ); - if (event.type === 'completion') { + if (event.type === 'completion' && !event.isTaskRecoverable) { shouldUnsubscribe = true; } } diff --git a/plugins/scaffolder-node/src/tasks/types.ts b/plugins/scaffolder-node/src/tasks/types.ts index 27a131c513..b830c0633a 100644 --- a/plugins/scaffolder-node/src/tasks/types.ts +++ b/plugins/scaffolder-node/src/tasks/types.ts @@ -76,6 +76,7 @@ export type TaskEventType = 'completion' | 'log' | 'cancelled' | 'recovered'; */ export type SerializedTaskEvent = { id: number; + isTaskRecoverable?: boolean; taskId: string; body: JsonObject; type: TaskEventType; @@ -163,6 +164,8 @@ export interface TaskContext { export interface TaskBroker { cancel?(taskId: string): Promise; + retry?(taskId: string): Promise; + claim(): Promise; recoverTasks?(): Promise; diff --git a/plugins/scaffolder-react/src/api/types.ts b/plugins/scaffolder-react/src/api/types.ts index eaef54ff1b..3816c4bcbf 100644 --- a/plugins/scaffolder-react/src/api/types.ts +++ b/plugins/scaffolder-react/src/api/types.ts @@ -161,6 +161,7 @@ export interface ScaffolderGetIntegrationsListResponse { * @public */ export interface ScaffolderStreamLogsOptions { + isTaskRecoverable?: boolean; taskId: string; after?: number; } @@ -213,6 +214,13 @@ export interface ScaffolderApi { */ cancelTask(taskId: string): Promise; + /** + * Starts the task again from the point where it failed. + * + * @param taskId - the id of the task + */ + retry?(taskId: string): Promise; + listTasks?(options: { filterByOwnership: 'owned' | 'all'; }): Promise<{ tasks: ScaffolderTask[] }>; diff --git a/plugins/scaffolder-react/src/hooks/useEventStream.ts b/plugins/scaffolder-react/src/hooks/useEventStream.ts index f63f6024e6..7c5ce09ab7 100644 --- a/plugins/scaffolder-react/src/hooks/useEventStream.ts +++ b/plugins/scaffolder-react/src/hooks/useEventStream.ts @@ -143,6 +143,11 @@ function reducer(draft: TaskStream, action: ReducerAction) { } case 'RECOVERED': { + draft.cancelled = false; + draft.completed = false; + draft.output = undefined; + draft.error = undefined; + for (const stepId in draft.steps) { if (draft.steps.hasOwnProperty(stepId)) { draft.steps[stepId].startedAt = undefined; @@ -185,12 +190,16 @@ export const useTaskEventStream = (taskId: string): TaskStream => { let subscription: Subscription | undefined; let logPusher: NodeJS.Timeout | undefined; let retryCount = 1; + let isTaskRecoverable = false; const startStreamLogProcess = () => scaffolderApi.getTask(taskId).then( task => { if (didCancel) { return; } + isTaskRecoverable = + task.spec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy === + 'startOver'; dispatch({ type: 'INIT', data: task }); // TODO(blam): Use a normal fetch to fetch the current log for the event stream @@ -199,7 +208,10 @@ export const useTaskEventStream = (taskId: string): TaskStream => { // stream logs. Without this, if you have a lot of logs, it can look like the // task is being rebuilt on load as it progresses through the steps at a slower // rate whilst it builds the status from the event logs - const observable = scaffolderApi.streamLogs({ taskId }); + const observable = scaffolderApi.streamLogs({ + isTaskRecoverable, + taskId, + }); const collectedLogEvents = new Array(); @@ -270,12 +282,14 @@ export const useTaskEventStream = (taskId: string): TaskStream => { ); void startStreamLogProcess(); return () => { - didCancel = true; - if (subscription) { - subscription.unsubscribe(); - } - if (logPusher) { - clearInterval(logPusher); + if (!isTaskRecoverable) { + didCancel = true; + if (subscription) { + subscription.unsubscribe(); + } + if (logPusher) { + clearInterval(logPusher); + } } }; }, [scaffolderApi, dispatch, taskId]); diff --git a/plugins/scaffolder/src/api.ts b/plugins/scaffolder/src/api.ts index 1809ed4c15..93db7386f3 100644 --- a/plugins/scaffolder/src/api.ts +++ b/plugins/scaffolder/src/api.ts @@ -217,12 +217,10 @@ export class ScaffolderClient implements ScaffolderApi { } private streamLogsEventStream({ + isTaskRecoverable, taskId, after, - }: { - taskId: string; - after?: number; - }): Observable { + }: ScaffolderStreamLogsOptions): Observable { return new ObservableImpl(subscriber => { const params = new URLSearchParams(); if (after !== undefined) { @@ -246,14 +244,14 @@ export class ScaffolderClient implements ScaffolderApi { }; const ctrl = new AbortController(); - fetchEventSource(url, { + void fetchEventSource(url, { fetch: this.fetchApi.fetch, signal: ctrl.signal, onmessage(e: EventSourceMessage) { if (e.event === 'log') { processEvent(e); return; - } else if (e.event === 'completion') { + } else if (e.event === 'completion' && !isTaskRecoverable) { processEvent(e); subscriber.complete(); ctrl.abort(); @@ -338,6 +336,21 @@ export class ScaffolderClient implements ScaffolderApi { return await response.json(); } + async retry?(taskId: string): Promise { + const baseUrl = await this.discoveryApi.getBaseUrl('scaffolder'); + const url = `${baseUrl}/v2/tasks/${encodeURIComponent(taskId)}/retry`; + + const response = await this.fetchApi.fetch(url, { + method: 'POST', + }); + + if (!response.ok) { + throw await ResponseError.fromResponse(response); + } + + return await response.json(); + } + async autocomplete({ token, resource, diff --git a/plugins/scaffolder/src/components/OngoingTask/ContextMenu.tsx b/plugins/scaffolder/src/components/OngoingTask/ContextMenu.tsx index e0ab01fc67..c3e9250d2b 100644 --- a/plugins/scaffolder/src/components/OngoingTask/ContextMenu.tsx +++ b/plugins/scaffolder/src/components/OngoingTask/ContextMenu.tsx @@ -41,8 +41,10 @@ import { scaffolderTranslationRef } from '../../translation'; type ContextMenuProps = { cancelEnabled?: boolean; + canRetry: boolean; logsVisible?: boolean; buttonBarVisible?: boolean; + onRetry?: () => void; onStartOver?: () => void; onToggleLogs?: (state: boolean) => void; onToggleButtonBar?: (state: boolean) => void; @@ -58,8 +60,10 @@ const useStyles = makeStyles(() => ({ export const ContextMenu = (props: ContextMenuProps) => { const { cancelEnabled, + canRetry, logsVisible, buttonBarVisible, + onRetry, onStartOver, onToggleLogs, onToggleButtonBar, @@ -151,6 +155,16 @@ export const ContextMenu = (props: ContextMenuProps) => { + + + + + + ({ cancelButton: { marginRight: theme.spacing(1), }, + retryButton: { + marginRight: theme.spacing(1), + }, logsVisibilityButton: { marginRight: theme.spacing(1), }, @@ -130,6 +133,12 @@ export const OngoingTask = (props: { return 0; }, [steps]); + const isRetryableTask = + taskStream.task?.spec.EXPERIMENTAL_recovery?.EXPERIMENTAL_strategy === + 'startOver'; + + const canRetry = canReadTask && canCreateTask && isRetryableTask; + const startOver = useCallback(() => { const { namespace, name } = taskStream.task?.spec.templateInfo?.entity?.metadata ?? {}; @@ -157,6 +166,13 @@ export const OngoingTask = (props: { templateRouteRef, ]); + const [{ status: _ }, { execute: triggerRetry }] = useAsync(async () => { + if (taskId) { + analytics.captureEvent('retried', 'Template has been retried'); + await scaffolderApi.retry?.(taskId); + } + }); + const [{ status: cancelStatus }, { execute: triggerCancel }] = useAsync( async () => { if (taskId) { @@ -190,9 +206,11 @@ export const OngoingTask = (props: { > {t('ongoingTask.cancelButtonTitle')} +