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;
  });
}