mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix(channels): keep deeper queued messages alive through long turns (#133248)
## Contributor description — original implementation
The contributor’s original description follows. Historical heads, scope and open decisions in this section refer to that implementation; the reviewed current implementation, compatibility decision and proof are recorded in the maintainer update below.
Closes #133245.
Related: #127950. This supersedes the current-main-incompatible queue-owned approach in #125385 while preserving today's canonical timeout retry disposition.
## What Problem This Solves
A durably claimed channel message can wait behind another follow-up turn for longer than the five-minute claim-to-adoption watchdog. Existing liveness renewal covers queue-head active-run admission checks, but not an ingress-backed lifecycle waiting deeper in the follow-up queue. The claim can therefore be retired before the message reaches the model.
## Why This Change Was Made
The follow-up queue now owns periodic deferred heartbeats for the exact ingress lifecycle while it remains in pending items, in-flight delivery, summary sources, or compacted summary elisions.
Ingress supplies a cadence derived from one third of its adoption-stall timeout. Debounce, Plugin SDK fan-in, and channel lifecycle wrappers preserve the shortest applicable cadence. Queue ownership triggers an immediate renewal, then periodic renewal; it stops after successful adoption, completion, ownership loss, or callback failure. A rejected adoption callback keeps renewal alive for the supported retry path.
The existing watchdog and canonical timeout retry policy remain unchanged, so silent handlers and orphaned claims still recover normally.
## User Impact
Messages accepted by durable channel ingress remain adoptable while they wait behind long-running work instead of silently disappearing before model execution. No setting or migration is required.
## Evidence
SDK documentation polish (2026-09-10, head `76097fef9397a04c10637e06dfaf64ebb6ca104a`): the public channel SDK guide now documents forwarding both heartbeat fields, the shortest-positive-finite fan-in cadence, renewal termination, and compatibility when older wrappers omit the optional cadence. The documentation-only addition passed changed-page formatting, MDX sanity, and `git diff --check`; production code and the previously tested runtime head below are unchanged. Current-head hosted CI has completed successfully; the final exact-head ClawSweeper review accepts the implementation and proof, leaving only SDK-owner acceptance of the documented contract.
Refreshed on 2026-09-08 for exact head `57c4faff0ed7bb2d9b4fd308aa46b19ae4095e61`, rebased onto `1ff45cad5a`.
- Preserved upstream's ingress-monitor type extraction and gateway-suspension repair; the optional cadence field follows its new type owner.
- Repaired the retry integration fixture: after real enqueue and initial renewal, only the abandoned message's heartbeat callback fails. The real watchdog then recovers it. Retry delivery, healthy sibling cleanup, and true duplicate assertions remain intact; no production behavior or timeout changed for this repair.
- 102 focused tests pass: dispatch ingress retry, queue in-flight/dedupe, ingress lifecycle/watchdog, Plugin SDK fan-in, Discord queue handling, and Slack handler.
- Changed-file gates pass: core/core-test/extension typechecks, formatting, changed core/extension lint, SDK boundaries/exports, dead-export scans, and repository guards.
- Fresh independent source review found no actionable P0–P2 defects. The standalone autoreview CLI failed at startup without a verdict and is not counted as review coverage.
- Current-head CI run [34520397963](https://github.com/openclaw/openclaw/actions/runs/34520397963) completed successfully. The earlier startup subprocess timeout is historical, not a current failing gate. Exact-head ClawSweeper review (September 10, 20:00 UTC) reports no actionable correctness or proof findings; explicit SDK-owner acceptance remains required before merge.
- Existing regressions cover immediate renewal after late handoff, periodic deeper-queue renewal, stopping after owner loss, and continuing across a rejected adoption callback.
Exact-head durable-ingress boundary proof used isolated SQLite with the production ingress drain/lifecycle binder, follow-up queue, and canonical adoption helper. The model callback was synthetic; this is not live Telegram or Discord transport proof. No production Gateway, channel state, or configuration was touched. Temporary state and queue ownership were cleaned up.
```json
{
"schema": "openclaw.pr133248.durable-ingress-proof.v1",
"exactHead": "57c4faff0ed7bb2d9b4fd308aa46b19ae4095e61",
"setup": {
"durableStore": "isolated OpenClaw SQLite state",
"executor": "production follow-up queue with synthetic model callback",
"transportLifecycle": "production ingress drain lifecycle",
"adoptionStallTimeoutMs": 180
},
"queuedBehindLongTurn": {
"heldMs": 620,
"formerDeadlineCrossed": true,
"heartbeatCount": 11,
"claimStillOwnedAtCheckpoint": [
"queued-event"
],
"retryRowsAtCheckpoint": 0,
"failedRowsAtCheckpoint": 0,
"executionCount": 1,
"duplicateAfterAdoption": "completed"
},
"orphanRecovery": {
"status": "released-for-retry",
"attempts": 1,
"lastErrorContainsHandlerTimeout": true
},
"productionTouched": false
}
```
Owner acceptance:
- Intended behavior: accepted messages remain adoptable while their exact lifecycle is owned by the queue; ownerless claims retain canonical timeout recovery.
- Boundary: one optional cadence field propagated through existing shared lifecycle and channel wrappers; no configuration, schema, or migration change.
- Maintainer decision remains open: accept the additive public Plugin SDK lifecycle field and its queue-owned cadence semantics. Bot review is not that acceptance.
- Rollback: revert the renewal fix and its companion retry-fixture adjustment.
- Scope: 23 files, +273/-2. The width is required lifecycle forwarding; renewal policy stays in the queue and ingress drain.
AI-assisted.
---
## Maintainer update — reviewed head `9325ad501aa4`
## What Problem This Solves
Messages accepted while a long reply is running can expire in the follow-up queue and consume a retry. The reproduction confirmed expiry and retry; all three messages eventually arrived. It did not reproduce the older permanent-loss report in #133245.
## Fix and impact
The queue starts one heartbeat when it accepts a lifecycle and stops it on adoption, completion, cancellation, or callback failure. Heartbeats no longer scan queue collections or copy the in-flight set. Ingress remains responsible for the watchdog and retry settlement, and a heartbeat cannot undo the watchdog pause during adoption finalization.
Ingress derives the cadence from its adoption timeout. The optional lifecycle metadata remains necessary because wrappers rebuild callbacks and combine abort signals, losing the original timeout. No channel setting, schema, migration, dependency, or protocol change is added.
The simplification removes 110 lines net from the initial implementation plus main merge. Total production growth is 40 lines. Existing lifecycle types are reused, and the queue cases share the existing lifecycle test fixtures.
## Evidence
- Real Gateway and Discord: three marked messages on one route, with a queued reply held for 330 seconds. Main expired the third claim after 300.004 seconds. This head retained the same claim at 322.257 seconds with zero attempts; replies arrived once, in order, at 17.804, 349.201, and 350.223 seconds.
- Real SQLite/drain/binder/queue controls: explicit abandonment releases the claim; callback failure stops renewal and lets the watchdog retry without model execution.
- 191 core owner/sibling tests and 190 channel tests passed. All five new regression cases failed on plain main for the intended reasons.
- The local preflight passed core/extension type-aware lint, production and test types, script types, and protocol checks in an isolated checkout of this head.
- External consumers compiled and ran against published 2026.9.4 and this head, including omitted metadata, forwarding, optional return fields, and asynchronous abandonment.
## Consumers
Shared ingress binding, batching, reply dispatch, and the Feishu, Slack, Telegram, and Twitch wrappers preserve the source cadence. The fan-in uses the shortest valid cadence. Legacy wrappers that omit the optional field retain their existing head-only heartbeat behavior. Adoption, queue clearing, overflow, cancellation, summaries, and drain replacement retain or finish the same lifecycle owner.
Closes #133245.
Original implementation and report by @PollyBot13 (#133245). The contributor's commits and authorship are retained.
AI-assisted.
Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
parent
15f63bf749
commit
dba4da1f8b
22 changed files with 208 additions and 47 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -183,6 +183,7 @@ export function createFeishuBroadcastIngressSettlement(params: {
|
|||
defer();
|
||||
},
|
||||
onDeferredHeartbeat: () => params.lifecycle?.onDeferredHeartbeat?.(),
|
||||
deferredHeartbeatIntervalMs: params.lifecycle?.deferredHeartbeatIntervalMs,
|
||||
onAdoptionFinalizing: beginFinalizing,
|
||||
onAbandoned: async () => {
|
||||
if (
|
||||
|
|
|
|||
|
|
@ -317,6 +317,7 @@ export function buildFeishuFlushIngressLifecycle(
|
|||
transportLifecycle.onDeferred();
|
||||
},
|
||||
onDeferredHeartbeat: () => transportLifecycle.onDeferredHeartbeat?.(),
|
||||
deferredHeartbeatIntervalMs: transportLifecycle.deferredHeartbeatIntervalMs,
|
||||
onAdoptionFinalizing: () => {
|
||||
transportLifecycle.onAdoptionFinalizing();
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<void>;
|
||||
onDeferred?: () => void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
onAbandoned?: () => void;
|
||||
abortSignal?: AbortSignal;
|
||||
};
|
||||
turnAdoptionLifecycle?: GetReplyOptions["turnAdoptionLifecycle"];
|
||||
};
|
||||
|
||||
export type TelegramDispatchResult =
|
||||
|
|
|
|||
|
|
@ -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<void>;
|
||||
onDeferred?: () => void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
onAbandoned?: () => void;
|
||||
abortSignal?: AbortSignal;
|
||||
};
|
||||
turnAdoptionLifecycle?: GetReplyOptions["turnAdoptionLifecycle"];
|
||||
}): Promise<TelegramMessageProcessingResult> => {
|
||||
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");
|
||||
|
|
|
|||
|
|
@ -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<void>;
|
||||
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<void>;
|
||||
};
|
||||
|
||||
type TelegramSpooledReplayFrame = {
|
||||
|
|
|
|||
|
|
@ -139,6 +139,7 @@ export function createTwitchIngress(options: {
|
|||
lifecycle.onDeferred();
|
||||
},
|
||||
onDeferredHeartbeat: () => lifecycle.onDeferredHeartbeat?.(),
|
||||
deferredHeartbeatIntervalMs: lifecycle.deferredHeartbeatIntervalMs,
|
||||
onAbandoned: async () => {
|
||||
handedOff = true;
|
||||
await lifecycle.onAbandoned();
|
||||
|
|
|
|||
|
|
@ -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. */
|
||||
|
|
|
|||
|
|
@ -47,6 +47,7 @@ type InboundDebounceAdmissionLifecycleInput = {
|
|||
onAdopted?: () => void | Promise<void>;
|
||||
onDeferred?: () => boolean | void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
deferredHeartbeatIntervalMs?: number;
|
||||
onAdoptionFinalizing?: () => void;
|
||||
onFailed?: (error: unknown) => void | Promise<void>;
|
||||
onAbandoned?: () => void | Promise<void>;
|
||||
|
|
@ -58,6 +59,7 @@ type InboundDebounceAdmissionLifecycle = {
|
|||
onAdopted: () => Promise<void>;
|
||||
onDeferred: () => boolean | void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
deferredHeartbeatIntervalMs?: number;
|
||||
onAdoptionFinalizing: () => void;
|
||||
onFailed?: (error: unknown) => Promise<void>;
|
||||
onAbandoned: () => Promise<void>;
|
||||
|
|
@ -98,6 +100,7 @@ function createInboundDebounceFlush(params: {
|
|||
return accepted;
|
||||
},
|
||||
onDeferredHeartbeat: () => source?.onDeferredHeartbeat?.(),
|
||||
deferredHeartbeatIntervalMs: source?.deferredHeartbeatIntervalMs,
|
||||
onAdoptionFinalizing: () => source?.onAdoptionFinalizing?.(),
|
||||
onFailed: source?.onFailed
|
||||
? async (error) => {
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<void>>()
|
||||
|
|
|
|||
|
|
@ -307,9 +307,43 @@ const admittingTurnAdoptionLifecycles = new WeakMap<TurnAdoptionLifecycle, Promi
|
|||
const retiredTurnAdoptionCancellationLifecycles = new WeakSet<TurnAdoptionLifecycle>();
|
||||
const completedTurnAdoptionLifecycles = new WeakSet<TurnAdoptionLifecycle>();
|
||||
const completedTurnAdoptionLifecycleCallbacks = new WeakSet<TurnAdoptionLifecycle>();
|
||||
const deferredHeartbeatStops = new WeakMap<TurnAdoptionLifecycle, () => void>();
|
||||
|
||||
type FollowupLifecycleRun = Pick<FollowupRun, "steerPending" | "turnAdoptionLifecycle">;
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<void>;
|
||||
onDeferred: () => void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
onAbandoned: () => void | Promise<void>;
|
||||
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,
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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<Payload>({
|
||||
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?.();
|
||||
|
|
|
|||
|
|
@ -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<TBody> = { 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<void>;
|
||||
onDeferred: () => void;
|
||||
onDeferredHeartbeat?: () => void;
|
||||
onAdoptionFinalizing: () => void;
|
||||
onFailed?: (error: unknown) => void | Promise<void>;
|
||||
onCancelled?: () => void | Promise<void>;
|
||||
onAbandoned: () => void | Promise<void>;
|
||||
};
|
||||
|
||||
/** Optional explicit outcome from a channel delivery. */
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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?.(),
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue