refactor(cron): schedule exit watcher retries (#160381)

## What Problem This Solves

Cron exit watcher retry deadlines now use the required Gateway scheduler. A scheduled retry joins startup and failed-start state publication, then the existing watcher owns the long-running child through exit. Cancellation and replacement still preserve exact slot ownership and drain watcher lifecycles.

The cutover deletes native timeout handles, `unref` calls, and the unused shell-injection option. No operator configuration or persisted schema changes.

Production **+32/-26/net +6**; tests **+101/-38/net +63**. This correctness-driven slice uses the closeout's real deletion offsets.

## Evidence

Final head `66d57e0dfa26e242a9c80f415c7f29d65a5721cd`, based on merged main `14f65626bd` after cron stream merge `2aaf6d41a7`. Only the exit commit is included. The rebase is patch-identical (`git range-diff` reports `=`), and all three changed files are byte-identical to the tested head below. Existing review and Testbox proof are reused; `git diff --check` passes and the worktree is clean.

- Independent Codex autoreview: clean through P2. The committed patch exactly matches the reviewed patch (SHA256 `3ca06c7e52888c92689d9dbc5cebcbb1a18199c1bf3290351d81a14693e81755`).
- Exact-head Blacksmith Testbox `tbx_01m3kjwyph0svfjebr22jy7w8q`, run [36399017112](https://github.com/openclaw/openclaw/actions/runs/36399017112), head `80fc468ac824da48ef94423d7a79da316d4ed976`: `pnpm test src/gateway/cron-exit-watchers.test.ts --maxWorkers=1` passed (**28 tests**, **3.75s wall**).
- Scoped `node scripts/check-changed.mjs --base 42975a258c -- src/gateway/cron-exit-watchers.test.ts src/gateway/cron-exit-watchers.ts src/gateway/server-cron.ts` passed on the same exact head: **1463.25s wall**, including production and all core test type graphs, format, lint, line ratchets, dead-export scan, and architecture guards.

Release-note context: Gateway-owned cron exit retries now share scheduler catch-up and shutdown ownership while child-process settlement remains with its watcher.


### Final inherited-CI classification

CI run [36411857535](https://github.com/openclaw/openclaw/actions/runs/36411857535) completed on reviewed head `66d57e0dfa26e242a9c80f415c7f29d65a5721cd`, testing merge tree `80ef6fe243fb9c0c148975007d087c81f1a1a52d` against main parent `fc168b2f6c`. All six failing test jobs have independent evidence in owners this change does not modify:

| Failure | This PR | Matching independent evidence |
| --- | --- | --- |
| Session rename remains Original session after ordinary Enter | [job 108894925328](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894925328) | Clean main at `251c610735`: [job 108908250303](https://github.com/openclaw/openclaw/actions/runs/36416233754/job/108908250303). The same parameterized test body fails at line 119 with the same expected committed title and received Original session. Matrix cells differ: this PR uses legacy/Escape; main uses modern/Enter. Both already passed the composition assertions, then execute the identical ordinary Enter, matching patch, list reconciliation, and final rendered-title assertion. |
| Manually unread UI state times out waiting for its unread dot | [job 108894925308](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894925308) | Independent MCP [job 108865636156](https://github.com/openclaw/openclaw/actions/runs/36402978963/job/108865636156): same `session-management.unread.e2e.test.ts` case, line 266, unread-dot locator, and 30000ms timeout. |
| Canonical skills refresh E2E exceeds its 90000ms timeout | [job 108894929439](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894929439) | Independent questions [job 108880022656](https://github.com/openclaw/openclaw/actions/runs/36407409443/job/108880022656): identical canonical-skills test and 90000ms timeout. |
| Subagent completion remains pending instead of suspended | [job 108894929700](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894929700) | Independent activity [job 108869549339](https://github.com/openclaw/openclaw/actions/runs/36404157978/job/108869549339): same ordinary-delivery-exhaustion test and pending/retryable versus suspended/permanent_failure assertion. |
| Requester serial continuation never reaches delivered cleanup | [job 108894933085](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894933085) | Config [job 108890604749](https://github.com/openclaw/openclaw/actions/runs/36410694894/job/108890604749) and companion [job 108888541543](https://github.com/openclaw/openclaw/actions/runs/36410052666/job/108888541543): identical CLI cell with next child accepted=false and requester final=false; same failed/intentional_non_delivery state, `worker task timed out`, and helper line 101. |
| Five OpenResponses HTTP fixture cases receive 500 instead of 200 | [job 108894933187](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108894933187) | Clean actual CI main parent `fc168b2f6c`, complete-file reproduction on Blacksmith Testbox `tbx_01m3kswm6aw0bxmdxv2e6qpfv6` ([lease run 36411554208](https://github.com/openclaw/openclaw/actions/runs/36411554208)): identical five failures, 29 passes, and missing Gateway context error. |

The remaining red is the aggregate [CI gate](https://github.com/openclaw/openclaw/actions/runs/36411857535/job/108903858681), derived from those six jobs.

For the HTTP baseline, a separate clean detached worktree passed exact-SHA, clean-tree, and frozen-install assertions. The complete `src/gateway/server-http.openai-compat.test.ts` ran once through `node --import tsx scripts/ci-run-node-test-shard.mts` with the CI `test/vitest/vitest.gateway-server.config.ts` config, `bun-compatible` runtime policy, empty prebuilt-dist setting, and two-worker group setting. Wall time: **73.419s**. Responses image-limit hot reload and the completed image/file requests with streaming enabled and disabled all fail at the same `OpenResponses bearer scope requires a current Gateway context` guard in `openresponses-http.ts:127`. The fixture omits `resolveGatewayContext`; the guard introduced on main by #159227 requires its `configRevisionProjector`.

This PR changes only cron exit watchers, their focused tests, and `server-cron.ts` wiring. No assertions were weakened, no baseline source was changed, and no hosted CI was rerun for this evidence.

The maintainer’s standing instruction authorizes pinned admin squash over these proven inherited failures. The protected native GraphQL route preserves the expected head, squash method, and reviewed message.
This commit is contained in:
Peter Steinberger 2026-09-28 05:26:50 -07:00 • committed by GitHub
parent 8aec7806a9
commit ec1b2b4602
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 133 additions and 64 deletions

View file

@ -1,10 +1,14 @@
import { AsyncLocalStorage } from "node:async_hooks";
import { setImmediate, setTimeout as delay } from "node:timers/promises";
import { setImmediate } from "node:timers/promises";
import { expectDefined } from "@openclaw/normalization-core";
import { afterEach, describe, expect, it, vi } from "vitest";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import type { CronJob } from "../cron/types.js";
import { AsyncWorkScope, trackAsyncWork } from "../shared/async-work-scope.js";
import {
createGatewaySchedulerClock,
createTestGatewayScheduler,
} from "../test-utils/gateway-scheduler-clock.js";
import { resolveExitWatchShell } from "./cron-exit-watch-shell.js";
import {
createCronExitWatchers,
@ -104,6 +108,8 @@ function onExitJob(id: string, command = "true", enabled = true): CronJob {
}
const noopLogger = { info: () => {}, warn: () => {} };
let clock: ReturnType<typeof createGatewaySchedulerClock>;
let scheduler: ReturnType<typeof createTestGatewayScheduler>;
type FixtureHandlers = Omit<CronExitWatcherHandlers, "fireOnExit"> & {
reserveExit: (job: CronJob) => Promise<void>;
@ -113,8 +119,7 @@ type FixtureHandlers = Omit<CronExitWatcherHandlers, "fireOnExit"> & {
};
function createWatcherFixture(
params: FixtureHandlers &
Pick<Parameters<typeof createCronExitWatchers>[0], "shell" | "retryBackoffMs">,
params: FixtureHandlers & NonNullable<Parameters<typeof createCronExitWatchers>[2]>,
) {
let jobs = new Map<string, CronJob>();
const handlers = (next: FixtureHandlers): CronExitWatcherHandlers => ({
@ -139,7 +144,7 @@ function createWatcherFixture(
await next.fireOnExit(current, exit);
},
});
const watchers = createCronExitWatchers({ ...params, ...handlers(params) });
const watchers = createCronExitWatchers(handlers(params), scheduler, params);
return {
...watchers,
reconcile: (current: CronJob[]) => {
@ -153,6 +158,14 @@ function createWatcherFixture(
const flush = () => setImmediate();
describe("createCronExitWatchers", () => {
beforeEach(() => {
clock = createGatewaySchedulerClock();
scheduler = createTestGatewayScheduler(clock.clock);
});
afterEach(async () => {
await scheduler.stop();
});
it("arms a watcher and fires on exit after the creating request closes", async () => {
const { supervisor, runs } = makeFakeSupervisor();
const creatorContext = new AsyncLocalStorage<string>();
@ -320,23 +333,26 @@ describe("createCronExitWatchers", () => {
const settled = createDeferred();
const fired = vi.fn();
let job = onExitJob("job-a");
const watchers = createCronExitWatchers({
getProcessSupervisor: () => supervisor as never,
logger: noopLogger,
fireOnExit: async (_job, _exit, controls) => {
try {
controls.commitGuard();
job = { ...job, enabled: false };
controls.onReserved();
reserved.resolve();
await release.promise;
controls.commitGuard();
fired();
} finally {
settled.resolve();
}
const watchers = createCronExitWatchers(
{
getProcessSupervisor: () => supervisor as never,
logger: noopLogger,
fireOnExit: async (_job, _exit, controls) => {
try {
controls.commitGuard();
job = { ...job, enabled: false };
controls.onReserved();
reserved.resolve();
await release.promise;
controls.commitGuard();
fired();
} finally {
settled.resolve();
}
},
},
});
scheduler,
);
try {
watchers.reconcile([job]);
await flush();
@ -367,7 +383,7 @@ describe("createCronExitWatchers", () => {
},
);
it("routes post-handoff retries through the replacement scheduler owner", async () => {
it("retains a pending retry across handler handoff and uses the replacement owner", async () => {
const { supervisor, runs } = makeFakeSupervisor();
const oldUpdateWatcherState = vi.fn(async () => {});
const newUpdateWatcherState = vi.fn(async () => {});
@ -377,11 +393,15 @@ describe("createCronExitWatchers", () => {
fireOnExit: vi.fn(async () => {}),
updateWatcherState: oldUpdateWatcherState,
logger: noopLogger,
retryBackoffMs: [0],
retryBackoffMs: [1_000],
});
watchers.reconcile([onExitJob("job-a")]);
await flush();
expectDefined(runs[0], "first watcher").deferred.reject(new Error("wait failed"));
await flush();
expect(oldUpdateWatcherState).toHaveBeenCalledOnce();
expect(scheduler.nextWakeAtMs).toBe(1_000);
await watchers.updateHandlers({
getProcessSupervisor: () => supervisor as never,
reserveExit: vi.fn(async () => {}),
@ -389,13 +409,59 @@ describe("createCronExitWatchers", () => {
updateWatcherState: newUpdateWatcherState,
logger: noopLogger,
});
expectDefined(runs[0], "runs[0] test invariant").deferred.reject(new Error("wait failed"));
await vi.waitFor(() => expect(newUpdateWatcherState).toHaveBeenCalledOnce());
expect(oldUpdateWatcherState).not.toHaveBeenCalled();
await delay(5);
await flush();
await clock.advanceBy(1_000);
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
expectDefined(runs[1], "replacement watcher").deferred.reject(new Error("replacement failed"));
await flush();
expect(newUpdateWatcherState).toHaveBeenCalledOnce();
expect(oldUpdateWatcherState).toHaveBeenCalledOnce();
await watchers.cancelAll();
await clock.advanceBy(1_000);
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
});
it("joins a scheduled retry's failed-start publication during scheduler shutdown", async () => {
const { supervisor } = makeFakeSupervisor();
supervisor.spawn.mockRejectedValue(new Error("spawn failed"));
const firstRetryScheduled = createDeferred();
const retryFailure = createDeferred();
const releaseRetry = createDeferred();
const watchers = createWatcherFixture({
getProcessSupervisor: () => supervisor as never,
reserveExit: vi.fn(async () => {}),
fireOnExit: vi.fn(async () => {}),
updateWatcherState: vi.fn(async (_job, patch) => {
if (patch.consecutiveErrors !== 1) {
retryFailure.resolve();
await releaseRetry.promise;
}
}),
logger: { ...noopLogger, warn: () => firstRetryScheduled.resolve() },
retryBackoffMs: [1_000],
});
watchers.reconcile([onExitJob("job-a")]);
await firstRetryScheduled.promise;
expect(scheduler.nextWakeAtMs).toBe(1_000);
const retry = clock.advanceBy(1_000);
await retryFailure.promise;
let stopped = false;
const stopping = scheduler.stop().then(() => {
stopped = true;
});
try {
for (let turn = 0; turn < 5; turn += 1) {
await Promise.resolve();
}
expect(stopped).toBe(false);
} finally {
releaseRetry.resolve();
await retry;
await stopping;
await watchers.cancelAll();
}
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
expect(scheduler.nextWakeAtMs).toBeNull();
expect(watchers.activeJobIds()).toEqual([]);
});
it("a fired job stays unarmed across a simulated restart (disabled in store → not re-run)", async () => {
@ -490,10 +556,8 @@ describe("createCronExitWatchers", () => {
expect.objectContaining({ consecutiveErrors: 1 }),
);
expect(w.activeJobIds()).toEqual(["job-a"]);
// The backoff timer re-arms without an external reconcile. The 0ms retry
// timer is armed only after async state persistence, so a fixed sleep races
// it on loaded CI workers — wait for the re-arm instead.
await vi.waitFor(() => expect(supervisor.spawn).toHaveBeenCalledTimes(2));
await clock.advanceBy(0);
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
expect(w.activeJobIds()).toEqual(["job-a"]);
});
@ -516,9 +580,8 @@ describe("createCronExitWatchers", () => {
expect.objectContaining({ id: "job-a" }),
expect.objectContaining({ consecutiveErrors: 1 }),
);
// Backoff re-arm succeeds on the second spawn. Same async-persist race as
// the wait-rejection case above: wait for the re-arm, not a fixed sleep.
await vi.waitFor(() => expect(supervisor.spawn).toHaveBeenCalledTimes(2));
await clock.advanceBy(0);
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
expect(w.activeJobIds()).toEqual(["job-a"]);
// Cancelling clears any pending retry timer; once the cancelled child
@ -528,8 +591,8 @@ describe("createCronExitWatchers", () => {
exitCode: null,
reason: "cancelled",
});
await delay(5);
await flush();
await w.cancelAll();
await clock.advanceBy(1_000);
expect(supervisor.spawn).toHaveBeenCalledTimes(2);
expect(w.activeJobIds()).toEqual([]);
});

View file

@ -1,4 +1,5 @@
import type { CronJob } from "../cron/types.js";
import type { GatewayScheduledJob, GatewayScheduler } from "../infra/gateway-scheduler.js";
import { markOpenClawExecEnv } from "../infra/openclaw-exec-env.js";
import type { ManagedRun, ProcessSupervisor } from "../process/supervisor/index.js";
import { runInDetachedAsyncContext } from "../shared/async-work-scope.js";
@ -68,12 +69,11 @@ function isWatchableExitJob(job: CronJob): job is OnExitCronJob {
}
export function createCronExitWatchers(
params: CronExitWatcherHandlers & {
shell?: { command: string; argsFor: (command: string) => string[] };
retryBackoffMs?: readonly number[];
},
initialHandlers: CronExitWatcherHandlers,
scheduler: GatewayScheduler,
options?: { retryBackoffMs?: readonly number[] },
): CronExitWatchers {
let handlers: CronExitWatcherHandlers = params;
let handlers = initialHandlers;
const ownerSettlements = new Set<Promise<void>>();
const settleOwnerCallback = async <T>(operation: Promise<T>): Promise<T> => {
const settlement = operation.then(
@ -87,10 +87,10 @@ export function createCronExitWatchers(
ownerSettlements.delete(settlement);
}
};
const shell = params.shell ?? resolveExitWatchShell();
const shell = resolveExitWatchShell();
const retryBackoffMs =
params.retryBackoffMs && params.retryBackoffMs.length > 0
? params.retryBackoffMs
options?.retryBackoffMs && options.retryBackoffMs.length > 0
? options.retryBackoffMs
: ON_EXIT_WATCH_RETRY_BACKOFF_MS;
// Reserving the slot before spawn lets cancel/replace retire an in-flight arm.
// Async continuations publish only while this exact slot remains current.
@ -106,7 +106,7 @@ export function createCronExitWatchers(
command: string;
cwd: string | undefined;
consecutiveFailures: number;
retryTimer: NodeJS.Timeout | undefined;
retryJob: GatewayScheduledJob | undefined;
};
const active = new Map<string, WatcherSlot>();
// A cancelled child can keep running until the supervisor observes exit.
@ -130,10 +130,8 @@ export function createCronExitWatchers(
if (!preserveReserved || !slot.fired) {
slot.admission?.abort();
}
if (slot.retryTimer) {
clearTimeout(slot.retryTimer);
slot.retryTimer = undefined;
}
slot.retryJob?.cancel();
slot.retryJob = undefined;
if (!slot.lifecycleSettled) {
settlingCancelledSlots.add(slot);
}
@ -150,7 +148,8 @@ export function createCronExitWatchers(
}
};
const arm = (job: OnExitCronJob, consecutiveFailures = 0) => {
const arm = (job: OnExitCronJob, consecutiveFailures = 0): Promise<void> => {
const armed = createDeferredCore();
const command = job.schedule.command;
const cwd = job.schedule.cwd;
const predecessors = Array.from(settlingCancelledSlots)
@ -170,7 +169,7 @@ export function createCronExitWatchers(
command,
cwd,
consecutiveFailures,
retryTimer: undefined,
retryJob: undefined,
};
active.set(job.id, slot);
const owns = () => active.get(job.id) === slot;
@ -208,15 +207,18 @@ export function createCronExitWatchers(
}
const delayMs =
retryBackoffMs[Math.min(slot.consecutiveFailures - 1, retryBackoffMs.length - 1)]!;
slot.retryTimer = setTimeout(() => {
slot.retryTimer = undefined;
if (!owns() || slot.cancelled) {
return;
}
active.delete(slot.job.id);
arm(slot.job, slot.consecutiveFailures);
}, delayMs);
slot.retryTimer.unref?.();
slot.retryJob = scheduler.schedule({
id: `${scopeKey(job.id)}:retry`,
delayMs,
run: () => {
slot.retryJob = undefined;
if (!owns() || slot.cancelled) {
return undefined;
}
active.delete(slot.job.id);
return arm(slot.job, slot.consecutiveFailures);
},
});
handlers.logger.warn(
{ err: String(error), jobId: slot.job.id, retryInMs: delayMs },
`cron-exit: watcher ${phase} failed; retry scheduled`,
@ -264,6 +266,8 @@ export function createCronExitWatchers(
{ jobId: job.id, runId: run.runId, command },
"cron-exit: watcher armed",
);
// The watcher now owns the child through exit; the scheduled arm only joins startup.
armed.resolve();
let exit: Awaited<ReturnType<ManagedRun["wait"]>>;
try {
exit = await run.wait();
@ -335,6 +339,7 @@ export function createCronExitWatchers(
}
}
})().finally(() => {
armed.resolve();
slot.lifecycleSettled = true;
settlingCancelledSlots.delete(slot);
if (slot.cancelled && active.get(job.id) === slot) {
@ -343,6 +348,7 @@ export function createCronExitWatchers(
slot.settlement.resolve(undefined);
}),
);
return armed.promise;
};
const reconcile = (jobs: CronJob[]) => {
@ -378,7 +384,7 @@ export function createCronExitWatchers(
}
cancel(jobId, slot.fired && slot.command === command && slot.cwd === cwd);
}
arm(job);
void arm(job);
}
};

View file

@ -1182,7 +1182,7 @@ export function buildGatewayCronService(params: {
}, "cron:watcher-state"),
logger: cronServiceLogger,
} satisfies CronExitWatcherHandlers;
exitWatchersRef.current = createCronExitWatchers(exitWatcherHandlers);
exitWatchersRef.current = createCronExitWatchers(exitWatcherHandlers, params.scheduler);
streamWatchersRef.current = createCronStreamWatchers({
scheduler: params.scheduler,
getProcessSupervisor,