Skip to content

Commit 13d06b5

Browse files
ghitafilalivincentkoc
authored andcommitted
fix: restart gateway after isolated cron setup timeout
1 parent 2f34d06 commit 13d06b5

28 files changed

Lines changed: 2367 additions & 114 deletions

src/agents/embedded-agent-runner/run.ts

Lines changed: 45 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ import { formatErrorMessage } from "../../infra/errors.js";
3030
import { buildAgentHookContextChannelFields } from "../../plugins/hook-agent-context.js";
3131
import { getGlobalHookRunner } from "../../plugins/hook-runner-global.js";
3232
import { resolveProviderAuthProfileId } from "../../plugins/provider-runtime.js";
33-
import { enqueueCommandInLane } from "../../process/command-queue.js";
33+
import { enqueueCommandInLane, getCommandLaneSnapshot } from "../../process/command-queue.js";
3434
import type { CommandQueueEnqueueOptions } from "../../process/command-queue.types.js";
3535
import { createAgentHarnessTaskRuntimeScope } from "../../tasks/agent-harness-task-runtime-scope.js";
3636
import { resolveUserPath } from "../../utils.js";
@@ -586,12 +586,38 @@ async function runEmbeddedAgentInternal(
586586
},
587587
laneTaskTimeoutMs,
588588
);
589+
const withRunLaneWait = (opts?: CommandQueueEnqueueOptions) => {
590+
if (!opts?.onWait && !params.onLaneWait) {
591+
return opts;
592+
}
593+
return {
594+
...opts,
595+
onWait: (waitMs, queuedAhead) => {
596+
opts?.onWait?.(waitMs, queuedAhead);
597+
params.onLaneWait?.({ waitMs, queuedAhead, waiting: true });
598+
},
599+
} satisfies CommandQueueEnqueueOptions;
600+
};
601+
const noteLaneWaitIfBusy = (lane: string) => {
602+
if (!params.onLaneWait) {
603+
return;
604+
}
605+
const snapshot = getCommandLaneSnapshot(lane);
606+
if (snapshot.queuedCount > 0 || snapshot.activeCount >= snapshot.maxConcurrent) {
607+
params.onLaneWait({
608+
waitMs: 0,
609+
queuedAhead: snapshot.queuedCount + snapshot.activeCount,
610+
waiting: true,
611+
});
612+
}
613+
};
589614
const enqueueGlobal = <T>(task: () => Promise<T>, opts?: CommandQueueEnqueueOptions) => {
590615
const globalOpts: CommandQueueEnqueueOptions = {
591616
...opts,
592617
priority: sessionQueuePriority,
593618
};
594619
const taskWithCurrentLifecycle = () => {
620+
params.onLaneWait?.({ waitMs: 0, queuedAhead: 0, waiting: false });
595621
throwIfAborted();
596622
const currentLifecycleGeneration = getAgentEventLifecycleGeneration();
597623
const existingContext = getAgentRunContext(params.runId);
@@ -618,15 +644,27 @@ async function runEmbeddedAgentInternal(
618644
});
619645
return withAgentRunLifecycleGeneration(lifecycleGeneration, task);
620646
};
621-
return params.enqueue
622-
? params.enqueue(taskWithCurrentLifecycle, withLaneTimeout(globalOpts))
623-
: enqueueCommandInLane(globalLane, taskWithCurrentLifecycle, withLaneTimeout(globalOpts));
647+
if (params.enqueue) {
648+
return params.enqueue(taskWithCurrentLifecycle, withLaneTimeout(withRunLaneWait(globalOpts)));
649+
}
650+
noteLaneWaitIfBusy(globalLane);
651+
return enqueueCommandInLane(
652+
globalLane,
653+
taskWithCurrentLifecycle,
654+
withLaneTimeout(withRunLaneWait(globalOpts)),
655+
);
624656
};
625657
const enqueueSession = <T>(task: () => Promise<T>, opts?: CommandQueueEnqueueOptions) => {
626658
const sessionOpts: CommandQueueEnqueueOptions = { ...opts, priority: sessionQueuePriority };
627-
return params.enqueue
628-
? params.enqueue(task, sessionOpts)
629-
: enqueueCommandInLane(sessionLane, task, sessionOpts);
659+
const taskWithLaneAdmission = () => {
660+
params.onLaneWait?.({ waitMs: 0, queuedAhead: 0, waiting: false });
661+
return task();
662+
};
663+
if (params.enqueue) {
664+
return params.enqueue(taskWithLaneAdmission, withRunLaneWait(sessionOpts));
665+
}
666+
noteLaneWaitIfBusy(sessionLane);
667+
return enqueueCommandInLane(sessionLane, taskWithLaneAdmission, withRunLaneWait(sessionOpts));
630668
};
631669
const channelHint = params.messageChannel ?? params.messageProvider;
632670
const resolvedToolResultFormat =

src/agents/embedded-agent-runner/run/params.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -187,6 +187,7 @@ export type RunEmbeddedAgentParams = {
187187
itemId?: string;
188188
firstModelCallStarted?: boolean;
189189
}) => void;
190+
onLaneWait?: (info: { waitMs: number; queuedAhead: number; waiting?: boolean }) => void;
190191
onRunProgress?: (info: {
191192
reason: string;
192193
provider?: string;

src/cli/gateway-cli/lifecycle.runtime.ts

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,16 @@ export { rotateAgentEventLifecycleGeneration } from "../../infra/agent-events.js
2929
export { markUpdateRestartSentinelFailure } from "../../infra/restart-sentinel.js";
3030
export { detectRespawnSupervisor } from "../../infra/supervisor-markers.js";
3131
export { writeDiagnosticStabilityBundleForFailureSync } from "../../logging/diagnostic-stability-bundle.js";
32+
export {
33+
advanceCronActiveJobGeneration,
34+
resetCronActiveJobs,
35+
waitForActiveCronJobs,
36+
} from "../../cron/active-jobs.js";
37+
export {
38+
abortActiveCronTaskRuns,
39+
retireActiveCronTaskRunTracking,
40+
waitForActiveCronTaskRuns,
41+
} from "../../tasks/cron-task-cancel.js";
3242
export {
3343
getActiveTaskCount,
3444
markGatewayDraining,

src/cli/gateway-cli/run-loop.test.ts

Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,18 @@ const getInspectableActiveTaskRestartBlockers = vi.fn(
5454
const markGatewayDraining = vi.fn();
5555
const waitForActiveTasks = vi.fn(async (_timeoutMs?: number) => ({ drained: true }));
5656
const resetAllLanes = vi.fn();
57+
const advanceCronActiveJobGeneration = vi.fn();
58+
const resetCronActiveJobs = vi.fn();
59+
const abortActiveCronTaskRuns = vi.fn((_reason?: string) => 0);
60+
const retireActiveCronTaskRunTracking = vi.fn();
61+
const waitForActiveCronTaskRuns = vi.fn(async (_timeoutMs?: number) => ({
62+
drained: true,
63+
active: 0,
64+
}));
65+
const waitForActiveCronJobs = vi.fn(async (_timeoutMs?: number) => ({
66+
drained: true,
67+
active: 0,
68+
}));
5769
const reloadTaskRegistryFromStore = vi.fn();
5870
const rotateAgentEventLifecycleGeneration = vi.fn();
5971
const clearRuntimeConfigSnapshot = vi.fn();
@@ -151,6 +163,18 @@ vi.mock("../../process/command-queue.js", () => ({
151163
resetAllLanes: () => resetAllLanes(),
152164
}));
153165

166+
vi.mock("../../cron/active-jobs.js", () => ({
167+
advanceCronActiveJobGeneration: () => advanceCronActiveJobGeneration(),
168+
resetCronActiveJobs: () => resetCronActiveJobs(),
169+
waitForActiveCronJobs: (timeoutMs: number) => waitForActiveCronJobs(timeoutMs),
170+
}));
171+
172+
vi.mock("../../tasks/cron-task-cancel.js", () => ({
173+
abortActiveCronTaskRuns: (reason?: string) => abortActiveCronTaskRuns(reason),
174+
retireActiveCronTaskRunTracking: () => retireActiveCronTaskRunTracking(),
175+
waitForActiveCronTaskRuns: (timeoutMs: number) => waitForActiveCronTaskRuns(timeoutMs),
176+
}));
177+
154178
vi.mock("../../tasks/runtime-internal.js", () => ({
155179
reloadTaskRegistryFromStore: () => reloadTaskRegistryFromStore(),
156180
}));
@@ -876,14 +900,32 @@ describe("runGatewayLoop", () => {
876900
expect(gatewayLog.warn).toHaveBeenCalledWith(DRAIN_TIMEOUT_LOG);
877901
expectRestartCloseCall(closeFirst, 1_234);
878902
expect(markGatewaySigusr1RestartHandled).toHaveBeenCalledTimes(1);
903+
expect(abortActiveCronTaskRuns).toHaveBeenCalledWith("Gateway restarting.");
904+
expect(waitForActiveCronTaskRuns).toHaveBeenCalledWith(1_000);
905+
expect(waitForActiveCronJobs).toHaveBeenCalledWith(1_000);
879906
expect(resetAllLanes).toHaveBeenCalledTimes(1);
907+
expect(advanceCronActiveJobGeneration).toHaveBeenCalledTimes(1);
908+
expect(retireActiveCronTaskRunTracking).toHaveBeenCalledTimes(1);
909+
expect(resetCronActiveJobs).toHaveBeenCalledTimes(1);
880910
expect(clearRuntimeConfigSnapshot).toHaveBeenCalledTimes(1);
881911
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(1);
882912
expect(rotateAgentEventLifecycleGeneration).toHaveBeenCalledTimes(1);
883913
expect(reloadTaskRegistryFromStore).toHaveBeenCalledTimes(1);
884914
expect(
885915
rotateAgentEventLifecycleGeneration.mock.invocationCallOrder[0] ?? Infinity,
886916
).toBeLessThan(resetAllLanes.mock.invocationCallOrder[0] ?? Infinity);
917+
expect(advanceCronActiveJobGeneration.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
918+
abortActiveCronTaskRuns.mock.invocationCallOrder[0] ?? Infinity,
919+
);
920+
expect(waitForActiveCronJobs.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
921+
retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity,
922+
);
923+
expect(retireActiveCronTaskRunTracking.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
924+
resetCronActiveJobs.mock.invocationCallOrder[0] ?? Infinity,
925+
);
926+
expect(resetCronActiveJobs.mock.invocationCallOrder[0] ?? Infinity).toBeLessThan(
927+
resetAllLanes.mock.invocationCallOrder[0] ?? Infinity,
928+
);
887929

888930
sigusr1();
889931

@@ -894,7 +936,13 @@ describe("runGatewayLoop", () => {
894936
expectRestartCloseCall(closeSecond, 1_234);
895937
expect(markGatewaySigusr1RestartHandled).toHaveBeenCalledTimes(2);
896938
expect(markGatewayDraining).toHaveBeenCalledTimes(2);
939+
expect(abortActiveCronTaskRuns).toHaveBeenCalledTimes(2);
940+
expect(waitForActiveCronTaskRuns).toHaveBeenCalledTimes(2);
941+
expect(waitForActiveCronJobs).toHaveBeenCalledTimes(2);
897942
expect(resetAllLanes).toHaveBeenCalledTimes(2);
943+
expect(advanceCronActiveJobGeneration).toHaveBeenCalledTimes(2);
944+
expect(retireActiveCronTaskRunTracking).toHaveBeenCalledTimes(2);
945+
expect(resetCronActiveJobs).toHaveBeenCalledTimes(2);
898946
expect(clearRuntimeConfigSnapshot).toHaveBeenCalledTimes(2);
899947
expect(resetGatewayRestartStateForInProcessRestart).toHaveBeenCalledTimes(2);
900948
expect(rotateAgentEventLifecycleGeneration).toHaveBeenCalledTimes(2);
@@ -910,6 +958,42 @@ describe("runGatewayLoop", () => {
910958
});
911959
});
912960

961+
it("advances stale cron active markers after bounded restart cron-run drain", async () => {
962+
vi.clearAllMocks();
963+
waitForActiveCronJobs.mockResolvedValueOnce({ drained: false, active: 1 });
964+
peekGatewaySigusr1RestartReason.mockReturnValue(undefined);
965+
respawnGatewayProcessForUpdate.mockReturnValue({
966+
mode: "disabled",
967+
detail: "OPENCLAW_NO_RESPAWN",
968+
});
969+
970+
await withIsolatedSignals(async ({ captureSignal }) => {
971+
const { start, exited } = await createSignaledLoopHarness();
972+
const sigusr1 = captureSignal("SIGUSR1");
973+
const sigint = captureSignal("SIGINT");
974+
975+
sigusr1();
976+
await waitForLoopCondition(
977+
() => start.mock.calls.length >= 2,
978+
"expected SIGUSR1 to trigger restart",
979+
);
980+
981+
expect(abortActiveCronTaskRuns).toHaveBeenCalledWith("Gateway restarting.");
982+
expect(waitForActiveCronTaskRuns).toHaveBeenCalledWith(1_000);
983+
expect(waitForActiveCronJobs).toHaveBeenCalledWith(1_000);
984+
expect(resetAllLanes).toHaveBeenCalledTimes(1);
985+
expect(advanceCronActiveJobGeneration).toHaveBeenCalledTimes(1);
986+
expect(retireActiveCronTaskRunTracking).toHaveBeenCalledTimes(1);
987+
expect(resetCronActiveJobs).toHaveBeenCalledTimes(1);
988+
expect(gatewayLog.warn).toHaveBeenCalledWith(
989+
"cron run drain timed out during restart lifecycle reset after retiring old cron admission; 0 task handle(s) and 1 active marker(s) remain after aborting old cron runs",
990+
);
991+
992+
sigint();
993+
await expect(exited).resolves.toBe(0);
994+
});
995+
});
996+
913997
it("queues SIGUSR1 received before the run-loop installs its restart waiter", async () => {
914998
vi.clearAllMocks();
915999
peekGatewaySigusr1RestartReason.mockReturnValue(undefined);

src/cli/gateway-cli/run-loop.ts

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -800,13 +800,30 @@ export async function runGatewayLoop(params: {
800800
// deferral timers and reloads the task registry from durable state so
801801
// cancelled/completed work is not kept alive by old in-memory maps.
802802
const {
803+
abortActiveCronTaskRuns,
804+
advanceCronActiveJobGeneration,
803805
reloadTaskRegistryFromStore,
806+
retireActiveCronTaskRunTracking,
807+
resetCronActiveJobs,
804808
resetAllLanes,
805809
resetGatewayRestartStateForInProcessRestart,
806810
rotateAgentEventLifecycleGeneration,
811+
waitForActiveCronJobs,
812+
waitForActiveCronTaskRuns,
807813
} = await loadGatewayLifecycleRuntimeModule();
808814
// Rotate ownership before reset pumps preserved queue entries.
809815
rotateAgentEventLifecycleGeneration();
816+
advanceCronActiveJobGeneration();
817+
abortActiveCronTaskRuns("Gateway restarting.");
818+
const cronTaskDrain = await waitForActiveCronTaskRuns(1_000);
819+
const cronDrain = await waitForActiveCronJobs(1_000);
820+
if (!cronTaskDrain.drained || !cronDrain.drained) {
821+
gatewayLog.warn(
822+
`cron run drain timed out during restart lifecycle reset after retiring old cron admission; ${cronTaskDrain.active} task handle(s) and ${cronDrain.active} active marker(s) remain after aborting old cron runs`,
823+
);
824+
}
825+
retireActiveCronTaskRunTracking();
826+
resetCronActiveJobs();
810827
resetAllLanes();
811828
clearRuntimeConfigSnapshot();
812829
resetGatewayRestartStateForInProcessRestart();

0 commit comments

Comments
 (0)