Adding manual scaffolder retry

Signed-off-by: bnechyporenko <bnechyporenko@bol.com>
This commit is contained in:
bnechyporenko
2024-07-21 19:49:02 +02:00
committed by blam
parent b304f78a31
commit 58410465bb
11 changed files with 163 additions and 18 deletions
@@ -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<RawDbTaskEventRow>('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<void> {
await this.db.transaction(async tx => {
const result = await tx<RawDbTaskRow>('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<RawDbTaskEventRow>('task_events')
.where('task_id', id)
.andWhere(q => q.whereIn('event_type', ['cancelled', 'completion']))
.del();
await tx<RawDbTaskEventRow>('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[] }> {
@@ -254,7 +254,9 @@ export interface CurrentClaimedTask {
* The creator of the task.
*/
createdBy?: string;
/**
* The workspace of the task.
*/
workspace?: Promise<Buffer>;
}
@@ -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<void> {
await this.storage.retryTask?.({ taskId });
this.signalDispatch();
}
}
@@ -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<void>;
@@ -74,6 +75,7 @@ export class DefaultWorkspaceService implements WorkspaceService {
targetPath: string;
}): Promise<void> {
if (this.isWorkspaceSerializationEnabled()) {
await fs.mkdirp(options.targetPath);
await this.workspaceProvider.rehydrateWorkspace(options);
}
}
@@ -119,6 +119,7 @@ export type TaskStoreEmitOptions<TBody = JsonObject> = {
* @public
*/
export type TaskStoreListEventsOptions = {
isTaskRecoverable?: boolean;
taskId: string;
after?: number | undefined;
};
@@ -170,6 +171,8 @@ export interface TaskStore {
options: TaskStoreCreateTaskOptions,
): Promise<TaskStoreCreateTaskResult>;
retryTask?(options: { taskId: string }): Promise<void>;
recoverTasks?(
options: TaskStoreRecoverTaskOptions,
): Promise<{ ids: string[] }>;
@@ -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;
}
}
@@ -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<void>;
retry?(taskId: string): Promise<void>;
claim(): Promise<TaskContext>;
recoverTasks?(): Promise<void>;
@@ -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<void>;
/**
* Starts the task again from the point where it failed.
*
* @param taskId - the id of the task
*/
retry?(taskId: string): Promise<void>;
listTasks?(options: {
filterByOwnership: 'owned' | 'all';
}): Promise<{ tasks: ScaffolderTask[] }>;
@@ -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<LogEvent>();
@@ -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]);
+19 -6
View File
@@ -217,12 +217,10 @@ export class ScaffolderClient implements ScaffolderApi {
}
private streamLogsEventStream({
isTaskRecoverable,
taskId,
after,
}: {
taskId: string;
after?: number;
}): Observable<LogEvent> {
}: ScaffolderStreamLogsOptions): Observable<LogEvent> {
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<void> {
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,
@@ -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<Theme, { fontColor: string }>(() => ({
export const ContextMenu = (props: ContextMenuProps) => {
const {
cancelEnabled,
canRetry,
logsVisible,
buttonBarVisible,
onRetry,
onStartOver,
onToggleLogs,
onToggleButtonBar,
@@ -151,6 +155,16 @@ export const ContextMenu = (props: ContextMenuProps) => {
</ListItemIcon>
<ListItemText primary={t('ongoingTask.contextMenu.startOver')} />
</MenuItem>
<MenuItem
onClick={onRetry}
disabled={cancelEnabled || !canRetry}
data-testid="retry-task"
>
<ListItemIcon>
<Retry fontSize="small" />
</ListItemIcon>
<ListItemText primary="Retry" />
</MenuItem>
<MenuItem
onClick={cancel}
disabled={
@@ -58,6 +58,9 @@ const useStyles = makeStyles(theme => ({
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: {
>
<ContextMenu
cancelEnabled={cancelEnabled}
canRetry={canRetry}
logsVisible={logsVisible}
buttonBarVisible={buttonBarVisible}
onStartOver={startOver}
onRetry={triggerRetry}
onToggleLogs={setLogVisibleState}
onToggleButtonBar={setButtonBarVisibleState}
taskId={taskId}
@@ -237,6 +255,14 @@ export const OngoingTask = (props: {
>
{t('ongoingTask.cancelButtonTitle')}
</Button>
<Button
className={classes.retryButton}
disabled={cancelEnabled || !canRetry}
onClick={triggerRetry}
data-testid="retry-button"
>
Retry
</Button>
<Button
className={classes.logsVisibilityButton}
color="primary"