mirror of
https://github.com/openclaw/openclaw.git
synced 2026-05-05 22:10:21 +00:00
refactor: unify gateway restart deferral and dispatcher cleanup
This commit is contained in:
61
src/auto-reply/dispatch.test.ts
Normal file
61
src/auto-reply/dispatch.test.ts
Normal file
@@ -0,0 +1,61 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import type { ReplyDispatcher } from "./reply/reply-dispatcher.js";
|
||||
import { withReplyDispatcher } from "./dispatch.js";
|
||||
|
||||
function createDispatcher(record: string[]): ReplyDispatcher {
|
||||
return {
|
||||
sendToolResult: () => true,
|
||||
sendBlockReply: () => true,
|
||||
sendFinalReply: () => true,
|
||||
getQueuedCounts: () => ({ tool: 0, block: 0, final: 0 }),
|
||||
markComplete: () => {
|
||||
record.push("markComplete");
|
||||
},
|
||||
waitForIdle: async () => {
|
||||
record.push("waitForIdle");
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
describe("withReplyDispatcher", () => {
|
||||
it("always marks complete and waits for idle after success", async () => {
|
||||
const order: string[] = [];
|
||||
const dispatcher = createDispatcher(order);
|
||||
|
||||
const result = await withReplyDispatcher({
|
||||
dispatcher,
|
||||
run: async () => {
|
||||
order.push("run");
|
||||
return "ok";
|
||||
},
|
||||
onSettled: () => {
|
||||
order.push("onSettled");
|
||||
},
|
||||
});
|
||||
|
||||
expect(result).toBe("ok");
|
||||
expect(order).toEqual(["run", "markComplete", "waitForIdle", "onSettled"]);
|
||||
});
|
||||
|
||||
it("still drains dispatcher after run throws", async () => {
|
||||
const order: string[] = [];
|
||||
const dispatcher = createDispatcher(order);
|
||||
const onSettled = vi.fn(() => {
|
||||
order.push("onSettled");
|
||||
});
|
||||
|
||||
await expect(
|
||||
withReplyDispatcher({
|
||||
dispatcher,
|
||||
run: async () => {
|
||||
order.push("run");
|
||||
throw new Error("boom");
|
||||
},
|
||||
onSettled,
|
||||
}),
|
||||
).rejects.toThrow("boom");
|
||||
|
||||
expect(onSettled).toHaveBeenCalledTimes(1);
|
||||
expect(order).toEqual(["run", "markComplete", "waitForIdle", "onSettled"]);
|
||||
});
|
||||
});
|
||||
@@ -14,6 +14,24 @@ import {
|
||||
|
||||
export type DispatchInboundResult = DispatchFromConfigResult;
|
||||
|
||||
export async function withReplyDispatcher<T>(params: {
|
||||
dispatcher: ReplyDispatcher;
|
||||
run: () => Promise<T>;
|
||||
onSettled?: () => void | Promise<void>;
|
||||
}): Promise<T> {
|
||||
try {
|
||||
return await params.run();
|
||||
} finally {
|
||||
// Ensure dispatcher reservations are always released on every exit path.
|
||||
params.dispatcher.markComplete();
|
||||
try {
|
||||
await params.dispatcher.waitForIdle();
|
||||
} finally {
|
||||
await params.onSettled?.();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export async function dispatchInboundMessage(params: {
|
||||
ctx: MsgContext | FinalizedMsgContext;
|
||||
cfg: OpenClawConfig;
|
||||
@@ -41,20 +59,23 @@ export async function dispatchInboundMessageWithBufferedDispatcher(params: {
|
||||
const { dispatcher, replyOptions, markDispatchIdle } = createReplyDispatcherWithTyping(
|
||||
params.dispatcherOptions,
|
||||
);
|
||||
|
||||
const result = await dispatchInboundMessage({
|
||||
ctx: params.ctx,
|
||||
cfg: params.cfg,
|
||||
return await withReplyDispatcher({
|
||||
dispatcher,
|
||||
replyResolver: params.replyResolver,
|
||||
replyOptions: {
|
||||
...params.replyOptions,
|
||||
...replyOptions,
|
||||
run: async () =>
|
||||
dispatchInboundMessage({
|
||||
ctx: params.ctx,
|
||||
cfg: params.cfg,
|
||||
dispatcher,
|
||||
replyResolver: params.replyResolver,
|
||||
replyOptions: {
|
||||
...params.replyOptions,
|
||||
...replyOptions,
|
||||
},
|
||||
}),
|
||||
onSettled: () => {
|
||||
markDispatchIdle();
|
||||
},
|
||||
});
|
||||
|
||||
markDispatchIdle();
|
||||
return result;
|
||||
}
|
||||
|
||||
export async function dispatchInboundMessageWithDispatcher(params: {
|
||||
@@ -65,13 +86,15 @@ export async function dispatchInboundMessageWithDispatcher(params: {
|
||||
replyResolver?: typeof import("./reply.js").getReplyFromConfig;
|
||||
}): Promise<DispatchInboundResult> {
|
||||
const dispatcher = createReplyDispatcher(params.dispatcherOptions);
|
||||
const result = await dispatchInboundMessage({
|
||||
ctx: params.ctx,
|
||||
cfg: params.cfg,
|
||||
return await withReplyDispatcher({
|
||||
dispatcher,
|
||||
replyResolver: params.replyResolver,
|
||||
replyOptions: params.replyOptions,
|
||||
run: async () =>
|
||||
dispatchInboundMessage({
|
||||
ctx: params.ctx,
|
||||
cfg: params.cfg,
|
||||
dispatcher,
|
||||
replyResolver: params.replyResolver,
|
||||
replyOptions: params.replyOptions,
|
||||
}),
|
||||
});
|
||||
await dispatcher.waitForIdle();
|
||||
return result;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user