From 2aab6b8e3726356b226159cbd2e2fa64f0784d8d Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 16 Aug 2026 10:03:29 -0700 Subject: [PATCH] refactor(reply): unify keyed FIFO leases (#124690) Share one lifecycle-owned reservation primitive between reply admission and foreground delivery ordering, and remove the obsolete admission-wait callback plumbing. --- ...est.ts => dispatch.delivery-order.test.ts} | 262 +++++------------- src/auto-reply/dispatch.ts | 101 ++----- .../foreground-reply-fence-state.ts | 26 -- src/auto-reply/reply/agent-runner-run.ts | 1 - .../reply/dispatch-from-config.lifecycle.ts | 2 - ...ispatch-from-config.stale-recovery.test.ts | 10 +- .../reply/followup-turn-admission.ts | 1 - src/auto-reply/reply/get-reply-run-execute.ts | 1 - src/auto-reply/reply/get-reply.types.ts | 2 - .../reply/queue/drain.identity-guard.test.ts | 44 --- src/auto-reply/reply/queue/drain.ts | 16 -- src/auto-reply/reply/queue/types.ts | 2 - .../reply/reply-admission-ticket.ts | 57 +--- .../reply/reply-turn-admission.test.ts | 14 - src/auto-reply/reply/reply-turn-admission.ts | 53 +--- .../reply/stranded-reply-recovery.test.ts | 2 - .../turn/run-channel-turn.pipeline.test.ts | 37 --- src/shared/keyed-fifo-lease.test.ts | 152 ++++++++++ src/shared/keyed-fifo-lease.ts | 87 ++++++ 19 files changed, 346 insertions(+), 524 deletions(-) rename src/auto-reply/{dispatch.freshness.test.ts => dispatch.delivery-order.test.ts} (67%) delete mode 100644 src/auto-reply/foreground-reply-fence-state.ts create mode 100644 src/shared/keyed-fifo-lease.test.ts create mode 100644 src/shared/keyed-fifo-lease.ts diff --git a/src/auto-reply/dispatch.freshness.test.ts b/src/auto-reply/dispatch.delivery-order.test.ts similarity index 67% rename from src/auto-reply/dispatch.freshness.test.ts rename to src/auto-reply/dispatch.delivery-order.test.ts index 977b91e9b71f..2f4c34908b72 100644 --- a/src/auto-reply/dispatch.freshness.test.ts +++ b/src/auto-reply/dispatch.delivery-order.test.ts @@ -155,12 +155,7 @@ describe("foreground reply delivery order", () => { } if (params.ctx.MessageSid === "waiting-message") { waitingSuccessorStarted.resolve(); - params.replyOptions?.onReplyAdmissionWaitChange?.(true); - try { - await releaseWaitingSuccessor.promise; - } finally { - params.replyOptions?.onReplyAdmissionWaitChange?.(false); - } + await releaseWaitingSuccessor.promise; return { queuedFinal: false, counts: { tool: 0, block: 0, final: 0 }, @@ -211,140 +206,65 @@ describe("foreground reply delivery order", () => { ]); }); - it("releases a WhatsApp-shaped lane after beforeDeliver times out", async () => { - vi.useFakeTimers(); - try { - const deliveries: Delivery[] = []; - const hookStarted = createDeferred(); - const onSettled = vi.fn(); - let hookCalls = 0; - const beforeDeliver = vi.fn((payload: ReplyPayload) => { - hookCalls += 1; - if (hookCalls === 1) { - hookStarted.resolve(); - return new Promise(() => {}); - } - return payload; - }); - hoisted.dispatchReplyFromConfigMock.mockImplementation( - async (params: DispatchReplyFromConfigParams) => { - params.dispatcher.sendFinalReply({ text: "stuck final" }); - params.dispatcher.sendFinalReply({ text: "follow-up final" }); - return { - queuedFinal: true, - counts: { tool: 0, block: 0, final: 2 }, - }; - }, - ); - - const dispatch = dispatchWithDeliveries(buildForegroundCtx(), deliveries, { - beforeDeliver, - onSettled, - }); - await hookStarted.promise; - await vi.advanceTimersByTimeAsync(15_000); - - await expect(dispatch).resolves.toEqual({ - queuedFinal: true, - counts: { tool: 0, block: 0, final: 1 }, - failedCounts: { tool: 0, block: 0, final: 1 }, - }); - expect(beforeDeliver).toHaveBeenCalledTimes(2); - expect(deliveries).toEqual([{ kind: "final", text: "follow-up final" }]); - expect(onSettled).toHaveBeenCalledOnce(); - expect(vi.getTimerCount()).toBe(0); - } finally { - vi.useRealTimers(); - } - }); - - it("honors a configured beforeDeliver budget inside the foreground fence", async () => { + it("does not charge predecessor waiting against the configured beforeDeliver budget", async () => { vi.useFakeTimers(); try { const deliveries: Delivery[] = []; + const olderHookStarted = createDeferred(); + const releaseOlderHook = createDeferred(); const hookStarted = createDeferred(); hoisted.dispatchReplyFromConfigMock.mockImplementation( async (params: DispatchReplyFromConfigParams) => { - params.dispatcher.sendFinalReply({ text: "budgeted final" }); + params.dispatcher.sendFinalReply({ text: `${params.ctx.MessageSid} final` }); return queuedFinalResult(); }, ); - const dispatch = dispatchWithDeliveries(buildForegroundCtx(), deliveries, { - beforeDeliver: async (payload) => { - hookStarted.resolve(); - await new Promise((resolve) => { - setTimeout(resolve, 16_000); - }); - return payload; + const olderDispatch = dispatchWithDeliveries( + buildForegroundCtx({ MessageSid: "older" }), + deliveries, + { + beforeDeliver: () => { + olderHookStarted.resolve(); + return releaseOlderHook.promise; + }, + beforeDeliverOptions: { timeoutMs: 40_000 }, }, - beforeDeliverOptions: { timeoutMs: 20_000 }, - }); - await hookStarted.promise; - await vi.advanceTimersByTimeAsync(15_000); + ); + await olderHookStarted.promise; + const newerDispatch = dispatchWithDeliveries( + buildForegroundCtx({ MessageSid: "newer" }), + deliveries, + { + beforeDeliver: async (payload) => { + hookStarted.resolve(); + await new Promise((resolve) => { + setTimeout(resolve, 16_000); + }); + return payload; + }, + beforeDeliverOptions: { timeoutMs: 20_000 }, + }, + ); + + await vi.advanceTimersByTimeAsync(20_000); expect(deliveries).toEqual([]); - await vi.advanceTimersByTimeAsync(1_000); + releaseOlderHook.resolve({ text: "older final" }); + await expect(olderDispatch).resolves.toEqual(queuedFinalResult()); + await hookStarted.promise; + await vi.advanceTimersByTimeAsync(16_000); - await expect(dispatch).resolves.toEqual(queuedFinalResult()); - expect(deliveries).toEqual([{ kind: "final", text: "budgeted final" }]); + await expect(newerDispatch).resolves.toEqual(queuedFinalResult()); + expect(deliveries).toEqual([ + { kind: "final", text: "older final" }, + { kind: "final", text: "newer final" }, + ]); expect(vi.getTimerCount()).toBe(0); } finally { vi.useRealTimers(); } }); - it("keeps an older foreground final when a newer inbound has no visible delivery while beforeDeliver is pending", async () => { - const deliveries: Delivery[] = []; - const beforeDeliverStarted = createDeferred(); - const releaseBeforeDeliver = createDeferred(); - const beforeDeliver = vi.fn(() => { - beforeDeliverStarted.resolve(); - return releaseBeforeDeliver.promise; - }); - - hoisted.dispatchReplyFromConfigMock.mockImplementation( - async (params: DispatchReplyFromConfigParams) => { - if (params.ctx.MessageSid === "old-message") { - params.dispatcher.sendFinalReply({ text: "old final" }); - return queuedFinalResult(); - } - if (params.ctx.MessageSid === "new-message") { - return { - queuedFinal: false, - counts: { tool: 0, block: 0, final: 0 }, - }; - } - throw new Error(`unexpected test message ${params.ctx.MessageSid ?? ""}`); - }, - ); - - const olderDispatch = dispatchWithDeliveries( - buildForegroundCtx({ MessageSid: "old-message" }), - deliveries, - { beforeDeliver }, - ); - await beforeDeliverStarted.promise; - - const newerResult = await dispatchWithDeliveries( - buildForegroundCtx({ MessageSid: "new-message" }), - deliveries, - ); - - releaseBeforeDeliver.resolve({ text: "old rewritten final" }); - const olderResult = await olderDispatch; - - expect(beforeDeliver).toHaveBeenCalledTimes(1); - expect(newerResult).toEqual({ - queuedFinal: false, - counts: { tool: 0, block: 0, final: 0 }, - }); - expect(olderResult).toEqual({ - queuedFinal: true, - counts: { tool: 0, block: 0, final: 1 }, - }); - expect(deliveries).toEqual([{ kind: "final", text: "old rewritten final" }]); - }); - it("does not fence an older final behind a newer inbound waiting for its delivery", async () => { const deliveries: Delivery[] = []; const olderStarted = createDeferred(); @@ -363,12 +283,7 @@ describe("foreground reply delivery order", () => { if (params.ctx.MessageSid === "new-message") { newerStarted.resolve(); // Same-session follow-up admission waits for the owning final delivery. - params.replyOptions?.onReplyAdmissionWaitChange?.(true); - try { - await olderDelivered.promise; - } finally { - params.replyOptions?.onReplyAdmissionWaitChange?.(false); - } + await olderDelivered.promise; return { queuedFinal: false, counts: { tool: 0, block: 0, final: 0 }, @@ -405,58 +320,6 @@ describe("foreground reply delivery order", () => { expect(deliveries).toEqual([{ kind: "final", text: "old final" }]); }); - it("does not make an older final wait for a newer independent turn", async () => { - const deliveries: Delivery[] = []; - const olderBeforeDeliverStarted = createDeferred(); - const releaseOlderBeforeDeliver = createDeferred(); - const newerStarted = createDeferred(); - const releaseNewerFinal = createDeferred(); - - hoisted.dispatchReplyFromConfigMock.mockImplementation( - async (params: DispatchReplyFromConfigParams) => { - if (params.ctx.MessageSid === "old-message") { - params.dispatcher.sendFinalReply({ text: "old final" }); - return queuedFinalResult(); - } - if (params.ctx.MessageSid === "new-message") { - newerStarted.resolve(); - await releaseNewerFinal.promise; - params.dispatcher.sendFinalReply({ text: "new final" }); - return queuedFinalResult(); - } - throw new Error(`unexpected test message ${params.ctx.MessageSid ?? ""}`); - }, - ); - - const olderDispatch = dispatchWithDeliveries( - buildForegroundCtx({ MessageSid: "old-message" }), - deliveries, - { - beforeDeliver: () => { - olderBeforeDeliverStarted.resolve(); - return releaseOlderBeforeDeliver.promise; - }, - }, - ); - await olderBeforeDeliverStarted.promise; - - const newerDispatch = dispatchWithDeliveries( - buildForegroundCtx({ MessageSid: "new-message" }), - deliveries, - ); - await newerStarted.promise; - releaseOlderBeforeDeliver.resolve({ text: "old final" }); - await expect(olderDispatch).resolves.toEqual(queuedFinalResult()); - expect(deliveries).toEqual([{ kind: "final", text: "old final" }]); - - releaseNewerFinal.resolve(); - await expect(newerDispatch).resolves.toEqual(queuedFinalResult()); - expect(deliveries).toEqual([ - { kind: "final", text: "old final" }, - { kind: "final", text: "new final" }, - ]); - }); - it.each(["onSettled", "onFreshSettledDelivery"] as const)( "orders %s delivery behind an earlier foreground final", async (settledHook) => { @@ -515,29 +378,42 @@ describe("foreground reply delivery order", () => { }, ); - it("runs the settled delivery hook when dispatch fails after queueing a reply", async () => { + it("releases a same-target successor when an earlier dispatch fails", async () => { const deliveries: Delivery[] = []; - let settled = false; + const olderStarted = createDeferred(); + const newerStarted = createDeferred(); + const releaseOlderFailure = createDeferred(); const error = new Error("resolver failed"); hoisted.dispatchReplyFromConfigMock.mockImplementation( async (params: DispatchReplyFromConfigParams) => { - params.dispatcher.sendFinalReply({ text: "queued final" }); - throw error; + if (params.ctx.MessageSid === "older") { + olderStarted.resolve(); + await releaseOlderFailure.promise; + throw error; + } + newerStarted.resolve(); + params.dispatcher.sendFinalReply({ text: "newer final" }); + return queuedFinalResult(); }, ); - await expect( - dispatchWithDeliveries(buildForegroundCtx(), deliveries, { - deliver: async () => ({ visibleReplySent: false }), - onSettled: () => { - settled = true; - return { visibleReplySent: true }; - }, - }), - ).rejects.toBe(error); + const olderDispatch = dispatchWithDeliveries( + buildForegroundCtx({ MessageSid: "older" }), + deliveries, + ); + await olderStarted.promise; + const newerDispatch = dispatchWithDeliveries( + buildForegroundCtx({ MessageSid: "newer" }), + deliveries, + ); + await newerStarted.promise; + expect(deliveries).toEqual([]); - expect(settled).toBe(true); + releaseOlderFailure.resolve(); + await expect(olderDispatch).rejects.toBe(error); + await expect(newerDispatch).resolves.toEqual(queuedFinalResult()); + expect(deliveries).toEqual([{ kind: "final", text: "newer final" }]); }); it("keeps concurrent foreground finals isolated for different targets sharing a session", async () => { diff --git a/src/auto-reply/dispatch.ts b/src/auto-reply/dispatch.ts index 3ebcf1f8380a..eb97abf06c3a 100644 --- a/src/auto-reply/dispatch.ts +++ b/src/auto-reply/dispatch.ts @@ -14,17 +14,13 @@ import { type ReplyPayloadSuppressedObserver, } from "../infra/outbound/deliver-hooks.js"; import { logMessageReceived } from "../logging/diagnostic.js"; +import { createKeyedFifoLeaseRegistry, type KeyedFifoLease } from "../shared/keyed-fifo-lease.js"; import type { SilentReplyConversationType } from "../shared/silent-reply-policy.js"; import { resolveCommandTurnContext, resolveCommandTurnTargetSessionKey, } from "./command-turn-context.js"; import { withReplyDispatcher } from "./dispatch-dispatcher.js"; -import { - foregroundReplyFenceByKey, - type ForegroundReplyFenceState, - notifyForegroundReplyFenceWaiters, -} from "./foreground-reply-fence-state.js"; import type { CommandSessionMetadataChange } from "./reply/command-session-metadata.js"; import { dispatchReplyFromConfig } from "./reply/dispatch-from-config.js"; import type { DispatchFromConfigResult } from "./reply/dispatch-from-config.types.js"; @@ -47,16 +43,14 @@ import type { FinalizedMsgContext, MsgContext } from "./templating.js"; type InternalDispatchReplyOptions = Omit; -type ForegroundReplyFenceSnapshot = { - key: string; - generation: number; -}; - type ReplyPayloadRunState = { runId?: string; }; const replyPayloadSendingDispatchers = new WeakSet(); +const foregroundReplyLeases = createKeyedFifoLeaseRegistry( + Symbol.for("openclaw.foregroundReplyFences"), +); function applyRuntimeToolsAllow( replyOptions: InternalDispatchReplyOptions | undefined, @@ -71,7 +65,7 @@ function applyRuntimeToolsAllow( }; } -function resolveForegroundReplyFenceKey(finalized: FinalizedMsgContext): string | undefined { +function resolveForegroundReplyOrderKey(finalized: FinalizedMsgContext): string | undefined { const sessionKey = normalizeOptionalString(finalized.SessionKey); const channel = normalizeOptionalString(finalized.OriginatingChannel) ?? @@ -98,83 +92,24 @@ function resolveForegroundReplyFenceKey(finalized: FinalizedMsgContext): string ]); } -function beginForegroundReplyFence( - finalized: FinalizedMsgContext, -): ForegroundReplyFenceSnapshot | undefined { - const key = resolveForegroundReplyFenceKey(finalized); - if (!key) { - return undefined; - } - const state = foregroundReplyFenceByKey.get(key) ?? { - generation: 0, - activeGenerations: new Set(), - waiters: new Set<() => void>(), - }; - // Keep every admitted generation until it settles so successors cannot overtake it. - state.generation += 1; - state.activeGenerations.add(state.generation); - foregroundReplyFenceByKey.set(key, state); - return { - key, - generation: state.generation, - }; +function reserveForegroundReplyLease(finalized: FinalizedMsgContext): KeyedFifoLease | undefined { + const key = resolveForegroundReplyOrderKey(finalized); + return key ? foregroundReplyLeases.reserve([key]) : undefined; } -function hasEarlierActiveForegroundReplyFenceGeneration( - state: ForegroundReplyFenceState, - generation: number, -): boolean { - for (const activeGeneration of state.activeGenerations) { - if (activeGeneration < generation) { - return true; - } - } - return false; -} - -async function waitForEarlierForegroundReplyFenceGenerations( - snapshot: ForegroundReplyFenceSnapshot | undefined, -): Promise { - if (!snapshot) { - return; - } - while (true) { - const state = foregroundReplyFenceByKey.get(snapshot.key); - if (!state || !hasEarlierActiveForegroundReplyFenceGeneration(state, snapshot.generation)) { - return; - } - // Delivery is FIFO; model work remains concurrent and only the visible boundary waits. - await new Promise((resolve) => { - state.waiters.add(resolve); - }); - } -} - -async function runForegroundReplyFenceSettledDeliveries( - snapshot: ForegroundReplyFenceSnapshot | undefined, +async function runOrderedForegroundReplySettledDeliveries( + lease: KeyedFifoLease | undefined, onSettled: (() => unknown) | undefined, onFreshSettledDelivery: (() => unknown) | undefined, ): Promise { if (!onSettled && !onFreshSettledDelivery) { return; } - await waitForEarlierForegroundReplyFenceGenerations(snapshot); + await lease?.wait(); await onSettled?.(); await onFreshSettledDelivery?.(); } -function endForegroundReplyFence(snapshot: ForegroundReplyFenceSnapshot): void { - const state = foregroundReplyFenceByKey.get(snapshot.key); - if (!state) { - return; - } - state.activeGenerations.delete(snapshot.generation); - notifyForegroundReplyFenceWaiters(state); - if (state.activeGenerations.size === 0) { - foregroundReplyFenceByKey.delete(snapshot.key); - } -} - function resolveDispatcherSilentReplyContext( ctx: MsgContext | FinalizedMsgContext, cfg: OpenClawConfig, @@ -381,7 +316,7 @@ async function dispatchInboundMessageWithBufferedDispatcherCore( }, ): Promise { const finalized = finalizeInboundContext(params.ctx); - const foregroundReplyFence = beginForegroundReplyFence(finalized); + const foregroundReplyLease = reserveForegroundReplyLease(finalized); const silentReplyContext = resolveDispatcherSilentReplyContext(finalized, params.cfg); const replyPayloadRunState = { runId: params.replyOptions?.runId, @@ -411,9 +346,9 @@ async function dispatchInboundMessageWithBufferedDispatcherCore( ) : globalBeforeDeliver; const beforeDeliver: ReplyDispatchBeforeDeliver | undefined = - foregroundReplyFence || configuredBeforeDeliver + foregroundReplyLease || configuredBeforeDeliver ? markReplyDispatchBeforeDeliverDeadlineOwned(async (payload, info) => { - await waitForEarlierForegroundReplyFenceGenerations(foregroundReplyFence); + await foregroundReplyLease?.wait(); return configuredBeforeDeliver ? await configuredBeforeDeliver(payload, info) : payload; }) : undefined; @@ -448,15 +383,13 @@ async function dispatchInboundMessageWithBufferedDispatcherCore( }); } finally { try { - await runForegroundReplyFenceSettledDeliveries( - foregroundReplyFence, + await runOrderedForegroundReplySettledDeliveries( + foregroundReplyLease, params.dispatcherOptions.onSettled, params.dispatcherOptions.onFreshSettledDelivery, ); } finally { - if (foregroundReplyFence) { - endForegroundReplyFence(foregroundReplyFence); - } + foregroundReplyLease?.release(); markRunComplete(); markDispatchIdle(); } diff --git a/src/auto-reply/foreground-reply-fence-state.ts b/src/auto-reply/foreground-reply-fence-state.ts deleted file mode 100644 index e6d8a533d238..000000000000 --- a/src/auto-reply/foreground-reply-fence-state.ts +++ /dev/null @@ -1,26 +0,0 @@ -import { resolveGlobalMap } from "../shared/global-singleton.js"; - -export type ForegroundReplyFenceState = { - generation: number; - activeGenerations: Set; - waiters: Set<() => void>; -}; - -export function notifyForegroundReplyFenceWaiters(state: ForegroundReplyFenceState): void { - const waiters = [...state.waiters]; - state.waiters.clear(); - for (const resolve of waiters) { - resolve(); - } -} - -export const foregroundReplyFenceByKey = resolveGlobalMap( - Symbol.for("openclaw.foregroundReplyFences"), - (fences) => { - for (const state of fences.values()) { - notifyForegroundReplyFenceWaiters(state); - } - fences.clear(); - }, - "close-only", -); diff --git a/src/auto-reply/reply/agent-runner-run.ts b/src/auto-reply/reply/agent-runner-run.ts index 92177fcdc4f4..e8bae5572814 100644 --- a/src/auto-reply/reply/agent-runner-run.ts +++ b/src/auto-reply/reply/agent-runner-run.ts @@ -436,7 +436,6 @@ export async function runReplyAgent( routeThreadId: replyRouteThreadId, originatingLeafEntryId: turnAdoptionLifecycle?.originatingLeafEntryId, upstreamAbortSignal: opts?.abortSignal, - onReplyAdmissionWaitChange: opts?.onReplyAdmissionWaitChange, }); if (replyOperationRunState) { replyOperationRunState.admission = diff --git a/src/auto-reply/reply/dispatch-from-config.lifecycle.ts b/src/auto-reply/reply/dispatch-from-config.lifecycle.ts index 516ed843175c..627157004776 100644 --- a/src/auto-reply/reply/dispatch-from-config.lifecycle.ts +++ b/src/auto-reply/reply/dispatch-from-config.lifecycle.ts @@ -212,7 +212,6 @@ export function createDispatchReplyOperationCoordinator(params: { waitForActive: !allowActivePreDispatch && !allowSlackRoutedThreadBypass, retainLifecycleAdmissionOnActive: allowActivePreDispatch || allowSlackRoutedThreadBypass, onLifecycleInterrupt, - onReplyAdmissionWaitChange: params.replyOptions?.onReplyAdmissionWaitChange, }); if ( admission.status === "skipped" && @@ -261,7 +260,6 @@ export function createDispatchReplyOperationCoordinator(params: { waitForActive: !allowActivePreDispatch && !allowSlackRoutedThreadBypass, retainLifecycleAdmissionOnActive: allowActivePreDispatch || allowSlackRoutedThreadBypass, onLifecycleInterrupt, - onReplyAdmissionWaitChange: params.replyOptions?.onReplyAdmissionWaitChange, }); } } diff --git a/src/auto-reply/reply/dispatch-from-config.stale-recovery.test.ts b/src/auto-reply/reply/dispatch-from-config.stale-recovery.test.ts index d60a53e89de6..6cda52d1cabe 100644 --- a/src/auto-reply/reply/dispatch-from-config.stale-recovery.test.ts +++ b/src/auto-reply/reply/dispatch-from-config.stale-recovery.test.ts @@ -85,14 +85,8 @@ describe("dispatchReplyFromConfig stale visible admission recovery", () => { activeOperation.abortSignal.addEventListener("abort", () => activeOperation.complete(), { once: true, }); - const waitChanges: boolean[] = []; const replyResolver = vi.fn(async () => ({ text: "telegram reply" }) satisfies ReplyPayload); - const dispatchParams = { - ...createVisibleDispatchParams(replyResolver), - replyOptions: { - onReplyAdmissionWaitChange: (waiting: boolean) => waitChanges.push(waiting), - }, - }; + const dispatchParams = createVisibleDispatchParams(replyResolver); let settled = false; const resultPromise = dispatchReplyFromConfig(dispatchParams).then((result) => { @@ -103,7 +97,6 @@ describe("dispatchReplyFromConfig stale visible admission recovery", () => { await vi.advanceTimersByTimeAsync(120_000); expect(settled).toBe(false); - expect(waitChanges).toEqual([true]); expect(replyResolver).not.toHaveBeenCalled(); activeOperation.complete(); @@ -115,7 +108,6 @@ describe("dispatchReplyFromConfig stale visible admission recovery", () => { }); expect(replyResolver).toHaveBeenCalledTimes(1); expect(dispatchParams.dispatcher.sendFinalReply).toHaveBeenCalledTimes(1); - expect(waitChanges).toEqual([true, false]); }); it("reclaims stale pre-backend work after bounded terminal settlement", async () => { diff --git a/src/auto-reply/reply/followup-turn-admission.ts b/src/auto-reply/reply/followup-turn-admission.ts index 82b6717b2fa7..98d9522a0e94 100644 --- a/src/auto-reply/reply/followup-turn-admission.ts +++ b/src/auto-reply/reply/followup-turn-admission.ts @@ -151,7 +151,6 @@ export async function admitFollowupTurn(params: { routeThreadId: params.queued.originatingThreadId, originatingLeafEntryId: params.queued.turnAdoptionLifecycle?.originatingLeafEntryId, upstreamAbortSignal: resolveFollowupAbortSignal(params.queued), - onReplyAdmissionWaitChange: params.queued.onReplyAdmissionWaitChange, }); if (admission.status === "skipped") { return admission.reason === "active-run" diff --git a/src/auto-reply/reply/get-reply-run-execute.ts b/src/auto-reply/reply/get-reply-run-execute.ts index 49a137435103..af2bca8e3c9c 100644 --- a/src/auto-reply/reply/get-reply-run-execute.ts +++ b/src/auto-reply/reply/get-reply-run-execute.ts @@ -362,7 +362,6 @@ export async function executePreparedReplyRun(state: PreparedReplyRunAdmission) ...(queuedFollowupAbortSignal ? { abortSignal: queuedFollowupAbortSignal } : {}), deliveryCorrelations: opts?.queuedDeliveryCorrelations, turnAdoptionLifecycle: opts?.turnAdoptionLifecycle, - onReplyAdmissionWaitChange: opts?.onReplyAdmissionWaitChange, ...(opts?.onFollowupQueueDisposition ? { onQueueDisposition: opts.onFollowupQueueDisposition } : {}), diff --git a/src/auto-reply/reply/get-reply.types.ts b/src/auto-reply/reply/get-reply.types.ts index ca03c31ce72a..7fa4101846d8 100644 --- a/src/auto-reply/reply/get-reply.types.ts +++ b/src/auto-reply/reply/get-reply.types.ts @@ -30,8 +30,6 @@ type InternalReplySessionOptions = { requestedSessionId?: string; resumeRequestedSession?: boolean; sessionPromptSourceReplyDeliveryMode?: GetReplyOptions["sourceReplyDeliveryMode"]; - /** Marks when this reply is waiting to own its session's reply lane. */ - onReplyAdmissionWaitChange?: (waiting: boolean) => void; /** Receives terminal queue-cap outcomes without widening the public reply API. */ onFollowupQueueDisposition?: (disposition: FollowupQueueDisposition) => void; /** Overrides persisted queue mode for this reply only. */ diff --git a/src/auto-reply/reply/queue/drain.identity-guard.test.ts b/src/auto-reply/reply/queue/drain.identity-guard.test.ts index 70b73b7652da..43a62afc4eea 100644 --- a/src/auto-reply/reply/queue/drain.identity-guard.test.ts +++ b/src/auto-reply/reply/queue/drain.identity-guard.test.ts @@ -117,48 +117,4 @@ describe("drain finally identity guard — late D1 must not orphan Q2", () => { expect(calls).toHaveLength(1); expect(calls[0]?.prompt).toBe("msg1"); }); - - it("keeps each inbound admission callback when it joins an active drain", async () => { - const key = `test-drain-callback-${Date.now()}-${Math.random()}`; - keysToCleanup.push(key); - const settings: QueueSettings = { mode: "followup", debounceMs: 0, cap: 50 }; - const gate = createDeferred(); - const firstEntered = createDeferred(); - const firstCallback = () => {}; - const secondCallback = () => {}; - const callbacks: Array = []; - const firstRunner = async (run: FollowupRun) => { - callbacks.push(run.onReplyAdmissionWaitChange); - if (callbacks.length === 1) { - firstEntered.resolve(); - await gate.promise; - } - }; - const secondRunner = async () => { - throw new Error("active drain must retain its current runner"); - }; - - enqueueFollowupRun( - key, - { ...createRun({ prompt: "msg1" }), onReplyAdmissionWaitChange: firstCallback }, - settings, - "message-id", - firstRunner, - ); - scheduleFollowupDrain(key, firstRunner); - await firstEntered.promise; - - enqueueFollowupRun( - key, - { ...createRun({ prompt: "msg2" }), onReplyAdmissionWaitChange: secondCallback }, - settings, - "message-id", - secondRunner, - ); - scheduleFollowupDrain(key, secondRunner); - gate.resolve(); - - await expect.poll(() => callbacks.length).toBe(2); - expect(callbacks).toEqual([firstCallback, secondCallback]); - }); }); diff --git a/src/auto-reply/reply/queue/drain.ts b/src/auto-reply/reply/queue/drain.ts index 11632b613c3a..6519f48c5445 100644 --- a/src/auto-reply/reply/queue/drain.ts +++ b/src/auto-reply/reply/queue/drain.ts @@ -412,7 +412,6 @@ type FollowupRuntimeMetadata = Pick< | "queueAbortSignal" | "deliveryCorrelations" | "turnAdoptionLifecycle" - | "onReplyAdmissionWaitChange" >; function hasCurrentTurnRuntimeMetadata(item: FollowupRun): boolean { @@ -642,11 +641,6 @@ function collectRuntimeMetadata( // Preserve the exact carrier (including hidden intersections); never derive it from identity evidence. const authoritySource = items.at(-1); const deliveryCorrelations = items.flatMap((item) => item.deliveryCorrelations ?? []); - const admissionWaitCallbacks = new Set( - items.flatMap((item) => - item.onReplyAdmissionWaitChange ? [item.onReplyAdmissionWaitChange] : [], - ), - ); const explicitSkillSelections = [ ...new Map( items @@ -669,14 +663,6 @@ function collectRuntimeMetadata( queueAbortSignal: items.find((item) => item.queueAbortSignal)?.queueAbortSignal, deliveryCorrelations: deliveryCorrelations.length > 0 ? deliveryCorrelations : undefined, turnAdoptionLifecycle: items.length === 1 ? items[0]?.turnAdoptionLifecycle : undefined, - onReplyAdmissionWaitChange: - admissionWaitCallbacks.size > 0 - ? (waiting) => { - for (const callback of admissionWaitCallbacks) { - callback(waiting); - } - } - : undefined, }; } @@ -1041,7 +1027,6 @@ export function createOverflowSummaryRetrySource(source: FollowupRun): FollowupR originatingChatType: source.originatingChatType, abortSignal: source.abortSignal, turnAdoptionLifecycle: source.turnAdoptionLifecycle, - onReplyAdmissionWaitChange: source.onReplyAdmissionWaitChange, ...(source.currentInboundEventKind === "room_event" ? { currentInboundEventKind: "room_event" } : {}), @@ -1103,7 +1088,6 @@ async function runSyntheticOverflowSummary(params: { run: resolveCollectedRun(params.sources, params.source.run), enqueuedAt: Date.now(), abortSignal: params.abortSignal, - onReplyAdmissionWaitChange: runtimeMetadata.onReplyAdmissionWaitChange, explicitSkillSelections: runtimeMetadata.explicitSkillSelections, channelAdmissionEvidence: runtimeMetadata.channelAdmissionEvidence, toolsAllow: runtimeMetadata.toolsAllow, diff --git a/src/auto-reply/reply/queue/types.ts b/src/auto-reply/reply/queue/types.ts index 46b71ab72756..ac5dbc9db387 100644 --- a/src/auto-reply/reply/queue/types.ts +++ b/src/auto-reply/reply/queue/types.ts @@ -98,8 +98,6 @@ export type FollowupRun = { deliveryCorrelations?: QueuedReplyDeliveryCorrelation[]; /** Canonical ownership lifecycle for durable ingress / reply-lane transfer. */ turnAdoptionLifecycle?: TurnAdoptionLifecycle; - /** Dispatch-scoped freshness owner for a queued delivery-barrier wait. */ - onReplyAdmissionWaitChange?: (waiting: boolean) => void; /** Records terminal queue-cap outcomes at the queue owner before lifecycle cleanup. */ onQueueDisposition?: (disposition: FollowupQueueDisposition) => void; /** Provider message ID, when available (for deduplication). */ diff --git a/src/auto-reply/reply/reply-admission-ticket.ts b/src/auto-reply/reply/reply-admission-ticket.ts index 55addb1a1195..55e54ef13707 100644 --- a/src/auto-reply/reply/reply-admission-ticket.ts +++ b/src/auto-reply/reply/reply-admission-ticket.ts @@ -1,61 +1,22 @@ import { normalizeStringifiedEntries } from "@openclaw/normalization-core/string-coerce"; -import { createDeferredCore } from "../../shared/deferred.js"; -import { resolveGlobalMap } from "../../shared/global-singleton.js"; +import { + createKeyedFifoLeaseRegistry, + type KeyedFifoLease, +} from "../../shared/keyed-fifo-lease.js"; export const REPLY_ADMISSION_TICKET = Symbol("openclaw.replyAdmissionTicket"); -type ReplyAdmissionTicket = { - wait(signal?: AbortSignal): Promise; - release(): void; -}; +type ReplyAdmissionTicket = KeyedFifoLease; export type ReplyOptionsWithAdmissionTicket = { [REPLY_ADMISSION_TICKET]?: ReplyAdmissionTicket; }; -const tails = resolveGlobalMap>(Symbol.for("openclaw.replyAdmissionTickets")); +const replyAdmissionTickets = createKeyedFifoLeaseRegistry( + Symbol.for("openclaw.replyAdmissionTickets"), +); /** Briefly orders queue publication across a command's source and target sessions. */ export function reserveReplyAdmissionTicket( sessionKeys: ReadonlyArray, ): ReplyAdmissionTicket | undefined { - const keys = [...new Set(normalizeStringifiedEntries(sessionKeys))].toSorted(); - if (keys.length === 0) { - return undefined; - } - const { promise: completed, resolve: finish } = createDeferredCore(); - const predecessors = keys.map((key) => tails.get(key) ?? Promise.resolve()); - const owned = keys.map((key, index) => { - const tail = predecessors[index]!.then(() => completed); - tails.set(key, tail); - return { key, tail }; - }); - let released = false; - return { - async wait(signal) { - if (signal?.aborted) { - return false; - } - const ready = Promise.all(predecessors).then(() => true); - if (!signal) { - return await ready; - } - return await new Promise((resolve) => { - const abort = () => resolve(false); - signal.addEventListener("abort", abort, { once: true }); - void ready.then((value) => { - signal.removeEventListener("abort", abort); - resolve(value); - }); - }); - }, - release() { - if (released) { - return; - } - released = true; - finish(); - for (const { key, tail } of owned) { - void tail.then(() => tails.get(key) === tail && tails.delete(key)); - } - }, - }; + return replyAdmissionTickets.reserve(normalizeStringifiedEntries(sessionKeys)); } diff --git a/src/auto-reply/reply/reply-turn-admission.test.ts b/src/auto-reply/reply/reply-turn-admission.test.ts index 9de3ac2d789c..3aeb1512e161 100644 --- a/src/auto-reply/reply/reply-turn-admission.test.ts +++ b/src/auto-reply/reply/reply-turn-admission.test.ts @@ -926,7 +926,6 @@ describe("reply turn admission", () => { }); it("waits for visible turns and reuses the active session id", async () => { - const waitChanges: boolean[] = []; const active = createTestReplyOperation({ sessionKey: "agent:main:telegram:topic:42", sessionId: "active-session", @@ -936,7 +935,6 @@ describe("reply turn admission", () => { const admitted = admitTestReplyTurn({ sessionKey: "agent:main:telegram:topic:42", sessionId: "new-session", - onReplyAdmissionWaitChange: (waiting) => waitChanges.push(waiting), }); let settled = false; @@ -947,11 +945,9 @@ describe("reply turn admission", () => { setImmediate(resolve); }); expect(settled).toBe(false); - expect(waitChanges).toEqual([true]); active.complete(); const result = await admitted; - expect(waitChanges).toEqual([true, false]); expect(result.status).toBe("owned"); if (result.status === "owned") { @@ -1024,7 +1020,6 @@ describe("reply turn admission", () => { }); it("keeps an already-waiting follow-up behind the delivery barrier", async () => { - const waitChanges: boolean[] = []; const active = createTestReplyOperation({ sessionKey: "agent:main:discord:channel:42", sessionId: "active-session", @@ -1037,7 +1032,6 @@ describe("reply turn admission", () => { sessionKey: "agent:main:discord:channel:42", sessionId: "queued-session", kind: "queued_followup", - onReplyAdmissionWaitChange: (waiting) => waitChanges.push(waiting), }); let settled = false; void admitted.then(() => { @@ -1049,13 +1043,9 @@ describe("reply turn admission", () => { await Promise.resolve(); expect(settled).toBe(false); - await vi.waitFor(() => { - expect(waitChanges).toEqual([true]); - }); releaseBarrier(); const result = await admitted; - expect(waitChanges).toEqual([true, false]); expect(result.status).toBe("owned"); if (result.status === "owned") { result.operation.complete(); @@ -1411,13 +1401,11 @@ describe("reply turn admission", () => { active.setPhase("running"); active.recordActivity(); const abortController = new AbortController(); - const waitChanges: boolean[] = []; let settled = false; const result = admitTestReplyTurn({ sessionKey: "agent:main:telegram:topic:fresh-visible", sessionId: "waiting-session", upstreamAbortSignal: abortController.signal, - onReplyAdmissionWaitChange: (waiting) => waitChanges.push(waiting), }).then((admission) => { settled = true; return admission; @@ -1425,7 +1413,6 @@ describe("reply turn admission", () => { await vi.advanceTimersByTimeAsync(REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS); expect(settled).toBe(false); - expect(waitChanges).toEqual([true]); expect(replyRunRegistry.get("agent:main:telegram:topic:fresh-visible")).toBe(active); abortController.abort(); @@ -1434,7 +1421,6 @@ describe("reply turn admission", () => { reason: "aborted", activeOperation: active, }); - expect(waitChanges).toEqual([true, false]); } finally { await vi.runOnlyPendingTimersAsync(); vi.useRealTimers(); diff --git a/src/auto-reply/reply/reply-turn-admission.ts b/src/auto-reply/reply/reply-turn-admission.ts index 86d9015217b5..ee6a28f428df 100644 --- a/src/auto-reply/reply/reply-turn-admission.ts +++ b/src/auto-reply/reply/reply-turn-admission.ts @@ -148,36 +148,11 @@ type ReplyTurnAdmissionParams = { waitForActive?: boolean; retainLifecycleAdmissionOnActive?: boolean; onLifecycleInterrupt?: () => void; - /** Reports one interval while blocked behind an older lane owner or its delivery barrier. */ - onReplyAdmissionWaitChange?: (waiting: boolean) => void; }; -type WaitForReplyAdmission = (wait: () => Promise) => Promise; - /** Waits for or claims the per-session reply run slot. */ export async function admitReplyTurn( params: ReplyTurnAdmissionParams, -): Promise { - let admissionWaitReported = false; - const waitForAdmission = async (wait: () => Promise): Promise => { - if (!admissionWaitReported) { - admissionWaitReported = true; - params.onReplyAdmissionWaitChange?.(true); - } - return await wait(); - }; - try { - return await admitReplyTurnWithWaitSignal(params, waitForAdmission); - } finally { - if (admissionWaitReported) { - params.onReplyAdmissionWaitChange?.(false); - } - } -} - -async function admitReplyTurnWithWaitSignal( - params: ReplyTurnAdmissionParams, - waitForAdmission: WaitForReplyAdmission, ): Promise { let sessionId = params.sessionId; let expectedSessionId = params.expectedSessionId; @@ -192,12 +167,10 @@ async function admitReplyTurnWithWaitSignal( if (params.kind === "heartbeat") { return { status: "skipped", reason: "active-run" }; } - const successorAdmission = await waitForAdmission(() => - waitForReplyRunSuccessorAdmission( - params.sessionKey, - params.kind === "visible" ? null : waitTimeoutMs, - { signal: params.upstreamAbortSignal }, - ), + const successorAdmission = await waitForReplyRunSuccessorAdmission( + params.sessionKey, + params.kind === "visible" ? null : waitTimeoutMs, + { signal: params.upstreamAbortSignal }, ); if (!successorAdmission.settled) { return { @@ -417,12 +390,10 @@ async function admitReplyTurnWithWaitSignal( if (params.kind === "heartbeat") { return { status: "skipped", reason: "active-run" }; } - const followupAdmission = await waitForAdmission(() => - waitForReplyRunFollowupAdmission( - params.sessionKey, - waitTimeoutMs ?? REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS, - { signal: params.upstreamAbortSignal }, - ), + const followupAdmission = await waitForReplyRunFollowupAdmission( + params.sessionKey, + waitTimeoutMs ?? REPLY_RUN_IDLE_SETTLE_TIMEOUT_MS, + { signal: params.upstreamAbortSignal }, ); if (!followupAdmission.settled) { return { @@ -457,11 +428,9 @@ async function admitReplyTurnWithWaitSignal( } const activeWaitTimeoutMs = params.kind === "visible" ? resolveVisibleActiveWaitMs(activeOperation) : waitTimeoutMs; - const ended = await waitForAdmission(() => - replyRunRegistry.waitForIdle(params.sessionKey, activeWaitTimeoutMs, { - signal: params.upstreamAbortSignal, - }), - ); + const ended = await replyRunRegistry.waitForIdle(params.sessionKey, activeWaitTimeoutMs, { + signal: params.upstreamAbortSignal, + }); if (!ended) { if (params.kind === "visible" && !isAbortSignalAborted(params.upstreamAbortSignal)) { // Visible turns block on active work like before, but in bounded wait diff --git a/src/auto-reply/reply/stranded-reply-recovery.test.ts b/src/auto-reply/reply/stranded-reply-recovery.test.ts index c0fb92660b7c..c6fa74f6ba7b 100644 --- a/src/auto-reply/reply/stranded-reply-recovery.test.ts +++ b/src/auto-reply/reply/stranded-reply-recovery.test.ts @@ -18,7 +18,6 @@ describe("buildStrandedReplyRetryFollowupRun lifecycle ownership", () => { onDeferred: onEnqueued, }, admissionSessionId: "sess-rotated", - onReplyAdmissionWaitChange: vi.fn(), }); const recovery = resolveStrandedReplyRecovery({ @@ -42,7 +41,6 @@ describe("buildStrandedReplyRetryFollowupRun lifecycle ownership", () => { expect(retry.summaryLine).toBe(STRANDED_REPLY_RETRY_MARKER); // Session routing stays; only the client-turn lifecycle identity is detached. expect(retry.admissionSessionId).toBe("sess-rotated"); - expect(retry.onReplyAdmissionWaitChange).toBe(parent.onReplyAdmissionWaitChange); expect(retry.run.sessionKey).toBe(parent.run.sessionKey); // mark/complete no-op when lifecycle is absent (drop-policy onDrop path too). diff --git a/src/channels/turn/run-channel-turn.pipeline.test.ts b/src/channels/turn/run-channel-turn.pipeline.test.ts index 48458ca0bb35..918f16b600a4 100644 --- a/src/channels/turn/run-channel-turn.pipeline.test.ts +++ b/src/channels/turn/run-channel-turn.pipeline.test.ts @@ -980,41 +980,4 @@ describe("channel turn pipeline", () => { expect.objectContaining({ reason: "zero-count-visible-dispatch" }), ]); }); - - it("still warns when a visible turn has zero counts and no observed delivery", async () => { - const events: string[] = []; - const log = vi.fn(); - const recordInboundSession = createRecordInboundSession(events); - // Guard against over-suppression: a genuinely empty visible dispatch must still warn. - const runDispatch = vi.fn(async () => ({ - queuedFinal: false, - counts: { tool: 0, block: 0, final: 0 }, - observedReplyDelivery: false, - })); - - const result = await runPreparedChannelTurn({ - channel: "test", - routeSessionKey: "agent:main:test:peer", - storePath: "/tmp/sessions.json", - ctxPayload: createCtx(), - recordInboundSession, - runDispatch, - log, - messageId: "msg-empty", - record: { - onRecordError: vi.fn(), - }, - }); - - expectDispatched(result); - expect(result.dispatchResult?.observedReplyDelivery).toBe(false); - expect(log.mock.calls).toContainEqual([ - expect.objectContaining({ - stage: "dispatch", - event: "warning", - messageId: "msg-empty", - reason: "zero-count-visible-dispatch", - }), - ]); - }); }); diff --git a/src/shared/keyed-fifo-lease.test.ts b/src/shared/keyed-fifo-lease.test.ts new file mode 100644 index 000000000000..4c7f693a8014 --- /dev/null +++ b/src/shared/keyed-fifo-lease.test.ts @@ -0,0 +1,152 @@ +import { importFreshModule } from "openclaw/plugin-sdk/test-fixtures"; +import { afterEach, describe, expect, it } from "vitest"; +import { drainGlobalSingletonLifecycleState } from "./global-singleton.js"; +import { createKeyedFifoLeaseRegistry } from "./keyed-fifo-lease.js"; + +const keysToDelete = new Set(); +const TEST_KEY = Symbol.for("openclaw.test.keyedFifoLease"); + +function createRegistry() { + keysToDelete.add(TEST_KEY); + return createKeyedFifoLeaseRegistry(TEST_KEY); +} + +afterEach(async () => { + await drainGlobalSingletonLifecycleState("close"); + for (const key of keysToDelete) { + delete (globalThis as Record)[key]; + } + keysToDelete.clear(); +}); + +describe("keyed FIFO leases", () => { + it("reserves order before reverse-completing work reaches wait", async () => { + const registry = createRegistry(); + const older = registry.reserve(["target"])!; + const newer = registry.reserve(["target"])!; + let newerReady = false; + + const newerWait = newer.wait().then((ready) => { + newerReady = ready; + }); + await Promise.resolve(); + expect(newerReady).toBe(false); + + older.release(); + await newerWait; + expect(newerReady).toBe(true); + newer.release(); + }); + + it("does not release an aborted waiter's reserved slot", async () => { + const registry = createRegistry(); + const older = registry.reserve(["target"])!; + const aborted = registry.reserve(["target"])!; + const newer = registry.reserve(["target"])!; + const controller = new AbortController(); + const abortedWait = aborted.wait(controller.signal); + const newerWait = newer.wait(); + let newerReady = false; + void newerWait.then(() => { + newerReady = true; + }); + + controller.abort(); + await expect(abortedWait).resolves.toBe(false); + older.release(); + await Promise.resolve(); + expect(newerReady).toBe(false); + + aborted.release(); + aborted.release(); + await expect(newerWait).resolves.toBe(true); + newer.release(); + }); + + it("reserves sorted unique multi-key leases without deadlock", async () => { + const registry = createRegistry(); + const first = registry.reserve(["b", "a", "b"])!; + const second = registry.reserve(["a", "b"])!; + let secondReady = false; + const secondWait = second.wait().then((ready) => { + secondReady = ready; + }); + + await expect(first.wait()).resolves.toBe(true); + await Promise.resolve(); + expect(secondReady).toBe(false); + first.release(); + await secondWait; + expect(secondReady).toBe(true); + second.release(); + }); + + it("keeps a newer tail when an older tail cleans up", async () => { + const registry = createRegistry(); + const first = registry.reserve(["target"])!; + first.release(); + const second = registry.reserve(["target"])!; + await expect(second.wait()).resolves.toBe(true); + const third = registry.reserve(["target"])!; + let thirdReady = false; + void third.wait().then(() => { + thirdReady = true; + }); + + await Promise.resolve(); + expect(thirdReady).toBe(false); + second.release(); + await expect(third.wait()).resolves.toBe(true); + third.release(); + }); + + it("shares reservations across duplicate runtime chunks", async () => { + const key = TEST_KEY; + keysToDelete.add(key); + const moduleA = await importFreshModule( + import.meta.url, + "./keyed-fifo-lease.js?scope=duplicate-a", + ); + const moduleB = await importFreshModule( + import.meta.url, + "./keyed-fifo-lease.js?scope=duplicate-b", + ); + const firstRegistry = moduleA.createKeyedFifoLeaseRegistry(key); + const secondRegistry = moduleB.createKeyedFifoLeaseRegistry(key); + const first = firstRegistry.reserve(["target"])!; + const second = secondRegistry.reserve(["target"])!; + let secondReady = false; + void second.wait().then(() => { + secondReady = true; + }); + + await Promise.resolve(); + expect(secondReady).toBe(false); + first.release(); + await expect(second.wait()).resolves.toBe(true); + second.release(); + }); + + it("releases live gates on full close but not restart", async () => { + const registry = createRegistry(); + registry.reserve(["target"]); + registry.reserve(["target"]); + const last = registry.reserve(["target"])!; + let ready = false; + const lastWait = last.wait().then((result) => { + ready = result; + return result; + }); + + await drainGlobalSingletonLifecycleState("restart"); + await Promise.resolve(); + expect(ready).toBe(false); + + await drainGlobalSingletonLifecycleState("close"); + await expect(lastWait).resolves.toBe(true); + + const fresh = registry.reserve(["target"])!; + await expect(fresh.wait()).resolves.toBe(true); + fresh.release(); + }); +}); diff --git a/src/shared/keyed-fifo-lease.ts b/src/shared/keyed-fifo-lease.ts new file mode 100644 index 000000000000..f44cb4fe37c2 --- /dev/null +++ b/src/shared/keyed-fifo-lease.ts @@ -0,0 +1,87 @@ +import { createDeferredCore } from "./deferred.js"; +import { resolveGlobalSingleton } from "./global-singleton.js"; + +export type KeyedFifoLease = { + wait(signal?: AbortSignal): Promise; + release(): void; +}; + +type KeyedFifoLeaseState = { + tails: Map>; + releases: Set<() => void>; +}; + +type KeyedFifoLeaseRegistry = { + reserve(keys: readonly string[]): KeyedFifoLease | undefined; +}; + +/** Creates a close-owned FIFO registry shared by every runtime chunk using globalKey. */ +export function createKeyedFifoLeaseRegistry(globalKey: symbol): KeyedFifoLeaseRegistry { + const state = resolveGlobalSingleton( + globalKey, + () => ({ tails: new Map(), releases: new Set() }), + (current) => { + // Full close must unblock every predecessor chain before forgetting its tails. + for (const release of current.releases) { + release(); + } + current.tails.clear(); + }, + "close-only", + ); + + return { + reserve(inputKeys) { + const keys = [...new Set(inputKeys)].toSorted(); + if (keys.length === 0) { + return undefined; + } + + const { promise: completed, resolve: complete } = createDeferredCore(); + const predecessors = keys.map((key) => state.tails.get(key) ?? Promise.resolve()); + const owned = keys.map((key, index) => { + const tail = predecessors[index]!.then(() => completed); + state.tails.set(key, tail); + return { key, tail }; + }); + let released = false; + + const release = () => { + if (released) { + return; + } + released = true; + state.releases.delete(release); + complete(); + for (const { key, tail } of owned) { + void tail.then(() => state.tails.get(key) === tail && state.tails.delete(key)); + } + }; + state.releases.add(release); + + return { + async wait(signal) { + if (signal?.aborted) { + return false; + } + const ready = Promise.all(predecessors).then(() => true); + if (!signal) { + return await ready; + } + return await new Promise((resolve) => { + const abort = () => resolve(false); + signal.addEventListener("abort", abort, { once: true }); + if (signal.aborted) { + abort(); + } + void ready.then((value) => { + signal.removeEventListener("abort", abort); + resolve(value); + }); + }); + }, + release, + }; + }, + }; +}