mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-25 20:05:46 -06:00
fix(cron): scope settlement grace to the cancelled run (#118393)
Co-authored-by: Peter Steinberger <steipete@macos.shared>
This commit is contained in:
committed by
GitHub
parent
b438b9605f
commit
57732647d1
@@ -1,8 +1,10 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
abortActiveCronTaskRuns,
|
||||
cancelActiveCronTaskRun,
|
||||
getSuspensionVisibleCronTaskRunCount,
|
||||
registerActiveCronTaskRun,
|
||||
retireActiveCronTaskRunTracking,
|
||||
startActiveCronTaskRunSettlementGrace,
|
||||
trackActiveCronTaskRunSettlement,
|
||||
waitForActiveCronTaskRuns,
|
||||
} from "./active-run-cancellation.js";
|
||||
@@ -38,16 +40,77 @@ describe("cron task cancellation tracking", () => {
|
||||
await vi.waitFor(() => expect(getSuspensionVisibleCronTaskRunCount()).toBe(0));
|
||||
});
|
||||
|
||||
it("drops never-settling cron promises after a bounded grace period", async () => {
|
||||
it.each([
|
||||
{ timing: "after tracking", abortBeforeTracking: false },
|
||||
{ timing: "before tracking", abortBeforeTracking: true },
|
||||
])(
|
||||
"drops never-settling cron promises after an abort $timing",
|
||||
async ({ abortBeforeTracking }) => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
resetActiveCronTaskRunsForTests();
|
||||
const controller = new AbortController();
|
||||
if (abortBeforeTracking) {
|
||||
controller.abort();
|
||||
}
|
||||
trackActiveCronTaskRunSettlement(new Promise<never>(() => {}), controller.signal);
|
||||
|
||||
await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({
|
||||
drained: false,
|
||||
active: 1,
|
||||
});
|
||||
|
||||
if (!abortBeforeTracking) {
|
||||
await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1);
|
||||
|
||||
await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({
|
||||
drained: false,
|
||||
active: 1,
|
||||
});
|
||||
|
||||
controller.abort();
|
||||
}
|
||||
await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1);
|
||||
|
||||
await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({
|
||||
drained: true,
|
||||
active: 0,
|
||||
});
|
||||
expect(getSuspensionVisibleCronTaskRunCount()).toBe(1);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
resetActiveCronTaskRunsForTests();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("does not retire an unrelated active run when one registered run is cancelled", async () => {
|
||||
vi.useFakeTimers();
|
||||
const cancelledController = new AbortController();
|
||||
const unaffectedController = new AbortController();
|
||||
let settleCancelled = () => {};
|
||||
let settleUnaffected = () => {};
|
||||
const cancelledRun = new Promise<void>((resolve) => {
|
||||
settleCancelled = resolve;
|
||||
});
|
||||
const unaffectedRun = new Promise<void>((resolve) => {
|
||||
settleUnaffected = resolve;
|
||||
});
|
||||
|
||||
try {
|
||||
resetActiveCronTaskRunsForTests();
|
||||
trackActiveCronTaskRunSettlement(new Promise<never>(() => {}));
|
||||
|
||||
await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({
|
||||
drained: false,
|
||||
active: 1,
|
||||
trackActiveCronTaskRunSettlement(cancelledRun, cancelledController.signal);
|
||||
trackActiveCronTaskRunSettlement(unaffectedRun, unaffectedController.signal);
|
||||
const release = registerActiveCronTaskRun({
|
||||
runId: "cancelled-run",
|
||||
controller: cancelledController,
|
||||
});
|
||||
const existingTimerCount = vi.getTimerCount();
|
||||
|
||||
expect(cancelActiveCronTaskRun({ runId: "cancelled-run" })).toBe(true);
|
||||
expect(vi.getTimerCount()).toBe(existingTimerCount + 1);
|
||||
expect(cancelActiveCronTaskRun({ runId: "cancelled-run" })).toBe(false);
|
||||
release?.();
|
||||
|
||||
await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1);
|
||||
|
||||
@@ -55,15 +118,51 @@ describe("cron task cancellation tracking", () => {
|
||||
drained: false,
|
||||
active: 1,
|
||||
});
|
||||
expect(getSuspensionVisibleCronTaskRunCount()).toBe(2);
|
||||
} finally {
|
||||
settleCancelled();
|
||||
settleUnaffected();
|
||||
await Promise.all([cancelledRun, unaffectedRun]);
|
||||
vi.useRealTimers();
|
||||
resetActiveCronTaskRunsForTests();
|
||||
}
|
||||
});
|
||||
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
it.each([
|
||||
{ scenario: "only handleless runs", cancellableRuns: 0, handlelessRuns: 1 },
|
||||
{ scenario: "mixed active runs", cancellableRuns: 1, handlelessRuns: 1 },
|
||||
{ scenario: "no active runs", cancellableRuns: 0, handlelessRuns: 0 },
|
||||
])("retires $scenario during a bulk shutdown", async ({ cancellableRuns, handlelessRuns }) => {
|
||||
vi.useFakeTimers();
|
||||
const controller = new AbortController();
|
||||
try {
|
||||
resetActiveCronTaskRunsForTests();
|
||||
if (cancellableRuns > 0) {
|
||||
trackActiveCronTaskRunSettlement(new Promise<never>(() => {}), controller.signal);
|
||||
}
|
||||
if (handlelessRuns > 0) {
|
||||
trackActiveCronTaskRunSettlement(new Promise<never>(() => {}));
|
||||
}
|
||||
const release =
|
||||
cancellableRuns > 0
|
||||
? registerActiveCronTaskRun({ runId: "shutdown-run", controller })
|
||||
: undefined;
|
||||
|
||||
expect(abortActiveCronTaskRuns()).toBe(cancellableRuns);
|
||||
const retirementTimerCount = vi.getTimerCount();
|
||||
expect(abortActiveCronTaskRuns()).toBe(0);
|
||||
expect(vi.getTimerCount()).toBe(retirementTimerCount);
|
||||
if (cancellableRuns + handlelessRuns === 0) {
|
||||
expect(retirementTimerCount).toBe(0);
|
||||
}
|
||||
release?.();
|
||||
await vi.advanceTimersByTimeAsync(CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS + 1);
|
||||
|
||||
await expect(waitForActiveCronTaskRuns(0)).resolves.toEqual({
|
||||
drained: true,
|
||||
active: 0,
|
||||
});
|
||||
expect(getSuspensionVisibleCronTaskRunCount()).toBe(1);
|
||||
expect(getSuspensionVisibleCronTaskRunCount()).toBe(cancellableRuns + handlelessRuns);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
resetActiveCronTaskRunsForTests();
|
||||
@@ -72,12 +171,13 @@ describe("cron task cancellation tracking", () => {
|
||||
|
||||
it("keeps suspension blocked until a timed-out core actually settles", async () => {
|
||||
resetActiveCronTaskRunsForTests();
|
||||
const controller = new AbortController();
|
||||
let settle = () => {};
|
||||
const core = new Promise<void>((resolve) => {
|
||||
settle = resolve;
|
||||
});
|
||||
trackActiveCronTaskRunSettlement(core);
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
trackActiveCronTaskRunSettlement(core, controller.signal);
|
||||
controller.abort();
|
||||
|
||||
expect(getSuspensionVisibleCronTaskRunCount()).toBe(1);
|
||||
settle();
|
||||
|
||||
@@ -1,33 +1,25 @@
|
||||
// Process-local cancellation handles for live cron task runs.
|
||||
|
||||
type CronTaskCancelHandle = {
|
||||
controller: AbortController;
|
||||
onCancel?: (reason: string) => void;
|
||||
};
|
||||
|
||||
type SettlingCronTaskRun = {
|
||||
retirementTimer?: NodeJS.Timeout;
|
||||
};
|
||||
|
||||
const activeCronTaskRunsByRunId = new Map<string, CronTaskCancelHandle>();
|
||||
const settlingCronTaskRuns = new Map<Promise<unknown>, SettlingCronTaskRun>();
|
||||
const activeCronTaskRunsByRunId = new Map<
|
||||
string,
|
||||
{ controller: AbortController; onCancel?: (reason: string) => void }
|
||||
>();
|
||||
const settlingCronTaskRuns = new Map<Promise<unknown>, { retirementTimer?: NodeJS.Timeout }>();
|
||||
// Restart drain may retire an abort-ignoring core after a bounded grace, but a
|
||||
// host snapshot must keep refusing readiness until that core actually settles.
|
||||
const suspensionVisibleCronTaskRuns = new Set<Promise<unknown>>();
|
||||
const DEFAULT_CRON_TASK_RUN_DRAIN_POLL_MS = 25;
|
||||
const CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS = 60_000;
|
||||
|
||||
export function startActiveCronTaskRunSettlementGrace(): void {
|
||||
for (const [promise, entry] of settlingCronTaskRuns) {
|
||||
if (entry.retirementTimer) {
|
||||
continue;
|
||||
}
|
||||
const retirementTimer = setTimeout(() => {
|
||||
settlingCronTaskRuns.delete(promise);
|
||||
}, CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS);
|
||||
retirementTimer.unref?.();
|
||||
entry.retirementTimer = retirementTimer;
|
||||
function startActiveCronTaskRunSettlementGrace(promise: Promise<unknown>): void {
|
||||
const entry = settlingCronTaskRuns.get(promise);
|
||||
if (!entry || entry.retirementTimer) {
|
||||
return;
|
||||
}
|
||||
entry.retirementTimer = setTimeout(() => {
|
||||
settlingCronTaskRuns.delete(promise);
|
||||
}, CRON_TASK_RUN_SETTLEMENT_TRACKING_MAX_MS);
|
||||
entry.retirementTimer.unref?.();
|
||||
}
|
||||
|
||||
export function registerActiveCronTaskRun(params: {
|
||||
@@ -60,18 +52,29 @@ export function abortActiveCronTaskRuns(reason = "Gateway restarting."): number
|
||||
handle.onCancel?.(reason);
|
||||
aborted += 1;
|
||||
}
|
||||
if (aborted > 0) {
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
// Shutdown also retires main-session runs without cancellation handles.
|
||||
for (const promise of settlingCronTaskRuns.keys()) {
|
||||
startActiveCronTaskRunSettlementGrace(promise);
|
||||
}
|
||||
return aborted;
|
||||
}
|
||||
|
||||
export function trackActiveCronTaskRunSettlement(promise: Promise<unknown>): void {
|
||||
export function trackActiveCronTaskRunSettlement(
|
||||
promise: Promise<unknown>,
|
||||
abortSignal?: AbortSignal,
|
||||
): void {
|
||||
settlingCronTaskRuns.set(promise, {});
|
||||
suspensionVisibleCronTaskRuns.add(promise);
|
||||
// Cancellation belongs to this core only; sibling jobs must remain drain-visible.
|
||||
const startSettlementGrace = () => startActiveCronTaskRunSettlementGrace(promise);
|
||||
abortSignal?.addEventListener("abort", startSettlementGrace, { once: true });
|
||||
if (abortSignal?.aborted) {
|
||||
startSettlementGrace();
|
||||
}
|
||||
void promise
|
||||
.catch(() => undefined)
|
||||
.finally(() => {
|
||||
abortSignal?.removeEventListener("abort", startSettlementGrace);
|
||||
const entry = settlingCronTaskRuns.get(promise);
|
||||
if (entry?.retirementTimer) {
|
||||
clearTimeout(entry.retirementTimer);
|
||||
@@ -131,7 +134,6 @@ export function cancelActiveCronTaskRun(params: {
|
||||
const reason = params.reason?.trim() || "Cancelled by operator.";
|
||||
handle.controller.abort(reason);
|
||||
handle.onCancel?.(reason);
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -6,7 +6,6 @@ import { createCronRunDiagnosticsFromError } from "../run-diagnostics.js";
|
||||
import type { CronAgentExecutionStarted, CronJob } from "../types.js";
|
||||
import {
|
||||
registerActiveCronTaskRun,
|
||||
startActiveCronTaskRunSettlementGrace,
|
||||
trackActiveCronTaskRunSettlement,
|
||||
} from "./active-run-cancellation.js";
|
||||
import {
|
||||
@@ -250,7 +249,7 @@ export async function executeJobCoreWithTimeout(
|
||||
progress.completedCoreResult = result;
|
||||
return await deliverPrimaryWebhook(state, job, result, runAbortController.signal, progress);
|
||||
});
|
||||
trackActiveCronTaskRunSettlement(runPromise);
|
||||
trackActiveCronTaskRunSettlement(runPromise, runAbortController.signal);
|
||||
void runPromise.catch((err: unknown) => {
|
||||
if (runAbortController.signal.aborted) {
|
||||
state.deps.log.warn(
|
||||
@@ -263,7 +262,6 @@ export async function executeJobCoreWithTimeout(
|
||||
if (first !== operatorCancellationMarker) {
|
||||
return first;
|
||||
}
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
const settled = resolveInterruptedRunProgress({
|
||||
progress,
|
||||
job,
|
||||
@@ -334,7 +332,7 @@ export async function executeJobCoreWithTimeout(
|
||||
watchdog.deadlineAtMs(),
|
||||
);
|
||||
});
|
||||
trackActiveCronTaskRunSettlement(runPromise);
|
||||
trackActiveCronTaskRunSettlement(runPromise, runAbortController.signal);
|
||||
void runPromise.catch((err: unknown) => {
|
||||
if (runAbortController.signal.aborted) {
|
||||
state.deps.log.warn(
|
||||
@@ -346,7 +344,6 @@ export async function executeJobCoreWithTimeout(
|
||||
try {
|
||||
const first = await Promise.race([runPromise, timeoutPromise, operatorCancellationPromise]);
|
||||
if (first === operatorCancellationMarker) {
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
const settled = resolveInterruptedRunProgress({
|
||||
progress,
|
||||
job,
|
||||
@@ -360,7 +357,6 @@ export async function executeJobCoreWithTimeout(
|
||||
if (first !== timeoutMarker) {
|
||||
return first;
|
||||
}
|
||||
startActiveCronTaskRunSettlementGrace();
|
||||
const activeExecution = watchdog.activeExecution();
|
||||
const settled = resolveInterruptedRunProgress({
|
||||
progress,
|
||||
|
||||
Reference in New Issue
Block a user