Download src/gateway/server/hooks-request-handler.ts from SaylorTwift/openclaw: direct link, hf CLI and curl.
- Browser
- Download file 27.1 kB
-
https://huggingface.co/SaylorTwift/openclaw/resolve/main/src/gateway/server/hooks-request-handler.ts
- Command line
-
hf download hf://SaylorTwift/openclaw/src/gateway/server/hooks-request-handler.ts
-
curl -L -o hooks-request-handler.ts https://huggingface.co/SaylorTwift/openclaw/resolve/main/src/gateway/server/hooks-request-handler.ts
27.1 kB
| import { createHash } from "node:crypto"; | |
| import type { IncomingMessage, ServerResponse } from "node:http"; | |
| import { asRecord } from "@openclaw/normalization-core/record-coerce"; | |
| import { sendHttpRequestRejection } from "../../infra/http-request-lifecycle.js"; | |
| import type { createSubsystemLogger } from "../../logging/subsystem.js"; | |
| import { resolveHookExternalContentSource as resolveHookExternalContentSourceFromSession } from "../../security/external-content.js"; | |
| import { safeEqualSecret } from "../../security/secret-equal.js"; | |
| import { | |
| AUTH_RATE_LIMIT_SCOPE_HOOK_AUTH, | |
| createAuthRateLimiter, | |
| normalizeRateLimitClientIp, | |
| } from "../auth-rate-limit.js"; | |
| import { applyHookMappings, HOOK_MAPPING_FAN_OUT_MAX_ITEMS } from "../hooks-mapping.js"; | |
| import { | |
| extractHookToken, | |
| getHookAgentPolicyError, | |
| getHookChannelError, | |
| getHookSessionKeyPrefixError, | |
| type HookAgentDispatchPayload, | |
| type HookSessionKeySource, | |
| type HookTargetAgentResolution, | |
| type HooksConfigResolved, | |
| isHookAgentAllowed, | |
| isSessionKeyAllowedByPrefix, | |
| normalizeAgentPayload, | |
| normalizeHookDispatchSessionKey, | |
| normalizeHookHeaders, | |
| normalizeWakePayload, | |
| readJsonBody, | |
| resolveEffectiveHookTargetAgentId, | |
| resolveHookChannel, | |
| resolveHookDeliver, | |
| resolveHookIdempotencyKey, | |
| resolveHookPathBodyLimit, | |
| resolveHookSessionKey, | |
| } from "../hooks.js"; | |
| import type { HookAgentDispatchResult, HookAgentDispatchSuccess } from "../hooks.types.js"; | |
| import { sendJson } from "../http-common.js"; | |
| import { readPreparedGatewayIngressAttribution } from "../ingress-attribution.js"; | |
| import { resolveRequestClientIpFromHeaders } from "../net.js"; | |
| import { DEDUPE_MAX, DEDUPE_TTL_MS } from "../server-constants.js"; | |
| import { | |
| HOOK_FAN_OUT_RESPONSE_DEADLINE_MS, | |
| sendAgentResult, | |
| sendFanOutResult, | |
| settleFanOutDispatches, | |
| type WakeResult, | |
| } from "./hooks-request-handler-response.js"; | |
| const HOOK_AUTH_FAILURE_LIMIT = 20; | |
| const HOOK_AUTH_FAILURE_WINDOW_MS = 60_000; | |
| // Marker for replay keys derived from item content when the producer supplies | |
| // no idempotency key; item identity lives in the dispatch-scope fingerprint. | |
| const HOOK_FAN_OUT_DERIVED_IDEMPOTENCY = "hook-fanout-item"; | |
| const HOOK_CONFIG_CHANGED_ERROR = "hook configuration changed; retry request"; | |
| const hashReplay = (value: string) => createHash("sha256").update(value, "utf8").digest("hex"); | |
| export type HookClientIpConfig = Readonly<{ | |
| trustedProxies?: string[]; | |
| allowRealIpFallback?: boolean; | |
| }>; | |
| export type HooksRequestHandler = (req: IncomingMessage, res: ServerResponse) => Promise<boolean>; | |
| type HookDispatchers = { | |
| dispatchWakeHook: ( | |
| value: { text: string; mode: "now" | "next-heartbeat"; sessionKey?: string }, | |
| agentId: string, | |
| ) => WakeResult; | |
| dispatchAgentHook: ( | |
| value: HookAgentDispatchPayload, | |
| ) => HookAgentDispatchResult | Promise<HookAgentDispatchResult>; | |
| }; | |
| type HookReplayEntry = | |
| | { state: "pending"; dispatch: Promise<HookAgentDispatchResult> } | |
| | { state: "active"; dispatch: HookAgentDispatchSuccess } | |
| | { state: "terminal"; ts: number; dispatch: HookAgentDispatchSuccess }; | |
| type HookReplayScope = { | |
| pathKey: string; | |
| token: string | undefined; | |
| idempotencyKey?: string; | |
| dispatchScope: Record<string, unknown>; | |
| }; | |
| function resolveMappedHookExternalContentSource(params: { subPath: string; sessionKey: string }) { | |
| if (params.subPath === "gmail") { | |
| return "gmail" as const; | |
| } | |
| return resolveHookExternalContentSourceFromSession(params.sessionKey) ?? "webhook"; | |
| } | |
| export function createHooksRequestHandler( | |
| opts: { | |
| /** Returns the stable resolved object for the current hooks-config generation. */ | |
| getHooksConfig: () => HooksConfigResolved | null; | |
| bindHost: string; | |
| port: number; | |
| logHooks: ReturnType<typeof createSubsystemLogger>; | |
| getClientIpConfig?: () => HookClientIpConfig; | |
| fanoutResponseDeadlineMs?: number; | |
| } & HookDispatchers, | |
| ): HooksRequestHandler { | |
| const { getHooksConfig, logHooks, dispatchAgentHook, dispatchWakeHook, getClientIpConfig } = opts; | |
| const fanoutResponseDeadlineMs = | |
| opts.fanoutResponseDeadlineMs ?? HOOK_FAN_OUT_RESPONSE_DEADLINE_MS; | |
| const hookReplayCache = new Map<string, HookReplayEntry>(); | |
| const hookAuthLimiter = createAuthRateLimiter({ | |
| maxAttempts: HOOK_AUTH_FAILURE_LIMIT, | |
| windowMs: HOOK_AUTH_FAILURE_WINDOW_MS, | |
| lockoutMs: HOOK_AUTH_FAILURE_WINDOW_MS, | |
| exemptLoopback: false, | |
| // Handler lifetimes are tied to gateway runtime/tests; skip background timer fanout. | |
| pruneIntervalMs: 0, | |
| }); | |
| const resolveHookClientKey = (req: IncomingMessage): string => { | |
| const attribution = readPreparedGatewayIngressAttribution(req); | |
| if (attribution && attribution.kind !== "unattributable-proxy") { | |
| return normalizeRateLimitClientIp(attribution.rateLimit.subject.key); | |
| } | |
| const clientIpConfig = getClientIpConfig?.(); | |
| const clientIp = | |
| resolveRequestClientIpFromHeaders( | |
| req, | |
| clientIpConfig?.trustedProxies, | |
| clientIpConfig?.allowRealIpFallback === true, | |
| ) ?? req.socket?.remoteAddress; | |
| return normalizeRateLimitClientIp(clientIp); | |
| }; | |
| const pruneHookReplayCache = (now: number) => { | |
| for (const [key, entry] of hookReplayCache) { | |
| if (entry.state === "terminal" && entry.ts < now - DEDUPE_TTL_MS) { | |
| hookReplayCache.delete(key); | |
| } | |
| } | |
| const terminal = [...hookReplayCache].filter(([, entry]) => entry.state === "terminal"); | |
| for (const [key] of terminal.slice(0, Math.max(0, terminal.length - DEDUPE_MAX))) { | |
| hookReplayCache.delete(key); | |
| } | |
| }; | |
| const buildHookReplayCacheKey = (params: HookReplayScope): string | undefined => { | |
| const idem = params.idempotencyKey?.trim(); | |
| if (!idem) { | |
| return undefined; | |
| } | |
| const scope = JSON.stringify({ | |
| pathKey: params.pathKey, | |
| dispatchScope: params.dispatchScope, | |
| }); | |
| return `${hashReplay(params.token ?? "")}:${hashReplay(scope)}:${hashReplay(idem)}`; | |
| }; | |
| const resolveHookReplay = (key: string | undefined) => { | |
| if (!key) { | |
| return undefined; | |
| } | |
| pruneHookReplayCache(Date.now()); | |
| const cached = hookReplayCache.get(key); | |
| if (!cached) { | |
| return undefined; | |
| } | |
| if (cached.state === "terminal") { | |
| hookReplayCache.delete(key); | |
| hookReplayCache.set(key, cached); | |
| } | |
| return cached.dispatch; | |
| }; | |
| const dispatchAgentHookWithReplay = ( | |
| key: string | undefined, | |
| dispatch: () => HookAgentDispatchResult | Promise<HookAgentDispatchResult>, | |
| ): HookAgentDispatchResult | Promise<HookAgentDispatchResult> => { | |
| if (!key) { | |
| return dispatch(); | |
| } | |
| const existing = resolveHookReplay(key); | |
| if (existing) { | |
| return existing; | |
| } | |
| const pending = Promise.resolve() | |
| .then(dispatch) | |
| .then((result) => { | |
| const current = hookReplayCache.get(key); | |
| if (current?.state === "pending" && current.dispatch === pending) { | |
| if (result.ok) { | |
| const active = { state: "active", dispatch: result } as const; | |
| hookReplayCache.set(key, active); | |
| const settle = () => { | |
| if (hookReplayCache.get(key) !== active) { | |
| return; | |
| } | |
| const terminal = { state: "terminal", ts: Date.now(), dispatch: result } as const; | |
| hookReplayCache.delete(key); | |
| hookReplayCache.set(key, terminal); | |
| pruneHookReplayCache(terminal.ts); | |
| }; | |
| void result.completion.then(settle, settle); | |
| } else { | |
| hookReplayCache.delete(key); | |
| } | |
| } | |
| return result; | |
| }) | |
| .catch((err: unknown) => { | |
| const current = hookReplayCache.get(key); | |
| if (current?.state === "pending" && current.dispatch === pending) { | |
| hookReplayCache.delete(key); | |
| } | |
| throw err; | |
| }); | |
| hookReplayCache.set(key, { state: "pending", dispatch: pending }); | |
| return pending; | |
| }; | |
| return async (req, res) => { | |
| const hooksConfig = getHooksConfig(); | |
| if (!hooksConfig) { | |
| return false; | |
| } | |
| // Only pathname/search are used here; keep the base host fixed so bind-host | |
| // representation (e.g. IPv6 wildcards) cannot break request parsing. | |
| const url = new URL(req.url ?? "/", "http://localhost"); | |
| const basePath = hooksConfig.basePath; | |
| if (url.pathname !== basePath && !url.pathname.startsWith(`${basePath}/`)) { | |
| return false; | |
| } | |
| if (url.searchParams.has("token")) { | |
| res.statusCode = 400; | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end( | |
| "Hook token must be provided via Authorization: Bearer <token> or X-OpenClaw-Token header (query parameters are not allowed).", | |
| ); | |
| return true; | |
| } | |
| if (req.method !== "POST") { | |
| res.statusCode = 405; | |
| res.setHeader("Allow", "POST"); | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end("Method Not Allowed"); | |
| return true; | |
| } | |
| const token = extractHookToken(req); | |
| const clientKey = resolveHookClientKey(req); | |
| if (!safeEqualSecret(token, hooksConfig.token)) { | |
| const throttle = hookAuthLimiter.check(clientKey, AUTH_RATE_LIMIT_SCOPE_HOOK_AUTH); | |
| if (!throttle.allowed) { | |
| const retryAfter = throttle.retryAfterMs > 0 ? Math.ceil(throttle.retryAfterMs / 1000) : 1; | |
| res.statusCode = 429; | |
| res.setHeader("Retry-After", String(retryAfter)); | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end("Too Many Requests"); | |
| logHooks.warn(`hook auth throttled for ${clientKey}; retry-after=${retryAfter}s`); | |
| return true; | |
| } | |
| hookAuthLimiter.recordFailure(clientKey, AUTH_RATE_LIMIT_SCOPE_HOOK_AUTH); | |
| res.statusCode = 401; | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end("Unauthorized"); | |
| return true; | |
| } | |
| hookAuthLimiter.reset(clientKey, AUTH_RATE_LIMIT_SCOPE_HOOK_AUTH); | |
| // The runtime state owns one resolved object per published hooks generation. | |
| // Object identity therefore acts as the generation token without re-running | |
| // mappings or comparing secret-bearing config values. | |
| const isHooksConfigCurrent = () => getHooksConfig() === hooksConfig; | |
| const rejectChangedHooksConfig = (): boolean => { | |
| if (isHooksConfigCurrent()) { | |
| return false; | |
| } | |
| sendJson(res, 409, { ok: false, error: HOOK_CONFIG_CHANGED_ERROR }); | |
| return true; | |
| }; | |
| const changedHooksConfigDispatchResult = (): HookAgentDispatchResult => ({ | |
| ok: false, | |
| statusCode: 409, | |
| error: HOOK_CONFIG_CHANGED_ERROR, | |
| }); | |
| const subPath = url.pathname.slice(basePath.length).replace(/^\/+/, ""); | |
| if (!subPath) { | |
| res.statusCode = 404; | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end("Not Found"); | |
| return true; | |
| } | |
| // gmail-path mappings carry a producer-derived bound (gog batch contract); | |
| // every other path keeps the shared default cap. | |
| const body = await readJsonBody(req, resolveHookPathBodyLimit(hooksConfig, subPath)); | |
| if (!body.ok) { | |
| const error = { ok: false, error: body.error }; | |
| if (body.error === "payload too large" || body.error === "request body timeout") { | |
| await sendHttpRequestRejection( | |
| req, | |
| res, | |
| body.error === "payload too large" ? 413 : 408, | |
| JSON.stringify(error), | |
| "application/json; charset=utf-8", | |
| ); | |
| } else { | |
| sendJson(res, 400, error); | |
| } | |
| return true; | |
| } | |
| if (rejectChangedHooksConfig()) { | |
| return true; | |
| } | |
| const payload = asRecord(body.value); | |
| const headers = normalizeHookHeaders(req); | |
| const idempotencyKey = resolveHookIdempotencyKey({ payload, headers }); | |
| // Later mapped validation errors must report any wake outcome that already occurred. | |
| let wakeResult: WakeResult | undefined; | |
| const sendHookError = (error: string) => | |
| sendJson(res, 400, { ok: false, error, ...wakeResult }); | |
| const resolveDispatchSessionKeyOrRespond = ( | |
| sessionKeyValue: string, | |
| targetAgentId: string, | |
| ): string | null => { | |
| const dispatchSessionKey = normalizeHookDispatchSessionKey({ | |
| sessionKey: sessionKeyValue, | |
| targetAgentId, | |
| }); | |
| const allowedPrefixes = hooksConfig.sessionPolicy.allowedSessionKeyPrefixes; | |
| if (allowedPrefixes && !isSessionKeyAllowedByPrefix(dispatchSessionKey, allowedPrefixes)) { | |
| sendHookError(getHookSessionKeyPrefixError(allowedPrefixes)); | |
| return null; | |
| } | |
| return dispatchSessionKey; | |
| }; | |
| const resolveTargetAgentOrRespond = ( | |
| agentId: string | undefined, | |
| source: "request" | "mapping", | |
| ): Extract<HookTargetAgentResolution, { ok: true }> | null => { | |
| const resolution = resolveEffectiveHookTargetAgentId(hooksConfig, agentId, source); | |
| if (!resolution.ok) { | |
| sendHookError(resolution.error); | |
| return null; | |
| } | |
| if (!isHookAgentAllowed(hooksConfig, resolution.effectiveAgentId)) { | |
| sendHookError(getHookAgentPolicyError()); | |
| return null; | |
| } | |
| return resolution; | |
| }; | |
| // Callers own the success response so mappings can dispatch several wakes first. | |
| const dispatchWake = ( | |
| value: Parameters<HookDispatchers["dispatchWakeHook"]>[0], | |
| targetAgentId: string, | |
| source: HookSessionKeySource, | |
| ): WakeResult | null => { | |
| let dispatchSessionKey: string | undefined; | |
| if (value.sessionKey) { | |
| const sessionKey = resolveHookSessionKey({ | |
| hooksConfig, | |
| source, | |
| sessionKey: value.sessionKey, | |
| }); | |
| if (!sessionKey.ok) { | |
| sendHookError(sessionKey.error); | |
| return null; | |
| } | |
| const resolvedSessionKey = resolveDispatchSessionKeyOrRespond( | |
| sessionKey.value, | |
| targetAgentId, | |
| ); | |
| if (resolvedSessionKey === null) { | |
| return null; | |
| } | |
| dispatchSessionKey = resolvedSessionKey; | |
| } | |
| const dispatchValue = { ...value, sessionKey: dispatchSessionKey }; | |
| if (rejectChangedHooksConfig()) { | |
| return null; | |
| } | |
| return dispatchWakeHook(dispatchValue, targetAgentId); | |
| }; | |
| if (subPath === "wake") { | |
| const normalized = normalizeWakePayload(payload); | |
| if (!normalized.ok) { | |
| sendJson(res, 400, { ok: false, error: normalized.error }); | |
| return true; | |
| } | |
| const target = resolveTargetAgentOrRespond(normalized.value.agentId, "request"); | |
| if (!target) { | |
| return true; | |
| } | |
| const directWakeResult = dispatchWake(normalized.value, target.effectiveAgentId, "request"); | |
| if (!directWakeResult) { | |
| return true; | |
| } | |
| sendJson(res, 200, { ok: true, mode: normalized.value.mode, ...directWakeResult }); | |
| return true; | |
| } | |
| if (subPath === "agent") { | |
| const waitForCompletion = payload.waitForCompletion; | |
| if (waitForCompletion !== undefined && typeof waitForCompletion !== "boolean") { | |
| sendJson(res, 400, { ok: false, error: "waitForCompletion must be boolean" }); | |
| return true; | |
| } | |
| const normalized = normalizeAgentPayload(payload); | |
| if (!normalized.ok) { | |
| sendJson(res, 400, { ok: false, error: normalized.error }); | |
| return true; | |
| } | |
| const target = resolveTargetAgentOrRespond(normalized.value.agentId, "request"); | |
| if (!target) { | |
| return true; | |
| } | |
| if (normalized.value.sessionMode === "persistent" && !normalized.value.sessionKey) { | |
| sendJson(res, 400, { | |
| ok: false, | |
| error: "sessionKey is required when sessionMode is persistent", | |
| }); | |
| return true; | |
| } | |
| const sessionKey = resolveHookSessionKey({ | |
| hooksConfig, | |
| source: "request", | |
| sessionKey: normalized.value.sessionKey, | |
| }); | |
| if (!sessionKey.ok) { | |
| sendJson(res, 400, { ok: false, error: sessionKey.error }); | |
| return true; | |
| } | |
| if ( | |
| normalized.value.sessionMode === "persistent" && | |
| !hooksConfig.sessionPolicy.allowedSessionKeyPrefixes?.length | |
| ) { | |
| sendJson(res, 400, { | |
| ok: false, | |
| error: | |
| "hooks.allowedSessionKeyPrefixes is required when direct hook sessionMode is persistent", | |
| }); | |
| return true; | |
| } | |
| const replayKey = buildHookReplayCacheKey({ | |
| pathKey: "agent", | |
| token, | |
| idempotencyKey, | |
| dispatchScope: { | |
| agentId: target.effectiveAgentId, | |
| sessionKey: | |
| normalized.value.sessionKey ?? hooksConfig.sessionPolicy.defaultSessionKey ?? null, | |
| message: normalized.value.message, | |
| name: normalized.value.name, | |
| wakeMode: normalized.value.wakeMode, | |
| sessionMode: normalized.value.sessionMode, | |
| deliver: normalized.value.deliver, | |
| channel: normalized.value.channel, | |
| to: normalized.value.to ?? null, | |
| accountId: normalized.value.accountId ?? null, | |
| model: normalized.value.model ?? null, | |
| thinking: normalized.value.thinking ?? null, | |
| timeoutSeconds: normalized.value.timeoutSeconds ?? null, | |
| }, | |
| }); | |
| const replay = resolveHookReplay(replayKey); | |
| if (replay) { | |
| await sendAgentResult(res, await replay, undefined, waitForCompletion === true); | |
| return true; | |
| } | |
| const dispatchSessionKey = resolveDispatchSessionKeyOrRespond( | |
| sessionKey.value, | |
| target.effectiveAgentId, | |
| ); | |
| if (dispatchSessionKey === null) { | |
| return true; | |
| } | |
| const dispatched = await dispatchAgentHookWithReplay(replayKey, () => { | |
| if (!isHooksConfigCurrent()) { | |
| return changedHooksConfigDispatchResult(); | |
| } | |
| return dispatchAgentHook({ | |
| ...normalized.value, | |
| effectiveAgentId: target.effectiveAgentId, | |
| idempotencyKey, | |
| sessionKey: dispatchSessionKey, | |
| sourcePath: `${basePath}/agent`, | |
| agentId: target.selectedAgentId, | |
| externalContentSource: "webhook", | |
| }); | |
| }); | |
| await sendAgentResult(res, dispatched, undefined, waitForCompletion === true); | |
| return true; | |
| } | |
| if (hooksConfig.mappings.length > 0) { | |
| try { | |
| const mapped = await applyHookMappings(hooksConfig.mappings, { | |
| payload, | |
| headers, | |
| url, | |
| path: subPath, | |
| }); | |
| if (rejectChangedHooksConfig()) { | |
| return true; | |
| } | |
| if (mapped) { | |
| if (!mapped.ok) { | |
| sendJson(res, 400, { ok: false, error: mapped.error }); | |
| return true; | |
| } | |
| if (mapped.dropped > 0) { | |
| logHooks.warn( | |
| `hook mapping ${subPath} fan-out dropped ${mapped.dropped} items beyond the ${HOOK_MAPPING_FAN_OUT_MAX_ITEMS}-item cap`, | |
| ); | |
| } | |
| if (mapped.actions.length === 0) { | |
| if (mapped.fanout) { | |
| logHooks.info(`hook mapping ${subPath} matched with no items to dispatch`); | |
| } | |
| res.statusCode = 204; | |
| res.end(); | |
| return true; | |
| } | |
| // Within-batch duplicates: content identity alone would collapse two | |
| // identical rendered items into one dispatch while the response | |
| // claims both ran. Numbering repeated scopes keeps one replay entry | |
| // per occurrence, and identical redeliveries renumber identically. | |
| const fanOutScopeOccurrences = new Map<string, number>(); | |
| // Resolves policy for one mapped agent action and returns its | |
| // dispatch closure; a null return means an error response was sent. | |
| const prepareMappedAgentDispatchOrRespond = ( | |
| action: Extract<(typeof mapped.actions)[number], { kind: "agent" }>, | |
| ): (() => HookAgentDispatchResult | Promise<HookAgentDispatchResult>) | null => { | |
| const channel = resolveHookChannel(action.channel); | |
| if (!channel) { | |
| sendHookError(getHookChannelError()); | |
| return null; | |
| } | |
| const deliver = resolveHookDeliver(action.deliver); | |
| const delivery = deliver | |
| ? { mode: "announce" as const, channel, to: action.to } | |
| : { mode: "none" as const }; | |
| const target = resolveTargetAgentOrRespond(action.agentId, "mapping"); | |
| if (!target) { | |
| return null; | |
| } | |
| if ( | |
| action.sessionMode === "persistent" && | |
| !action.sessionKey && | |
| !hooksConfig.sessionPolicy.defaultSessionKey | |
| ) { | |
| sendHookError( | |
| "sessionKey or hooks.defaultSessionKey is required when mapped hook sessionMode is persistent", | |
| ); | |
| return null; | |
| } | |
| const sessionKey = resolveHookSessionKey({ | |
| hooksConfig, | |
| source: action.sessionKeySource === "static" ? "mapping-static" : "mapping-templated", | |
| sessionKey: action.sessionKey, | |
| }); | |
| if (!sessionKey.ok) { | |
| sendHookError(sessionKey.error); | |
| return null; | |
| } | |
| const dispatchSessionKey = resolveDispatchSessionKeyOrRespond( | |
| sessionKey.value, | |
| target.effectiveAgentId, | |
| ); | |
| if (dispatchSessionKey === null) { | |
| return null; | |
| } | |
| const dispatchScope: Record<string, unknown> = { | |
| agentId: target.effectiveAgentId, | |
| sessionKey: action.sessionKey ?? hooksConfig.sessionPolicy.defaultSessionKey ?? null, | |
| message: action.message, | |
| name: action.name ?? "Hook", | |
| wakeMode: action.wakeMode, | |
| sessionMode: action.sessionMode, | |
| deliver, | |
| channel, | |
| to: action.to ?? null, | |
| model: action.model ?? null, | |
| thinking: action.thinking ?? null, | |
| timeoutSeconds: action.timeoutSeconds ?? null, | |
| }; | |
| if (mapped.fanout) { | |
| const fingerprint = JSON.stringify(dispatchScope); | |
| const occurrence = fanOutScopeOccurrences.get(fingerprint) ?? 0; | |
| fanOutScopeOccurrences.set(fingerprint, occurrence + 1); | |
| dispatchScope.occurrence = occurrence; | |
| } | |
| const replayKey = buildHookReplayCacheKey({ | |
| pathKey: subPath || "mapping", | |
| token, | |
| // Fan-out producers (gog gmail) send no idempotency key, yet a | |
| // non-2xx batch response makes them redeliver the same batch. | |
| // Deriving item identity from the dispatch scope lets retries | |
| // replay already-dispatched items instead of duplicating them. | |
| idempotencyKey: mapped.fanout | |
| ? (idempotencyKey ?? HOOK_FAN_OUT_DERIVED_IDEMPOTENCY) | |
| : idempotencyKey, | |
| dispatchScope, | |
| }); | |
| return () => | |
| dispatchAgentHookWithReplay(replayKey, () => { | |
| if (!isHooksConfigCurrent()) { | |
| return changedHooksConfigDispatchResult(); | |
| } | |
| return dispatchAgentHook({ | |
| message: action.message, | |
| name: action.name ?? "Hook", | |
| idempotencyKey, | |
| agentId: target.selectedAgentId, | |
| effectiveAgentId: target.effectiveAgentId, | |
| wakeMode: action.wakeMode, | |
| sessionKey: dispatchSessionKey, | |
| sessionMode: action.sessionMode, | |
| sourcePath: `${basePath}/${subPath}`, | |
| deliver, | |
| channel, | |
| to: action.to, | |
| delivery, | |
| model: action.model, | |
| thinking: action.thinking, | |
| timeoutSeconds: action.timeoutSeconds, | |
| mappingId: action.mappingId, | |
| allowUnsafeExternalContent: action.allowUnsafeExternalContent, | |
| ...(mapped.fanout ? { admissionMode: "background" as const } : {}), | |
| externalContentSource: resolveMappedHookExternalContentSource({ | |
| subPath, | |
| sessionKey: sessionKey.value, | |
| }), | |
| }); | |
| }); | |
| }; | |
| // One pass over every action so a per-item transform emitting mixed | |
| // kinds loses nothing: wakes dispatch immediately (no replay | |
| // identity — a producer retry after a partial agent failure | |
| // dispatches them again), agents collect for dispatch. | |
| const dispatches: Array< | |
| () => HookAgentDispatchResult | Promise<HookAgentDispatchResult> | |
| > = []; | |
| let wakeMode: "now" | "next-heartbeat" | undefined; | |
| for (const action of mapped.actions) { | |
| if (action.kind === "wake") { | |
| const target = resolveTargetAgentOrRespond(action.agentId, "mapping"); | |
| if (!target) { | |
| return true; | |
| } | |
| const dispatched = dispatchWake( | |
| { text: action.text, mode: action.mode, sessionKey: action.sessionKey }, | |
| target.effectiveAgentId, | |
| action.sessionKeySource === "static" ? "mapping-static" : "mapping-templated", | |
| ); | |
| if (!dispatched) { | |
| return true; | |
| } | |
| if (!wakeResult || dispatched.eventOutcome === "queued") { | |
| wakeResult = dispatched; | |
| } | |
| wakeMode = action.mode; | |
| continue; | |
| } | |
| const prepared = prepareMappedAgentDispatchOrRespond(action); | |
| if (!prepared) { | |
| return true; | |
| } | |
| dispatches.push(prepared); | |
| } | |
| if (dispatches.length === 0) { | |
| sendJson(res, 200, { ok: true, mode: wakeMode ?? "now", ...wakeResult }); | |
| return true; | |
| } | |
| if (!mapped.fanout) { | |
| // Non-fanout mappings produce exactly one action. | |
| const dispatched = await dispatches[0]!(); | |
| void sendAgentResult(res, dispatched, wakeResult); | |
| return true; | |
| } | |
| const settled = await settleFanOutDispatches( | |
| dispatches.map((dispatch) => Promise.resolve(dispatch())), | |
| fanoutResponseDeadlineMs, | |
| ); | |
| sendFanOutResult(res, settled, wakeResult); | |
| return true; | |
| } | |
| } catch (err) { | |
| logHooks.warn(`hook mapping failed: ${String(err)}`); | |
| sendJson(res, 500, { ok: false, error: "hook mapping failed" }); | |
| return true; | |
| } | |
| } | |
| res.statusCode = 404; | |
| res.setHeader("Content-Type", "text/plain; charset=utf-8"); | |
| res.end("Not Found"); | |
| return true; | |
| }; | |
| } | |