From 7fcddedaaf63aaf904a1744fdd19abc4190b3fb2 Mon Sep 17 00:00:00 2001 From: Ayaan Zaidi Date: Fri, 14 Aug 2026 17:51:05 +0530 Subject: [PATCH] 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 --- src/agents/embedded-agent-runner.ts | 1 + src/agents/embedded-agent-runner/run-state.ts | 3 + .../run/attempt-stream-prepare.test.ts | 26 +++ .../run/attempt-stream-prepare.ts | 5 + .../runs.steering.test.ts | 62 +++++++ src/agents/embedded-agent-runner/runs.ts | 53 ++++-- src/agents/embedded-agent.runtime.ts | 1 + src/agents/embedded-agent.ts | 1 + src/auto-reply/reply/agent-runner-core.ts | 4 + ...agent-runner-direct-runtime-config.test.ts | 2 + src/auto-reply/reply/agent-runner-execute.ts | 12 +- .../reply/agent-runner-memory.test.ts | 2 + src/auto-reply/reply/agent-runner-run.ts | 7 +- .../agent-runner.misc.runreplyagent.test.ts | 160 ++++++++++-------- .../dispatch-from-config.base.test-utils.ts | 64 +++++++ .../reply/dispatch-from-config.lifecycle.ts | 3 +- .../reply/get-reply-run-admission.ts | 24 ++- .../reply/get-reply-run.media-only.test.ts | 153 ++++++++++++++++- src/auto-reply/reply/reply-operation-abort.ts | 22 ++- .../reply/reply-operation-agent-turn-state.ts | 20 ++- .../reply/reply-run-registry.contracts.ts | 9 +- .../reply/reply-run-registry.operation.ts | 41 ++++- .../reply/reply-run-registry.test.ts | 20 +++ src/auto-reply/reply/reply-run-registry.ts | 1 + src/auto-reply/reply/reply-turn-admission.ts | 10 +- src/auto-reply/reply/test-helpers.ts | 2 + src/cron/service/timer-execution.ts | 18 +- src/cron/service/timer.regression.test.ts | 61 ++++++- src/infra/heartbeat-runner-execution.ts | 6 +- src/infra/heartbeat-runner-run.ts | 14 +- src/infra/heartbeat-runner-scheduler.ts | 24 +-- src/infra/heartbeat-runner.scheduler.test.ts | 45 ++--- ...eat-runner.skips-busy-session-lane.test.ts | 52 ++++++ .../heartbeat-runner.tool-response.test.ts | 36 +++- src/infra/heartbeat-wake.preemption.test.ts | 136 +++++++++++++++ src/infra/heartbeat-wake.test.ts | 33 ++-- src/infra/heartbeat-wake.ts | 85 +++++++--- 37 files changed, 1013 insertions(+), 205 deletions(-) create mode 100644 src/infra/heartbeat-wake.preemption.test.ts diff --git a/src/agents/embedded-agent-runner.ts b/src/agents/embedded-agent-runner.ts index 025046b1641e..6c753625fc12 100644 --- a/src/agents/embedded-agent-runner.ts +++ b/src/agents/embedded-agent-runner.ts @@ -7,6 +7,7 @@ export { runEmbeddedAgent } from "./embedded-agent-runner/run.js"; export { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, + preemptAndDrainEmbeddedHeartbeatRun, isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunActive, isEmbeddedAgentRunHandleActive, diff --git a/src/agents/embedded-agent-runner/run-state.ts b/src/agents/embedded-agent-runner/run-state.ts index 547dbb5a490e..4ff072929b72 100644 --- a/src/agents/embedded-agent-runner/run-state.ts +++ b/src/agents/embedded-agent-runner/run-state.ts @@ -38,6 +38,8 @@ export type EmbeddedAgentQueueHandle = { ) => Promise; /** Cancels this run's pending user-input request before an image is queued as a later turn. */ cancelPendingUserInput?: (resolvedBy: string) => Promise; + /** 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; }; diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-prepare.test.ts b/src/agents/embedded-agent-runner/run/attempt-stream-prepare.test.ts index bf0964b1ac12..f777483458d6 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream-prepare.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream-prepare.test.ts @@ -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( diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-prepare.ts b/src/agents/embedded-agent-runner/run/attempt-stream-prepare.ts index 59a2dca5347b..1e6b31de7e3b 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream-prepare.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream-prepare.ts @@ -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: () => diff --git a/src/agents/embedded-agent-runner/runs.steering.test.ts b/src/agents/embedded-agent-runner/runs.steering.test.ts index 57ab68a85cfa..7808b299e7e8 100644 --- a/src/agents/embedded-agent-runner/runs.steering.test.ts +++ b/src/agents/embedded-agent-runner/runs.steering.test.ts @@ -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", { diff --git a/src/agents/embedded-agent-runner/runs.ts b/src/agents/embedded-agent-runner/runs.ts index 485b5184b5bb..8b8c2ccb18aa 100644 --- a/src/agents/embedded-agent-runner/runs.ts +++ b/src/agents/embedded-agent-runner/runs.ts @@ -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 { + 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 { - 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(); diff --git a/src/agents/embedded-agent.runtime.ts b/src/agents/embedded-agent.runtime.ts index 1ce6abf8b4bf..f5558a907f92 100644 --- a/src/agents/embedded-agent.runtime.ts +++ b/src/agents/embedded-agent.runtime.ts @@ -7,6 +7,7 @@ export { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, + preemptAndDrainEmbeddedHeartbeatRun, isEmbeddedAgentRunActive, isEmbeddedAgentRunStreaming, resolveActiveEmbeddedRunSessionId, diff --git a/src/agents/embedded-agent.ts b/src/agents/embedded-agent.ts index 778f7f8bf7df..356fce3f38ac 100644 --- a/src/agents/embedded-agent.ts +++ b/src/agents/embedded-agent.ts @@ -9,6 +9,7 @@ export type { export { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, + preemptAndDrainEmbeddedHeartbeatRun, compactEmbeddedAgentSession, isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunActive, diff --git a/src/auto-reply/reply/agent-runner-core.ts b/src/auto-reply/reply/agent-runner-core.ts index 1475d6b9e1aa..2aa366a5b692 100644 --- a/src/auto-reply/reply/agent-runner-core.ts +++ b/src/auto-reply/reply/agent-runner-core.ts @@ -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" diff --git a/src/auto-reply/reply/agent-runner-direct-runtime-config.test.ts b/src/auto-reply/reply/agent-runner-direct-runtime-config.test.ts index 9b6d172ef88a..ef47490ea5de 100644 --- a/src/auto-reply/reply/agent-runner-direct-runtime-config.test.ts +++ b/src/auto-reply/reply/agent-runner-direct-runtime-config.test.ts @@ -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(), diff --git a/src/auto-reply/reply/agent-runner-execute.ts b/src/auto-reply/reply/agent-runner-execute.ts index a51f9c301d2b..92028a491e09 100644 --- a/src/auto-reply/reply/agent-runner-execute.ts +++ b/src/auto-reply/reply/agent-runner-execute.ts @@ -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")); diff --git a/src/auto-reply/reply/agent-runner-memory.test.ts b/src/auto-reply/reply/agent-runner-memory.test.ts index a68a5f55f5da..e7a4f744892b 100644 --- a/src/auto-reply/reply/agent-runner-memory.test.ts +++ b/src/auto-reply/reply/agent-runner-memory.test.ts @@ -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(), diff --git a/src/auto-reply/reply/agent-runner-run.ts b/src/auto-reply/reply/agent-runner-run.ts index a1ba3491ca8f..212fbd443af7 100644 --- a/src/auto-reply/reply/agent-runner-run.ts +++ b/src/auto-reply/reply/agent-runner-run.ts @@ -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, diff --git a/src/auto-reply/reply/agent-runner.misc.runreplyagent.test.ts b/src/auto-reply/reply/agent-runner.misc.runreplyagent.test.ts index 78b242c7c0fb..86ab4cf7505d 100644 --- a/src/auto-reply/reply/agent-runner.misc.runreplyagent.test.ts +++ b/src/auto-reply/reply/agent-runner.misc.runreplyagent.test.ts @@ -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((resolve) => { - markCandidateSettled = resolve; - }); - const fallbackRelease = new Promise((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((resolve) => { + markCandidateSettled = resolve; + }); + const fallbackRelease = new Promise((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({ diff --git a/src/auto-reply/reply/dispatch-from-config.base.test-utils.ts b/src/auto-reply/reply/dispatch-from-config.base.test-utils.ts index 1ae4847874a8..4d0b5396cdbf 100644 --- a/src/auto-reply/reply/dispatch-from-config.base.test-utils.ts +++ b/src/auto-reply/reply/dispatch-from-config.base.test-utils.ts @@ -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(); diff --git a/src/auto-reply/reply/dispatch-from-config.lifecycle.ts b/src/auto-reply/reply/dispatch-from-config.lifecycle.ts index 993fa31baeea..516ed843175c 100644 --- a/src/auto-reply/reply/dispatch-from-config.lifecycle.ts +++ b/src/auto-reply/reply/dispatch-from-config.lifecycle.ts @@ -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. diff --git a/src/auto-reply/reply/get-reply-run-admission.ts b/src/auto-reply/reply/get-reply-run-admission.ts index 1ee0e5d14135..7a2839644d37 100644 --- a/src/auto-reply/reply/get-reply-run-admission.ts +++ b/src/auto-reply/reply/get-reply-run-admission.ts @@ -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" || diff --git a/src/auto-reply/reply/get-reply-run.media-only.test.ts b/src/auto-reply/reply/get-reply-run.media-only.test.ts index ff95a1e1abf8..e2136c26117f 100644 --- a/src/auto-reply/reply/get-reply-run.media-only.test.ts +++ b/src/auto-reply/reply/get-reply-run.media-only.test.ts @@ -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((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); diff --git a/src/auto-reply/reply/reply-operation-abort.ts b/src/auto-reply/reply/reply-operation-abort.ts index 724a424388e4..227537d4422a 100644 --- a/src/auto-reply/reply/reply-operation-abort.ts +++ b/src/auto-reply/reply/reply-operation-abort.ts @@ -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 { diff --git a/src/auto-reply/reply/reply-operation-agent-turn-state.ts b/src/auto-reply/reply/reply-operation-agent-turn-state.ts index bca35bc359e3..3fafec8ac058 100644 --- a/src/auto-reply/reply/reply-operation-agent-turn-state.ts +++ b/src/auto-reply/reply/reply-operation-agent-turn-state.ts @@ -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(); +type ReplyOperationAgentTurn = { + status: ReplyOperationAgentTurnStatus; + owner?: ReplyOperation; +}; + +const agentTurns = new WeakMap(); 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; } diff --git a/src/auto-reply/reply/reply-run-registry.contracts.ts b/src/auto-reply/reply/reply-run-registry.contracts.ts index fd820625cf4c..04f8f2684a65 100644 --- a/src/auto-reply/reply/reply-run-registry.contracts.ts +++ b/src/auto-reply/reply/reply-run-registry.contracts.ts @@ -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, cause?: unknown): void; abortByUser(): boolean; abortForRestart(): boolean; + supersede(): boolean; }; export type ReplyRunRegistry = { diff --git a/src/auto-reply/reply/reply-run-registry.operation.ts b/src/auto-reply/reply/reply-run-registry.operation.ts index 02e168ccb940..96be79d17596 100644 --- a/src/auto-reply/reply/reply-run-registry.operation.ts +++ b/src/auto-reply/reply/reply-run-registry.operation.ts @@ -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; -type ReplyOperationStaleReason = replyRunSettle.ReplyOperationStaleReason; +type ReplyOperationAbortCode = Extract["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 | 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; diff --git a/src/auto-reply/reply/reply-run-registry.test.ts b/src/auto-reply/reply/reply-run-registry.test.ts index eaf69bf60569..e80bf0b38829 100644 --- a/src/auto-reply/reply/reply-run-registry.test.ts +++ b/src/auto-reply/reply/reply-run-registry.test.ts @@ -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); diff --git a/src/auto-reply/reply/reply-run-registry.ts b/src/auto-reply/reply/reply-run-registry.ts index f0f21cc548f6..951c7c34f5b2 100644 --- a/src/auto-reply/reply/reply-run-registry.ts +++ b/src/auto-reply/reply/reply-run-registry.ts @@ -16,6 +16,7 @@ export type { ReplyOperation, ReplyOperationPhase, ReplyToolAuthorityOverlay, + ReplyTurnKind, } from "./reply-run-registry.contracts.js"; export { abortReplyMessageInjectionTarget, diff --git a/src/auto-reply/reply/reply-turn-admission.ts b/src/auto-reply/reply/reply-turn-admission.ts index ba5fd4de1d33..86d9015217b5 100644 --- a/src/auto-reply/reply/reply-turn-admission.ts +++ b/src/auto-reply/reply/reply-turn-admission.ts @@ -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; } diff --git a/src/auto-reply/reply/test-helpers.ts b/src/auto-reply/reply/test-helpers.ts index 1081c53a25ec..5f36ec65ba5c 100644 --- a/src/auto-reply/reply/test-helpers.ts +++ b/src/auto-reply/reply/test-helpers.ts @@ -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, diff --git a/src/cron/service/timer-execution.ts b/src/cron/service/timer-execution.ts index c602a1954395..f750853f2502 100644 --- a/src/cron/service/timer-execution.ts +++ b/src/cron/service/timer-execution.ts @@ -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") { diff --git a/src/cron/service/timer.regression.test.ts b/src/cron/service/timer.regression.test.ts index cfd292000689..7218749eedaf 100644 --- a/src/cron/service/timer.regression.test.ts +++ b/src/cron/service/timer.regression.test.ts @@ -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>(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 { diff --git a/src/infra/heartbeat-runner-execution.ts b/src/infra/heartbeat-runner-execution.ts index e611077943f1..81c210631a56 100644 --- a/src/infra/heartbeat-runner-execution.ts +++ b/src/infra/heartbeat-runner-execution.ts @@ -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 && diff --git a/src/infra/heartbeat-runner-run.ts b/src/infra/heartbeat-runner-run.ts index dea656c30438..afbd6f91ea56 100644 --- a/src/infra/heartbeat-runner-run.ts +++ b/src/infra/heartbeat-runner-run.ts @@ -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 { 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(); }); }); diff --git a/src/infra/heartbeat-runner.skips-busy-session-lane.test.ts b/src/infra/heartbeat-runner.skips-busy-session-lane.test.ts index 3dbdc0f2be1c..b571cc0789e7 100644 --- a/src/infra/heartbeat-runner.skips-busy-session-lane.test.ts +++ b/src/infra/heartbeat-runner.skips-busy-session-lane.test.ts @@ -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 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); diff --git a/src/infra/heartbeat-runner.tool-response.test.ts b/src/infra/heartbeat-runner.tool-response.test.ts index 3612f037f580..b530532c6db4 100644 --- a/src/infra/heartbeat-runner.tool-response.test.ts +++ b/src/infra/heartbeat-runner.tool-response.test.ts @@ -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 }); diff --git a/src/infra/heartbeat-wake.preemption.test.ts b/src/infra/heartbeat-wake.preemption.test.ts new file mode 100644 index 000000000000..1c5dd04c40ec --- /dev/null +++ b/src/infra/heartbeat-wake.preemption.test.ts @@ -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[0]; + type WakeRequest = Parameters[0]; + let disposeHandler: (() => void) | undefined; + + function setHeartbeatWakeHandler(handler: HeartbeatWakeHandler) { + disposeHandler = setRuntimeHeartbeatWakeHandler(handler); + } + + function wake(reason: "interval" | "manual" | "exec-event", opts: Partial = {}) { + 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 }); + }); +}); diff --git a/src/infra/heartbeat-wake.test.ts b/src/infra/heartbeat-wake.test.ts index 485f3114ee7f..3c22f511497b 100644 --- a/src/infra/heartbeat-wake.test.ts +++ b/src/infra/heartbeat-wake.test.ts @@ -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 () => { diff --git a/src/infra/heartbeat-wake.ts b/src/infra/heartbeat-wake.ts index 231a90b16ad7..4a553469a287 100644 --- a/src/infra/heartbeat-wake.ts +++ b/src/infra/heartbeat-wake.ts @@ -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, +): { 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; } } }