fix(heartbeat): prevent replies disappearing during active heartbeats (#123458)

* fix(heartbeat): prioritize visible turns over active runs

* fix(heartbeat): preserve embedded run ownership

* fix(heartbeat): drain preempted run ownership

* test(heartbeat): cover retained preemption work

* fix(heartbeat): retain work across visible preemption

* fix(heartbeat): fence finalizing supersession
This commit is contained in:
Ayaan Zaidi
2026-08-14 17:51:05 +05:30
committed by GitHub
parent c9152fef4e
commit 7fcddedaaf
37 changed files with 1013 additions and 205 deletions
+1
View File
@@ -7,6 +7,7 @@ export { runEmbeddedAgent } from "./embedded-agent-runner/run.js";
export {
abortAndDrainEmbeddedAgentRun,
abortEmbeddedAgentRun,
preemptAndDrainEmbeddedHeartbeatRun,
isEmbeddedAgentRunAbortableForCompaction,
isEmbeddedAgentRunActive,
isEmbeddedAgentRunHandleActive,
@@ -38,6 +38,8 @@ export type EmbeddedAgentQueueHandle = {
) => Promise<boolean>;
/** Cancels this run's pending user-input request before an image is queued as a later turn. */
cancelPendingUserInput?: (resolvedBy: string) => Promise<boolean>;
/** Exact heartbeat owner retained after its reply-operation registration clears. */
readonly preemptByVisibleTurn?: () => boolean;
queueMessage: (
text: string,
options?: EmbeddedAgentQueueMessageOptions,
@@ -70,6 +72,7 @@ export type ActiveEmbeddedRunSnapshot = {
export type EmbeddedRunWaiter = {
resolve: (ended: boolean) => void;
handle?: EmbeddedAgentQueueHandle;
timer?: NodeJS.Timeout;
};
@@ -108,6 +108,32 @@ describe("prepareEmbeddedAttemptStream", () => {
mocks.runBeforeFinalizeHook.mockResolvedValue({ action: "continue" });
});
it("retains exact heartbeat preemption on the embedded queue handle", () => {
const operation = createReplyOperation({
sessionKey: "agent:main:main",
sessionId: "session-output-schema",
turnKind: "heartbeat",
resetTriggered: false,
});
try {
const prepared = prepareCatalogExecutor([], { replyOperation: operation });
expect(prepared.queueHandle.preemptByVisibleTurn?.()).toBe(true);
expect(operation.result).toEqual({
kind: "aborted",
code: "aborted_for_supersession",
});
expect(mocks.setActiveRun).toHaveBeenCalledWith(
"session-output-schema",
expect.objectContaining({ preemptByVisibleTurn: expect.any(Function) }),
"agent:main:main",
undefined,
);
} finally {
operation.complete();
}
});
it("uses the persisted assistant entry id and closes steering during revision settlement", async () => {
let resolveHook: ((value: { action: "revise"; reason: string }) => void) | undefined;
mocks.runBeforeFinalizeHook.mockImplementation(
@@ -427,6 +427,8 @@ export function prepareEmbeddedAttemptStream(input: {
activeQueueAdmissions--;
}
};
const heartbeatReplyOperation =
attempt.replyOperation?.turnKind === "heartbeat" ? attempt.replyOperation : undefined;
const queueHandle: AttemptStreamQueueHandle = {
kind: "embedded",
runId: attempt.runId,
@@ -437,6 +439,9 @@ export function prepareEmbeddedAttemptStream(input: {
claimEmbeddedPendingUserInputAnswer(text, options, attempt.sessionKey),
cancelPendingUserInput: (resolvedBy) =>
cancelPendingAgentQuestionForSession({ sessionKey: attempt.sessionKey, resolvedBy }),
preemptByVisibleTurn: heartbeatReplyOperation
? () => heartbeatReplyOperation.supersede()
: undefined,
queueMessage,
messageInjection: {
isAvailable: () =>
@@ -8,10 +8,13 @@ import { resetDiagnosticSessionStateForTest } from "../../logging/diagnostic-ses
import { createUserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.js";
import { createTestUserTurnTranscriptTarget } from "../../sessions/user-turn-transcript.test-support.js";
import {
clearActiveEmbeddedRun,
formatEmbeddedAgentQueueFailureSummary,
preemptAndDrainEmbeddedHeartbeatRun,
queueEmbeddedAgentMessageWithOutcome,
queueEmbeddedAgentMessageWithOutcomeAsync,
setActiveEmbeddedRun,
type EmbeddedAgentQueueHandle,
} from "./runs.js";
import { createEmbeddedRunHandle, testing } from "./runs.test-support.js";
@@ -24,6 +27,65 @@ describe("embedded-agent active-run steering", () => {
vi.restoreAllMocks();
});
it("aborts and drains the exact heartbeat handle through session replacement", async () => {
const heartbeatPreempt = vi.fn(() => true);
const finalizingHeartbeatPreempt = vi.fn(() => true);
const visibleAbort = vi.fn();
const heartbeatHandle: EmbeddedAgentQueueHandle = {
...createEmbeddedRunHandle(),
preemptByVisibleTurn: heartbeatPreempt,
};
const replacementHandle: EmbeddedAgentQueueHandle = {
...createEmbeddedRunHandle({ abort: visibleAbort }),
};
const finalizingHeartbeatHandle: EmbeddedAgentQueueHandle = {
...createEmbeddedRunHandle({ isAbortable: false }),
preemptByVisibleTurn: finalizingHeartbeatPreempt,
};
setActiveEmbeddedRun("heartbeat-session", heartbeatHandle);
setActiveEmbeddedRun("visible-session", replacementHandle);
setActiveEmbeddedRun("finalizing-heartbeat-session", finalizingHeartbeatHandle);
const heartbeatPreemption = preemptAndDrainEmbeddedHeartbeatRun("heartbeat-session", 1_000);
const finalizingPreemption = preemptAndDrainEmbeddedHeartbeatRun(
"finalizing-heartbeat-session",
1_000,
);
await expect(preemptAndDrainEmbeddedHeartbeatRun("visible-session", 1_000)).resolves.toBe(
"not-heartbeat",
);
let heartbeatDrained = false;
void heartbeatPreemption.then(() => {
heartbeatDrained = true;
});
setActiveEmbeddedRun("heartbeat-session", replacementHandle);
await Promise.resolve();
expect(heartbeatDrained).toBe(false);
clearActiveEmbeddedRun("heartbeat-session", heartbeatHandle);
clearActiveEmbeddedRun("finalizing-heartbeat-session", finalizingHeartbeatHandle);
await expect(heartbeatPreemption).resolves.toBe("drained");
await expect(finalizingPreemption).resolves.toBe("drained");
expect(heartbeatPreempt).toHaveBeenCalledOnce();
expect(finalizingHeartbeatPreempt).toHaveBeenCalledOnce();
expect(visibleAbort).not.toHaveBeenCalled();
});
it("leaves a heartbeat in another session running", async () => {
const preemptIsolatedHeartbeat = vi.fn(() => true);
setActiveEmbeddedRun("base-session:heartbeat", {
...createEmbeddedRunHandle(),
preemptByVisibleTurn: preemptIsolatedHeartbeat,
});
await expect(preemptAndDrainEmbeddedHeartbeatRun("base-session", 1_000)).resolves.toBe(
"not-heartbeat",
);
expect(preemptIsolatedHeartbeat).not.toHaveBeenCalled();
});
it("passes steering options to active embedded runs", () => {
const queueMessage = vi.fn(async () => {});
setActiveEmbeddedRun("session-steer", {
+42 -11
View File
@@ -730,6 +730,25 @@ export function abortEmbeddedAgentRun(
return false;
}
type EmbeddedHeartbeatPreemptionResult = "not-heartbeat" | "drained" | "timed-out";
export async function preemptAndDrainEmbeddedHeartbeatRun(
sessionId: string,
timeoutMs: number,
): Promise<EmbeddedHeartbeatPreemptionResult> {
const handle = ACTIVE_EMBEDDED_RUNS.get(sessionId);
if (!handle?.preemptByVisibleTurn) {
return "not-heartbeat";
}
const drainPromise = waitForCurrentEmbeddedAgentRunEnd(sessionId, timeoutMs, handle);
try {
handle.preemptByVisibleTurn();
} catch (err) {
diag.warn(`heartbeat preemption failed: sessionId=${sessionId} err=${String(err)}`);
}
return (await drainPromise) ? "drained" : "timed-out";
}
export function isEmbeddedAgentRunActive(sessionId: string): boolean {
const active = ACTIVE_EMBEDDED_RUNS.has(sessionId) || isReplyRunActiveForSessionId(sessionId);
if (active) {
@@ -866,8 +885,14 @@ export async function waitForActiveEmbeddedRuns(
function waitForCurrentEmbeddedAgentRunEnd(
sessionId: string,
timeoutMs: number | null,
handle?: EmbeddedAgentQueueHandle,
): Promise<boolean> {
if (!ACTIVE_EMBEDDED_RUNS.has(sessionId)) {
const isHandleActive = () =>
handle ? ACTIVE_EMBEDDED_RUNS.get(sessionId) === handle : ACTIVE_EMBEDDED_RUNS.has(sessionId);
if (!isHandleActive()) {
if (handle) {
return Promise.resolve(true);
}
return waitForReplyRunEndBySessionId(sessionId, timeoutMs);
}
const timeoutLabel = timeoutMs === null ? "none" : String(timeoutMs);
@@ -876,6 +901,7 @@ function waitForCurrentEmbeddedAgentRunEnd(
const waiters = EMBEDDED_RUN_WAITERS.get(sessionId) ?? new Set();
const waiter: EmbeddedRunWaiter = {
resolve,
handle,
};
if (timeoutMs !== null) {
waiter.timer = setTimeout(
@@ -892,7 +918,7 @@ function waitForCurrentEmbeddedAgentRunEnd(
}
waiters.add(waiter);
EMBEDDED_RUN_WAITERS.set(sessionId, waiters);
if (!ACTIVE_EMBEDDED_RUNS.has(sessionId)) {
if (!isHandleActive()) {
waiters.delete(waiter);
if (waiters.size === 0) {
EMBEDDED_RUN_WAITERS.delete(sessionId);
@@ -1111,19 +1137,26 @@ async function persistForceClearedEmbeddedRunTerminalState(params: {
}
}
function notifyEmbeddedRunEnded(sessionId: string) {
function notifyEmbeddedRunEnded(sessionId: string, endedHandle: EmbeddedAgentQueueHandle) {
const waiters = EMBEDDED_RUN_WAITERS.get(sessionId);
if (!waiters || waiters.size === 0) {
return;
}
EMBEDDED_RUN_WAITERS.delete(sessionId);
const sessionIdle = !ACTIVE_EMBEDDED_RUNS.has(sessionId);
diag.debug(`notifying waiters: sessionId=${sessionId} waiterCount=${waiters.size}`);
for (const waiter of waiters) {
if (waiter.handle ? waiter.handle !== endedHandle : !sessionIdle) {
continue;
}
waiters.delete(waiter);
if (waiter.timer) {
clearTimeout(waiter.timer);
}
waiter.resolve(true);
}
if (waiters.size === 0) {
EMBEDDED_RUN_WAITERS.delete(sessionId);
}
}
export function setActiveEmbeddedRun(
@@ -1193,9 +1226,6 @@ export function clearActiveEmbeddedRun(
reason = "run_completed",
) {
const activeHandle = ACTIVE_EMBEDDED_RUNS.get(sessionId);
if (activeHandle === undefined) {
return;
}
if (activeHandle === handle) {
ACTIVE_EMBEDDED_RUNS.delete(sessionId);
clearEmbeddedRunAbortability(handle, { retainFinalizing: true });
@@ -1213,10 +1243,11 @@ export function clearActiveEmbeddedRun(
if (!sessionId.startsWith("probe-")) {
diag.debug(`run cleared: sessionId=${sessionId} totalActive=${ACTIVE_EMBEDDED_RUNS.size}`);
}
notifyEmbeddedRunEnded(sessionId);
} else {
} else if (activeHandle !== undefined) {
diag.debug(`run clear skipped: sessionId=${sessionId} reason=handle_mismatch`);
}
// Exact-handle waiters own teardown even after another run takes the session slot.
notifyEmbeddedRunEnded(sessionId, handle);
}
function forceClearEmbeddedAgentRun(
@@ -1236,7 +1267,7 @@ function forceClearEmbeddedAgentRun(
clearActiveRunSessionFiles(sessionId);
logSessionStateChange({ sessionId, sessionKey, state: "idle", reason });
markDiagnosticEmbeddedRunEnded({ sessionId, sessionKey });
notifyEmbeddedRunEnded(sessionId);
notifyEmbeddedRunEnded(sessionId, handle);
cleared = true;
}
const cause = new Error(`Embedded run force-cleared by ${reason}`);
@@ -1254,7 +1285,7 @@ const testing = {
if (waiter.timer) {
clearTimeout(waiter.timer);
}
waiter.resolve(true);
waiter.resolve(!waiter.handle);
}
}
EMBEDDED_RUN_WAITERS.clear();
+1
View File
@@ -7,6 +7,7 @@
export {
abortAndDrainEmbeddedAgentRun,
abortEmbeddedAgentRun,
preemptAndDrainEmbeddedHeartbeatRun,
isEmbeddedAgentRunActive,
isEmbeddedAgentRunStreaming,
resolveActiveEmbeddedRunSessionId,
+1
View File
@@ -9,6 +9,7 @@ export type {
export {
abortAndDrainEmbeddedAgentRun,
abortEmbeddedAgentRun,
preemptAndDrainEmbeddedHeartbeatRun,
compactEmbeddedAgentSession,
isEmbeddedAgentRunAbortableForCompaction,
isEmbeddedAgentRunActive,
@@ -39,6 +39,7 @@ import { normalizeReplyPayload } from "./normalize-reply.js";
import { sanitizePendingFinalDeliveryText } from "./pending-final-delivery.js";
import { type FollowupRun, type QueueSettings, scheduleFollowupDrain } from "./queue.js";
import { normalizeReplyPayloadDirectives } from "./reply-delivery.js";
import { isReplyOperationSuperseded } from "./reply-operation-abort.js";
import { type ReplyOperation, runAfterReplyOperationClear } from "./reply-run-registry.js";
import { resolveSourceReplyVisibilityPolicy } from "./source-reply-delivery-mode.js";
import type { TypingController } from "./typing.js";
@@ -312,6 +313,9 @@ export async function handleReplyAgentRunError(
sessionCtx,
} = context;
if (isReplyOperationSuperseded(replyOperation)) {
return { text: SILENT_REPLY_TOKEN };
}
if (
replyOperation.result?.kind === "aborted" &&
replyOperation.result.code === "aborted_by_user"
@@ -131,6 +131,7 @@ function createReplyOperation(): TestReplyOperation {
get sessionId() {
return sessionId;
},
turnKind: "visible",
abortSignal: new AbortController().signal,
resetTriggered: false,
phase: "queued",
@@ -164,6 +165,7 @@ function createReplyOperation(): TestReplyOperation {
freezeAbort: vi.fn(),
abortByUser: vi.fn(),
abortForRestart: vi.fn(),
supersede: vi.fn(),
terminalRecovery: false,
acceptedSteeredInboundAudio: false,
markTerminalRecovery: vi.fn(),
+11 -1
View File
@@ -27,6 +27,7 @@ import {
} from "./pending-final-delivery.js";
import type { FollowupRun } from "./queue.js";
import type { ReplyMediaContext } from "./reply-media-paths.js";
import { isReplyOperationSuperseded } from "./reply-operation-abort.js";
import { recordReplyOperationAgentTurn } from "./reply-operation-agent-turn-state.js";
import { resolveReplyOperationRunState } from "./reply-operation-run-state.js";
import type { ReplyOperation } from "./reply-run-registry.js";
@@ -393,13 +394,22 @@ export async function executePreparedReplyAgentRun(
}),
),
);
const operationSuperseded = isReplyOperationSuperseded(replyOperation);
recordReplyOperationAgentTurn(
resolveReplyOperationRunState(opts),
runOutcome.outcome.kind === "settled" ? runOutcome.outcome.status : "failed",
operationSuperseded
? "superseded"
: runOutcome.outcome.kind === "settled"
? runOutcome.outcome.status
: "failed",
replyOperation,
);
activeSessionEntry = getActiveSessionEntry();
activeIsNewSession = getActiveIsNewSession();
if (operationSuperseded) {
return { text: SILENT_REPLY_TOKEN };
}
if (runOutcome.outcome.kind !== "settled") {
if (runOutcome.outcome.kind === "rejected" && !replyOperation.result) {
replyOperation.fail("run_failed", new Error("reply operation exited with final payload"));
@@ -81,6 +81,7 @@ function createReplyOperation(): TestReplyOperation {
return {
key: "test",
sessionId: "session",
turnKind: "visible",
abortSignal: new AbortController().signal,
staleExpiryReason: undefined,
resetTriggered: false,
@@ -111,6 +112,7 @@ function createReplyOperation(): TestReplyOperation {
fail: vi.fn(),
abortByUser: vi.fn(() => true),
abortForRestart: vi.fn(() => true),
supersede: vi.fn(() => true),
markTerminalRecovery: vi.fn(),
markAcceptedSteeredInboundAudio: vi.fn(),
markWaitingForDeferredMaintenance: vi.fn(),
+6 -1
View File
@@ -43,6 +43,7 @@ import { resolveActiveRunQueueAction } from "./queue-policy.js";
import { enqueueFollowupRun, scheduleFollowupDrain } from "./queue.js";
import { REPLY_ADMISSION_TICKET } from "./reply-admission-ticket.js";
import { createReplyMediaContext } from "./reply-media-paths.js";
import { isReplyOperationSuperseded } from "./reply-operation-abort.js";
import { recordReplyOperationAgentTurn } from "./reply-operation-agent-turn-state.js";
import * as replyRunState from "./reply-operation-run-state.js";
import { type ReplyOperation, replyRunRegistry } from "./reply-run-registry.js";
@@ -609,7 +610,11 @@ export async function runReplyAgent(
typingSignals,
});
} catch (error) {
recordReplyOperationAgentTurn(replyOperationRunState, "failed");
recordReplyOperationAgentTurn(
replyOperationRunState,
isReplyOperationSuperseded(replyOperation) ? "superseded" : "failed",
replyOperation,
);
return await handleReplyAgentRunError(error, {
blockReplyPipeline,
cfg,
@@ -813,80 +813,98 @@ describe("runReplyAgent auto-compaction token update", () => {
expect(scheduleFollowupDrain).toHaveBeenCalledTimes(1);
});
it("records a settled fallback cancelled by its upstream signal as aborted", async () => {
const upstreamAbort = new AbortController();
const sessionKey = "upstream-cancelled-settled-fallback";
const sessionEntry = {
sessionId: "session-upstream-cancelled",
updatedAt: Date.now(),
totalTokens: 50_000,
};
const replyOperation = createReplyOperation({
sessionKey,
sessionId: sessionEntry.sessionId,
resetTriggered: false,
upstreamAbortSignal: upstreamAbort.signal,
});
let releaseFallback: () => void = () => undefined;
let markCandidateSettled: () => void = () => undefined;
const candidateSettled = new Promise<void>((resolve) => {
markCandidateSettled = resolve;
});
const fallbackRelease = new Promise<void>((resolve) => {
releaseFallback = resolve;
});
runEmbeddedAgentMock.mockResolvedValueOnce({
payloads: [{ text: "late reply" }],
meta: { agentMeta: {} },
});
runWithModelFallbackMock.mockImplementationOnce(
async ({ provider, model, run }: RunWithModelFallbackParams) => {
const result = await run(provider, model);
markCandidateSettled();
await fallbackRelease;
return { result, provider, model };
},
);
const { typing, sessionCtx, resolvedQueue, followupRun } = createBaseRun({
storePath: "",
sessionEntry,
});
followupRun.run.sessionKey = sessionKey;
try {
const pending = runReplyAgent({
commandBody: "hello",
followupRun,
queueKey: sessionKey,
resolvedQueue,
shouldSteer: false,
shouldFollowup: false,
isActive: false,
typing,
sessionCtx,
sessionEntry,
sessionStore: { [sessionKey]: sessionEntry },
it.each([
{
label: "its upstream signal",
superseded: false,
expectedCode: "aborted_by_user" as const,
},
{
label: "a visible-turn supersession",
superseded: true,
expectedCode: "aborted_for_supersession" as const,
},
])(
"records a settled fallback cancelled by $label as aborted",
async ({ superseded, expectedCode }) => {
const upstreamAbort = new AbortController();
const sessionKey = `${superseded ? "superseded" : "upstream-cancelled"}-settled-fallback`;
const sessionEntry = {
sessionId: "session-upstream-cancelled",
updatedAt: Date.now(),
totalTokens: 50_000,
};
const replyOperation = createReplyOperation({
sessionKey,
defaultModel: "anthropic/claude-opus-4-6",
agentCfgContextTokens: 200_000,
resolvedVerboseLevel: "off",
isNewSession: false,
blockStreamingEnabled: false,
resolvedBlockStreamingBreak: "message_end",
shouldInjectGroupIntro: false,
typingMode: "instant",
replyOperation,
sessionId: sessionEntry.sessionId,
resetTriggered: false,
upstreamAbortSignal: upstreamAbort.signal,
});
await candidateSettled;
upstreamAbort.abort(new Error("caller cancelled"));
releaseFallback();
let releaseFallback: () => void = () => undefined;
let markCandidateSettled: () => void = () => undefined;
const candidateSettled = new Promise<void>((resolve) => {
markCandidateSettled = resolve;
});
const fallbackRelease = new Promise<void>((resolve) => {
releaseFallback = resolve;
});
runEmbeddedAgentMock.mockResolvedValueOnce({
payloads: [{ text: "late reply" }],
meta: { agentMeta: {} },
});
runWithModelFallbackMock.mockImplementationOnce(
async ({ provider, model, run }: RunWithModelFallbackParams) => {
const result = await run(provider, model);
markCandidateSettled();
await fallbackRelease;
return { result, provider, model };
},
);
const { typing, sessionCtx, resolvedQueue, followupRun } = createBaseRun({
storePath: "",
sessionEntry,
});
followupRun.run.sessionKey = sessionKey;
expectReplyText(await pending, SILENT_REPLY_TOKEN);
expect(replyOperation.result).toEqual({ kind: "aborted", code: "aborted_by_user" });
} finally {
replyOperation.complete();
}
});
try {
const pending = runReplyAgent({
commandBody: "hello",
followupRun,
queueKey: sessionKey,
resolvedQueue,
shouldSteer: false,
shouldFollowup: false,
isActive: false,
typing,
sessionCtx,
sessionEntry,
sessionStore: { [sessionKey]: sessionEntry },
sessionKey,
defaultModel: "anthropic/claude-opus-4-6",
agentCfgContextTokens: 200_000,
resolvedVerboseLevel: "off",
isNewSession: false,
blockStreamingEnabled: false,
resolvedBlockStreamingBreak: "message_end",
shouldInjectGroupIntro: false,
typingMode: "instant",
replyOperation,
});
await candidateSettled;
if (superseded) {
replyOperation.supersede();
} else {
upstreamAbort.abort(new Error("caller cancelled"));
}
releaseFallback();
expectReplyText(await pending, SILENT_REPLY_TOKEN);
expect(replyOperation.result).toEqual({ kind: "aborted", code: expectedCode });
} finally {
replyOperation.complete();
}
},
);
it("reports live diagnostic context from promptTokens, not provider usage totals", async () => {
const { sessionKey, stored, usageEvent } = await runBaseReplyWithAgentMeta({
@@ -54,6 +54,7 @@ import {
} from "./dispatch-from-config.test-harness.js";
import { getPreparedReplyDispatchRuntime } from "./prepared-reply-dispatch-context.js";
import { createReplyDispatcher } from "./reply-dispatcher.js";
import { admitReplyTurn } from "./reply-turn-admission.js";
import { buildChannelSourceTurnId } from "./source-turn-id.js";
import { buildTestCtx } from "./test-ctx.js";
@@ -393,6 +394,69 @@ describe("dispatchReplyFromConfig", () => {
activeOperation.complete();
});
it("preempts a heartbeat before resolving a visible Telegram turn", async () => {
setNoAbort();
const sessionKey = "agent:main:telegram:direct:heartbeat-preemption";
const heartbeatAdmission = await admitReplyTurn({
sessionKey,
sessionId: "heartbeat-session",
kind: "heartbeat",
resetTriggered: false,
});
expect(heartbeatAdmission.status).toBe("owned");
if (heartbeatAdmission.status !== "owned") {
return;
}
const heartbeatOperation = heartbeatAdmission.operation;
const cancel = vi.fn(() => heartbeatOperation.complete());
heartbeatOperation.attachBackend({
kind: "embedded",
cancel,
isStreaming: () => true,
});
heartbeatOperation.setPhase("running");
sessionStoreMocks.currentEntry = {
sessionId: "heartbeat-session",
updatedAt: Date.now(),
};
let heartbeatWasAbortedBeforeReply = false;
const replyResolver = vi.fn(async () => {
heartbeatWasAbortedBeforeReply = heartbeatOperation.abortSignal.aborted;
return { text: "visible reply" } satisfies ReplyPayload;
});
const result = await dispatchReplyFromConfig({
ctx: buildTestCtx({
Provider: "telegram",
Surface: "telegram",
OriginatingChannel: "telegram",
OriginatingTo: "user:1",
ChatType: "direct",
SessionKey: sessionKey,
BodyForAgent: "answer this now",
}),
cfg: automaticDirectReplyConfig,
dispatcher: createDispatcher(),
replyOptions: {
turnAdoptionLifecycle: {
onAdopted: async () => {},
onDeferred: vi.fn(),
onSettled: vi.fn(),
},
},
replyResolver,
});
expect(result.queuedFinal).toBe(true);
expect(heartbeatWasAbortedBeforeReply).toBe(true);
expect(heartbeatOperation.result).toEqual({
kind: "aborted",
code: "aborted_for_supersession",
});
expect(cancel).toHaveBeenCalledWith("superseded");
expect(replyResolver).toHaveBeenCalledOnce();
});
it("does not route when Provider matches OriginatingChannel (even if Surface is missing)", async () => {
setNoAbort();
mocks.routeReply.mockClear();
@@ -178,7 +178,8 @@ export function createDispatchReplyOperationCoordinator(params: {
phase === "dispatch" &&
replyTurnKind === "visible" &&
params.replyOptions?.turnAdoptionLifecycle !== undefined &&
activeReplyOperation !== undefined;
activeReplyOperation !== undefined &&
activeReplyOperation.turnKind !== "heartbeat";
if (allowGatewayQueueResolution) {
// Gateway turns need to reach getReplyFromConfig while the owner is active;
// that layer applies the session's steer/followup/collect/drop policy.
@@ -30,7 +30,10 @@ import {
loadSessionUpdatesRuntime,
routeThreadIdsMatch,
} from "./get-reply-run-helpers.js";
import { resolvePreparedReplyQueueState } from "./get-reply-run-queue.js";
import {
REPLY_RUN_STILL_SHUTTING_DOWN_TEXT,
resolvePreparedReplyQueueState,
} from "./get-reply-run-queue.js";
import { buildReplyPromptEnvelope } from "./prompt-prelude.js";
import { resolveActiveRunQueueAction } from "./queue-policy.js";
import { resolveQueueSettings } from "./queue/settings-runtime.js";
@@ -362,6 +365,23 @@ export async function prepareReplyRunAdmission(context: PreparedReplyRunContext)
)
? undefined
: rawActiveSessionIdForInterrupt;
const shouldPreemptHeartbeat =
!isRoomEvent && !context.isHeartbeat && rawActiveSessionIdForInterrupt !== undefined;
const heartbeatPreemption =
shouldPreemptHeartbeat && embeddedAgentRuntime
? await embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun(
rawActiveSessionIdForInterrupt,
REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS,
)
: "not-heartbeat";
if (heartbeatPreemption === "timed-out") {
typing.cleanup();
return {
kind: "reply",
reply: { text: REPLY_RUN_STILL_SHUTTING_DOWN_TEXT },
} as const;
}
const visibleTurnPreemptsHeartbeat = heartbeatPreemption === "drained";
if (
activeRunQueueMode === "interrupt" &&
!isRoomEvent &&
@@ -522,9 +542,11 @@ export async function prepareReplyRunAdmission(context: PreparedReplyRunContext)
activeRunAcceptsCurrentThread &&
!context.isHeartbeat &&
!effectiveResetTriggered &&
!visibleTurnPreemptsHeartbeat &&
resolvedQueue.mode === "steer";
const shouldFollowup =
!effectiveResetTriggered &&
!visibleTurnPreemptsHeartbeat &&
((isRoomEvent && isActive) ||
resolvedQueue.mode === "steer" ||
resolvedQueue.mode === "followup" ||
@@ -24,7 +24,11 @@ import {
buildInboundUserContextPrefix,
resolveInboundUserContextPromptJoiner,
} from "./inbound-meta.js";
import { createReplyOperation, getActiveReplyRunCount } from "./reply-run-registry.js";
import {
REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS,
createReplyOperation,
getActiveReplyRunCount,
} from "./reply-run-registry.js";
import { testing as replyRunTesting } from "./reply-run-registry.test-support.js";
import { routeReply } from "./route-reply.runtime.js";
import { drainFormattedSystemEvents } from "./session-system-events.js";
@@ -45,6 +49,7 @@ vi.mock("../../agents/embedded-agent.runtime.js", () => ({
abortEmbeddedAgentRun: vi.fn().mockReturnValue(false),
isEmbeddedAgentRunActive: vi.fn().mockReturnValue(false),
isEmbeddedAgentRunStreaming: vi.fn().mockReturnValue(false),
preemptAndDrainEmbeddedHeartbeatRun: vi.fn().mockResolvedValue("not-heartbeat"),
resolveActiveEmbeddedRunSessionId: vi.fn().mockReturnValue(undefined),
resolveActiveEmbeddedRunSessionIdBySessionFile: vi.fn().mockReturnValue(undefined),
resolveEmbeddedSessionLane: vi.fn().mockReturnValue("session:session-key"),
@@ -2023,6 +2028,147 @@ describe("runPreparedReply media-only handling", () => {
);
expect(vi.mocked(runReplyAgent)).toHaveBeenCalledOnce();
});
it("interrupts an embedded-only heartbeat before running a visible Telegram turn", async () => {
const queueSettings = await import("./queue/settings-runtime.js");
const embeddedAgentRuntime = await import("../../agents/embedded-agent.runtime.js");
let embeddedRunActive = true;
vi.mocked(queueSettings.resolveQueueSettings).mockReturnValueOnce({ mode: "steer" });
vi.mocked(embeddedAgentRuntime.resolveActiveEmbeddedRunSessionId).mockImplementation(() =>
embeddedRunActive ? "session-embedded-heartbeat" : undefined,
);
vi.mocked(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun).mockImplementation(
async (sessionId) => {
if (sessionId !== "session-embedded-heartbeat") {
return "not-heartbeat";
}
embeddedRunActive = false;
return "drained";
},
);
vi.mocked(embeddedAgentRuntime.isEmbeddedAgentRunActive).mockImplementation(
() => embeddedRunActive,
);
vi.mocked(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd).mockImplementation(async () => {
embeddedRunActive = false;
return true;
});
try {
await expect(
runPrepared({
isNewSession: false,
sessionId: "session-embedded-heartbeat",
ctx: {
...createInboundTurn("answer this now", "telegram", "direct"),
OriginatingChannel: "telegram",
OriginatingTo: "user:1",
},
sessionCtx: {
...createSessionTurn("answer this now", "telegram", "direct"),
OriginatingChannel: "telegram",
OriginatingTo: "user:1",
},
}),
).resolves.toEqual({ text: "ok" });
} finally {
vi.mocked(embeddedAgentRuntime.resolveActiveEmbeddedRunSessionId).mockReturnValue(undefined);
vi.mocked(embeddedAgentRuntime.isEmbeddedAgentRunActive).mockReturnValue(false);
vi.mocked(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun).mockResolvedValue(
"not-heartbeat",
);
vi.mocked(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd).mockResolvedValue(true);
}
expect(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun).toHaveBeenCalledWith(
"session-embedded-heartbeat",
REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS,
);
expect(embeddedAgentRuntime.abortEmbeddedAgentRun).not.toHaveBeenCalled();
expect(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd).not.toHaveBeenCalled();
expect(vi.mocked(runReplyAgent)).toHaveBeenCalledOnce();
});
it("drains an embedded heartbeat hidden by the visible pre-dispatch operation", async () => {
const queueSettings = await import("./queue/settings-runtime.js");
const embeddedAgentRuntime = await import("../../agents/embedded-agent.runtime.js");
const operation = createReplyOperation({
sessionId: "session-pre-dispatch-heartbeat",
sessionKey: "session-key",
turnKind: "visible",
resetTriggered: false,
});
let embeddedRunActive = true;
let releaseDrain: (() => void) | undefined;
const drainBarrier = new Promise<void>((resolve) => {
releaseDrain = resolve;
});
vi.mocked(queueSettings.resolveQueueSettings).mockReturnValueOnce({ mode: "steer" });
vi.mocked(embeddedAgentRuntime.resolveActiveEmbeddedRunSessionId).mockImplementation(() =>
embeddedRunActive ? "session-pre-dispatch-heartbeat" : undefined,
);
vi.mocked(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun).mockImplementation(
async () => {
await drainBarrier;
embeddedRunActive = false;
return "drained";
},
);
vi.mocked(embeddedAgentRuntime.isEmbeddedAgentRunActive).mockImplementation(
() => embeddedRunActive,
);
vi.mocked(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd).mockImplementation(async () => {
await drainBarrier;
embeddedRunActive = false;
return true;
});
try {
const runPromise = runPrepared({
isNewSession: false,
sessionId: "session-pre-dispatch-heartbeat",
opts: { replyOperation: operation } as never,
ctx: {
...createInboundTurn("answer this now", "telegram", "direct"),
OriginatingChannel: "telegram",
OriginatingTo: "user:1",
},
sessionCtx: {
...createSessionTurn("answer this now", "telegram", "direct"),
OriginatingChannel: "telegram",
OriginatingTo: "user:1",
},
});
await vi.waitFor(
() => {
expect(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun).toHaveBeenCalledWith(
"session-pre-dispatch-heartbeat",
REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS,
);
},
{ timeout: 1_000 },
);
expect(vi.mocked(runReplyAgent)).not.toHaveBeenCalled();
expect(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd).not.toHaveBeenCalled();
releaseDrain?.();
await expect(runPromise).resolves.toEqual({ text: "ok" });
} finally {
releaseDrain?.();
operation.complete();
vi.mocked(embeddedAgentRuntime.resolveActiveEmbeddedRunSessionId)
.mockReset()
.mockReturnValue(undefined);
vi.mocked(embeddedAgentRuntime.preemptAndDrainEmbeddedHeartbeatRun)
.mockReset()
.mockResolvedValue("not-heartbeat");
vi.mocked(embeddedAgentRuntime.isEmbeddedAgentRunActive).mockReset().mockReturnValue(false);
vi.mocked(embeddedAgentRuntime.waitForEmbeddedAgentRunEnd)
.mockReset()
.mockResolvedValue(true);
}
expect(vi.mocked(runReplyAgent)).toHaveBeenCalledOnce();
});
it("refreshes goal context after interrupt admission waits", async () => {
const queueSettings = await import("./queue/settings-runtime.js");
const inboundMeta = await import("./inbound-meta.js");
@@ -2120,7 +2266,10 @@ describe("runPreparedReply media-only handling", () => {
expect(result).toEqual({ text: "ok" });
expect(commandQueue.clearCommandLane).toHaveBeenCalledWith("session:session-key");
expect(embeddedAgentRuntime.abortEmbeddedAgentRun).toHaveBeenCalledWith("session-active");
expect(activeOperation.result).toEqual({ kind: "aborted", code: "aborted_by_user" });
expect(activeOperation.result).toEqual({
kind: "aborted",
code: "aborted_by_user",
});
expect(vi.mocked(runReplyAgent)).toHaveBeenCalledOnce();
const call = requireRunReplyAgentCall();
expect(call?.shouldSteer).toBe(false);
+20 -2
View File
@@ -1,5 +1,8 @@
import { isFallbackSummaryError } from "../../agents/model-fallback-attempt.js";
import { isAgentRunRestartAbortReason } from "../../agents/run-termination.js";
import {
isAgentRunRestartAbortReason,
isAgentRunSupersededAbortReason,
} from "../../agents/run-termination.js";
import { CommandLaneClearedError, GatewayDrainingError } from "../../process/command-queue.js";
import type { ReplyOperation } from "./reply-run-registry.js";
@@ -15,7 +18,11 @@ export function isReplyOperationUserAbort(replyOperation?: ReplyOperation): bool
return true;
}
const abortSignal = replyOperation?.abortSignal;
return abortSignal?.aborted === true && !isAgentRunRestartAbortReason(abortSignal.reason);
return (
abortSignal?.aborted === true &&
!isAgentRunRestartAbortReason(abortSignal.reason) &&
!isAgentRunSupersededAbortReason(abortSignal.reason)
);
}
export function isReplyOperationRestartAbort(replyOperation?: ReplyOperation): boolean {
@@ -29,6 +36,17 @@ export function isReplyOperationRestartAbort(replyOperation?: ReplyOperation): b
return abortSignal?.aborted === true && isAgentRunRestartAbortReason(abortSignal.reason);
}
export function isReplyOperationSuperseded(replyOperation?: ReplyOperation): boolean {
if (
replyOperation?.result?.kind === "aborted" &&
replyOperation.result.code === "aborted_for_supersession"
) {
return true;
}
const abortSignal = replyOperation?.abortSignal;
return abortSignal?.aborted === true && isAgentRunSupersededAbortReason(abortSignal.reason);
}
export function resolveRestartLifecycleError(
error: unknown,
): GatewayDrainingError | CommandLaneClearedError | undefined {
@@ -1,20 +1,32 @@
import { isReplyOperationSuperseded } from "./reply-operation-abort.js";
import type { ReplyOperationRunState } from "./reply-operation-run-state.js";
import type { ReplyOperation } from "./reply-run-registry.js";
type ReplyOperationAgentTurnStatus = "ok" | "failed";
type ReplyOperationAgentTurnStatus = "ok" | "failed" | "superseded";
const agentTurns = new WeakMap<ReplyOperationRunState, ReplyOperationAgentTurnStatus>();
type ReplyOperationAgentTurn = {
status: ReplyOperationAgentTurnStatus;
owner?: ReplyOperation;
};
const agentTurns = new WeakMap<ReplyOperationRunState, ReplyOperationAgentTurn>();
export function recordReplyOperationAgentTurn(
state: ReplyOperationRunState | undefined,
status: ReplyOperationAgentTurnStatus,
owner?: ReplyOperation,
): void {
if (state) {
agentTurns.set(state, status);
agentTurns.set(state, { status, owner });
}
}
export function resolveReplyOperationAgentTurn(
state: ReplyOperationRunState | undefined,
): ReplyOperationAgentTurnStatus | undefined {
return state ? agentTurns.get(state) : undefined;
if (!state) {
return undefined;
}
const turn = agentTurns.get(state);
return isReplyOperationSuperseded(turn?.owner) ? "superseded" : turn?.status;
}
@@ -22,6 +22,8 @@ type ReplyBackendKind = "embedded" | "cli";
type ReplyBackendCancelReason = "user_abort" | "restart" | "superseded";
export type ReplyTurnKind = "visible" | "heartbeat" | "queued_followup";
export type ReplyBackendQueueMessageOptions = {
steeringMode?: "all";
/** True when this queue item came from the channel's current user turn. */
@@ -198,7 +200,10 @@ type ReplyOperationFailureCode =
| "run_stalled"
| "run_failed";
type ReplyOperationAbortCode = "aborted_by_user" | "aborted_for_restart";
type ReplyOperationAbortCode =
| "aborted_by_user"
| "aborted_for_restart"
| "aborted_for_supersession";
type ReplyOperationResult =
| { kind: "completed" }
@@ -208,6 +213,7 @@ type ReplyOperationResult =
export type ReplyOperation = {
readonly key: ReplyRunKey;
readonly sessionId: string;
readonly turnKind: ReplyTurnKind;
/** Gateway lifecycle that admitted this process-local owner. */
readonly lifecycleGeneration?: string;
readonly routeThreadId?: string | number;
@@ -305,6 +311,7 @@ export type ReplyOperation = {
fail(code: Exclude<ReplyOperationFailureCode, "aborted_by_user">, cause?: unknown): void;
abortByUser(): boolean;
abortForRestart(): boolean;
supersede(): boolean;
};
export type ReplyRunRegistry = {
@@ -1,7 +1,9 @@
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import {
createAgentRunRestartAbortError,
createAgentRunSupersededAbortError as createSupersededError,
isAgentRunRestartAbortReason,
isAgentRunSupersededAbortReason,
} from "../../agents/run-termination.js";
import { createAbortError } from "../../infra/abort-signal.js";
import { getAgentEventLifecycleGeneration } from "../../infra/agent-events.js";
@@ -17,6 +19,7 @@ import {
type ReplyOperation,
type ReplyOperationPhase,
type ReplyToolAuthorityProjector,
type ReplyTurnKind,
} from "./reply-run-registry.contracts.js";
import {
abortFrozenOperations,
@@ -46,13 +49,13 @@ import {
} from "./reply-run-registry.state.js";
type ReplyBackendCancelReason = "user_abort" | "restart" | "superseded";
type ReplyOperationAbortCode = "aborted_by_user" | "aborted_for_restart";
type ReplyOperationResult = NonNullable<ReplyOperation["result"]>;
type ReplyOperationStaleReason = replyRunSettle.ReplyOperationStaleReason;
type ReplyOperationAbortCode = Extract<ReplyOperationResult, { kind: "aborted" }>["code"];
export function createReplyOperation(params: {
sessionKey: string;
sessionId: string;
turnKind?: ReplyTurnKind;
resetTriggered: boolean;
routeThreadId?: string | number;
originatingLeafEntryId?: string | null;
@@ -87,7 +90,7 @@ export function createReplyOperation(params: {
let currentSessionId = sessionId;
let phase: ReplyOperationPhase = "queued";
let phaseBeforeGlobalLaneWait: "queued" | "running" | undefined;
let staleExpiryReason: ReplyOperationStaleReason | undefined;
let staleExpiryReason: replyRunSettle.ReplyOperationStaleReason | undefined;
let result: ReplyOperationResult | null = null;
let stateCleared = false;
let clearBarrierSettlement: Promise<void> | undefined;
@@ -222,6 +225,7 @@ export function createReplyOperation(params: {
get sessionId() {
return currentSessionId;
},
turnKind: params.turnKind ?? "visible",
lifecycleGeneration,
get routeThreadId() {
return params.routeThreadId;
@@ -444,7 +448,9 @@ export function createReplyOperation(params: {
result.kind === "aborted"
? result.code === "aborted_for_restart"
? "restart"
: "user_abort"
: result.code === "aborted_for_supersession"
? "superseded"
: "user_abort"
: "superseded",
);
return;
@@ -537,6 +543,22 @@ export function createReplyOperation(params: {
abortOperation("restart", createAgentRunRestartAbortError(), "aborted_for_restart");
return true;
},
supersede() {
if (result || stateCleared) {
return false;
}
if (abortFrozenOperations.has(operation)) {
setResult({ kind: "aborted", code: "aborted_for_supersession" });
phase = "aborted";
scheduleTerminalSettle();
return true;
}
if (!isReplyOperationAbortable(operation)) {
return false;
}
abortOperation("superseded", createSupersededError(), "aborted_for_supersession");
return true;
},
};
expireReplyOperationByOperation.set(operation, (reason, options) => {
@@ -680,10 +702,15 @@ export function createReplyOperation(params: {
return;
}
const restart = isAgentRunRestartAbortReason(upstreamAbortSignal.reason);
const superseded = isAgentRunSupersededAbortReason(upstreamAbortSignal.reason);
abortOperation(
restart ? "restart" : "user_abort",
restart ? "restart" : superseded ? "superseded" : "user_abort",
upstreamAbortSignal.reason,
restart ? "aborted_for_restart" : "aborted_by_user",
restart
? "aborted_for_restart"
: superseded
? "aborted_for_supersession"
: "aborted_by_user",
);
};
if (upstreamAbortSignal.aborted) {
@@ -699,7 +726,7 @@ export function createReplyOperation(params: {
export function expireStaleReplyOperation(
operation: ReplyOperation,
reason: ReplyOperationStaleReason,
reason: replyRunSettle.ReplyOperationStaleReason,
options?: ReplyOperationStaleExpiryOptions,
): boolean {
return expireReplyOperationByOperation.get(operation)?.(reason, options) ?? false;
@@ -1449,6 +1449,26 @@ describe("reply run registry", () => {
expect(replyRunRegistry.get("agent:main:reentrant-expire")).toBeUndefined();
});
it("keeps supersession attribution when backend cancellation re-enters user abort", () => {
const operation = createTestReplyOperation({
sessionKey: "agent:main:heartbeat-preemption",
sessionId: "heartbeat-preemption-session",
turnKind: "heartbeat",
});
const cancel = vi.fn(() => {
operation.abortByUser();
});
operation.attachBackend({ kind: "embedded", cancel, isStreaming: () => true });
operation.setPhase("running");
expect(operation.supersede()).toBe(true);
expect(cancel).toHaveBeenCalledWith("superseded");
expect(operation.result).toEqual({
kind: "aborted",
code: "aborted_for_supersession",
});
});
it("cancels terminal settle when the owner clears state first", async () => {
await withFakeReplyTimers(async () => {
const warnSpy = vi.spyOn(diagnosticLogger, "warn").mockImplementation(() => undefined);
@@ -16,6 +16,7 @@ export type {
ReplyOperation,
ReplyOperationPhase,
ReplyToolAuthorityOverlay,
ReplyTurnKind,
} from "./reply-run-registry.contracts.js";
export {
abortReplyMessageInjectionTarget,
+7 -3
View File
@@ -36,13 +36,11 @@ import {
retainReplyOperationUntilComplete,
runAfterReplyOperationClear,
type ReplyOperation,
type ReplyTurnKind,
waitForReplyRunFollowupAdmission,
waitForReplyRunSuccessorAdmission,
} from "./reply-run-registry.js";
/** Kinds of turns that compete for one reply run slot per session. */
type ReplyTurnKind = "visible" | "heartbeat" | "queued_followup";
/** Admission result for a reply turn attempting to own the session run slot. */
type ReplyTurnAdmission =
| { status: "owned"; operation: ReplyOperation; sessionEntry?: SessionEntry }
@@ -336,6 +334,7 @@ async function admitReplyTurnWithWaitSignal(
operation = createReplyOperation({
sessionKey: params.sessionKey,
sessionId,
turnKind: params.kind,
resetTriggered: params.resetTriggered,
routeThreadId: params.routeThreadId,
originatingLeafEntryId: params.originatingLeafEntryId,
@@ -441,6 +440,11 @@ async function admitReplyTurnWithWaitSignal(
throw error;
}
const activeOperation = replyRunRegistry.get(params.sessionKey);
if (params.kind === "visible" && activeOperation?.turnKind === "heartbeat") {
// Background heartbeats must yield before queue policy can steer this
// user turn into the heartbeat's model run and lose its visible reply.
activeOperation.supersede();
}
if (params.kind === "visible" && expireVisibleStaleOperation(activeOperation)) {
continue;
}
+2
View File
@@ -26,6 +26,7 @@ export function createMockReplyOperation(
const replyOperation: ReplyOperation = {
key: overrides.key ?? "main",
sessionId,
turnKind: "visible",
abortSignal: overrides.abortSignal ?? new AbortController().signal,
resetTriggered: false,
terminalRecovery: false,
@@ -76,6 +77,7 @@ export function createMockReplyOperation(
fail: failMock,
abortByUser: vi.fn(() => true),
abortForRestart: vi.fn(() => true),
supersede: vi.fn(() => true),
};
return {
replyOperation,
+14 -4
View File
@@ -1,7 +1,9 @@
import {
HEARTBEAT_IDLE_RETRY_GRACE_MS,
HEARTBEAT_SKIP_CRON_IN_PROGRESS,
HEARTBEAT_SKIP_PREEMPTED,
type HeartbeatRunResult,
isRetryableHeartbeatBusySkipReason,
isRetryableHeartbeatSkipReason,
} from "../../infra/heartbeat-wake.js";
import type { CommandLaneTaskMarker } from "../../process/command-queue.js";
import { type CronActiveJobMarker, isCronActiveJobMarkerCurrent } from "../active-jobs.js";
@@ -297,7 +299,7 @@ async function executeMainSessionCronJob(
}
if (
heartbeatResult.status !== "skipped" ||
!isRetryableHeartbeatBusySkipReason(heartbeatResult.reason)
!isRetryableHeartbeatSkipReason(heartbeatResult.reason)
) {
break;
}
@@ -318,7 +320,8 @@ async function executeMainSessionCronJob(
removeQueuedSystemEventHandle(state, job, queuedSystemEvent);
return { status: "error", error: timeoutErrorMessage() };
}
if (state.deps.nowMs() - waitStartedAt > maxWaitMs) {
const elapsedMs = state.deps.nowMs() - waitStartedAt;
if (elapsedMs >= maxWaitMs) {
if (abortSignal?.aborted) {
removeQueuedSystemEventHandle(state, job, queuedSystemEvent);
return { status: "error", error: timeoutErrorMessage() };
@@ -333,7 +336,14 @@ async function executeMainSessionCronJob(
});
return { status: "ok", summary: text, sessionKey: cronRunSessionKey };
}
await waitWithAbort(retryDelayMs);
await waitWithAbort(
Math.min(
heartbeatResult.reason === HEARTBEAT_SKIP_PREEMPTED
? HEARTBEAT_IDLE_RETRY_GRACE_MS
: retryDelayMs,
maxWaitMs - elapsedMs,
),
);
}
if (heartbeatResult.status === "ran") {
+60 -1
View File
@@ -11,7 +11,12 @@ import {
} from "../../../test/helpers/cron/service-regression-fixtures.js";
import { createDeferred } from "../../../test/helpers/promise.js";
import { DEFAULT_CRON_MAX_CONCURRENT_RUNS } from "../../config/cron-limits.js";
import { HEARTBEAT_SKIP_LANES_BUSY, type HeartbeatRunResult } from "../../infra/heartbeat-wake.js";
import {
HEARTBEAT_IDLE_RETRY_GRACE_MS,
HEARTBEAT_SKIP_LANES_BUSY,
HEARTBEAT_SKIP_PREEMPTED,
type HeartbeatRunResult,
} from "../../infra/heartbeat-wake.js";
import { openOpenClawStateDatabase } from "../../state/openclaw-state-db.js";
import { CRON_TASK_KIND } from "../../tasks/cron-task-contract.js";
import { cancelTaskById, listTaskRecords } from "../../tasks/task-registry.js";
@@ -1278,6 +1283,60 @@ describe("cron service timer regressions", () => {
expect(requestHeartbeat).not.toHaveBeenCalled();
});
it.each([
{ label: "retries after idle grace", staysPreempted: false, expectedCalls: 2 },
{ label: "requeues at the wait budget", staysPreempted: true, expectedCalls: 3 },
])("$label for a preempted wake-now heartbeat", async ({ staysPreempted, expectedCalls }) => {
vi.useFakeTimers();
try {
let heartbeatAttempt = 0;
const runHeartbeatOnce = vi.fn<() => Promise<HeartbeatRunResult>>(async () =>
staysPreempted || ++heartbeatAttempt === 1
? { status: "skipped", reason: HEARTBEAT_SKIP_PREEMPTED }
: { status: "ran", durationMs: 1 },
);
const requestHeartbeat = vi.fn();
const mainJob: CronJob = {
id: "main-preempted",
name: "main preempted",
enabled: true,
createdAtMs: Date.now(),
updatedAtMs: Date.now(),
schedule: { kind: "at", at: new Date(Date.now() + 60_000).toISOString() },
sessionTarget: "main",
wakeMode: "now",
payload: { kind: "systemEvent", text: "tick" },
state: {},
};
const state = createCronServiceState({
cronEnabled: true,
storePath: "/tmp/openclaw-cron-preempted-test/jobs.json",
log: noopLogger,
nowMs: () => Date.now(),
enqueueSystemEvent: vi.fn(),
requestHeartbeat,
runHeartbeatOnce,
wakeNowHeartbeatBusyMaxWaitMs: 2 * HEARTBEAT_IDLE_RETRY_GRACE_MS,
wakeNowHeartbeatBusyRetryDelayMs: 1,
runIsolatedAgentJob: createDefaultIsolatedRunner(),
});
const resultPromise = executeJobCore(state, mainJob);
await vi.advanceTimersByTimeAsync(HEARTBEAT_IDLE_RETRY_GRACE_MS - 1);
expect(runHeartbeatOnce).toHaveBeenCalledTimes(1);
await vi.advanceTimersByTimeAsync(1);
if (staysPreempted) {
await vi.advanceTimersByTimeAsync(HEARTBEAT_IDLE_RETRY_GRACE_MS);
}
await expect(resultPromise).resolves.toMatchObject({ status: "ok" });
expect(runHeartbeatOnce).toHaveBeenCalledTimes(expectedCalls);
expect(requestHeartbeat).toHaveBeenCalledTimes(staysPreempted ? 1 : 0);
} finally {
vi.useRealTimers();
}
});
it("keeps user cancellation disabled for main-session cron wrappers", async () => {
vi.useFakeTimers();
try {
+5 -1
View File
@@ -676,13 +676,17 @@ export async function invokeHeartbeatAgentRun(
replyOpts,
cfg,
);
const agentTurnStatus = resolveReplyOperationAgentTurn(replyOperationRunState);
if (agentTurnStatus === "superseded") {
return { kind: "preempted" } as const;
}
const heartbeatToolResponse = resolveHeartbeatToolResponseFromReplyResult(replyResult);
const heartbeatScratchProposal = resolveHeartbeatScratchProposalFromReplyResult(replyResult);
const heartbeatTerminalToolFailure: HeartbeatTerminalToolFailure | undefined =
resolveHeartbeatTerminalToolFailure(replyResult);
const selectedReplyPayload = resolveHeartbeatReplyPayload(replyResult);
const replyPayload = selectedReplyPayload;
const agentRunFailed = resolveReplyOperationAgentTurn(replyOperationRunState) === "failed";
const agentRunFailed = agentTurnStatus === "failed";
if (
heartbeatScratchProposal !== undefined &&
heartbeatToolResponse &&
+13 -1
View File
@@ -22,7 +22,11 @@ import {
type HeartbeatRunOptions,
} from "./heartbeat-runner-execution.js";
import { createHeartbeatTypingCallbacks } from "./heartbeat-typing.js";
import { HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT, type HeartbeatRunResult } from "./heartbeat-wake.js";
import {
HEARTBEAT_SKIP_PREEMPTED,
HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT,
type HeartbeatRunResult,
} from "./heartbeat-wake.js";
import { resolveAgentOutboundIdentity } from "./outbound/identity.js";
import { buildOutboundSessionContext } from "./outbound/session-context.js";
@@ -147,6 +151,14 @@ export async function runHeartbeatOnce(opts: HeartbeatRunOptions): Promise<Heart
});
return { status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT };
}
if (agentRun.kind === "preempted") {
emitHeartbeatEvent({
status: "skipped",
reason: HEARTBEAT_SKIP_PREEMPTED,
durationMs: Date.now() - startedAt,
});
return { status: "skipped", reason: HEARTBEAT_SKIP_PREEMPTED };
}
const outcome = classifyHeartbeatAgentOutcome({
agentRun,
hasRelayableExecCompletion,
+12 -12
View File
@@ -35,7 +35,7 @@ import {
type HeartbeatWakeHandler,
type HeartbeatWakeIntent,
type HeartbeatWakeRequest,
isRetryableHeartbeatBusySkipReason,
isRetryableHeartbeatSkipReason,
setHeartbeatWakeHandler,
} from "./heartbeat-wake.js";
@@ -353,7 +353,7 @@ export function startHeartbeatRunner(opts: {
// closures are safe to fan out under `Promise.all`.
type AgentWakeOutcome = {
ran: boolean;
retryableBusySkip?: HeartbeatRunResult;
retryableSkip?: HeartbeatRunResult;
// Terminal per-agent result so targeted callers can report the real
// skip reason instead of collapsing everything to not-due.
result?: HeartbeatRunResult;
@@ -407,10 +407,10 @@ export function startHeartbeatRunner(opts: {
advanceAgentSchedule(agent, now, reason);
return { ran: false, result: { status: "failed", reason: formatErrorMessage(err) } };
}
if (res.status === "skipped" && isRetryableHeartbeatBusySkipReason(res.reason)) {
if (res.status === "skipped" && isRetryableHeartbeatSkipReason(res.reason)) {
// Do not advance the schedule or record run bookkeeping for this
// agent — its target runtime is busy and the wake layer retries.
return { ran: false, retryableBusySkip: res };
return { ran: false, retryableSkip: res };
}
if (
params.source === "exec-event" &&
@@ -455,8 +455,8 @@ export function startHeartbeatRunner(opts: {
// replaced broadcast timer — not resolveHeartbeatForWake, which only
// ever served override-carrying targeted event wakes.
const outcome = await runOneAgent(targetAgent, authoritativeScheduledTick);
if (outcome.retryableBusySkip) {
return outcome.retryableBusySkip;
if (outcome.retryableSkip) {
return outcome.retryableSkip;
}
if (outcome.ran) {
return { status: "ran", durationMs: Date.now() - startedAt };
@@ -500,7 +500,7 @@ export function startHeartbeatRunner(opts: {
tasks: requestedTasks,
deps: { runtime: state.runtime },
});
if (res.status === "skipped" && isRetryableHeartbeatBusySkipReason(res.reason)) {
if (res.status === "skipped" && isRetryableHeartbeatSkipReason(res.reason)) {
// Retryable busy — do NOT record run bookkeeping. The wake layer
// retries the same reason shortly; if we recorded `lastRunStartedAtMs`
// here, the retry would falsely defer with `not-due`/`min-spacing`
@@ -548,20 +548,20 @@ export function startHeartbeatRunner(opts: {
runOneAgent(agent, authoritativeScheduledTick),
),
);
let firstRetryableBusy: HeartbeatRunResult | undefined;
let firstRetryableSkip: HeartbeatRunResult | undefined;
for (const outcome of agentOutcomes) {
if (outcome.ran) {
ran = true;
}
if (outcome.retryableBusySkip && !firstRetryableBusy) {
firstRetryableBusy = outcome.retryableBusySkip;
if (outcome.retryableSkip && !firstRetryableSkip) {
firstRetryableSkip = outcome.retryableSkip;
}
}
if (firstRetryableBusy) {
if (firstRetryableSkip) {
// At least one agent's runtime was busy. The wake layer schedules a
// retry; on retry, agents that already advanced their schedule will
// defer via cooldown, so only the still-busy agent actually re-runs.
return firstRetryableBusy;
return firstRetryableSkip;
}
if (ran) {
+12 -33
View File
@@ -13,6 +13,7 @@ import { startHeartbeatRunner } from "./heartbeat-runner.js";
import { computeNextHeartbeatPhaseDueMs, resolveHeartbeatPhaseMs } from "./heartbeat-schedule.js";
import {
getHeartbeatWakeAbortSignal,
HEARTBEAT_SKIP_PREEMPTED,
HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT,
requestHeartbeat,
setHeartbeatWakeHandler,
@@ -662,11 +663,8 @@ describe("startHeartbeatRunner", () => {
await pokeIntervalWake();
expect(runSpy).toHaveBeenCalledTimes(1);
// The wake layer auto-retries the busy interval wake every 1s; the busy
// skips must not advance nextDueMs, so each retry reaches runOnce until
// the 6th attempt succeeds.
for (let i = 0; i < 5; i++) {
await vi.advanceTimersByTimeAsync(1_000);
await vi.advanceTimersByTimeAsync(60_000);
}
expect(runSpy).toHaveBeenCalledTimes(6);
const scheduledSlotCallsBeforeInterval = callTimes.filter(
@@ -1157,26 +1155,16 @@ describe("startHeartbeatRunner", () => {
runner.stop();
});
it("retryable busy skip does not poison the cooldown for the next retry", async () => {
// Reproduces P2 finding from #75439 review: if a targeted exec-event wake
// hits requests-in-flight on its first attempt, the wake layer retries the
// same reason. The cooldown must NOT have been advanced by the busy attempt
// — otherwise the retry would falsely defer with `not-due`/`min-spacing`.
it.each([
[HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT, 1_500],
[HEARTBEAT_SKIP_PREEMPTED, 60_500],
])("retryable %s does not poison the next retry", async (skipReason, retryDelayMs) => {
useFakeHeartbeatTime();
let attempt = 0;
const runSpy = vi.fn().mockImplementation(async () => {
attempt += 1;
if (attempt === 1) {
return { status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT } as const;
}
return { status: "ran", durationMs: 1 } as const;
});
const runner = startHeartbeatRunner({
cfg: heartbeatConfig(),
runOnce: runSpy,
stableSchedulerSeed: TEST_SCHEDULER_SEED,
});
const runSpy = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: skipReason })
.mockResolvedValue({ status: "ran", durationMs: 1 });
const runner = startDefaultRunner(runSpy);
requestHeartbeat({
source: "exec-event",
@@ -1188,22 +1176,13 @@ describe("startHeartbeatRunner", () => {
await vi.advanceTimersByTimeAsync(1);
expect(runSpy).toHaveBeenCalledTimes(1);
// Wake layer retries via DEFAULT_RETRY_MS (1s). Advance past it.
await vi.advanceTimersByTimeAsync(1500);
await vi.advanceTimersByTimeAsync(retryDelayMs);
// The retry must NOT be deferred to `not-due` or `min-spacing`. Since the
// first attempt was a retryable busy skip, the cooldown bookkeeping was
// never recorded — so the retry should reach runOnce normally.
expect(runSpy).toHaveBeenCalledTimes(2);
expectRunCallFields(runSpy, 1, {
reason: "exec-event",
sessionKey: "agent:main:main",
});
await expect(runSpy.mock.results[1]?.value).resolves.toEqual({
status: "ran",
durationMs: 1,
});
runner.stop();
});
});
@@ -1,6 +1,16 @@
// Covers heartbeat skipping while session lanes or cron jobs are busy.
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
import {
clearActiveEmbeddedRun,
preemptAndDrainEmbeddedHeartbeatRun,
setActiveEmbeddedRun,
} from "../agents/embedded-agent-runner/runs.js";
import {
createEmbeddedRunHandle,
testing as embeddedRunTesting,
} from "../agents/embedded-agent-runner/runs.test-support.js";
import { resolveNestedAgentLaneForSession } from "../agents/lanes.js";
import { recordReplyOperationAgentTurn } from "../auto-reply/reply/reply-operation-agent-turn-state.js";
import { resolveReplyOperationRunState } from "../auto-reply/reply/reply-operation-run-state.js";
import { createReplyOperation } from "../auto-reply/reply/reply-run-registry.js";
import { testing as replyRunRegistryTesting } from "../auto-reply/reply/reply-run-registry.test-support.js";
@@ -45,6 +55,7 @@ afterAll(() => {
});
beforeEach(() => {
embeddedRunTesting.resetActiveEmbeddedRuns();
resetSystemEventsForTest();
resetCronActiveJobs();
replyRunRegistryTesting.resetReplyRunRegistry();
@@ -481,6 +492,47 @@ describe("heartbeat runner skips when target session lane is busy", () => {
});
});
it("suppresses delivery when a visible turn supersedes a finalizing heartbeat", async () => {
await withTempHeartbeatSandbox(async ({ storePath, replySpy }) => {
const cfg = createHeartbeatTelegramConfig(storePath);
const sessionKey = await seedHeartbeatTelegramSession(storePath, cfg);
const sessionId = "finalizing-heartbeat-session";
let preempt: ReturnType<typeof vi.fn<() => boolean>> | undefined;
replySpy.mockImplementationOnce(async (_ctx, options) => {
const operation = createReplyOperation({
sessionKey,
sessionId,
turnKind: "heartbeat",
resetTriggered: false,
});
preempt = vi.fn(() => operation.supersede());
const handle = {
...createEmbeddedRunHandle({ isAbortable: false }),
preemptByVisibleTurn: preempt,
};
const runState = resolveReplyOperationRunState(options);
if (!runState) {
throw new Error("Expected heartbeat reply operation run state");
}
recordReplyOperationAgentTurn(runState, "ok", operation);
operation.freezeAbort();
setActiveEmbeddedRun(sessionId, handle, sessionKey);
const drained = preemptAndDrainEmbeddedHeartbeatRun(sessionId, 1_000);
clearActiveEmbeddedRun(sessionId, handle, sessionKey);
await expect(drained).resolves.toBe("drained");
operation.complete();
return { text: "Background work finished." };
});
const sendTelegram = vi.fn().mockResolvedValue({ messageId: "m1", chatId: "123" });
const result = await runHeartbeat(cfg, replySpy, {}, { telegram: sendTelegram });
expect(result).toEqual({ status: "skipped", reason: "preempted" });
expect(preempt).toHaveBeenCalledOnce();
expect(sendTelegram).not.toHaveBeenCalled();
});
});
it("does not infer admission rejection from a replacement run after an empty heartbeat", async () => {
await withTempHeartbeatSandbox(async ({ storePath }) => {
const cfg = createHeartbeatTelegramConfig(storePath);
@@ -176,7 +176,7 @@ describe("runHeartbeatOnce heartbeat response tool", () => {
return call;
}
function setAgentTurnStatus(options: object | undefined, status: "ok" | "failed") {
function setAgentTurnStatus(options: object | undefined, status: "ok" | "failed" | "superseded") {
const runState = resolveReplyOperationRunState(options);
if (!runState) {
throw new Error("Expected heartbeat reply operation run state");
@@ -332,6 +332,40 @@ describe("runHeartbeatOnce heartbeat response tool", () => {
});
});
it("retains heartbeat work and scratch when a visible turn supersedes the run", async () => {
await withTempTelegramHeartbeatSandbox(async ({ tmpDir, storePath, replySpy }) => {
const cfg = createConfig({ tmpDir, storePath });
const jobId = await seedHeartbeatScratchForTest({ content: "old scratch" });
const sessionKey = await seedTelegramSession(storePath, cfg);
enqueueSystemEvent("exec finished: backup completed", { sessionKey });
const inspectedEvents = peekSystemEventEntries(sessionKey);
replySpy.mockImplementationOnce(async (_ctx, options) => {
setAgentTurnStatus(options, "superseded");
return createHeartbeatToolResponsePayload({
outcome: "progress",
notify: true,
summary: "Backup completed.",
notificationText: "Backup completed successfully.",
scratch: "new private scratch",
});
});
const sendTelegram = vi.fn().mockResolvedValue({ messageId: "m1" });
const result = await runHeartbeat(cfg, replySpy, sendTelegram, {
source: "exec-event",
intent: "event",
reason: "exec-event",
});
expect(result).toEqual({ status: "skipped", reason: "preempted" });
expect(sendTelegram).not.toHaveBeenCalled();
expect(peekSystemEventEntries(sessionKey)).toEqual(inspectedEvents);
expect(readCronJobScratchState(resolveCronJobsStorePath(), jobId).scratch?.content).toBe(
"old scratch",
);
});
});
it("does not recreate scratch when its monitor is deleted while the heartbeat runs", async () => {
await withTempTelegramHeartbeatSandbox(async ({ tmpDir, storePath, replySpy }) => {
const cfg = createConfig({ tmpDir, storePath });
+136
View File
@@ -0,0 +1,136 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { resetGatewayWorkAdmission } from "../process/gateway-work-admission.js";
import {
HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT,
requestHeartbeat,
setHeartbeatWakeHandler as setRuntimeHeartbeatWakeHandler,
} from "./heartbeat-wake.js";
describe("heartbeat wake preemption retry", () => {
type HeartbeatWakeHandler = Parameters<typeof setRuntimeHeartbeatWakeHandler>[0];
type WakeRequest = Parameters<typeof requestHeartbeat>[0];
let disposeHandler: (() => void) | undefined;
function setHeartbeatWakeHandler(handler: HeartbeatWakeHandler) {
disposeHandler = setRuntimeHeartbeatWakeHandler(handler);
}
function wake(reason: "interval" | "manual" | "exec-event", opts: Partial<WakeRequest> = {}) {
const source =
reason === "interval" ? "interval" : reason === "manual" ? "manual" : "exec-event";
const intent = reason === "interval" ? "scheduled" : reason === "manual" ? "manual" : "event";
return { source, intent, reason, ...opts } satisfies WakeRequest;
}
beforeEach(() => {
resetGatewayWorkAdmission();
vi.useFakeTimers();
});
afterEach(async () => {
resetGatewayWorkAdmission();
disposeHandler?.();
const disposeDrain = setRuntimeHeartbeatWakeHandler(async () => ({
status: "skipped",
reason: "disabled",
}));
await vi.runAllTimersAsync();
disposeDrain();
disposeHandler = undefined;
vi.useRealTimers();
vi.restoreAllMocks();
});
it("gives scheduled requests-in-flight a 60-second idle grace", async () => {
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT })
.mockResolvedValueOnce({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat(wake("interval", { coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(59_999);
expect(handler).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(1);
expect(handler).toHaveBeenCalledTimes(2);
expect(handler.mock.calls[1]?.[0]).toEqual({ ...wake("interval"), retainedWork: true });
});
it("keeps manual requests-in-flight on the default retry delay", async () => {
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT })
.mockResolvedValueOnce({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat(wake("manual", { coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(999);
expect(handler).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(1);
expect(handler).toHaveBeenCalledTimes(2);
});
it("retries preempted task work after idle grace without losing its payload", async () => {
const tasks = [{ jobId: "job-backup", name: "backup", prompt: "Check backup" }];
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: "preempted" })
.mockResolvedValueOnce({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat({
source: "background-task",
intent: "task",
reason: "heartbeat-task:job-backup",
agentId: "main",
sessionKey: "agent:main:main",
tasks,
coalesceMs: 0,
});
await vi.advanceTimersByTimeAsync(59_999);
expect(handler).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(1);
expect(handler).toHaveBeenCalledTimes(2);
expect(handler.mock.calls[1]?.[0]).toMatchObject({ tasks, retainedWork: true });
});
it("lets a fresh manual wake bypass a scheduled idle grace", async () => {
const target = { agentId: "main", sessionKey: "agent:main:main" };
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT })
.mockResolvedValue({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat(wake("interval", { ...target, coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(1);
requestHeartbeat(wake("manual", { ...target, coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(1);
expect(handler.mock.calls[1]?.[0]).toEqual(wake("manual", target));
await vi.advanceTimersByTimeAsync(59_998);
expect(handler.mock.calls[2]?.[0]).toEqual({
...wake("interval", target),
retainedWork: true,
});
});
it("keeps guarded event work retained through preemption", async () => {
const handler = vi
.fn()
.mockResolvedValueOnce({
status: "skipped",
reason: "not-due",
retryAtMs: Date.now() + 30_000,
})
.mockResolvedValueOnce({ status: "skipped", reason: "preempted" })
.mockResolvedValueOnce({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat(wake("exec-event", { coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(30_000);
expect(handler.mock.calls[1]?.[0]).toMatchObject({ retainedWork: true });
await vi.advanceTimersByTimeAsync(60_000);
expect(handler.mock.calls[2]?.[0]).toMatchObject({ retainedWork: true });
});
});
+12 -21
View File
@@ -359,11 +359,11 @@ describe("heartbeat-wake", () => {
requestHeartbeat({ ...request, coalesceMs: 0 });
await vi.advanceTimersByTimeAsync(1);
await vi.advanceTimersByTimeAsync(1_000);
await vi.advanceTimersByTimeAsync(60_000);
expect(handler).toHaveBeenCalledTimes(2);
expect(handler).toHaveBeenNthCalledWith(1, request);
expect(handler).toHaveBeenNthCalledWith(2, request);
expect(handler).toHaveBeenNthCalledWith(2, { ...request, retainedWork: true });
});
it("runs equal-period tasks at staggered anchors by retaining the spaced task", async () => {
@@ -639,19 +639,6 @@ describe("heartbeat-wake", () => {
});
});
it("retries requests-in-flight after the default retry delay", async () => {
vi.useFakeTimers();
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT })
.mockResolvedValueOnce({ status: "ran", durationMs: 1 });
await expectRetryAfterDefaultDelay({
handler,
initialReason: "interval",
expectedRetryReason: "interval",
});
});
it.each([HEARTBEAT_SKIP_CRON_IN_PROGRESS, HEARTBEAT_SKIP_LANES_BUSY])(
"retries %s after the default retry delay",
async (reason) => {
@@ -668,22 +655,26 @@ describe("heartbeat-wake", () => {
},
);
it("keeps retry cooldown even when a sooner request arrives", async () => {
it("lets a fresh event run while a scheduled retry observes idle grace", async () => {
vi.useFakeTimers();
const handler = setRetryOnceHeartbeatHandler();
const handler = vi
.fn()
.mockResolvedValueOnce({ status: "skipped", reason: HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT })
.mockResolvedValue({ status: "ran", durationMs: 1 });
setHeartbeatWakeHandler(handler);
requestHeartbeat(wake("interval", { coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(1);
expect(handler).toHaveBeenCalledTimes(1);
// Retry is now waiting for 1000ms. This should not preempt cooldown.
requestHeartbeat(wake("hook:wake", { coalesceMs: 0 }));
await vi.advanceTimersByTimeAsync(998);
expect(handler).toHaveBeenCalledTimes(1);
await vi.advanceTimersByTimeAsync(1);
expect(handler).toHaveBeenCalledTimes(2);
expectWakeCall(handler, 1, wake("hook:wake"));
await vi.advanceTimersByTimeAsync(59_998);
expect(handler).toHaveBeenCalledTimes(3);
expect(handler.mock.calls[2]?.[0]).toEqual({ ...wake("interval"), retainedWork: true });
});
it("retries thrown handler errors after the default retry delay", async () => {
+59 -26
View File
@@ -37,15 +37,17 @@ export const HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT = "requests-in-flight";
export const HEARTBEAT_SKIP_CRON_IN_PROGRESS = "cron-in-progress";
export const HEARTBEAT_SKIP_LANES_BUSY = "lanes-busy";
export const HEARTBEAT_SKIP_NO_PENDING_EVENT = "no-pending-event";
const RETRYABLE_BUSY_SKIP_REASONS = new Set([
export const HEARTBEAT_SKIP_PREEMPTED = "preempted";
const RETRYABLE_HEARTBEAT_SKIP_REASONS = new Set([
HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT,
HEARTBEAT_SKIP_CRON_IN_PROGRESS,
HEARTBEAT_SKIP_LANES_BUSY,
HEARTBEAT_SKIP_PREEMPTED,
]);
const RETRYABLE_GUARD_SKIP_REASONS = new Set(["not-due", "min-spacing", "flood"]);
export function isRetryableHeartbeatBusySkipReason(reason: string): boolean {
return RETRYABLE_BUSY_SKIP_REASONS.has(reason);
export function isRetryableHeartbeatSkipReason(reason: string): boolean {
return RETRYABLE_HEARTBEAT_SKIP_REASONS.has(reason);
}
let heartbeatsEnabled = true;
@@ -78,8 +80,8 @@ type PendingWakeReason = {
tasks?: HeartbeatScheduledTask[];
/** Earliest instant at which this retained wake class may be dispatched. */
notBeforeMs?: number;
/** The wake was retained after a spacing/cooldown guard deferred its work. */
guardRetry?: boolean;
/** The wake was deferred with real work that must survive later retries. */
retainedWork?: boolean;
};
type PendingWakeGroup = {
@@ -107,6 +109,7 @@ let wakeEnqueueSequence = 0;
const DEFAULT_COALESCE_MS = 250;
const DEFAULT_RETRY_MS = 1_000;
export const HEARTBEAT_IDLE_RETRY_GRACE_MS = 60_000;
// Heartbeat turns can start model/provider work; bound cross-target fan-out so
// one aligned monitor tick cannot exhaust gateway or provider capacity.
const MAX_CONCURRENT_HEARTBEAT_WAKE_TARGETS = 4;
@@ -167,12 +170,12 @@ function mergePendingWakeReasons(
? next
: previous;
const other = preferred === previous ? next : previous;
// Explicit wakes bypass a retained spacing guard, but busy backoff remains
// Explicit wakes bypass deferred background work, but busy backoff remains
// target-owned in PendingWakeGroup.blockedUntilMs.
const bypassGuardRetry =
const bypassRetainedWork =
(preferred.intent === "manual" || preferred.intent === "immediate") &&
preferred.guardRetry !== true &&
(previous.guardRetry === true || next.guardRetry === true);
preferred.retainedWork !== true &&
(previous.retainedWork === true || next.retainedWork === true);
const scheduledEveryMs = preferred.scheduledEveryMs ?? other.scheduledEveryMs;
const scheduledAnchorMs = preferred.scheduledAnchorMs ?? other.scheduledAnchorMs;
const immediateBarrierSequences = [
@@ -187,7 +190,8 @@ function mergePendingWakeReasons(
...preferred,
enqueueSequence: Math.min(previous.enqueueSequence, next.enqueueSequence),
readyAtMs,
...(!bypassGuardRetry && (previous.notBeforeMs !== undefined || next.notBeforeMs !== undefined)
...(!bypassRetainedWork &&
(previous.notBeforeMs !== undefined || next.notBeforeMs !== undefined)
? {
requestedAt: Math.min(previous.requestedAt, next.requestedAt),
notBeforeMs: Math.max(previous.notBeforeMs ?? 0, next.notBeforeMs ?? 0),
@@ -200,10 +204,10 @@ function mergePendingWakeReasons(
...(scheduledAnchorMs !== undefined ? { scheduledAnchorMs } : {}),
...(mergedTasks.length ? { tasks: mergedTasks } : {}),
};
if (!bypassGuardRetry && (previous.guardRetry || next.guardRetry)) {
merged.guardRetry = true;
if (!bypassRetainedWork && (previous.retainedWork || next.retainedWork)) {
merged.retainedWork = true;
} else {
delete merged.guardRetry;
delete merged.retainedWork;
}
if (immediateBarrierSequences.length > 0) {
merged.immediateBarrierSequence = Math.min(...immediateBarrierSequences);
@@ -315,8 +319,8 @@ function takePendingWakeBatch(maxGroups: number, now = Date.now()): ReadyWakeGro
// prevents a periodic task stream from starving an older event forever.
wakes.push(
...[taskWake, group.event].toSorted((left, right) => {
if (left.guardRetry !== right.guardRetry) {
return left.guardRetry ? -1 : 1;
if (left.retainedWork !== right.retainedWork) {
return left.retainedWork ? -1 : 1;
}
if (left.requestedAt !== right.requestedAt) {
return left.requestedAt - right.requestedAt;
@@ -355,7 +359,7 @@ function queuePendingWakeReason(params: {
tasks?: readonly HeartbeatScheduledTask[];
notBeforeMs?: number;
blockTargetUntilMs?: number;
guardRetry?: boolean;
retainedWork?: boolean;
}) {
const requestedAt = params.requestedAt ?? Date.now();
const enqueueSequence = params.enqueueSequence ?? ++wakeEnqueueSequence;
@@ -391,7 +395,7 @@ function queuePendingWakeReason(params: {
scheduledAnchorMs: params.scheduledAnchorMs,
...(params.tasks?.length ? { tasks: [...params.tasks] } : {}),
...(params.notBeforeMs === undefined ? {} : { notBeforeMs: params.notBeforeMs }),
...(params.guardRetry ? { guardRetry: true } : {}),
...(params.retainedWork ? { retainedWork: true } : {}),
};
const group = pendingWakes.get(wakeTargetKey) ?? {};
if (params.blockTargetUntilMs !== undefined) {
@@ -409,9 +413,36 @@ function queuePendingWakeReason(params: {
pendingWakes.set(wakeTargetKey, group);
}
function retryPendingWake(pendingWake: PendingWakeReason) {
function resolveHeartbeatRetrySchedule(
pendingWake: PendingWakeReason,
result: Extract<HeartbeatRunResult, { status: "skipped" }>,
): { delayMs: number; deferWakeOnly: boolean } {
const now = Date.now();
const deferWakeOnly =
result.reason === HEARTBEAT_SKIP_PREEMPTED ||
(result.reason === HEARTBEAT_SKIP_REQUESTS_IN_FLIGHT &&
(pendingWake.intent === "scheduled" || pendingWake.intent === "task"));
return {
delayMs:
result.retryAtMs !== undefined
? Math.max(0, result.retryAtMs - now)
: deferWakeOnly
? HEARTBEAT_IDLE_RETRY_GRACE_MS
: DEFAULT_RETRY_MS,
deferWakeOnly,
};
}
function retryPendingWake(
pendingWake: PendingWakeReason,
retrySchedule: { delayMs: number; deferWakeOnly: boolean } = {
delayMs: DEFAULT_RETRY_MS,
deferWakeOnly: false,
},
) {
// A thrown or busy wake owns only its target; replaying the whole batch
// duplicates completed reminders and stalls unrelated agents.
const retryAtMs = Date.now() + retrySchedule.delayMs;
queuePendingWakeReason({
source: pendingWake.source,
intent: pendingWake.intent,
@@ -425,9 +456,11 @@ function retryPendingWake(pendingWake: PendingWakeReason) {
requestedAt: pendingWake.requestedAt,
enqueueSequence: pendingWake.enqueueSequence,
immediateBarrierSequence: pendingWake.immediateBarrierSequence,
blockTargetUntilMs: Date.now() + DEFAULT_RETRY_MS,
...(retrySchedule.deferWakeOnly
? { notBeforeMs: retryAtMs, retainedWork: true }
: { blockTargetUntilMs: retryAtMs, retainedWork: pendingWake.retainedWork }),
});
schedule(DEFAULT_RETRY_MS);
schedule(retrySchedule.delayMs);
}
function handOffPendingWakeBatch(pendingBatch: PendingWakeReason[], startIndex: number) {
@@ -469,7 +502,7 @@ async function dispatchPendingWakeGroup(params: {
? { scheduledAnchorMs: pendingWake.scheduledAnchorMs }
: {}),
...(pendingWake.tasks ? { tasks: pendingWake.tasks } : {}),
...(pendingWake.guardRetry ? { retainedWork: true } : {}),
...(pendingWake.retainedWork ? { retainedWork: true } : {}),
};
let result: HeartbeatRunResult;
try {
@@ -488,7 +521,7 @@ async function dispatchPendingWakeGroup(params: {
if (handlerGeneration !== generation) {
const retainWake =
result.status === "skipped" &&
(isRetryableHeartbeatBusySkipReason(result.reason) ||
(isRetryableHeartbeatSkipReason(result.reason) ||
(RETRYABLE_GUARD_SKIP_REASONS.has(result.reason) &&
(pendingWake.tasks?.length ||
pendingWake.intent === "task" ||
@@ -497,8 +530,8 @@ async function dispatchPendingWakeGroup(params: {
handOffPendingWakeBatch(wakes, wakeIndex + (retainWake ? 0 : 1));
return;
}
if (result.status === "skipped" && isRetryableHeartbeatBusySkipReason(result.reason)) {
retryPendingWake(pendingWake);
if (result.status === "skipped" && isRetryableHeartbeatSkipReason(result.reason)) {
retryPendingWake(pendingWake, resolveHeartbeatRetrySchedule(pendingWake, result));
} else if (
result.status === "skipped" &&
RETRYABLE_GUARD_SKIP_REASONS.has(result.reason) &&
@@ -523,7 +556,7 @@ async function dispatchPendingWakeGroup(params: {
enqueueSequence: pendingWake.enqueueSequence,
immediateBarrierSequence: pendingWake.immediateBarrierSequence,
notBeforeMs: retryAtMs,
guardRetry: true,
retainedWork: true,
});
schedule(retryAtMs - Date.now());
}
@@ -668,7 +701,7 @@ function clearPendingWakeRetryState() {
continue;
}
delete pending.notBeforeMs;
delete pending.guardRetry;
delete pending.retainedWork;
}
}
}