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.
This commit is contained in:
Peter Steinberger
2026-08-16 10:03:29 -07:00
committed by GitHub
parent 3ca7b038f0
commit 2aab6b8e37
19 changed files with 346 additions and 524 deletions
@@ -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<ReplyPayload>(() => {});
}
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<ReplyPayload | null>();
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<ReplyPayload | null>();
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 ?? "<missing>"}`);
},
);
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<ReplyPayload | null>();
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 ?? "<missing>"}`);
},
);
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 () => {
+17 -84
View File
@@ -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<InternalGetReplyOptions, "onBlockReply">;
type ForegroundReplyFenceSnapshot = {
key: string;
generation: number;
};
type ReplyPayloadRunState = {
runId?: string;
};
const replyPayloadSendingDispatchers = new WeakSet<ReplyDispatcher>();
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<number>(),
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<void> {
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<void>((resolve) => {
state.waiters.add(resolve);
});
}
}
async function runForegroundReplyFenceSettledDeliveries(
snapshot: ForegroundReplyFenceSnapshot | undefined,
async function runOrderedForegroundReplySettledDeliveries(
lease: KeyedFifoLease | undefined,
onSettled: (() => unknown) | undefined,
onFreshSettledDelivery: (() => unknown) | undefined,
): Promise<void> {
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<DispatchInboundResult> {
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();
}
@@ -1,26 +0,0 @@
import { resolveGlobalMap } from "../shared/global-singleton.js";
export type ForegroundReplyFenceState = {
generation: number;
activeGenerations: Set<number>;
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<string, ForegroundReplyFenceState>(
Symbol.for("openclaw.foregroundReplyFences"),
(fences) => {
for (const state of fences.values()) {
notifyForegroundReplyFenceWaiters(state);
}
fences.clear();
},
"close-only",
);
-1
View File
@@ -436,7 +436,6 @@ export async function runReplyAgent(
routeThreadId: replyRouteThreadId,
originatingLeafEntryId: turnAdoptionLifecycle?.originatingLeafEntryId,
upstreamAbortSignal: opts?.abortSignal,
onReplyAdmissionWaitChange: opts?.onReplyAdmissionWaitChange,
});
if (replyOperationRunState) {
replyOperationRunState.admission =
@@ -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,
});
}
}
@@ -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 () => {
@@ -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"
@@ -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 }
: {}),
-2
View File
@@ -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. */
@@ -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<FollowupRun["onReplyAdmissionWaitChange"]> = [];
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]);
});
});
-16
View File
@@ -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,
-2
View File
@@ -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). */
+9 -48
View File
@@ -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<boolean>;
release(): void;
};
type ReplyAdmissionTicket = KeyedFifoLease;
export type ReplyOptionsWithAdmissionTicket = {
[REPLY_ADMISSION_TICKET]?: ReplyAdmissionTicket;
};
const tails = resolveGlobalMap<string, Promise<void>>(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<string | undefined>,
): 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<boolean>((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));
}
@@ -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();
+11 -42
View File
@@ -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 = <T>(wait: () => Promise<T>) => Promise<T>;
/** Waits for or claims the per-session reply run slot. */
export async function admitReplyTurn(
params: ReplyTurnAdmissionParams,
): Promise<ReplyTurnAdmission> {
let admissionWaitReported = false;
const waitForAdmission = async <T>(wait: () => Promise<T>): Promise<T> => {
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<ReplyTurnAdmission> {
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
@@ -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).
@@ -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",
}),
]);
});
});
+152
View File
@@ -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<symbol>();
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<PropertyKey, unknown>)[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<typeof import("./keyed-fifo-lease.js")>(
import.meta.url,
"./keyed-fifo-lease.js?scope=duplicate-a",
);
const moduleB = await importFreshModule<typeof import("./keyed-fifo-lease.js")>(
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();
});
});
+87
View File
@@ -0,0 +1,87 @@
import { createDeferredCore } from "./deferred.js";
import { resolveGlobalSingleton } from "./global-singleton.js";
export type KeyedFifoLease = {
wait(signal?: AbortSignal): Promise<boolean>;
release(): void;
};
type KeyedFifoLeaseState = {
tails: Map<string, Promise<void>>;
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<KeyedFifoLeaseState>(
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<boolean>((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,
};
},
};
}