diff --git a/docs/plugins/sdk-channel-outbound.md b/docs/plugins/sdk-channel-outbound.md index f7004eae8d27..327990f9ab60 100644 --- a/docs/plugins/sdk-channel-outbound.md +++ b/docs/plugins/sdk-channel-outbound.md @@ -114,6 +114,20 @@ classification, and persisted payload shape in the plugin. Webhook transports should acknowledge only after `admit` resolves; non-replay transports should surface durable append exhaustion rather than silently dispatching. +### Deferred claim heartbeats + +Forward both `onDeferredHeartbeat` and `deferredHeartbeatIntervalMs` when a +plugin wraps the ingress lifecycle or maps it to `turnAdoptionLifecycle`. +`bindIngressLifecycleToReplyOptions(...)` forwards both. The drain derives the +optional cadence from its adoption-stall timeout; fan-in uses the shortest +positive, finite source cadence. The queue renews only while it owns the +lifecycle, stopping after adoption, completion, ownership loss, or callback +failure. A heartbeat does not adopt or complete a claim. + +Wrappers that omit the cadence remain valid but do not enable periodic renewal; +their deferred claims can still reach the adoption watchdog timeout. Plugins +must not run independent timers that keep abandoned work alive. + ## Adapter Most plugins define one `message` adapter: diff --git a/extensions/feishu/src/bot-broadcast.ts b/extensions/feishu/src/bot-broadcast.ts index 4f2b81f48fe7..161ae2f52a15 100644 --- a/extensions/feishu/src/bot-broadcast.ts +++ b/extensions/feishu/src/bot-broadcast.ts @@ -183,6 +183,7 @@ export function createFeishuBroadcastIngressSettlement(params: { defer(); }, onDeferredHeartbeat: () => params.lifecycle?.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: params.lifecycle?.deferredHeartbeatIntervalMs, onAdoptionFinalizing: beginFinalizing, onAbandoned: async () => { if ( diff --git a/extensions/feishu/src/feishu-ingress.ts b/extensions/feishu/src/feishu-ingress.ts index 83d8aa878e40..b6b726c3ac28 100644 --- a/extensions/feishu/src/feishu-ingress.ts +++ b/extensions/feishu/src/feishu-ingress.ts @@ -317,6 +317,7 @@ export function buildFeishuFlushIngressLifecycle( transportLifecycle.onDeferred(); }, onDeferredHeartbeat: () => transportLifecycle.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: transportLifecycle.deferredHeartbeatIntervalMs, onAdoptionFinalizing: () => { transportLifecycle.onAdoptionFinalizing(); }, diff --git a/extensions/slack/src/monitor/message-handler.ts b/extensions/slack/src/monitor/message-handler.ts index 66069fa004d9..2e31864b0658 100644 --- a/extensions/slack/src/monitor/message-handler.ts +++ b/extensions/slack/src/monitor/message-handler.ts @@ -282,6 +282,13 @@ export function createSlackMessageHandler(params: { return; } await turnAdoptionLifecycle?.onSessionRouted?.(prepared.route.sessionKey); + const deferredHeartbeatIntervals = [ + turnAdoptionLifecycle?.deferredHeartbeatIntervalMs, + admissionLifecycle.deferredHeartbeatIntervalMs, + ].filter( + (interval): interval is number => + interval !== undefined && Number.isFinite(interval) && interval > 0, + ); // Commit at adoption (durable turn ownership), release on abandonment; // deferred turns hand settlement to the reply lane with the claim held. prepared.turnAdoptionLifecycle = { @@ -308,6 +315,9 @@ export function createSlackMessageHandler(params: { turnAdoptionLifecycle?.onDeferredHeartbeat?.(); admissionLifecycle.onDeferredHeartbeat?.(); }, + ...(deferredHeartbeatIntervals.length > 0 + ? { deferredHeartbeatIntervalMs: Math.min(...deferredHeartbeatIntervals) } + : {}), onAbandoned: () => { settlementHandedOff = true; releaseClaims(); diff --git a/extensions/telegram/src/bot-message-dispatch-turn.ts b/extensions/telegram/src/bot-message-dispatch-turn.ts index 0b14ff11e071..688b1ff6234f 100644 --- a/extensions/telegram/src/bot-message-dispatch-turn.ts +++ b/extensions/telegram/src/bot-message-dispatch-turn.ts @@ -156,12 +156,8 @@ export async function runTelegramDispatchTurn(turn: Turn) { abortSignal: turn.turnAdoptionLifecycle?.abortSignal, turnAdoptionLifecycle: turn.turnAdoptionLifecycle ? { + ...turn.turnAdoptionLifecycle, admission: turn.turnAdoptionLifecycle.admission ?? "exclusive", - onAdopted: turn.turnAdoptionLifecycle.onAdopted, - onDeferred: turn.turnAdoptionLifecycle.onDeferred, - onDeferredHeartbeat: turn.turnAdoptionLifecycle.onDeferredHeartbeat, - onAbandoned: turn.turnAdoptionLifecycle.onAbandoned, - abortSignal: turn.turnAdoptionLifecycle.abortSignal, } : undefined, sourceReplyDeliveryMode: isRoomEvent ? "message_tool_only" : undefined, diff --git a/extensions/telegram/src/bot-message-dispatch.types.ts b/extensions/telegram/src/bot-message-dispatch.types.ts index 53d3eace2a66..2fa989e6622e 100644 --- a/extensions/telegram/src/bot-message-dispatch.types.ts +++ b/extensions/telegram/src/bot-message-dispatch.types.ts @@ -11,6 +11,7 @@ import type { TelegramAccountConfig, } from "openclaw/plugin-sdk/config-contracts"; import type { ReplyPayload } from "openclaw/plugin-sdk/reply-payload"; +import type { GetReplyOptions } from "openclaw/plugin-sdk/reply-runtime"; import type { RuntimeEnv } from "openclaw/plugin-sdk/runtime-env"; import type { SessionEntry } from "openclaw/plugin-sdk/session-store-runtime"; import type { TelegramBotDeps } from "./bot-deps.js"; @@ -45,14 +46,7 @@ export type DispatchTelegramMessageParams = { * Canonical turn ownership lifecycle from the durable ingress drain * (or a test double). Pre-adoption abort + adopt/defer/abandon. */ - turnAdoptionLifecycle?: { - admission?: "exclusive" | "cancel-only"; - onAdopted: () => void | Promise; - onDeferred?: () => void; - onDeferredHeartbeat?: () => void; - onAbandoned?: () => void; - abortSignal?: AbortSignal; - }; + turnAdoptionLifecycle?: GetReplyOptions["turnAdoptionLifecycle"]; }; export type TelegramDispatchResult = diff --git a/extensions/telegram/src/bot-message.ts b/extensions/telegram/src/bot-message.ts index c49ac9cf61ac..7f0e38367730 100644 --- a/extensions/telegram/src/bot-message.ts +++ b/extensions/telegram/src/bot-message.ts @@ -2,6 +2,7 @@ import type { OpenClawConfig, TelegramAccountConfig } from "openclaw/plugin-sdk/config-contracts"; import { resolveTextChunkLimit } from "openclaw/plugin-sdk/reply-chunking"; import { DEFAULT_GROUP_HISTORY_LIMIT } from "openclaw/plugin-sdk/reply-history"; +import type { GetReplyOptions } from "openclaw/plugin-sdk/reply-runtime"; import { createSubsystemLogger, danger, @@ -272,14 +273,7 @@ export const createTelegramMessageProcessor = (deps: TelegramMessageProcessorDep await turnContext.onDispatchStart?.(); } const runTelegramDispatch = async (params: { - turnAdoptionLifecycle?: { - admission?: "exclusive" | "cancel-only"; - onAdopted: () => void | Promise; - onDeferred?: () => void; - onDeferredHeartbeat?: () => void; - onAbandoned?: () => void; - abortSignal?: AbortSignal; - }; + turnAdoptionLifecycle?: GetReplyOptions["turnAdoptionLifecycle"]; }): Promise => { try { const dispatchResult = await dispatchTelegramMessage({ @@ -438,6 +432,7 @@ export const createTelegramMessageProcessor = (deps: TelegramMessageProcessorDep drainLifecycle?.onDeferred(); }, onDeferredHeartbeat: () => drainLifecycle?.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: drainLifecycle?.deferredHeartbeatIntervalMs, onAbandoned: () => { if (!adopted) { void settle({ kind: "failed-retryable", error: "turn-abandoned" }, "terminal"); diff --git a/extensions/telegram/src/bot-processing-outcome.ts b/extensions/telegram/src/bot-processing-outcome.ts index d6dbdb206350..d7a90b29c913 100644 --- a/extensions/telegram/src/bot-processing-outcome.ts +++ b/extensions/telegram/src/bot-processing-outcome.ts @@ -1,5 +1,6 @@ // Telegram plugin module tracks per-update processing outcomes. import { AsyncLocalStorage } from "node:async_hooks"; +import type { ChannelIngressMonitorLifecycle } from "openclaw/plugin-sdk/channel-outbound"; export type TelegramMessageProcessingResult = | { kind: "completed" } @@ -10,14 +11,12 @@ type TelegramUpdateProcessingFrame = { result?: TelegramMessageProcessingResult; }; -type TelegramSpooledReplayLifecycle = { - abortSignal: AbortSignal; - onAdopted: () => void | Promise; - onDeferred: () => void; - onDeferredHeartbeat?: () => void; +type TelegramSpooledReplayLifecycle = Omit< + ChannelIngressMonitorLifecycle, + "admission" | "onFailed" | "onCancelled" | "onAdoptionFinalizing" +> & { /** Clears pre-adoption stall while durable adoption finalization is held. */ onAdoptionFinalizing?: () => void; - onAbandoned: () => void | Promise; }; type TelegramSpooledReplayFrame = { diff --git a/extensions/twitch/src/twitch-ingress.ts b/extensions/twitch/src/twitch-ingress.ts index c8451db1a4f7..a1cdce9121f2 100644 --- a/extensions/twitch/src/twitch-ingress.ts +++ b/extensions/twitch/src/twitch-ingress.ts @@ -139,6 +139,7 @@ export function createTwitchIngress(options: { lifecycle.onDeferred(); }, onDeferredHeartbeat: () => lifecycle.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: lifecycle.deferredHeartbeatIntervalMs, onAbandoned: async () => { handedOff = true; await lifecycle.onAbandoned(); diff --git a/src/auto-reply/get-reply-options.types.ts b/src/auto-reply/get-reply-options.types.ts index 09663139b399..9758831b0dfb 100644 --- a/src/auto-reply/get-reply-options.types.ts +++ b/src/auto-reply/get-reply-options.types.ts @@ -94,6 +94,8 @@ export type TurnAdoptionLifecycle = { onDeferred?: () => boolean | void; /** Reports that a deferred turn is still queued behind an active turn. */ onDeferredHeartbeat?: () => void; + /** Requested cadence for queue-owned deferred heartbeats. */ + deferredHeartbeatIntervalMs?: number; /** Deferred turn finished without owning the reply lane. */ onAbandoned?: () => void; /** Always fires when the followup ownership cycle ends (admitted or not). Gateway cleanup. */ diff --git a/src/auto-reply/inbound-debounce.ts b/src/auto-reply/inbound-debounce.ts index c08ec214bb05..f71a989237f0 100644 --- a/src/auto-reply/inbound-debounce.ts +++ b/src/auto-reply/inbound-debounce.ts @@ -47,6 +47,7 @@ type InboundDebounceAdmissionLifecycleInput = { onAdopted?: () => void | Promise; onDeferred?: () => boolean | void; onDeferredHeartbeat?: () => void; + deferredHeartbeatIntervalMs?: number; onAdoptionFinalizing?: () => void; onFailed?: (error: unknown) => void | Promise; onAbandoned?: () => void | Promise; @@ -58,6 +59,7 @@ type InboundDebounceAdmissionLifecycle = { onAdopted: () => Promise; onDeferred: () => boolean | void; onDeferredHeartbeat?: () => void; + deferredHeartbeatIntervalMs?: number; onAdoptionFinalizing: () => void; onFailed?: (error: unknown) => Promise; onAbandoned: () => Promise; @@ -98,6 +100,7 @@ function createInboundDebounceFlush(params: { return accepted; }, onDeferredHeartbeat: () => source?.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: source?.deferredHeartbeatIntervalMs, onAdoptionFinalizing: () => source?.onAdoptionFinalizing?.(), onFailed: source?.onFailed ? async (error) => { diff --git a/src/auto-reply/reply/dispatch-from-config.ingress-retry.test.ts b/src/auto-reply/reply/dispatch-from-config.ingress-retry.test.ts index 9787988a143c..2bc484591a20 100644 --- a/src/auto-reply/reply/dispatch-from-config.ingress-retry.test.ts +++ b/src/auto-reply/reply/dispatch-from-config.ingress-retry.test.ts @@ -87,6 +87,17 @@ describe("dispatch retry after queued ingress abandonment", () => { }); run.turnAdoptionLifecycle = options?.turnAdoptionLifecycle; run.abortSignal = options?.turnAdoptionLifecycle?.abortSignal; + let renewalFailed = false; + const lifecycle = run.turnAdoptionLifecycle; + const heartbeat = lifecycle?.onDeferredHeartbeat; + if (lifecycle) { + lifecycle.onDeferredHeartbeat = () => { + if (renewalFailed) { + throw new Error("deferred heartbeat owner failed"); + } + heartbeat?.(); + }; + } expect( enqueueFollowupRun( key, @@ -103,6 +114,7 @@ describe("dispatch retry after queued ingress abandonment", () => { ), ).toBe(true); if (abandonment === "watchdog-after-commit") { + renewalFailed = true; expect( enqueueFollowupRun( key, diff --git a/src/auto-reply/reply/queue.collect.test.ts b/src/auto-reply/reply/queue.collect.test.ts index 733820ba6c63..6069a1133107 100644 --- a/src/auto-reply/reply/queue.collect.test.ts +++ b/src/auto-reply/reply/queue.collect.test.ts @@ -211,6 +211,66 @@ describe("followup queue collect routing", () => { expect(cancelA).toEqual(cancelB); }); + it.each(["admission", "abandonment", "abort", "callback failure"] as const)( + "renews a deeper queued lifecycle until %s", + async (transition) => { + vi.useFakeTimers(); + const key = `test-deferred-heartbeat-${transition}`; + const abort = new AbortController(); + let lastHeartbeat = -Infinity; + let failHeartbeat = false; + const heartbeat = vi.fn(() => { + if (failHeartbeat) { + throw new Error("heartbeat unavailable"); + } + lastHeartbeat = Date.now(); + }); + const pending = createRun({ prompt: "deeper queued turn" }); + pending.turnAdoptionLifecycle = { + admission: "exclusive", + abortSignal: abort.signal, + onAdopted: async () => {}, + onDeferredHeartbeat: heartbeat, + deferredHeartbeatIntervalMs: 1_000, + }; + try { + const settings = createQueueSettings({ mode: "followup" }); + enqueueFollowupRun(key, createRun({ prompt: "earlier turn" }), settings); + enqueueFollowupRun(key, pending, settings); + await vi.advanceTimersByTimeAsync(3_000); + expect(Date.now() - lastHeartbeat).toBeLessThan(1_000); + + if (transition === "admission") { + await admitFollowupRunLifecycle(pending); + } else if (transition === "abandonment") { + clearFollowupQueue(key); + } else if (transition === "abort") { + abort.abort(); + } else { + failHeartbeat = true; + await vi.advanceTimersByTimeAsync(1_000); + } + const callsAtTransition = heartbeat.mock.calls.length; + await vi.advanceTimersByTimeAsync(3_000); + expect(heartbeat).toHaveBeenCalledTimes(callsAtTransition); + + if (transition === "admission" || transition === "callback failure") { + const delivered: string[] = []; + scheduleFollowupDrain(key, async (run) => { + await admitFollowupRunLifecycle(run); + delivered.push(run.prompt); + completeFollowupRunLifecycle(run); + }); + await vi.runAllTimersAsync(); + expect(delivered).toEqual(["earlier turn", "deeper queued turn"]); + } + } finally { + clearFollowupQueue(key); + vi.useRealTimers(); + } + }, + ); + it("retries lifecycle admission after a callback rejection", async () => { const onAdmitted = vi .fn<() => Promise>() diff --git a/src/auto-reply/reply/queue/types.ts b/src/auto-reply/reply/queue/types.ts index 6aa997c12cb8..dac135f3f557 100644 --- a/src/auto-reply/reply/queue/types.ts +++ b/src/auto-reply/reply/queue/types.ts @@ -307,9 +307,43 @@ const admittingTurnAdoptionLifecycles = new WeakMap(); const completedTurnAdoptionLifecycles = new WeakSet(); const completedTurnAdoptionLifecycleCallbacks = new WeakSet(); +const deferredHeartbeatStops = new WeakMap void>(); type FollowupLifecycleRun = Pick; +function startFollowupRunDeferredHeartbeat(lifecycle: TurnAdoptionLifecycle): void { + const intervalMs = lifecycle.deferredHeartbeatIntervalMs; + const heartbeat = lifecycle.onDeferredHeartbeat; + if ( + !heartbeat || + intervalMs === undefined || + !Number.isFinite(intervalMs) || + intervalMs <= 0 || + lifecycle.abortSignal?.aborted || + admittedTurnAdoptionLifecycles.has(lifecycle) || + completedTurnAdoptionLifecycles.has(lifecycle) + ) { + return; + } + const pulse = () => { + try { + heartbeat(); + } catch { + // Leave recovery to the ingress watchdog when its liveness callback fails. + deferredHeartbeatStops.get(lifecycle)?.(); + } + }; + const timer = setInterval(pulse, intervalMs).unref(); + const stop = () => { + clearInterval(timer); + lifecycle.abortSignal?.removeEventListener("abort", stop); + deferredHeartbeatStops.delete(lifecycle); + }; + deferredHeartbeatStops.set(lifecycle, stop); + lifecycle.abortSignal?.addEventListener("abort", stop, { once: true }); + pulse(); +} + export function markFollowupRunEnqueued(run: FollowupLifecycleRun): boolean { const lifecycle = run.turnAdoptionLifecycle; if (lifecycle && !enqueuedTurnAdoptionLifecycles.has(lifecycle)) { @@ -317,6 +351,7 @@ export function markFollowupRunEnqueued(run: FollowupLifecycleRun): boolean { return false; } enqueuedTurnAdoptionLifecycles.add(lifecycle); + startFollowupRunDeferredHeartbeat(lifecycle); } return true; } @@ -348,6 +383,7 @@ export async function admitFollowupRunLifecycle(run: FollowupLifecycleRun): Prom if (!admittedTurnAdoptionLifecycles.has(lifecycle)) { await lifecycle.onAdopted(); admittedTurnAdoptionLifecycles.add(lifecycle); + deferredHeartbeatStops.get(lifecycle)?.(); } }); @@ -383,6 +419,7 @@ export function completeFollowupRunLifecycle( }; if (lifecycle && !completedTurnAdoptionLifecycles.has(lifecycle)) { + deferredHeartbeatStops.get(lifecycle)?.(); completedTurnAdoptionLifecycles.add(lifecycle); } diff --git a/src/channels/message/ingress-drain-lifecycle.test.ts b/src/channels/message/ingress-drain-lifecycle.test.ts index 9cec5f32693d..e7c30757ca0a 100644 --- a/src/channels/message/ingress-drain-lifecycle.test.ts +++ b/src/channels/message/ingress-drain-lifecycle.test.ts @@ -22,6 +22,10 @@ describe("channel ingress drain lifecycle", () => { onDeferred: () => { calls.push("deferred"); }, + onDeferredHeartbeat: () => { + calls.push("heartbeat"); + }, + deferredHeartbeatIntervalMs: 1_234, onAbandoned: () => { calls.push("abandoned"); }, @@ -30,14 +34,16 @@ describe("channel ingress drain lifecycle", () => { expect(bound.turnAdoptionLifecycle).toMatchObject({ admission: "exclusive", abortSignal: abort.signal, + deferredHeartbeatIntervalMs: 1_234, }); expect("onFailed" in bound.turnAdoptionLifecycle).toBe(false); expect("onCancelled" in bound.turnAdoptionLifecycle).toBe(false); expect("onAdopted" in bound).toBe(false); expect(Object.keys(bound)).toEqual(["turnAdoptionLifecycle"]); bound.turnAdoptionLifecycle.onDeferred(); + bound.turnAdoptionLifecycle.onDeferredHeartbeat?.(); await bound.turnAdoptionLifecycle.onAbandoned(); - expect(calls).toEqual(["deferred", "abandoned"]); + expect(calls).toEqual(["deferred", "heartbeat", "abandoned"]); calls.length = 0; bound.turnAdoptionLifecycle.onDeferred(); await bound.turnAdoptionLifecycle.onAdopted(); diff --git a/src/channels/message/ingress-drain-lifecycle.ts b/src/channels/message/ingress-drain-lifecycle.ts index cba8e4cfa308..0c068e166b4f 100644 --- a/src/channels/message/ingress-drain-lifecycle.ts +++ b/src/channels/message/ingress-drain-lifecycle.ts @@ -14,6 +14,7 @@ export type ChannelIngressDispatchLifecycle = { onDeferred: () => void; /** Deferred reply-lane admission is still waiting behind an active turn. */ onDeferredHeartbeat?: () => void; + deferredHeartbeatIntervalMs?: number; /** * Durable adoption finalization is in progress (e.g. settlement hold while * committing dedupe). Clears the pre-adoption stall watchdog so a timeout @@ -34,14 +35,10 @@ export type ChannelIngressDispatchLifecycle = { /** Maps a drain lifecycle onto the reply-lane ownership surface. */ export function bindIngressLifecycleToReplyOptions(lifecycle: ChannelIngressDispatchLifecycle): { - turnAdoptionLifecycle: { - admission: "exclusive"; - onAdopted: () => void | Promise; - onDeferred: () => void; - onDeferredHeartbeat?: () => void; - onAbandoned: () => void | Promise; - abortSignal: AbortSignal; - }; + turnAdoptionLifecycle: Omit< + ChannelIngressDispatchLifecycle, + "onAdoptionFinalizing" | "onFailed" | "onCancelled" + > & { admission: "exclusive" }; } { return { turnAdoptionLifecycle: { @@ -49,6 +46,7 @@ export function bindIngressLifecycleToReplyOptions(lifecycle: ChannelIngressDisp onAdopted: lifecycle.onAdopted, onDeferred: lifecycle.onDeferred, onDeferredHeartbeat: lifecycle.onDeferredHeartbeat, + deferredHeartbeatIntervalMs: lifecycle.deferredHeartbeatIntervalMs, onAbandoned: lifecycle.onAbandoned, abortSignal: lifecycle.abortSignal, }, diff --git a/src/channels/message/ingress-drain.ts b/src/channels/message/ingress-drain.ts index a46100159866..760172768b30 100644 --- a/src/channels/message/ingress-drain.ts +++ b/src/channels/message/ingress-drain.ts @@ -365,11 +365,12 @@ export function createChannelIngressDrain< } }, onDeferredHeartbeat: () => { - // Abort also covers disposal; retired callbacks cannot restart the watchdog. - if (state.phase === "deferred" && !state.abortController.signal.aborted) { + // A cleared watchdog marks adoption finalization or retired ownership. + if (state.phase === "deferred" && state.stallTimer) { armStallWatchdog(state); } }, + deferredHeartbeatIntervalMs: Math.max(1, Math.floor(adoptionStallTimeoutMs / 3)), onAdoptionFinalizing: () => { if (state.phase !== "dispatching" && state.phase !== "deferred") { return; diff --git a/src/channels/message/ingress-drain.watchdog.test.ts b/src/channels/message/ingress-drain.watchdog.test.ts index 4af242f54d01..71747b676051 100644 --- a/src/channels/message/ingress-drain.watchdog.test.ts +++ b/src/channels/message/ingress-drain.watchdog.test.ts @@ -92,12 +92,34 @@ describe("channel ingress drain watchdog", () => { }); }); + it("keeps adoption finalization paused across deferred heartbeats", async () => { + await withTempState(async (stateDir) => { + const queue = createTestIngressQueue(stateDir); + await queue.enqueue("finalizing", { text: "x" }, { laneKey: "l1" }); + const { drain, lifecycle, heartbeat } = await deferNext(queue); + lifecycle.onAdoptionFinalizing(); + try { + await vi.advanceTimersByTimeAsync(333); + heartbeat(); + await vi.advanceTimersByTimeAsync(1_100); + expect(lifecycle.abortSignal.aborted).toBe(false); + expect(await queue.listClaims()).toMatchObject([{ id: "finalizing", attempts: 0 }]); + await lifecycle.onAdopted(); + expect(await queue.listClaims()).toEqual([]); + expect(await queue.listPending({ limit: "all" })).toEqual([]); + } finally { + drain.dispose(); + } + }); + }); + it("rearms a live deferred wait, then guillotines silence", async () => { await withTempState(async (stateDir) => { let clock = 30_000; const queue = createTestIngressQueue(stateDir, { now: () => clock }); await queue.enqueue("evt-def-stall", { text: "x" }, { laneKey: "l1" }); let heartbeat: (() => void) | undefined; + let heartbeatIntervalMs: number | undefined; const drain = createChannelIngressDrain({ queue, @@ -106,6 +128,7 @@ describe("channel ingress drain watchdog", () => { dispatchClaimedEvent: async (_event, lifecycle) => { lifecycle.onDeferred(); heartbeat = lifecycle.onDeferredHeartbeat; + heartbeatIntervalMs = lifecycle.deferredHeartbeatIntervalMs; // Stay deferred without adoption -- watchdog must still fire. await new Promise(() => {}); }, @@ -113,6 +136,7 @@ describe("channel ingress drain watchdog", () => { await drain.drainOnce(); expect(await queue.listClaims()).toHaveLength(1); + expect(heartbeatIntervalMs).toBe(1_666); clock += 4_000; await vi.advanceTimersByTimeAsync(4_000); heartbeat?.(); diff --git a/src/channels/message/ingress-monitor-types.ts b/src/channels/message/ingress-monitor-types.ts index d5824a7335c0..580b85d4eab8 100644 --- a/src/channels/message/ingress-monitor-types.ts +++ b/src/channels/message/ingress-monitor-types.ts @@ -1,3 +1,4 @@ +import type { ChannelIngressDispatchLifecycle } from "./ingress-drain-lifecycle.js"; import type { CreateChannelIngressDrainOptions } from "./ingress-drain.js"; import type { ChannelIngressQueue, ChannelIngressQueueClaim } from "./ingress-queue.js"; @@ -8,16 +9,8 @@ export type ChannelIngressMonitorFacts = { eventId: string; laneKey: string }; type ChannelIngressPayloadEnvelope = { version: number; body: TBody }; /** Claim ownership lifecycle handed to one channel delivery. */ -export type ChannelIngressMonitorLifecycle = { +export type ChannelIngressMonitorLifecycle = ChannelIngressDispatchLifecycle & { admission: "exclusive"; - abortSignal: AbortSignal; - onAdopted: () => void | Promise; - onDeferred: () => void; - onDeferredHeartbeat?: () => void; - onAdoptionFinalizing: () => void; - onFailed?: (error: unknown) => void | Promise; - onCancelled?: () => void | Promise; - onAbandoned: () => void | Promise; }; /** Optional explicit outcome from a channel delivery. */ diff --git a/src/plugin-sdk/channel-ingress-runtime.test.ts b/src/plugin-sdk/channel-ingress-runtime.test.ts index 69b73213de47..4623f0a65043 100644 --- a/src/plugin-sdk/channel-ingress-runtime.test.ts +++ b/src/plugin-sdk/channel-ingress-runtime.test.ts @@ -47,6 +47,7 @@ describe("plugin-sdk/channel-ingress-runtime", () => { abortSignal: new AbortController().signal, onAdopted: vi.fn(async () => {}), onDeferred: vi.fn(), + deferredHeartbeatIntervalMs: 3_000, onAdoptionFinalizing: vi.fn(), onFailed: vi.fn(async () => {}), onCancelled: vi.fn(async () => {}), @@ -54,10 +55,13 @@ describe("plugin-sdk/channel-ingress-runtime", () => { }); const first = createLifecycle(); const second = createLifecycle(); + second.deferredHeartbeatIntervalMs = 1_000; const cancellation = fanInChannelIngressLifecycles([first, second]); await cancellation.lifecycle?.onCancelled?.(); const combined = fanInChannelIngressLifecycles([undefined, first, second]); + expect(combined.lifecycle?.deferredHeartbeatIntervalMs).toBe(1_000); + combined.lifecycle?.onAdoptionFinalizing(); await combined.lifecycle?.onAdopted(); await combined.settle(); diff --git a/src/plugin-sdk/channel-ingress-runtime.ts b/src/plugin-sdk/channel-ingress-runtime.ts index 1bfc3993ac35..19dccdabaf13 100644 --- a/src/plugin-sdk/channel-ingress-runtime.ts +++ b/src/plugin-sdk/channel-ingress-runtime.ts @@ -193,6 +193,12 @@ export function fanInChannelIngressLifecycles( } }; const supportsCancellation = lifecycles.every((lifecycle) => lifecycle.onCancelled !== undefined); + const deferredHeartbeatIntervals = lifecycles + .map((lifecycle) => lifecycle.deferredHeartbeatIntervalMs) + .filter( + (interval): interval is number => + interval !== undefined && Number.isFinite(interval) && interval > 0, + ); // Omit aggregate cancellation unless every durable source supports it. Callers // can then use settle/abandon without an acknowledged-but-unsettled claim. const cancelAll = () => @@ -222,6 +228,9 @@ export function fanInChannelIngressLifecycles( lifecycle.onDeferredHeartbeat?.(); } }, + ...(deferredHeartbeatIntervals.length > 0 + ? { deferredHeartbeatIntervalMs: Math.min(...deferredHeartbeatIntervals) } + : {}), onAdoptionFinalizing: () => { for (const lifecycle of lifecycles) { lifecycle.onAdoptionFinalizing(); diff --git a/src/plugin-sdk/test-helpers/plugin-runtime-mock.ts b/src/plugin-sdk/test-helpers/plugin-runtime-mock.ts index 16d507f76743..6ce2f565c0bc 100644 --- a/src/plugin-sdk/test-helpers/plugin-runtime-mock.ts +++ b/src/plugin-sdk/test-helpers/plugin-runtime-mock.ts @@ -33,6 +33,7 @@ export const createTestInboundDebounceFlush: InboundDebounceFlushFactory = (para onAdopted: async () => await source?.onAdopted?.(), onDeferred: () => source?.onDeferred?.(), onDeferredHeartbeat: () => source?.onDeferredHeartbeat?.(), + deferredHeartbeatIntervalMs: source?.deferredHeartbeatIntervalMs, onAdoptionFinalizing: () => source?.onAdoptionFinalizing?.(), onFailed: source?.onFailed ? async (error) => await source.onFailed?.(error) : undefined, onAbandoned: async () => await source?.onAbandoned?.(),