mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
fix(feishu): preserve webhooks on supported 9.6 hosts (#160455)
* fix(feishu): preserve webhooks on supported 9.6 hosts * docs(feishu): clarify legacy webhook update order * test(feishu): retain webhook fixture port claims
This commit is contained in:
parent
5c30c9afcb
commit
dba4c4f0e1
4 changed files with 345 additions and 13 deletions
|
|
@ -103,19 +103,31 @@ to the same plugin route and signature verifier. Set
|
|||
An omitted object `host` binds to `127.0.0.1`; explicit hosts, including wildcard
|
||||
addresses, are preserved. Account entries inherit the root setting, and
|
||||
`accounts.<id>.legacyWebhook: false` disables forwarding for that account.
|
||||
A shared legacy socket stays open while another account still uses that endpoint.
|
||||
On supported 2026.9.6 hosts that predate Gateway-owned forwarding, Feishu keeps
|
||||
an account-owned compatibility listener at that endpoint, using the same
|
||||
signature checks and dispatch path. Those hosts require distinct legacy endpoints
|
||||
for separate accounts. Newer hosts keep listener ownership in the Gateway, where
|
||||
a shared legacy socket stays open while another account still uses that endpoint.
|
||||
On account shutdown, authenticated responses may finish for up to five seconds,
|
||||
matching the previous listener's close grace period. Unfinished responses close
|
||||
at that deadline; other accounts keep their routes and listeners.
|
||||
During that grace period, correctly signed callbacks for the stopping account
|
||||
receive a retryable `503` unless a live successor already accepts their signature.
|
||||
|
||||
On update, the plugin's Doctor migration moves `webhookPort` and `webhookHost`
|
||||
The plugin's Doctor migration moves `webhookPort` and `webhookHost`
|
||||
into `legacyWebhook: { port, host }`, preserving the effective old defaults when
|
||||
only one key was set. Doctor's normal config backup protects the original
|
||||
only one key was set. The normal config backup protects the original
|
||||
settings. Existing canonical `legacyWebhook` settings, including `false`, win.
|
||||
An install that omitted both old settings keeps receiving traffic on port `3000`.
|
||||
|
||||
When updating from a 2026.9.6 host with these old keys, first update OpenClaw core
|
||||
to a release containing the [plugin-update migration repair](https://github.com/openclaw/openclaw/pull/160682).
|
||||
Then explicitly update any pinned Feishu package to your chosen release. The core
|
||||
updater preserves explicit plugin version pins. The updated installer applies
|
||||
Feishu's migration before activating the replacement package; the published
|
||||
2026.9.6 installer rejects the new schema before it can run that repair. A refused
|
||||
plugin-only update leaves the previous installation and settings intact.
|
||||
|
||||
The deprecated TypeScript `webhookPort` and `webhookHost` input fields remain
|
||||
source-compatible until the next Plugin SDK major. Runtime config uses
|
||||
`legacyWebhook`; run `openclaw doctor --fix` to migrate the old keys.
|
||||
|
|
|
|||
|
|
@ -9,7 +9,6 @@ import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
|
|||
import { truncateUtf16Safe } from "openclaw/plugin-sdk/text-utility-runtime";
|
||||
import {
|
||||
applyBasicWebhookRequestGuards,
|
||||
getWebhookLegacyListener,
|
||||
resolveRequestClientIp,
|
||||
} from "openclaw/plugin-sdk/webhook-ingress";
|
||||
import {
|
||||
|
|
@ -39,6 +38,7 @@ import {
|
|||
recordWebhookStatus,
|
||||
wsClients,
|
||||
} from "./monitor.state.js";
|
||||
import { feishuWebhookHost, startFeishuLegacyWebhookListener } from "./monitor.webhook-legacy.js";
|
||||
import type { ResolvedFeishuAccount } from "./types.js";
|
||||
import { DEFAULT_FEISHU_WEBHOOK_PATH, normalizeFeishuWebhookPath } from "./webhook-path.js";
|
||||
import {
|
||||
|
|
@ -401,10 +401,10 @@ async function handleFeishuWebhook(
|
|||
req: http.IncomingMessage,
|
||||
res: http.ServerResponse,
|
||||
webhookTargets: Map<string, FeishuWebhookTarget[]>,
|
||||
legacyListener = feishuWebhookHost.getWebhookLegacyListener?.(req),
|
||||
): Promise<void> {
|
||||
const requestUrl = req.url ?? "/";
|
||||
const requestPath = requestUrl.split("?", 1)[0];
|
||||
const legacyListener = getWebhookLegacyListener(req);
|
||||
const targets = (
|
||||
webhookTargets.get(canonicalizeWebhookRouteKey(requestPath ?? "/")) ?? []
|
||||
).filter(
|
||||
|
|
@ -609,8 +609,13 @@ async function handleFeishuWebhook(
|
|||
}
|
||||
|
||||
export async function monitorWebhook(params: MonitorTransportParams): Promise<void> {
|
||||
const { account, accountId, runtime, abortSignal, statusSink } = params;
|
||||
const { account, accountId, runtime, statusSink } = params;
|
||||
const stopped = new AbortController();
|
||||
const abortSignal = params.abortSignal
|
||||
? AbortSignal.any([params.abortSignal, stopped.signal])
|
||||
: stopped.signal;
|
||||
const legacyListener = resolveFeishuLegacyWebhookListener(account.config);
|
||||
const gatewayOwnsLegacyListeners = feishuWebhookHost.getWebhookLegacyListener !== undefined;
|
||||
const encryptKey = account.encryptKey?.trim();
|
||||
if (!encryptKey) {
|
||||
throw new Error(`Feishu account "${accountId}" webhook mode requires encryptKey`);
|
||||
|
|
@ -639,6 +644,7 @@ export async function monitorWebhook(params: MonitorTransportParams): Promise<vo
|
|||
});
|
||||
const registration = registerWebhookTarget(webhookTargets, {
|
||||
...params,
|
||||
abortSignal,
|
||||
path,
|
||||
rawPath,
|
||||
legacyListener,
|
||||
|
|
@ -647,6 +653,8 @@ export async function monitorWebhook(params: MonitorTransportParams): Promise<vo
|
|||
pendingResponses: new Map<http.ServerResponse, Promise<void>>(),
|
||||
});
|
||||
let unregisterRoute: (() => void) | undefined;
|
||||
let ownedListener: Awaited<ReturnType<typeof startFeishuLegacyWebhookListener>> | undefined;
|
||||
let listenerError: Error | undefined;
|
||||
let cleanupStarted = false;
|
||||
let pendingDrain: Promise<void> | undefined;
|
||||
const cleanup = () => {
|
||||
|
|
@ -654,6 +662,8 @@ export async function monitorWebhook(params: MonitorTransportParams): Promise<vo
|
|||
return pendingDrain;
|
||||
}
|
||||
cleanupStarted = true;
|
||||
stopped.abort();
|
||||
ownedListener?.stopAccepting();
|
||||
const identityRevision = readFeishuBotIdentityRevision(accountId);
|
||||
pendingDrain = (async () => {
|
||||
const pendingResponses = registration.target.pendingResponses;
|
||||
|
|
@ -681,12 +691,16 @@ export async function monitorWebhook(params: MonitorTransportParams): Promise<vo
|
|||
) {
|
||||
clearFeishuBotIdentityState(accountId);
|
||||
}
|
||||
registration.unregister();
|
||||
unregisterRoute?.();
|
||||
try {
|
||||
registration.unregister();
|
||||
unregisterRoute?.();
|
||||
} finally {
|
||||
await ownedListener?.close();
|
||||
}
|
||||
})();
|
||||
return pendingDrain;
|
||||
};
|
||||
if (abortSignal?.aborted) {
|
||||
if (abortSignal.aborted) {
|
||||
await cleanup();
|
||||
return;
|
||||
}
|
||||
|
|
@ -700,18 +714,44 @@ export async function monitorWebhook(params: MonitorTransportParams): Promise<vo
|
|||
handler: (req, res) => handleFeishuWebhook(req, res, webhookTargets),
|
||||
reuseExistingSameOwner: true,
|
||||
throwOnFailure: true,
|
||||
legacyListener,
|
||||
legacyListener: gatewayOwnsLegacyListeners ? legacyListener : undefined,
|
||||
log: runtime?.log,
|
||||
});
|
||||
if (!gatewayOwnsLegacyListeners && legacyListener && !abortSignal.aborted) {
|
||||
const accountTargets = new Map([[path, [registration.target]]]);
|
||||
ownedListener = await startFeishuLegacyWebhookListener({
|
||||
endpoint: legacyListener,
|
||||
handleRequest: (req, res) => handleFeishuWebhook(req, res, accountTargets, legacyListener),
|
||||
onFailure: (error) => {
|
||||
listenerError = error;
|
||||
stopped.abort(error);
|
||||
},
|
||||
onRequestError: (error) =>
|
||||
(runtime?.error ?? console.error)(
|
||||
`feishu[${accountId}]: legacy webhook request failed: ${formatFeishuWsErrorForLog(error)}`,
|
||||
),
|
||||
});
|
||||
}
|
||||
if (abortSignal.aborted) {
|
||||
if (listenerError) {
|
||||
throw listenerError;
|
||||
}
|
||||
return;
|
||||
}
|
||||
const connectedAt = Date.now();
|
||||
statusSink?.(channelReadyPatch({ lastConnectedAt: connectedAt, lastEventAt: connectedAt }));
|
||||
runtime?.log?.(
|
||||
pathConflict
|
||||
? `feishu[${accountId}]: ${pathConflict} The legacy listener keeps the old path working; move the path and callback before setting legacyWebhook:false.`
|
||||
: `feishu[${accountId}]: webhook registered on Gateway port ${params.gatewayPort ?? 18789} at ${rawPath}; point the Feishu callback URL or reverse-proxy upstream to this Gateway route. ${legacyListener ? `The legacy listener on ${legacyListener.host}:${legacyListener.port} forwards here; set legacyWebhook:false after verifying delivery through the Gateway to disable legacy forwarding for this account.` : "legacyWebhook:false disables legacy forwarding for this account."}`,
|
||||
ownedListener
|
||||
? `feishu[${accountId}]: 2026.9.6 compatibility listener ${legacyListener?.host}:${legacyListener?.port} serves this account directly. This host requires distinct legacy endpoints for separate accounts; set legacyWebhook:false after verifying delivery on the Gateway route.`
|
||||
: pathConflict
|
||||
? `feishu[${accountId}]: ${pathConflict} The legacy listener keeps the old path working; move the path and callback before setting legacyWebhook:false.`
|
||||
: `feishu[${accountId}]: webhook registered on Gateway port ${params.gatewayPort ?? 18789} at ${rawPath}; point the Feishu callback URL or reverse-proxy upstream to this Gateway route. ${legacyListener ? `The legacy listener on ${legacyListener.host}:${legacyListener.port} forwards here; set legacyWebhook:false after verifying delivery through the Gateway to disable legacy forwarding for this account.` : "legacyWebhook:false disables legacy forwarding for this account."}`,
|
||||
);
|
||||
// Stopping targets retain only signature recognition until their responses finish.
|
||||
await waitUntilAbort(abortSignal, cleanup);
|
||||
if (listenerError) {
|
||||
throw listenerError;
|
||||
}
|
||||
} finally {
|
||||
await cleanup();
|
||||
}
|
||||
|
|
|
|||
217
extensions/feishu/src/monitor.webhook-host-floor.test.ts
Normal file
217
extensions/feishu/src/monitor.webhook-host-floor.test.ts
Normal file
|
|
@ -0,0 +1,217 @@
|
|||
import { once } from "node:events";
|
||||
import { createServer } from "node:http";
|
||||
import * as Lark from "@larksuiteoapi/node-sdk";
|
||||
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
|
||||
import { getActivePluginRegistry } from "openclaw/plugin-sdk/plugin-test-runtime";
|
||||
import { acquireTestPortBlock } from "openclaw/plugin-sdk/test-env";
|
||||
import { afterAll, afterEach, beforeEach, expect, it, vi } from "vitest";
|
||||
import { createRuntimeSpies } from "../../test-support/runtime-spies.js";
|
||||
import { cleanupFeishuMonitorStateForTests } from "./monitor.cleanup.test-helpers.js";
|
||||
import { monitorWebhook } from "./monitor.transport.js";
|
||||
import {
|
||||
createFeishuWebhookTestAccount,
|
||||
getGatewayPort,
|
||||
postSignedPayload,
|
||||
} from "./monitor.webhook.test-helpers.js";
|
||||
|
||||
const host = vi.hoisted(() => ({ ownsLegacyListeners: false }));
|
||||
vi.mock("openclaw/plugin-sdk/webhook-ingress", async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof import("openclaw/plugin-sdk/webhook-ingress")>();
|
||||
return {
|
||||
...actual,
|
||||
get getWebhookLegacyListener() {
|
||||
return host.ownsLegacyListeners ? actual.getWebhookLegacyListener : undefined;
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
const running: Array<{ abort: AbortController; monitor: Promise<void> }> = [];
|
||||
const portClaims: Array<Awaited<ReturnType<typeof acquireTestPortBlock>>> = [];
|
||||
beforeEach(() => {
|
||||
host.ownsLegacyListeners = false;
|
||||
});
|
||||
afterEach(async () => {
|
||||
for (const entry of running) {
|
||||
entry.abort.abort();
|
||||
}
|
||||
await Promise.allSettled(running.splice(0).map((entry) => entry.monitor));
|
||||
await using claims = new AsyncDisposableStack();
|
||||
for (const claim of portClaims.splice(0)) {
|
||||
claims.defer(() => claim.release());
|
||||
}
|
||||
await cleanupFeishuMonitorStateForTests();
|
||||
});
|
||||
|
||||
afterAll(() => {
|
||||
vi.doUnmock("openclaw/plugin-sdk/webhook-ingress");
|
||||
vi.resetModules();
|
||||
});
|
||||
|
||||
async function reservePort(port: number) {
|
||||
const server = createServer();
|
||||
server.listen(port, "127.0.0.1");
|
||||
await once(server, "listening");
|
||||
const address = server.address();
|
||||
if (!address || typeof address === "string") {
|
||||
throw new Error("Expected a loopback port");
|
||||
}
|
||||
return {
|
||||
port: address.port,
|
||||
close: () =>
|
||||
new Promise<void>((resolve) => {
|
||||
server.close(() => resolve());
|
||||
}),
|
||||
};
|
||||
}
|
||||
|
||||
async function claimPort() {
|
||||
const claim = await acquireTestPortBlock({ offsets: [0] });
|
||||
portClaims.push(claim);
|
||||
return claim.port;
|
||||
}
|
||||
|
||||
async function start(
|
||||
port: number,
|
||||
options: {
|
||||
accountId?: string;
|
||||
disabled?: boolean;
|
||||
abort?: AbortController;
|
||||
invoke?: () => Promise<{ kind: "non-durable"; value: unknown }>;
|
||||
} = {},
|
||||
) {
|
||||
const gatewayPort = await getGatewayPort();
|
||||
const accountId = options.accountId ?? "floor";
|
||||
const account = createFeishuWebhookTestAccount(accountId, `/hook-${accountId}`);
|
||||
account.config.legacyWebhook = options.disabled ? false : { port, host: "127.0.0.1" };
|
||||
const abort = options.abort ?? new AbortController();
|
||||
const ready = createDeferred<void>();
|
||||
const invoked = vi.fn(
|
||||
options.invoke ??
|
||||
(async (): Promise<{ kind: "non-durable"; value: unknown }> => ({
|
||||
kind: "non-durable",
|
||||
value: { accountId },
|
||||
})),
|
||||
);
|
||||
const monitor = monitorWebhook({
|
||||
account,
|
||||
accountId,
|
||||
gatewayPort,
|
||||
abortSignal: abort.signal,
|
||||
eventDispatcher: new Lark.EventDispatcher({ encryptKey: "encrypt_key" }),
|
||||
invokeWebhookEvent: invoked,
|
||||
runtime: createRuntimeSpies(),
|
||||
statusSink: (patch) => {
|
||||
if (patch.lifecycle === "ready") {
|
||||
ready.resolve();
|
||||
}
|
||||
},
|
||||
});
|
||||
running.push({ abort, monitor });
|
||||
if (abort.signal.aborted) {
|
||||
await monitor;
|
||||
} else {
|
||||
await Promise.race([
|
||||
ready.promise,
|
||||
monitor.then(() => {
|
||||
throw new Error("Monitor stopped before becoming ready");
|
||||
}),
|
||||
]);
|
||||
}
|
||||
return { abort, monitor, invoked, gatewayPort, path: `/hook-${accountId}` };
|
||||
}
|
||||
|
||||
it("serves the shipped account endpoint with the existing signature and path checks", async () => {
|
||||
const port = await claimPort();
|
||||
const entry = await start(port);
|
||||
const url = `http://127.0.0.1:${port}${entry.path}`;
|
||||
const wrongPath = await postSignedPayload(`${url}/other`, { schema: "2.0", event: {} });
|
||||
expect(wrongPath.status).toBe(404);
|
||||
await wrongPath.text();
|
||||
const unsigned = await fetch(url, {
|
||||
method: "POST",
|
||||
headers: { "content-type": "application/json", connection: "close" },
|
||||
body: JSON.stringify({ schema: "2.0", event: {} }),
|
||||
});
|
||||
expect(unsigned.status).toBe(401);
|
||||
await unsigned.text();
|
||||
expect(entry.invoked).not.toHaveBeenCalled();
|
||||
const accepted = await postSignedPayload(url, { schema: "2.0", event: {} });
|
||||
expect(accepted.status).toBe(200);
|
||||
await expect(accepted.json()).resolves.toEqual({ accountId: "floor" });
|
||||
expect(entry.invoked).toHaveBeenCalledOnce();
|
||||
expect(getActivePluginRegistry()?.httpRoutes[0]?.legacyListeners).toBeUndefined();
|
||||
});
|
||||
|
||||
it("stops listener admission before draining an authenticated response and permits rebinding", async () => {
|
||||
const port = await claimPort();
|
||||
const entered = createDeferred<void>();
|
||||
const release = createDeferred<void>();
|
||||
const entry = await start(port, {
|
||||
invoke: async () => {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
return { kind: "non-durable", value: { completed: true } };
|
||||
},
|
||||
});
|
||||
const url = `http://127.0.0.1:${port}${entry.path}`;
|
||||
const response = postSignedPayload(url, { schema: "2.0", event: {} });
|
||||
let stopped = false;
|
||||
void entry.monitor.then(() => {
|
||||
stopped = true;
|
||||
});
|
||||
try {
|
||||
await entered.promise;
|
||||
entry.abort.abort();
|
||||
await expect(fetch(url, { headers: { connection: "close" } })).rejects.toMatchObject({
|
||||
cause: { code: "ECONNREFUSED" },
|
||||
});
|
||||
expect(stopped).toBe(false);
|
||||
release.resolve();
|
||||
const accepted = await response;
|
||||
expect(accepted.status).toBe(200);
|
||||
await expect(accepted.json()).resolves.toEqual({ completed: true });
|
||||
await entry.monitor;
|
||||
const rebound = await reservePort(port);
|
||||
await rebound.close();
|
||||
} finally {
|
||||
release.resolve();
|
||||
entry.abort.abort();
|
||||
await response.catch(() => {});
|
||||
await entry.monitor;
|
||||
}
|
||||
});
|
||||
|
||||
it("keeps the shipped per-account bind refusal without disturbing the live account", async () => {
|
||||
const port = await claimPort();
|
||||
const first = await start(port, { accountId: "first" });
|
||||
await expect(start(port, { accountId: "second" })).rejects.toMatchObject({ code: "EADDRINUSE" });
|
||||
const response = await postSignedPayload(`http://127.0.0.1:${port}${first.path}`, {
|
||||
schema: "2.0",
|
||||
event: {},
|
||||
});
|
||||
expect(response.status).toBe(200);
|
||||
await expect(response.json()).resolves.toEqual({ accountId: "first" });
|
||||
expect(getActivePluginRegistry()?.httpRoutes.some((route) => route.path === "/hook-second")).toBe(
|
||||
false,
|
||||
);
|
||||
});
|
||||
|
||||
it.each(["capable-host", "disabled", "already-aborted"] as const)(
|
||||
"does not own a listener for %s",
|
||||
async (mode) => {
|
||||
const port = await claimPort();
|
||||
host.ownsLegacyListeners = mode === "capable-host";
|
||||
const abort = new AbortController();
|
||||
if (mode === "already-aborted") {
|
||||
abort.abort();
|
||||
}
|
||||
await start(port, { disabled: mode === "disabled", abort });
|
||||
const unbound = await reservePort(port);
|
||||
await unbound.close();
|
||||
const routes = getActivePluginRegistry()?.httpRoutes ?? [];
|
||||
expect(routes).toHaveLength(mode === "already-aborted" ? 0 : 1);
|
||||
expect(routes[0]?.legacyListeners).toEqual(
|
||||
mode === "capable-host" ? [{ port, host: "127.0.0.1" }] : undefined,
|
||||
);
|
||||
},
|
||||
);
|
||||
63
extensions/feishu/src/monitor.webhook-legacy.ts
Normal file
63
extensions/feishu/src/monitor.webhook-legacy.ts
Normal file
|
|
@ -0,0 +1,63 @@
|
|||
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
|
||||
import { extractErrorCode } from "openclaw/plugin-sdk/error-runtime";
|
||||
import * as webhookIngressSdk from "openclaw/plugin-sdk/webhook-ingress";
|
||||
|
||||
// Shipped 2026.9.6 hosts lack Gateway-owned forwarding. Retire this adapter when
|
||||
// Feishu's declared host floor includes that listener capability.
|
||||
export const feishuWebhookHost: Partial<
|
||||
Pick<typeof webhookIngressSdk, "getWebhookLegacyListener">
|
||||
> = webhookIngressSdk;
|
||||
|
||||
export async function startFeishuLegacyWebhookListener(params: {
|
||||
endpoint: { port: number; host: string };
|
||||
handleRequest: (req: IncomingMessage, res: ServerResponse) => Promise<void>;
|
||||
onFailure: (error: Error) => void;
|
||||
onRequestError: (error: unknown) => void;
|
||||
}) {
|
||||
const server = createServer((req, res) => {
|
||||
void params.handleRequest(req, res).catch((error: unknown) => {
|
||||
params.onRequestError(error);
|
||||
if (!res.headersSent) {
|
||||
res.statusCode = 500;
|
||||
}
|
||||
res.end();
|
||||
});
|
||||
});
|
||||
let closing: Promise<void> | undefined;
|
||||
const stopAccepting = () => {
|
||||
closing ??= new Promise<void>((resolve, reject) => {
|
||||
server.close((error) => {
|
||||
if (error && extractErrorCode(error) !== "ERR_SERVER_NOT_RUNNING") {
|
||||
reject(error);
|
||||
} else {
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
});
|
||||
// The transport joins any close error after its authenticated response drain.
|
||||
void closing.catch(() => {});
|
||||
};
|
||||
const close = async () => {
|
||||
stopAccepting();
|
||||
server.closeAllConnections();
|
||||
try {
|
||||
await closing;
|
||||
} finally {
|
||||
server.off("error", params.onFailure);
|
||||
}
|
||||
};
|
||||
server.on("error", params.onFailure);
|
||||
try {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
server.once("error", reject);
|
||||
server.listen(params.endpoint.port, params.endpoint.host, () => {
|
||||
server.off("error", reject);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
} catch (error) {
|
||||
await close();
|
||||
throw error;
|
||||
}
|
||||
return { stopAccepting, close };
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue