File size: 3,229 Bytes
fcd8223 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 | import { AsyncLocalStorage } from "node:async_hooks";
import type { ReplyPayload } from "../auto-reply/reply-payload.js";
import { SILENT_REPLY_TOKEN } from "../auto-reply/tokens.js";
import { runOncePerAgentRun } from "../infra/agent-events.js";
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
import { getGlobalHookRunner } from "./hook-runner-global.js";
import type {
PluginHookAgentContext,
PluginHookBeforeAgentReplyEvent,
PluginHookBeforeAgentReplyResult,
} from "./hook-types.js";
import { isPluginHookAgentTrigger } from "./hook-types.js";
const BEFORE_AGENT_REPLY_OBSERVER_KEY = Symbol.for("openclaw.beforeAgentReply.observer");
type BeforeAgentReplyObserver = {
beforeDispatch: () => Promise<boolean | void>;
afterDispatch: (
result: PluginHookBeforeAgentReplyResult | undefined,
) => Promise<PluginHookBeforeAgentReplyResult | undefined>;
};
type BeforeAgentReplyObserverScope = BeforeAgentReplyObserver & { runId?: string };
const beforeAgentReplyObserver = resolveGlobalSingleton<
AsyncLocalStorage<BeforeAgentReplyObserverScope>
>(BEFORE_AGENT_REPLY_OBSERVER_KEY, () => new AsyncLocalStorage());
/** Attaches durable admission bookkeeping without moving hook ownership out of the runner. */
export function withBeforeAgentReplyObserver<T>(
observer: BeforeAgentReplyObserver,
run: () => T,
): T {
return beforeAgentReplyObserver.run({ ...observer }, run);
}
/** Preserves the full plugin reply contract, including private payload metadata. */
export function buildHandledBeforeAgentReplyPayloads(reply?: ReplyPayload): ReplyPayload[] {
return [reply ?? { text: SILENT_REPLY_TOKEN }];
}
/** Runs the reply claim hook once for one admitted turn, across model fallbacks. */
export function runBeforeAgentReplyForTurn(params: {
runId: string;
trigger?: string;
event: PluginHookBeforeAgentReplyEvent;
context: PluginHookAgentContext;
onDispatch?: () => void;
onDeclined?: () => void;
}): Promise<PluginHookBeforeAgentReplyResult | undefined> {
const trigger = params.trigger;
if (!isPluginHookAgentTrigger(trigger)) {
return Promise.resolve(undefined);
}
const context = { ...params.context, trigger };
return runOncePerAgentRun(params.runId, "before_agent_reply", async () => {
const hookRunner = getGlobalHookRunner();
if (!hookRunner?.hasHooks("before_agent_reply", context)) {
return undefined;
}
const observerScope = beforeAgentReplyObserver.getStore();
// Nested agent runs inherit async context. Bind recovery to the first runner
// so a hook-spawned child cannot checkpoint its parent's admitted turn.
const observer =
observerScope && (!observerScope.runId || observerScope.runId === params.runId)
? observerScope
: undefined;
if (observer && !observer.runId) {
observer.runId = params.runId;
}
if ((await observer?.beforeDispatch()) === false) {
return undefined;
}
params.onDispatch?.();
let result = await hookRunner.runBeforeAgentReply(params.event, context);
if (!result?.handled) {
params.onDeclined?.();
}
if (observer) {
result = await observer.afterDispatch(result);
}
return result;
});
}
|