refactor: controller doesn't shudown abort down to the agent
This commit is contained in:
parent
c27e717e0f
commit
729f6822d5
2 changed files with 1 additions and 57 deletions
|
|
@ -270,34 +270,6 @@ test("runBackup rejects before sending when the abort signal is already aborted"
|
||||||
await stopAgentController();
|
await stopAgentController();
|
||||||
});
|
});
|
||||||
|
|
||||||
test("runBackup requests cancellation when the abort signal fires while sending", async () => {
|
|
||||||
resetAgentRuntime();
|
|
||||||
const abortController = new AbortController();
|
|
||||||
controllerMock.sendBackup.mockImplementation(() =>
|
|
||||||
Effect.sync(() => {
|
|
||||||
abortController.abort();
|
|
||||||
return true;
|
|
||||||
}),
|
|
||||||
);
|
|
||||||
controllerMock.cancelBackup.mockImplementation(() => Effect.succeed(false));
|
|
||||||
const { agentManager, startAgentController, stopAgentController } = await import("../agents-manager");
|
|
||||||
|
|
||||||
await startAgentController();
|
|
||||||
const result = await agentManager.runBackup("local", {
|
|
||||||
scheduleId: 42,
|
|
||||||
payload: backupPayload,
|
|
||||||
signal: abortController.signal,
|
|
||||||
onProgress: vi.fn(),
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(result).toEqual({ status: "cancelled" });
|
|
||||||
expect(controllerMock.cancelBackup).toHaveBeenCalledWith("local", {
|
|
||||||
jobId: "job-1",
|
|
||||||
scheduleId: "schedule-1",
|
|
||||||
});
|
|
||||||
await stopAgentController();
|
|
||||||
});
|
|
||||||
|
|
||||||
test("restore events are delivered to the running restore callbacks", async () => {
|
test("restore events are delivered to the running restore callbacks", async () => {
|
||||||
resetAgentRuntime();
|
resetAgentRuntime();
|
||||||
controllerMock.sendRestore.mockImplementation(() => Effect.succeed(true));
|
controllerMock.sendRestore.mockImplementation(() => Effect.succeed(true));
|
||||||
|
|
|
||||||
|
|
@ -431,16 +431,8 @@ export const agentManager = {
|
||||||
});
|
});
|
||||||
getActiveBackupScheduleIdsByJobId().set(request.payload.jobId, request.scheduleId);
|
getActiveBackupScheduleIdsByJobId().set(request.payload.jobId, request.scheduleId);
|
||||||
});
|
});
|
||||||
const cancelOnAbort = () => {
|
|
||||||
void requestBackupCancellation(agentId, request.scheduleId).catch((error) => {
|
|
||||||
logger.error(`Failed to cancel backup ${request.scheduleId} after abort: ${String(error)}`);
|
|
||||||
});
|
|
||||||
};
|
|
||||||
request.signal.addEventListener("abort", cancelOnAbort, { once: true });
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
if (!(await Effect.runPromise(runtime.sendBackup(agentId, request.payload)))) {
|
if (!(await Effect.runPromise(runtime.sendBackup(agentId, request.payload)))) {
|
||||||
request.signal.removeEventListener("abort", cancelOnAbort);
|
|
||||||
clearActiveBackupRun(request.scheduleId);
|
clearActiveBackupRun(request.scheduleId);
|
||||||
return {
|
return {
|
||||||
status: "unavailable",
|
status: "unavailable",
|
||||||
|
|
@ -448,17 +440,10 @@ export const agentManager = {
|
||||||
} satisfies BackupExecutionResult;
|
} satisfies BackupExecutionResult;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (request.signal.aborted) {
|
|
||||||
await requestBackupCancellation(agentId, request.scheduleId);
|
|
||||||
}
|
|
||||||
|
|
||||||
return await completion;
|
return await completion;
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
request.signal.removeEventListener("abort", cancelOnAbort);
|
|
||||||
clearActiveBackupRun(request.scheduleId);
|
clearActiveBackupRun(request.scheduleId);
|
||||||
throw error;
|
throw error;
|
||||||
} finally {
|
|
||||||
request.signal.removeEventListener("abort", cancelOnAbort);
|
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
cancelBackup: async (agentId: string, scheduleId: number) => {
|
cancelBackup: async (agentId: string, scheduleId: number) => {
|
||||||
|
|
@ -504,16 +489,8 @@ export const agentManager = {
|
||||||
cancellationRequested: false,
|
cancellationRequested: false,
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
const cancelOnAbort = () => {
|
|
||||||
void requestRestoreCancellation(agentId, request.payload.restoreId).catch((error) => {
|
|
||||||
logger.error(`Failed to cancel restore ${request.payload.restoreId} after abort: ${String(error)}`);
|
|
||||||
});
|
|
||||||
};
|
|
||||||
request.signal.addEventListener("abort", cancelOnAbort, { once: true });
|
|
||||||
|
|
||||||
try {
|
try {
|
||||||
if (!(await Effect.runPromise(runtime.sendRestore(agentId, request.payload)))) {
|
if (!(await Effect.runPromise(runtime.sendRestore(agentId, request.payload)))) {
|
||||||
request.signal.removeEventListener("abort", cancelOnAbort);
|
|
||||||
clearActiveRestoreRun(request.payload.restoreId);
|
clearActiveRestoreRun(request.payload.restoreId);
|
||||||
return {
|
return {
|
||||||
status: "unavailable",
|
status: "unavailable",
|
||||||
|
|
@ -521,16 +498,11 @@ export const agentManager = {
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
if (request.signal.aborted) {
|
|
||||||
await requestRestoreCancellation(agentId, request.payload.restoreId);
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
status: "started",
|
status: "started",
|
||||||
result: completion.finally(() => request.signal.removeEventListener("abort", cancelOnAbort)),
|
result: completion,
|
||||||
};
|
};
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
request.signal.removeEventListener("abort", cancelOnAbort);
|
|
||||||
clearActiveRestoreRun(request.payload.restoreId);
|
clearActiveRestoreRun(request.payload.restoreId);
|
||||||
throw error;
|
throw error;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue