mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-21 01:51:39 -06:00
0b8aabe864
* docs: document markdown marker renderer * docs: document rendered markdown chunking * docs: document markdown text chunking * docs: document shared text chunking * docs: document plugin text chunking exports * docs: document avatar policy constants * docs: document node match candidates * docs: document scoped expiring id cache * docs: document runtime import normalization * docs: document string sample summaries * docs: document session usage timeseries types * docs: document session usage response types * docs: document manifest frontmatter shapes * docs: document channel route input metadata * docs: document pair loop guard settings * docs: document migration config patch helpers * docs: document api provider registry * docs: document tool call repair payloads * docs: document plugin tool payload helpers * docs: document lazy promise loader * docs: document store writer queue state * docs: document thread binding lifecycle * docs: document concurrency helper contract * docs: document gateway client info contract * docs: document delivery context contracts * docs: document secret ref defaults contract * docs: document command gating contract * docs: document avatar policy contract * docs: document node match policy * docs: document message channel normalization * docs: document boolean parsing contract * docs: document zod parse helpers * docs: document direct dm guard policy * docs: document fixed window limiter contract * docs: document node presence event contract * docs: document secret normalization contract * docs: document progress draft line removal * docs: document usage formatting contracts * docs: document agent run status contract * docs: document runtime import helpers * docs: document provider utility ownership * docs: document invalid config helpers * docs: document json compat parser * docs: document channel config metadata ownership * docs: document channel logging helpers * docs: document sender identity validation ownership * docs: document string sampling helper * docs: document global singleton helpers * docs: document transcript tool helpers * docs: document exec safe-bin normalization * docs: document reaction level resolver * docs: document account snapshot redaction boundary * docs: document messaging target helpers * docs: document thread binding messages * docs: document conversation binding context * docs: document conversation resolution helper * docs: document owner display secret retention * docs: document provider request config types * docs: document skills config types * docs: document memory config types * docs: document imessage config types * docs: document crestodian config types * docs: document tools config policies * docs: document shared config base types * docs: document channel config contracts * docs: document openclaw config state types * docs: document model config contracts * docs: document shared agent config types * docs: document agent defaults config types * docs: document secret input contracts * docs: document auth config contracts * docs: document gateway config contracts * docs: document tool call stream repair contracts * docs: document memory host facades * docs: document llm core contracts * docs: document markdown core contracts * docs: document gateway connect error contracts * docs: document gateway protocol primitives * docs: document gateway frame schemas * docs: document gateway device schemas * docs: document gateway environment schemas * docs: document gateway push schemas * docs: document gateway plugin schemas * docs: document gateway artifact schemas * docs: document gateway command schemas * docs: document gateway task schemas * docs: document gateway exec approval schemas * docs: document gateway secret schemas * docs: document gateway config schemas * docs: document gateway snapshot schemas * docs: document gateway chat schemas * docs: document gateway wizard schemas * docs: document gateway node schemas * docs: document gateway plugin approval schemas * docs: document gateway talk schemas * docs: document gateway agent schemas * docs: document gateway session schemas * docs: document gateway cron schemas * docs: document gateway agent model skill schemas * docs: document gateway skill proposal tool schemas * docs: document gateway protocol registry * docs: document gateway channel status schemas * docs: document gateway schema regression tests * docs: document gateway schema barrel * docs: document gateway validator tests * docs: document gateway primitive push tests * docs: document gateway contract tests * docs: document native protocol guard * docs: document channel schema tests * docs: document gateway protocol smoke tests * docs: document gateway protocol entrypoint * docs: document gateway protocol type exports * docs: document gateway error codes * docs: document protocol schema registry * docs: document talk audio codec * docs: document talk activation names * docs: document talk consult questions * docs: document talk consult tool * docs: document talk run control contracts * docs: document talk run control adapter * docs: document talkback consult queue * docs: document talk consult transcript guard * docs: document talk fast context runtime * docs: document forced talk consult coordinator * docs: document talk output activity tracker * docs: document talk event metrics * docs: document talk diagnostics * docs: document talk observability hook * docs: document talk provider resolver * docs: document talk provider registry * docs: document talk runtime primitives * docs: document talk consult controller logs * docs: document channel identity helpers * docs: document channel account allowlist helpers * docs: document channel metadata draft controls * docs: document channel ingress policy * docs: document channel sender access gates * docs: document channel catalog message contracts * docs: document channel account plugin helpers * docs: document configured binding helpers * docs: document channel acp approval config helpers * docs: document channel bundled config write helpers * docs: document channel plugin utility contracts * docs: document channel config access helpers * docs: document channel message action helpers * docs: document channel outbound runtime helpers * docs: document channel pairing promotion helpers * docs: document channel registry helpers * docs: document channel setup wizard helpers * docs: document channel lifecycle status helpers * docs: document channel target thread helpers * docs: document channel session binding helpers * docs: document channel package module probes * docs: document channel setup wizard contracts * docs: document channel plugin API barrels * docs: document channel contract test helpers * docs: document channel core helpers * docs: document small core facades * docs: document provider runtime helpers * docs: document persistence and realtime helpers * docs: document mcp and state helpers * docs: document tool planner contracts * docs: document music generation runtime * docs: document crestodian command flow * docs: document utility helpers * docs: document node host helpers * docs: document transcript contracts * docs: document trajectory export contracts * docs: document image generation contracts * docs: document routing helper contracts * docs: document session helper contracts * docs: document video generation contracts * docs: document model catalog contracts * docs: document proxy capture contracts * docs: document status rendering contracts * docs: document test helper contracts * docs: document wizard setup contracts * docs: document process contracts * docs: document memory host sdk contracts * docs: document tts contracts * docs: document secrets runtime contracts * docs: document shared helper contracts * docs: document hook runtime contracts * docs: document security audit contracts * docs: document flow contracts * docs: document media understanding contracts * docs: document tui contracts * docs: document logging contracts * docs: document llm contracts * docs: document cron contracts * docs: document daemon contracts * docs: document task contracts * docs: document acp contracts * docs: document test utility contracts * docs: document skill contracts * docs: document config contracts * docs: document outbound infra contracts * docs: document command analysis contracts * docs: document provider usage infra contracts * docs: document file safety infra contracts * docs: document exec approval infra contracts * docs: document gateway runtime infra contracts * docs: document infra utility contracts * docs: document infra queue storage contracts * docs: document heartbeat infra contracts * docs: document remaining infra contracts * docs: document gateway auth contracts * docs: document gateway display helpers * docs: document gateway http helpers * docs: document gateway node helpers * docs: document gateway mcp helpers * docs: document gateway support helpers * docs: document gateway server runtime helpers * docs: document gateway runtime bootstrap helpers * docs: document gateway session events * docs: document gateway utility helpers * docs: document gateway talk helpers * docs: document gateway helper contracts * docs: document gateway server method helpers * docs: document gateway server auth helpers * docs: document gateway server tests * docs: document gateway test helpers * docs: document gateway node tests * docs: document gateway channel tests * docs: document gateway session tests * docs: document gateway server startup tests * docs: document gateway tool test helpers * docs: document gateway server test helpers * docs: document gateway server method tests * docs: document remaining gateway tests * docs: document plugin sdk public subpaths * docs: document plugin sdk runtime helpers * docs: document plugin sdk memory provider helpers * docs: document plugin sdk runtime facades * docs: document plugin sdk command approval helpers * docs: document plugin sdk runtime types * docs: document plugin sdk browser account helpers * docs: document plugin sdk media memory helpers * docs: document plugin sdk core tests * docs: document plugin sdk contract helpers * docs: document plugin sdk test helpers * docs: document remaining plugin sdk tests * docs: document cli utility helpers * docs: document cli runtime helpers * docs: document cli command registration helpers * docs: document node cli helpers * docs: document cli program registration * docs: document message cli registration * docs: document daemon cli helpers * docs: document cli route parsers
340 lines
13 KiB
TypeScript
340 lines
13 KiB
TypeScript
import type { IncomingMessage, ServerResponse } from "node:http";
|
|
import { registerPluginHttpRoute } from "../plugins/http-registry.js";
|
|
import type { FixedWindowRateLimiter } from "./webhook-memory-guards.js";
|
|
import { normalizeWebhookPath } from "./webhook-path.js";
|
|
import {
|
|
beginWebhookRequestPipelineOrReject,
|
|
type WebhookInFlightLimiter,
|
|
} from "./webhook-request-guards.js";
|
|
|
|
/** Registration handle returned for one live webhook target. */
|
|
export type RegisteredWebhookTarget<T> = {
|
|
/** Normalized target stored in the caller-owned path registry. */
|
|
target: T;
|
|
/** Idempotently remove this target and run path teardown when it was the last target. */
|
|
unregister: () => void;
|
|
};
|
|
|
|
/** Lifecycle hooks for path-level webhook target registration. */
|
|
export type RegisterWebhookTargetOptions<T extends { path: string }> = {
|
|
/** Called before the first target for a normalized path is stored; may return path teardown. */
|
|
onFirstPathTarget?: (params: { path: string; target: T }) => void | (() => void);
|
|
/** Called after the last target for a normalized path has been removed. */
|
|
onLastPathTargetRemoved?: (params: { path: string }) => void;
|
|
};
|
|
|
|
type RegisterPluginHttpRouteParams = Parameters<typeof registerPluginHttpRoute>[0];
|
|
|
|
export { registerPluginHttpRoute };
|
|
|
|
/** Plugin HTTP route options supplied when webhook paths are registered lazily. */
|
|
export type RegisterWebhookPluginRouteOptions = Omit<
|
|
RegisterPluginHttpRouteParams,
|
|
"path" | "fallbackPath"
|
|
>;
|
|
|
|
/** Register a webhook target and lazily install the matching plugin HTTP route on first use. */
|
|
export function registerWebhookTargetWithPluginRoute<T extends { path: string }>(params: {
|
|
/** Caller-owned normalized path registry shared by all targets for this plugin/runtime. */
|
|
targetsByPath: Map<string, T[]>;
|
|
/** Target to normalize, store, and later return from the registration handle. */
|
|
target: T;
|
|
/** Plugin HTTP route configuration used when the first target for a path is registered. */
|
|
route: RegisterWebhookPluginRouteOptions;
|
|
/** Optional last-target hook forwarded to `registerWebhookTarget`. */
|
|
onLastPathTargetRemoved?: RegisterWebhookTargetOptions<T>["onLastPathTargetRemoved"];
|
|
}): RegisteredWebhookTarget<T> {
|
|
return registerWebhookTarget(params.targetsByPath, params.target, {
|
|
onFirstPathTarget: ({ path }) =>
|
|
registerPluginHttpRoute({
|
|
...params.route,
|
|
path,
|
|
// Webhook targets own this path while registered; default replacement lets
|
|
// plugin reload/setup refresh the handler without accumulating stale routes.
|
|
replaceExisting: params.route.replaceExisting ?? true,
|
|
}),
|
|
onLastPathTargetRemoved: params.onLastPathTargetRemoved,
|
|
});
|
|
}
|
|
|
|
const pathTeardownByTargetMap = new WeakMap<Map<string, unknown[]>, Map<string, () => void>>();
|
|
|
|
function getPathTeardownMap<T>(targetsByPath: Map<string, T[]>): Map<string, () => void> {
|
|
const mapKey = targetsByPath as unknown as Map<string, unknown[]>;
|
|
const existing = pathTeardownByTargetMap.get(mapKey);
|
|
if (existing) {
|
|
return existing;
|
|
}
|
|
const created = new Map<string, () => void>();
|
|
// Teardown is scoped to the caller-owned registry map so independent plugins using the same
|
|
// path do not unregister each other's HTTP routes.
|
|
pathTeardownByTargetMap.set(mapKey, created);
|
|
return created;
|
|
}
|
|
|
|
/** Add a normalized target to a path bucket and clean up route state when the last target leaves. */
|
|
export function registerWebhookTarget<T extends { path: string }>(
|
|
targetsByPath: Map<string, T[]>,
|
|
target: T,
|
|
opts?: RegisterWebhookTargetOptions<T>,
|
|
): RegisteredWebhookTarget<T> {
|
|
const key = normalizeWebhookPath(target.path);
|
|
const normalizedTarget = { ...target, path: key };
|
|
const existing = targetsByPath.get(key) ?? [];
|
|
|
|
if (existing.length === 0) {
|
|
const onFirstPathResult = opts?.onFirstPathTarget?.({
|
|
path: key,
|
|
target: normalizedTarget,
|
|
});
|
|
if (typeof onFirstPathResult === "function") {
|
|
getPathTeardownMap(targetsByPath).set(key, onFirstPathResult);
|
|
}
|
|
}
|
|
|
|
targetsByPath.set(key, [...existing, normalizedTarget]);
|
|
|
|
let isActive = true;
|
|
const unregister = () => {
|
|
if (!isActive) {
|
|
return;
|
|
}
|
|
isActive = false;
|
|
|
|
const updated = (targetsByPath.get(key) ?? []).filter((entry) => entry !== normalizedTarget);
|
|
if (updated.length > 0) {
|
|
targetsByPath.set(key, updated);
|
|
return;
|
|
}
|
|
targetsByPath.delete(key);
|
|
|
|
const teardown = getPathTeardownMap(targetsByPath).get(key);
|
|
if (teardown) {
|
|
getPathTeardownMap(targetsByPath).delete(key);
|
|
teardown();
|
|
}
|
|
opts?.onLastPathTargetRemoved?.({ path: key });
|
|
};
|
|
return { target: normalizedTarget, unregister };
|
|
}
|
|
|
|
/** Resolve all registered webhook targets for the incoming request path. */
|
|
export function resolveWebhookTargets<T>(
|
|
req: IncomingMessage,
|
|
targetsByPath: Map<string, T[]>,
|
|
): { path: string; targets: T[] } | null {
|
|
const url = new URL(req.url ?? "/", "http://localhost");
|
|
const path = normalizeWebhookPath(url.pathname);
|
|
const targets = targetsByPath.get(path);
|
|
if (!targets || targets.length === 0) {
|
|
return null;
|
|
}
|
|
return { path, targets };
|
|
}
|
|
|
|
/** Run common webhook guards, then dispatch only when the request path resolves to live targets. */
|
|
export async function withResolvedWebhookRequestPipeline<T>(params: {
|
|
/** Incoming HTTP request whose pathname selects the target bucket. */
|
|
req: IncomingMessage;
|
|
/** HTTP response used by guard failures before handler dispatch. */
|
|
res: ServerResponse;
|
|
/** Caller-owned target registry keyed by normalized webhook path. */
|
|
targetsByPath: Map<string, T[]>;
|
|
/** Allowed methods for the common request guard. */
|
|
allowMethods?: readonly string[];
|
|
/** Optional per-key fixed-window limiter shared across requests. */
|
|
rateLimiter?: FixedWindowRateLimiter;
|
|
/** Explicit rate-limit key; defaults are owned by the request guard. */
|
|
rateLimitKey?: string;
|
|
/** Clock override for deterministic limiter tests. */
|
|
nowMs?: number;
|
|
/** Require JSON content type before dispatching to the webhook handler. */
|
|
requireJsonContentType?: boolean;
|
|
/** Optional in-flight limiter to cap concurrent handling for a key. */
|
|
inFlightLimiter?: WebhookInFlightLimiter;
|
|
/** Explicit or derived key for concurrent request limiting. */
|
|
inFlightKey?: string | ((args: { req: IncomingMessage; path: string; targets: T[] }) => string);
|
|
/** Status code returned when the in-flight guard rejects. */
|
|
inFlightLimitStatusCode?: number;
|
|
/** Response body returned when the in-flight guard rejects. */
|
|
inFlightLimitMessage?: string;
|
|
/** Handler invoked only after target resolution and common guards succeed. */
|
|
handle: (args: { path: string; targets: T[] }) => Promise<boolean | void> | boolean | void;
|
|
}): Promise<boolean> {
|
|
const resolved = resolveWebhookTargets(params.req, params.targetsByPath);
|
|
if (!resolved) {
|
|
return false;
|
|
}
|
|
|
|
const inFlightKey =
|
|
typeof params.inFlightKey === "function"
|
|
? params.inFlightKey({ req: params.req, path: resolved.path, targets: resolved.targets })
|
|
: (params.inFlightKey ?? `${resolved.path}:${params.req.socket?.remoteAddress ?? "unknown"}`);
|
|
const requestLifecycle = beginWebhookRequestPipelineOrReject({
|
|
req: params.req,
|
|
res: params.res,
|
|
allowMethods: params.allowMethods,
|
|
rateLimiter: params.rateLimiter,
|
|
rateLimitKey: params.rateLimitKey,
|
|
nowMs: params.nowMs,
|
|
requireJsonContentType: params.requireJsonContentType,
|
|
inFlightLimiter: params.inFlightLimiter,
|
|
inFlightKey,
|
|
inFlightLimitStatusCode: params.inFlightLimitStatusCode,
|
|
inFlightLimitMessage: params.inFlightLimitMessage,
|
|
});
|
|
if (!requestLifecycle.ok) {
|
|
return true;
|
|
}
|
|
|
|
try {
|
|
await params.handle(resolved);
|
|
return true;
|
|
} finally {
|
|
// Release even when the handler throws; otherwise one failed webhook can pin the in-flight
|
|
// slot and permanently reject later deliveries for the same key.
|
|
requestLifecycle.release();
|
|
}
|
|
}
|
|
|
|
/** Result of matching a request against zero, one, or multiple webhook targets. */
|
|
export type WebhookTargetMatchResult<T> =
|
|
| { kind: "none" }
|
|
| { kind: "single"; target: T }
|
|
| { kind: "ambiguous" };
|
|
|
|
function updateMatchedWebhookTarget<T>(
|
|
matched: T | undefined,
|
|
target: T,
|
|
): { ok: true; matched: T } | { ok: false; result: WebhookTargetMatchResult<T> } {
|
|
if (matched) {
|
|
return { ok: false, result: { kind: "ambiguous" } };
|
|
}
|
|
return { ok: true, matched: target };
|
|
}
|
|
|
|
function finalizeMatchedWebhookTarget<T>(matched: T | undefined): WebhookTargetMatchResult<T> {
|
|
if (!matched) {
|
|
return { kind: "none" };
|
|
}
|
|
return { kind: "single", target: matched };
|
|
}
|
|
|
|
/** Match exactly one synchronous target or report whether resolution was empty or ambiguous. */
|
|
export function resolveSingleWebhookTarget<T>(
|
|
targets: readonly T[],
|
|
isMatch: (target: T) => boolean,
|
|
): WebhookTargetMatchResult<T> {
|
|
let matched: T | undefined;
|
|
for (const target of targets) {
|
|
if (!isMatch(target)) {
|
|
continue;
|
|
}
|
|
// Stop at the second match so auth callers can reject ambiguous secrets without inspecting
|
|
// or accidentally selecting a later target.
|
|
const updated = updateMatchedWebhookTarget(matched, target);
|
|
if (!updated.ok) {
|
|
return updated.result;
|
|
}
|
|
matched = updated.matched;
|
|
}
|
|
return finalizeMatchedWebhookTarget(matched);
|
|
}
|
|
|
|
/** Async variant of single-target resolution for auth checks that need I/O. */
|
|
export async function resolveSingleWebhookTargetAsync<T>(
|
|
targets: readonly T[],
|
|
isMatch: (target: T) => Promise<boolean>,
|
|
): Promise<WebhookTargetMatchResult<T>> {
|
|
let matched: T | undefined;
|
|
for (const target of targets) {
|
|
if (!(await isMatch(target))) {
|
|
continue;
|
|
}
|
|
const updated = updateMatchedWebhookTarget(matched, target);
|
|
if (!updated.ok) {
|
|
return updated.result;
|
|
}
|
|
matched = updated.matched;
|
|
}
|
|
return finalizeMatchedWebhookTarget(matched);
|
|
}
|
|
|
|
/** Resolve an authorized target and send the standard unauthorized or ambiguous response on failure. */
|
|
export async function resolveWebhookTargetWithAuthOrReject<T>(params: {
|
|
/** Candidate targets for the already-resolved webhook path. */
|
|
targets: readonly T[];
|
|
/** HTTP response used to send unauthorized or ambiguous failures. */
|
|
res: ServerResponse;
|
|
/** Auth or routing predicate; exactly one target must match. */
|
|
isMatch: (target: T) => boolean | Promise<boolean>;
|
|
/** Status code for no matching target. Defaults to 401. */
|
|
unauthorizedStatusCode?: number;
|
|
/** Response body for no matching target. */
|
|
unauthorizedMessage?: string;
|
|
/** Status code for multiple matching targets. Defaults to 401. */
|
|
ambiguousStatusCode?: number;
|
|
/** Response body for multiple matching targets. */
|
|
ambiguousMessage?: string;
|
|
}): Promise<T | null> {
|
|
const match = await resolveSingleWebhookTargetAsync(params.targets, async (target) =>
|
|
params.isMatch(target),
|
|
);
|
|
return resolveWebhookTargetMatchOrReject(params, match);
|
|
}
|
|
|
|
/** Synchronous variant of webhook auth resolution for cheap in-memory match checks. */
|
|
export function resolveWebhookTargetWithAuthOrRejectSync<T>(params: {
|
|
/** Candidate targets for the already-resolved webhook path. */
|
|
targets: readonly T[];
|
|
/** HTTP response used to send unauthorized or ambiguous failures. */
|
|
res: ServerResponse;
|
|
/** Synchronous auth or routing predicate; exactly one target must match. */
|
|
isMatch: (target: T) => boolean;
|
|
/** Status code for no matching target. Defaults to 401. */
|
|
unauthorizedStatusCode?: number;
|
|
/** Response body for no matching target. */
|
|
unauthorizedMessage?: string;
|
|
/** Status code for multiple matching targets. Defaults to 401. */
|
|
ambiguousStatusCode?: number;
|
|
/** Response body for multiple matching targets. */
|
|
ambiguousMessage?: string;
|
|
}): T | null {
|
|
const match = resolveSingleWebhookTarget(params.targets, params.isMatch);
|
|
return resolveWebhookTargetMatchOrReject(params, match);
|
|
}
|
|
|
|
function resolveWebhookTargetMatchOrReject<T>(
|
|
params: {
|
|
res: ServerResponse;
|
|
unauthorizedStatusCode?: number;
|
|
unauthorizedMessage?: string;
|
|
ambiguousStatusCode?: number;
|
|
ambiguousMessage?: string;
|
|
},
|
|
match: WebhookTargetMatchResult<T>,
|
|
): T | null {
|
|
if (match.kind === "single") {
|
|
return match.target;
|
|
}
|
|
if (match.kind === "ambiguous") {
|
|
params.res.statusCode = params.ambiguousStatusCode ?? 401;
|
|
params.res.end(params.ambiguousMessage ?? "ambiguous webhook target");
|
|
return null;
|
|
}
|
|
params.res.statusCode = params.unauthorizedStatusCode ?? 401;
|
|
params.res.end(params.unauthorizedMessage ?? "unauthorized");
|
|
return null;
|
|
}
|
|
|
|
/** Reject non-POST webhook requests with the conventional Allow header. */
|
|
export function rejectNonPostWebhookRequest(req: IncomingMessage, res: ServerResponse): boolean {
|
|
if (req.method === "POST") {
|
|
return false;
|
|
}
|
|
res.statusCode = 405;
|
|
res.setHeader("Allow", "POST");
|
|
res.end("Method Not Allowed");
|
|
return true;
|
|
}
|