mirror of
https://github.com/openclaw/openclaw.git
synced 2026-07-20 02:41:41 +00:00
fix(tasks): fail closed when task-flow state cannot restore (#105088)
* fix(tasks): fail closed on task-flow restore errors * fix(tasks): keep maintenance checks behind owner seam * fix(tasks): gate runtime startup on flow restore * test(tasks): prove runtime restore gate failures * fix(tasks): gate cancellation on flow readiness
This commit is contained in:
@@ -47,5 +47,5 @@ export {
|
||||
waitForActiveTasks,
|
||||
} from "../../process/command-queue.js";
|
||||
export { getInspectableActiveTaskRestartBlockers } from "../../tasks/task-registry.maintenance.js";
|
||||
export { reloadTaskRegistryFromStore } from "../../tasks/runtime-internal.js";
|
||||
export { reloadTaskRuntimeStateFromStore } from "../../tasks/runtime-internal.js";
|
||||
export { abortPendingChannelReloads } from "../../gateway/server-reload-handlers.js";
|
||||
|
||||
@@ -67,7 +67,7 @@ const waitForActiveCronJobs = vi.fn(async (_timeoutMs?: number) => ({
|
||||
drained: true,
|
||||
active: 0,
|
||||
}));
|
||||
const reloadTaskRegistryFromStore = vi.fn();
|
||||
const reloadTaskRuntimeStateFromStore = vi.fn();
|
||||
const rotateAgentEventLifecycleGeneration = vi.fn();
|
||||
const clearRuntimeConfigSnapshot = vi.fn();
|
||||
const restartGatewayProcessWithFreshPid = vi.fn<
|
||||
@@ -180,7 +180,7 @@ vi.mock("../../tasks/cron-task-cancel.js", () => ({
|
||||
}));
|
||||
|
||||
vi.mock("../../tasks/runtime-internal.js", () => ({
|
||||
reloadTaskRegistryFromStore: () => reloadTaskRegistryFromStore(),
|
||||
reloadTaskRuntimeStateFromStore: () => reloadTaskRuntimeStateFromStore(),
|
||||
}));
|
||||
|
||||
vi.mock("../../infra/agent-events.js", () => ({
|
||||
@@ -947,7 +947,10 @@ describe("runGatewayLoop", () => {
|
||||
expect(clearRuntimeConfigSnapshot).toHaveBeenCalledTimes(1);
|
||||
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(1);
|
||||
expect(rotateAgentEventLifecycleGeneration).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRuntimeStateFromStore.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
|
||||
start.mock.invocationCallOrder[1] ?? Infinity,
|
||||
);
|
||||
expect(
|
||||
rotateAgentEventLifecycleGeneration.mock.invocationCallOrder[0] ?? Infinity,
|
||||
).toBeLessThan(resetAllLanes.mock.invocationCallOrder[0] ?? Infinity);
|
||||
@@ -983,7 +986,7 @@ describe("runGatewayLoop", () => {
|
||||
expect(clearRuntimeConfigSnapshot).toHaveBeenCalledTimes(2);
|
||||
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(2);
|
||||
expect(rotateAgentEventLifecycleGeneration).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(acquireGatewayLock).toHaveBeenCalledTimes(3);
|
||||
|
||||
sigterm();
|
||||
@@ -1095,7 +1098,7 @@ describe("runGatewayLoop", () => {
|
||||
expect(markGatewayDraining).toHaveBeenCalledTimes(1);
|
||||
expect(resetAllLanes).toHaveBeenCalledTimes(1);
|
||||
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(1);
|
||||
} finally {
|
||||
sigterm();
|
||||
await expect(exited).resolves.toBe(0);
|
||||
@@ -1284,7 +1287,7 @@ describe("runGatewayLoop", () => {
|
||||
expect(markGatewayDraining).toHaveBeenCalledTimes(2);
|
||||
expect(resetAllLanes).toHaveBeenCalledTimes(2);
|
||||
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(acquireGatewayLock).toHaveBeenCalledTimes(3);
|
||||
expect(gatewayLog.error).toHaveBeenCalledWith(
|
||||
expect.stringContaining("gateway startup failed: restart startup failed."),
|
||||
@@ -1356,7 +1359,7 @@ describe("runGatewayLoop", () => {
|
||||
expect(markGatewayDraining).toHaveBeenCalledTimes(2);
|
||||
expect(resetAllLanes).toHaveBeenCalledTimes(2);
|
||||
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(acquireGatewayLock).toHaveBeenCalledTimes(3);
|
||||
} finally {
|
||||
sigterm();
|
||||
@@ -1365,12 +1368,16 @@ describe("runGatewayLoop", () => {
|
||||
});
|
||||
});
|
||||
|
||||
it("keeps the process alive and retries after task-registry restore fails during restart", async () => {
|
||||
it("keeps the process alive and retries after task runtime state restores fail", async () => {
|
||||
vi.clearAllMocks();
|
||||
reloadTaskRegistryFromStore.mockReset();
|
||||
reloadTaskRegistryFromStore.mockImplementationOnce(() => {
|
||||
throw new Error("task registry restore failed");
|
||||
});
|
||||
reloadTaskRuntimeStateFromStore.mockReset();
|
||||
reloadTaskRuntimeStateFromStore
|
||||
.mockImplementationOnce(() => {
|
||||
throw new Error("task-flow registry restore failed");
|
||||
})
|
||||
.mockImplementationOnce(() => {
|
||||
throw new Error("task registry restore failed");
|
||||
});
|
||||
peekGatewaySigusr1RestartReason.mockReturnValue(undefined);
|
||||
respawnGatewayProcessForUpdate.mockReturnValue({
|
||||
mode: "disabled",
|
||||
@@ -1407,6 +1414,22 @@ describe("runGatewayLoop", () => {
|
||||
const sigterm = captureSignal("SIGTERM");
|
||||
|
||||
try {
|
||||
sigusr1();
|
||||
await waitForLoopCondition(
|
||||
() =>
|
||||
gatewayLog.error.mock.calls.some(([message]) =>
|
||||
String(message).includes(
|
||||
"gateway startup failed: task-flow registry restore failed.",
|
||||
),
|
||||
),
|
||||
"expected failed task-flow registry restore to be logged",
|
||||
);
|
||||
|
||||
expectRestartCloseCall(closeFirst, 90_000);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(1);
|
||||
expect(start).toHaveBeenCalledTimes(1);
|
||||
expect(runtime.exit).not.toHaveBeenCalled();
|
||||
|
||||
sigusr1();
|
||||
await waitForLoopCondition(
|
||||
() =>
|
||||
@@ -1416,15 +1439,14 @@ describe("runGatewayLoop", () => {
|
||||
"expected failed task-registry restore to be logged",
|
||||
);
|
||||
|
||||
expectRestartCloseCall(closeFirst, 90_000);
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(1);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(start).toHaveBeenCalledTimes(1);
|
||||
expect(runtime.exit).not.toHaveBeenCalled();
|
||||
|
||||
sigusr1();
|
||||
await startedSecond;
|
||||
|
||||
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(2);
|
||||
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(3);
|
||||
expect(start).toHaveBeenCalledTimes(2);
|
||||
expect(runtime.exit).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
@@ -1438,7 +1460,7 @@ describe("runGatewayLoop", () => {
|
||||
});
|
||||
});
|
||||
} finally {
|
||||
reloadTaskRegistryFromStore.mockReset();
|
||||
reloadTaskRuntimeStateFromStore.mockReset();
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
@@ -885,7 +885,7 @@ export async function runGatewayLoop(params: {
|
||||
const {
|
||||
abortActiveCronTaskRuns,
|
||||
advanceCronActiveJobGeneration,
|
||||
reloadTaskRegistryFromStore,
|
||||
reloadTaskRuntimeStateFromStore,
|
||||
retireActiveCronTaskRunTracking,
|
||||
resetCronActiveJobs,
|
||||
resetAllLanes,
|
||||
@@ -910,7 +910,7 @@ export async function runGatewayLoop(params: {
|
||||
resetAllLanes();
|
||||
clearRuntimeConfigSnapshot();
|
||||
resetGatewayRestartStateForInProcessRestart();
|
||||
reloadTaskRegistryFromStore();
|
||||
reloadTaskRuntimeStateFromStore();
|
||||
markGatewayRestartTrace("restart.next-start");
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user