fix: rearm restart drain after lifecycle reset

This commit is contained in:
Shakker
2026-08-28 05:13:44 +02:00
parent b0d69b069c
commit 98bfeb480e
2 changed files with 35 additions and 37 deletions
+32 -34
View File
@@ -1202,6 +1202,13 @@ describe("runGatewayLoop", () => {
detail: "OPENCLAW_NO_RESPAWN",
});
markUpdateRestartSentinelFailure.mockClear();
let releaseFirstCronTaskDrain: (() => void) | undefined;
waitForActiveCronTaskRuns.mockImplementationOnce(
async () =>
await new Promise<{ drained: true; active: 0 }>((resolve) => {
releaseFirstCronTaskDrain = () => resolve({ drained: true, active: 0 });
}),
);
await withIsolatedSignals(async ({ captureSignal }) => {
const timedOutSnapshot = createActiveWorkSnapshot({ activeTasks: 2, embeddedRuns: 1 }, [
@@ -1244,11 +1251,17 @@ describe("runGatewayLoop", () => {
});
let resolveSecond: (() => void) | null = null;
let secondAgentEventGeneration: string | undefined;
let secondRestartDrainSignal: AbortSignal | undefined;
const startedSecond = new Promise<void>((resolve) => {
resolveSecond = resolve;
});
start.mockImplementationOnce(async () => {
expect(lifecycleSlot.size).toBe(0);
secondAgentEventGeneration = agentEventsActual.getAgentEventLifecycleGeneration();
secondRestartDrainSignal = gatewayWorkAdmissionActual.getGatewayRestartDrainSignal();
expect(secondRestartDrainSignal.aborted).toBe(false);
lifecycleSlot.set("second", 2);
resolveSecond?.();
return createGatewayServer(closeSecond);
});
@@ -1280,16 +1293,15 @@ describe("runGatewayLoop", () => {
sigusr1();
await waitForLoopCondition(
() => waitForActiveCronTaskRuns.mock.calls.length === 1,
"expected first restart to reach cron task drain",
);
sigusr1();
releaseFirstCronTaskDrain?.();
await startedSecond;
const secondAgentEventGeneration = agentEventsActual.getAgentEventLifecycleGeneration();
expect(secondAgentEventGeneration).not.toBe(firstAgentEventGeneration);
expect(gatewayWorkAdmissionActual.isGatewayWorkAdmissionClosed()).toBe(false);
expect(start).toHaveBeenCalledTimes(2);
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
expect(waitForGatewayActiveWork).toHaveBeenCalledOnce();
expect(waitForGatewayActiveWork.mock.calls[0]?.[0]).toBeLessThanOrEqual(
DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS,
);
@@ -1297,34 +1309,8 @@ describe("runGatewayLoop", () => {
"active-work drain timeout reached; proceeding with restart: 2 active background task run(s); 1 active embedded run(s)",
);
expectRestartCloseCall(closeFirst, DEFAULT_RESTART_DEFERRAL_TIMEOUT_MS);
expect(markGatewaySigusr1RestartHandled).toHaveBeenCalledTimes(1);
expect(abortActiveCronTaskRuns).toHaveBeenCalledWith("Gateway restarting.");
expect(waitForActiveCronTaskRuns).toHaveBeenCalledWith(1_000);
expect(waitForActiveCronJobs).toHaveBeenCalledWith(1_000);
expect(advanceCronActiveJobGeneration).toHaveBeenCalledTimes(1);
expect(retireActiveCronTaskRunTracking).toHaveBeenCalledTimes(1);
expect(resetCronActiveJobs).toHaveBeenCalledTimes(1);
expect(clearRuntimeConfigSnapshot).toHaveBeenCalledTimes(1);
expect(resetGatewaySuspendCoordinatorForLifecycleRestart).toHaveBeenCalledTimes(1);
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(1);
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(1);
expect(reloadTaskRuntimeStateFromStore.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
start.mock.invocationCallOrder[1] ?? Infinity,
);
expect(advanceCronActiveJobGeneration.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
abortActiveCronTaskRuns.mock.invocationCallOrder[0] ?? Infinity,
);
expect(waitForActiveCronJobs.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity,
);
expect(retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
resetCronActiveJobs.mock.invocationCallOrder[0] ?? Infinity,
);
lifecycleSlot.set("second", 2);
sigusr1();
await startedThird;
expect(secondRestartDrainSignal?.aborted).toBe(true);
const thirdAgentEventGeneration = agentEventsActual.getAgentEventLifecycleGeneration();
expect(thirdAgentEventGeneration).not.toBe(secondAgentEventGeneration);
expect(gatewayWorkAdmissionActual.isGatewayWorkAdmissionClosed()).toBe(false);
@@ -1344,6 +1330,18 @@ describe("runGatewayLoop", () => {
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(2);
expect(reloadTaskRuntimeStateFromStore).toHaveBeenCalledTimes(2);
expect(acquireGatewayLock).toHaveBeenCalledTimes(3);
expect(reloadTaskRuntimeStateFromStore.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
start.mock.invocationCallOrder[1] ?? Infinity,
);
expect(advanceCronActiveJobGeneration.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
abortActiveCronTaskRuns.mock.invocationCallOrder[0] ?? Infinity,
);
expect(waitForActiveCronJobs.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity,
);
expect(retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
resetCronActiveJobs.mock.invocationCallOrder[0] ?? Infinity,
);
sigterm();
await expect(exited).resolves.toBe(0);
+3 -3
View File
@@ -1006,6 +1006,9 @@ export async function runGatewayLoop(params: {
// suspension admission callback and discards the coordinator entry.
resetGatewaySuspendCoordinatorForLifecycleRestart();
resetAllLanes();
// resetAllLanes installs the next admission generation. Keep the local
// mirror aligned so a restart queued during cleanup closes that generation.
restartDrainingMarked = false;
clearRuntimeConfigSnapshot();
resetGatewayRestartStateForInProcessRestart();
// Rent: a failed startup has no server close handle, and restart hooks can
@@ -1023,9 +1026,6 @@ export async function runGatewayLoop(params: {
// SIGTERM/SIGINT still exit after a graceful shutdown.
let isFirstIteration = true;
for (;;) {
// The restart hook reopens admission before reloading durable state. Clear
// its local mirror first so a failed reload cannot skip the next drain.
restartDrainingMarked = false;
let startupFailedBeforeServerHandle = false;
const isRestartIteration = !isFirstIteration;
isFirstIteration = false;