mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-12 21:53:00 -06:00
feat(gateway): expose public worker ingress (#122578)
* feat(gateway): expose public worker ingress * docs(plan): link public worker ingress PR
This commit is contained in:
committed by
GitHub
parent
c5ba4efbd7
commit
6e71e9b156
@@ -1 +1 @@
|
||||
{"contentHash":"ada2636485ebedd6593c0223b66d75222498468af082773aad002a69e7253592","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
|
||||
{"contentHash":"592d7a63d16191c80337a99b54d1bd7abc90dec80ec5619a33aa0127f8065cd9","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"fc00024b4d58f04ce99f258213f597126bcfc8adae9ad389376c23ae8598ae3a","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
|
||||
{"contentHash":"9966a1fe3a54b7d4441c797b1278fdefbbdc570ad2a16c54a7e1d59d79b0d78c","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"ae7ebc2ff1de2a97232b992dd6cf494bc80f4666fa6f92bf05411c4412d0e089","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
|
||||
{"contentHash":"235411e28d230c2601c89729616e91bdf467e9aafead3986d6b3944f8a3793e2","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"15e8d88af3ded7a56f56d3eca4f59dbcc14d2d14dd4f5d291919eb204528082d","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
|
||||
{"contentHash":"64ab8862f4dc8bb89cba1508704737dfdb760e74de3330308ee33160d109e085","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"5ae94502f0098ffd3c91160440795c3955bc63558e034688ab890f0edac27e74","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
|
||||
{"contentHash":"1bee58753249c7720eb218b9a160d58c460c9051d2e55c5c16e6b560565b1670","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"9b17e70764b958845e9b5ad547b8a58c730acd76c450137b7063de4d84724d91","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
|
||||
{"contentHash":"5751ca5bba6c92d75ab169cdb8614be5e7a65e3a15cc68873a9c73b8e9e59f2a","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"628a6a5b9674566c5956e04e505be93079c8324bd2701368c6c47adbe93956cd","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
|
||||
{"contentHash":"9a3dae76ff23015d94e5b1946236534ea2fa92b6d7b51cf9e9e409cac134a972","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"a0a3832ec0e4a2529333b18d4474940e773b9afe01e8cd13c1d79edf75f583f4","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
|
||||
{"contentHash":"bf1d2b91991796d3bf82a1802d2997e4c8692d5ee68ca9e07b7dab99bfc3c5fb","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"2365d560ccc2c217c971c57918a096b2921727cc1df4357c1af486251c126048","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
|
||||
{"contentHash":"b45ce74762fee5267b16b339f8f97ccce0f4f5df92608a56bc48dc9f46824055","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"7c27ccd856e944332b06e5c79ee0275bad1c3e2d1a2174f108d5ead742b32783","entrypoint":"gateway-runtime","importSpecifier":"openclaw/plugin-sdk/gateway-runtime"}
|
||||
{"contentHash":"bb7cea1bf66d810319e1a1b1b18d8612121dba3ca6da4584dcc3faba53a4a1b6","entrypoint":"gateway-runtime","importSpecifier":"openclaw/plugin-sdk/gateway-runtime"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"0b9df99990acff17f2a8980c944ec6dbd6e6a1797d82238b0684e61cae7c227b","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
|
||||
{"contentHash":"36d41a6be6c6d5e523ae164662a8080bfb7c34e37d1b66e00964ec078db22c93","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"558d33dffe62849474af6ba9706c57b2b26ebf71187ee04f4dcfae039caa0678","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
|
||||
{"contentHash":"0b85bf09a339ce32b3f8e43159eb88fbfb1691b6a8346b97d0798e5188d2d32a","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"27b07f72e16f1a0d0bbb143d3c19ee80064914971909bc72099051f68ceb7b89","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
|
||||
{"contentHash":"73cf4e87d121ed568cf4339436d3fa629144c86902060bbf6e1e8fe06f261fa6","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"71cb620f8b936e6d83b260f45e9ebb16e55363b6fd7bccd8944632e78b14d70d","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
|
||||
{"contentHash":"5dd6717af8284afcf6ef88f19593a4d6f5363cb3be26f693157605bc64a29ea0","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"a5c6beede9403833929c5491019c92d691df2f092c1db13c9caf125da4013aaa","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
|
||||
{"contentHash":"a860f6e24630ad2c4f1964f10f1de3300e01029cf40dae90586c0913a7bf0867","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"contentHash":"3bd16cfa9be68c8d517d1a51956f588482e45ef70ea39fb192b2129e9f433ffe","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
|
||||
{"contentHash":"9c9485971934a4223d97838ab69a6d8cb381092d2e552933c9bf98e6b5d643c6","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
|
||||
|
||||
@@ -1 +1 @@
|
||||
{"contentHash":"c8b52887a5faf418866ba16c3b71d8f62509f7eb9231e30cd39f4cd331d3e7f1","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
|
||||
{"contentHash":"14e9274585375d9e517da0d7a45724407c76a56e57c2ad2a0c12c32e868d31ee","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
|
||||
|
||||
+16
-11
@@ -275,17 +275,22 @@ go through normal pairing and scope-upgrade checks.
|
||||
|
||||
### Worker role and closed protocol
|
||||
|
||||
Cloud workers use a dedicated loopback ingress through the gateway-owned,
|
||||
host-key-pinned SSH tunnel. It accepts only worker identity and never dispatches
|
||||
general auth, node events, operator RPCs, or plugin methods. A strict `connect`
|
||||
verifies a hash-at-rest, short-lived credential bound to the environment, bundle
|
||||
hash, owner epoch, RPC-set version, expiry, and one nullable session; it
|
||||
separately checks the current version and feature set. Success returns minimal
|
||||
`worker-hello-ok`; feature negotiation is independent of the general protocol
|
||||
version. Frames stay under 64 KiB, except a negotiated `worker.inference.start`
|
||||
frame may be up to 25 MiB. The closed allowlist contains `worker.heartbeat`,
|
||||
`worker.transcript.commit`, `worker.live-event`, `worker.inference.start`, and
|
||||
`worker.inference.cancel`.
|
||||
Workers use a closed protocol through either the public
|
||||
`/__openclaw__/worker` WebSocket path on the main TLS endpoint or the dedicated
|
||||
loopback ingress reached through the gateway-owned, host-key-pinned SSH tunnel.
|
||||
The route selects worker mode before reading frames, so it never dispatches
|
||||
general auth, node events, operator RPCs, or plugin methods. Public admission
|
||||
shares the main per-client pre-auth budget and authentication rate limiter; its
|
||||
wire errors collapse credential and environment details to
|
||||
`admission-rejected`, while trusted gateway diagnostics retain the internal
|
||||
reason. A strict `connect` verifies a hash-at-rest, short-lived credential bound
|
||||
to the environment, bundle hash, owner epoch, RPC-set version, expiry, and one
|
||||
nullable session; it separately checks the current version and feature set.
|
||||
Success returns minimal `worker-hello-ok`; feature negotiation is independent of
|
||||
the general protocol version. Frames stay under 64 KiB, except a negotiated
|
||||
`worker.inference.start` frame may be up to 25 MiB. The closed allowlist contains
|
||||
`worker.heartbeat`, `worker.transcript.commit`, `worker.live-event`,
|
||||
`worker.inference.start`, and `worker.inference.cancel`.
|
||||
|
||||
Transcript commits use owner-epoch fencing, a gateway-owned session binding,
|
||||
base-leaf compare-and-swap, and durable sequence replay; the gateway generates
|
||||
|
||||
@@ -18,11 +18,12 @@ advances a milestone.
|
||||
| 0 | This plan (revision 2) | landed | #122454 |
|
||||
| 1a | Naming: session copy revert | landed | #120667 |
|
||||
| 1b | Naming: devices consolidation | landed | #120689 |
|
||||
| 1c | Cleanup: node-pairing → device-pairing merge | not started | — |
|
||||
| 1c | Cleanup: node-pairing → device-pairing merge | landed | #120726 |
|
||||
| 2 | `openclaw resume` + web Continue in terminal | in progress | #120664 |
|
||||
| 3 | `openclaw connect` one-paste onboarding + `/j/` join route | in progress | #122499 |
|
||||
| 3 | `openclaw connect` one-paste onboarding + `/j/` join route | in progress | #120768, #122499 |
|
||||
| 4 | Picker: grouping, placement, liveness, enrichment | in progress | #120804, #122531 |
|
||||
| 5 | Public worker ingress path | not started | — |
|
||||
| F | Real-wire session boundary harness | landed | #121212 |
|
||||
| 5 | Public worker ingress path | in progress | #122578 |
|
||||
| 6 | Node worker provider (device runners) | not started | — |
|
||||
| 7 | Bundle push consent + runner updates | not started | — |
|
||||
| 8 | Stop-and-continue moves | not started | — |
|
||||
|
||||
@@ -626,6 +626,7 @@ describe("worker protocol schemas", () => {
|
||||
});
|
||||
|
||||
it("keeps worker close reasons closed", () => {
|
||||
expect(Value.Check(WorkerProtocolCloseReasonSchema, "admission-rejected")).toBe(true);
|
||||
expect(Value.Check(WorkerProtocolCloseReasonSchema, "credential-replaced")).toBe(true);
|
||||
expect(Value.Check(WorkerProtocolCloseReasonSchema, "placement-mismatch")).toBe(true);
|
||||
expect(Value.Check(WorkerProtocolCloseReasonSchema, "not-a-worker-reason")).toBe(false);
|
||||
|
||||
@@ -19,6 +19,7 @@ import {
|
||||
} from "./worker-protocol-primitives.js";
|
||||
|
||||
export {
|
||||
WORKER_PUBLIC_INGRESS_PATH,
|
||||
WORKER_PROTOCOL_MAX_FRAME_ID_LENGTH,
|
||||
WORKER_PROTOCOL_MAX_IDENTIFIER_LENGTH,
|
||||
WORKER_PROTOCOL_MAX_PAYLOAD_BYTES,
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { Type } from "typebox";
|
||||
import { closedObject } from "./closed-object.js";
|
||||
|
||||
export const WORKER_PUBLIC_INGRESS_PATH = "/__openclaw__/worker";
|
||||
export const WORKER_PROTOCOL_MAX_IDENTIFIER_LENGTH = 256;
|
||||
export const WORKER_PROTOCOL_MAX_FRAME_ID_LENGTH = 128;
|
||||
export const WORKER_PROTOCOL_MAX_PAYLOAD_BYTES = 64 * 1024;
|
||||
@@ -32,6 +33,7 @@ export const WorkerAdmissionFailureReasonSchema = Type.Union([
|
||||
|
||||
export const WorkerProtocolCloseReasonSchema = Type.Union([
|
||||
WorkerAdmissionFailureReasonSchema,
|
||||
Type.Literal("admission-rejected"),
|
||||
Type.Literal("invalid-handshake"),
|
||||
Type.Literal("protocol-mismatch"),
|
||||
Type.Literal("gateway-unavailable"),
|
||||
|
||||
@@ -59,6 +59,9 @@ export const AUTH_RATE_LIMIT_SCOPE_BOOTSTRAP_TOKEN = "bootstrap-token";
|
||||
// Public join-code exchange burns SQLite state, so misses are serialized and
|
||||
// throttled before they can queue unbounded writes behind the shared DB lock.
|
||||
export const AUTH_RATE_LIMIT_SCOPE_DEVICE_JOIN = "device-join";
|
||||
// Public worker admission performs store-backed credential verification before
|
||||
// the worker is authenticated, so it gets an independent per-IP guess budget.
|
||||
export const AUTH_RATE_LIMIT_SCOPE_WORKER_ADMISSION = "worker-admission";
|
||||
// Public watchOS challenge issuance is throttled separately from credential
|
||||
// failures so challenge floods cannot displace legitimate device handshakes.
|
||||
export const AUTH_RATE_LIMIT_SCOPE_WATCH_CHALLENGE = "watch-challenge";
|
||||
|
||||
@@ -9,6 +9,7 @@ import {
|
||||
import { createServer as createHttpsServer } from "node:https";
|
||||
import type { TlsOptions } from "node:tls";
|
||||
import type { WebSocketServer } from "ws";
|
||||
import { WORKER_PUBLIC_INGRESS_PATH } from "../../packages/gateway-protocol/src/schema/worker-admission.js";
|
||||
import { isCoreCanvasHostEnabled } from "../canvas/config.js";
|
||||
import { isCanvasDocumentHttpPath } from "../canvas/constants.js";
|
||||
import { resolveBundledChannelGatewayAuthBypassPaths } from "../channels/plugins/gateway-auth-bypass.js";
|
||||
@@ -76,6 +77,7 @@ import type { ReadinessChecker, StartupChecker, StartupResult } from "./server/r
|
||||
import {
|
||||
GATEWAY_WS_CONNECTION_KIND_PROPERTY,
|
||||
GATEWAY_WS_PREAUTH_BUDGET_PROPERTY,
|
||||
GATEWAY_WS_WORKER_INGRESS_PROPERTY,
|
||||
type GatewayIngressWebSocket,
|
||||
type GatewayWsClient,
|
||||
} from "./server/ws-types.js";
|
||||
@@ -926,6 +928,7 @@ export function attachGatewayUpgradeHandler(opts: {
|
||||
rateLimiter?: AuthRateLimiter;
|
||||
/** Optional logger for error diagnostics. */
|
||||
log?: { warn: (msg: string) => void };
|
||||
workerIngressEnabled?: boolean;
|
||||
desktopSessionRegistry?: DesktopSessionRegistry;
|
||||
}) {
|
||||
const {
|
||||
@@ -958,6 +961,31 @@ export function attachGatewayUpgradeHandler(opts: {
|
||||
}
|
||||
const resolvedAuthLocal = getResolvedAuth();
|
||||
const requestPath = scopedNodeCapability.pathname;
|
||||
if (requestPath === WORKER_PUBLIC_INGRESS_PATH) {
|
||||
if (!opts.workerIngressEnabled) {
|
||||
writeGatewayUpgradeServiceUnavailable(socket, "Worker websocket ingress unavailable");
|
||||
socket.destroy();
|
||||
return;
|
||||
}
|
||||
try {
|
||||
handleBudgetedGatewayWebSocketUpgrade({
|
||||
req,
|
||||
socket,
|
||||
head,
|
||||
wss,
|
||||
preauthConnectionBudget,
|
||||
preauthBudgetKey: requestClientIp,
|
||||
ingressName: "Worker",
|
||||
prepareSocket: (workerSocket) => {
|
||||
workerSocket[GATEWAY_WS_CONNECTION_KIND_PROPERTY] = "worker";
|
||||
workerSocket[GATEWAY_WS_WORKER_INGRESS_PROPERTY] = "public";
|
||||
},
|
||||
});
|
||||
} catch {
|
||||
throw new Error("public worker websocket upgrade failed");
|
||||
}
|
||||
return;
|
||||
}
|
||||
const pathContext = resolvePluginRoutePathContext(requestPath);
|
||||
const nodeCapability = resolvePluginNodeCapabilityRoute?.(pathContext);
|
||||
if (nodeCapability) {
|
||||
@@ -1090,6 +1118,7 @@ export function attachWorkerGatewayUpgradeHandler(params: {
|
||||
prepareSocket: (workerSocket) => {
|
||||
workerSocket[GATEWAY_WS_CONNECTION_KIND_PROPERTY] = "worker";
|
||||
workerSocket[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY] = params.preauthConnectionBudget;
|
||||
workerSocket[GATEWAY_WS_WORKER_INGRESS_PROPERTY] = "loopback";
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
|
||||
@@ -304,6 +304,7 @@ export async function createGatewayHttpTransport(params: {
|
||||
getResolvedAuth: params.getResolvedAuth,
|
||||
rateLimiter: params.rateLimiter,
|
||||
log: params.log,
|
||||
workerIngressEnabled: params.workerIngressEnabled,
|
||||
desktopSessionRegistry: params.desktopSessionRegistry,
|
||||
});
|
||||
gatewayHttpServers.push(httpServer);
|
||||
|
||||
@@ -2,8 +2,15 @@
|
||||
* Gateway pre-auth hardening tests.
|
||||
*/
|
||||
import http from "node:http";
|
||||
import { afterEach, describe, expect, it } from "vitest";
|
||||
import { rawDataToString } from "@openclaw/gateway-client/websocket-data";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { WebSocket, WebSocketServer } from "ws";
|
||||
import {
|
||||
GATEWAY_CLIENT_IDS,
|
||||
GATEWAY_CLIENT_MODES,
|
||||
} from "../../packages/gateway-protocol/src/client-info.js";
|
||||
import { PROTOCOL_VERSION } from "../../packages/gateway-protocol/src/index.js";
|
||||
import { WORKER_PUBLIC_INGRESS_PATH } from "../../packages/gateway-protocol/src/schema/worker-admission.js";
|
||||
import {
|
||||
onDiagnosticEvent,
|
||||
resetDiagnosticEventsForTest,
|
||||
@@ -23,9 +30,16 @@ import {
|
||||
createGatewayHttpServer,
|
||||
} from "./server-http.js";
|
||||
import { createPreauthConnectionBudget } from "./server/preauth-connection-budget.js";
|
||||
import { attachGatewayWsConnectionHandler } from "./server/ws-connection.js";
|
||||
import {
|
||||
createGatewayWsTestLogger,
|
||||
createGatewayWsTestRequestContext,
|
||||
} from "./server/ws-connection.test-helpers.js";
|
||||
import type { WorkerConnectionService } from "./server/ws-connection/worker-connection.js";
|
||||
import {
|
||||
GATEWAY_WS_CONNECTION_KIND_PROPERTY,
|
||||
GATEWAY_WS_PREAUTH_BUDGET_PROPERTY,
|
||||
GATEWAY_WS_WORKER_INGRESS_PROPERTY,
|
||||
type GatewayIngressWebSocket,
|
||||
type GatewayWsClient,
|
||||
} from "./server/ws-types.js";
|
||||
@@ -67,12 +81,15 @@ function setGatewayAuthNoneForTest() {
|
||||
});
|
||||
}
|
||||
|
||||
async function requestUpgradeRejection(port: number): Promise<{ status: number; body: string }> {
|
||||
async function requestUpgradeRejection(
|
||||
port: number,
|
||||
path = "/",
|
||||
): Promise<{ status: number; body: string }> {
|
||||
return await new Promise<{ status: number; body: string }>((resolve, reject) => {
|
||||
const req = http.request({
|
||||
host: "127.0.0.1",
|
||||
port,
|
||||
path: "/",
|
||||
path,
|
||||
headers: {
|
||||
Connection: "Upgrade",
|
||||
Upgrade: "websocket",
|
||||
@@ -143,6 +160,7 @@ describe("gateway pre-auth hardening", () => {
|
||||
const socket = await accepted;
|
||||
expect(socket[GATEWAY_WS_CONNECTION_KIND_PROPERTY]).toBe("worker");
|
||||
expect(socket[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY]).toBe(workerBudget);
|
||||
expect(socket[GATEWAY_WS_WORKER_INGRESS_PROPERTY]).toBe("loopback");
|
||||
} finally {
|
||||
client.close();
|
||||
await new Promise<void>((resolve) => {
|
||||
@@ -157,6 +175,239 @@ describe("gateway pre-auth hardening", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("reserves the public worker path before plugin upgrade routing", async () => {
|
||||
const clients = new Set<GatewayWsClient>();
|
||||
const resolvedAuth: ResolvedGatewayAuth = { mode: "none", allowTailscale: false };
|
||||
const httpServer = createGatewayHttpServer({
|
||||
clients,
|
||||
controlUiEnabled: false,
|
||||
controlUiBasePath: "/__control__",
|
||||
openAiChatCompletionsEnabled: false,
|
||||
openResponsesEnabled: false,
|
||||
handleHooksRequest: async () => false,
|
||||
resolvedAuth,
|
||||
});
|
||||
const wss = new WebSocketServer({ maxPayload: 1024, noServer: true });
|
||||
const pluginUpgrade = vi.fn(async () => false);
|
||||
const accepted = new Promise<GatewayIngressWebSocket>((resolve) => {
|
||||
wss.once("connection", (socket) => resolve(socket as GatewayIngressWebSocket));
|
||||
});
|
||||
attachGatewayUpgradeHandler({
|
||||
httpServer,
|
||||
wss,
|
||||
handlePluginUpgrade: pluginUpgrade,
|
||||
clients,
|
||||
preauthConnectionBudget: createPreauthConnectionBudget(1),
|
||||
resolvedAuth,
|
||||
workerIngressEnabled: true,
|
||||
});
|
||||
await new Promise<void>((resolve) => {
|
||||
httpServer.listen(0, "127.0.0.1", resolve);
|
||||
});
|
||||
const address = httpServer.address();
|
||||
const port = typeof address === "object" && address ? address.port : 0;
|
||||
const client = new WebSocket(`ws://127.0.0.1:${port}${WORKER_PUBLIC_INGRESS_PATH}`);
|
||||
|
||||
try {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
client.once("open", resolve);
|
||||
client.once("error", reject);
|
||||
});
|
||||
const socket = await accepted;
|
||||
expect(socket[GATEWAY_WS_CONNECTION_KIND_PROPERTY]).toBe("worker");
|
||||
expect(socket[GATEWAY_WS_WORKER_INGRESS_PROPERTY]).toBe("public");
|
||||
expect(socket[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY]).toBeUndefined();
|
||||
expect(pluginUpgrade).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
client.close();
|
||||
await new Promise<void>((resolve) => {
|
||||
client.once("close", () => resolve());
|
||||
});
|
||||
await new Promise<void>((resolve) => {
|
||||
wss.close(() => resolve());
|
||||
});
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
httpServer.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
it("admits a valid worker over the public path without a gateway challenge", async () => {
|
||||
const clients = new Set<GatewayWsClient>();
|
||||
const resolvedAuth: ResolvedGatewayAuth = { mode: "none", allowTailscale: false };
|
||||
const httpServer = createGatewayHttpServer({
|
||||
clients,
|
||||
controlUiEnabled: false,
|
||||
controlUiBasePath: "/__control__",
|
||||
openAiChatCompletionsEnabled: false,
|
||||
openResponsesEnabled: false,
|
||||
handleHooksRequest: async () => false,
|
||||
resolvedAuth,
|
||||
});
|
||||
const wss = new WebSocketServer({ maxPayload: 64 * 1024, noServer: true });
|
||||
const preauthConnectionBudget = createPreauthConnectionBudget(1);
|
||||
const workerConnectionService: WorkerConnectionService = {
|
||||
admitWorker: vi.fn(async () => ({
|
||||
ok: true as const,
|
||||
identity: {
|
||||
environmentId: "worker-public",
|
||||
credentialHash: "h".repeat(43),
|
||||
bundleHash: "a".repeat(64),
|
||||
sessionId: null,
|
||||
runId: null,
|
||||
ownerEpoch: 1,
|
||||
rpcSetVersion: 1,
|
||||
protocolFeatures: [],
|
||||
credentialExpiresAtMs: Date.now() + 60_000,
|
||||
},
|
||||
})),
|
||||
validateWorkerConnection: vi.fn(() => null),
|
||||
commitTranscript: vi.fn(async () => {
|
||||
throw new Error("unexpected transcript commit");
|
||||
}),
|
||||
pushLiveEvent: vi.fn(async () => {
|
||||
throw new Error("unexpected live event");
|
||||
}),
|
||||
};
|
||||
attachGatewayUpgradeHandler({
|
||||
httpServer,
|
||||
wss,
|
||||
clients,
|
||||
preauthConnectionBudget,
|
||||
resolvedAuth,
|
||||
workerIngressEnabled: true,
|
||||
});
|
||||
const logGateway = createGatewayWsTestLogger();
|
||||
const logHealth = createGatewayWsTestLogger();
|
||||
const logWsControl = createGatewayWsTestLogger();
|
||||
attachGatewayWsConnectionHandler({
|
||||
wss,
|
||||
clients,
|
||||
preauthConnectionBudget,
|
||||
port: 0,
|
||||
getResolvedAuth: () => resolvedAuth,
|
||||
preauthHandshakeTimeoutMs: 2_000,
|
||||
gatewayMethods: [],
|
||||
events: [],
|
||||
refreshHealthSnapshot: vi.fn(async () => ({}) as never),
|
||||
logGateway: logGateway as never,
|
||||
logHealth: logHealth as never,
|
||||
logWsControl: logWsControl as never,
|
||||
extraHandlers: {},
|
||||
broadcast: vi.fn(),
|
||||
buildRequestContext: () => createGatewayWsTestRequestContext() as never,
|
||||
workerConnectionService,
|
||||
});
|
||||
await new Promise<void>((resolve) => {
|
||||
httpServer.listen(0, "127.0.0.1", resolve);
|
||||
});
|
||||
const address = httpServer.address();
|
||||
const port = typeof address === "object" && address ? address.port : 0;
|
||||
const client = new WebSocket(`ws://127.0.0.1:${port}${WORKER_PUBLIC_INGRESS_PATH}`);
|
||||
const received: unknown[] = [];
|
||||
client.on("message", (data) => received.push(JSON.parse(rawDataToString(data))));
|
||||
|
||||
try {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
client.once("open", resolve);
|
||||
client.once("error", reject);
|
||||
});
|
||||
client.send(
|
||||
JSON.stringify({
|
||||
type: "req",
|
||||
id: "connect-public-worker",
|
||||
method: "connect",
|
||||
params: {
|
||||
minProtocol: PROTOCOL_VERSION,
|
||||
maxProtocol: PROTOCOL_VERSION,
|
||||
client: {
|
||||
id: GATEWAY_CLIENT_IDS.WORKER,
|
||||
version: "2026.8.12",
|
||||
platform: "linux",
|
||||
mode: GATEWAY_CLIENT_MODES.WORKER,
|
||||
},
|
||||
role: "worker",
|
||||
admission: {
|
||||
environmentId: "worker-public",
|
||||
credential: "public-worker-credential",
|
||||
sessionId: null,
|
||||
runId: null,
|
||||
ownerEpoch: 1,
|
||||
rpcSetVersion: 1,
|
||||
handshake: {
|
||||
bundleHash: "a".repeat(64),
|
||||
openclawVersion: "2026.8.12",
|
||||
protocolFeatures: [],
|
||||
},
|
||||
},
|
||||
},
|
||||
}),
|
||||
);
|
||||
await vi.waitFor(() => expect(received).toHaveLength(1));
|
||||
expect(received[0]).toMatchObject({
|
||||
type: "res",
|
||||
id: "connect-public-worker",
|
||||
ok: true,
|
||||
payload: { type: "worker-hello-ok", environmentId: "worker-public" },
|
||||
});
|
||||
expect(received).not.toContainEqual(
|
||||
expect.objectContaining({ type: "event", event: "connect.challenge" }),
|
||||
);
|
||||
expect(workerConnectionService.admitWorker).toHaveBeenCalledOnce();
|
||||
} finally {
|
||||
client.close();
|
||||
await new Promise<void>((resolve) => {
|
||||
client.once("close", () => resolve());
|
||||
});
|
||||
await new Promise<void>((resolve) => {
|
||||
wss.close(() => resolve());
|
||||
});
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
httpServer.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
it("rejects the reserved worker path when worker admission is unavailable", async () => {
|
||||
const clients = new Set<GatewayWsClient>();
|
||||
const resolvedAuth: ResolvedGatewayAuth = { mode: "none", allowTailscale: false };
|
||||
const httpServer = createGatewayHttpServer({
|
||||
clients,
|
||||
controlUiEnabled: false,
|
||||
controlUiBasePath: "/__control__",
|
||||
openAiChatCompletionsEnabled: false,
|
||||
openResponsesEnabled: false,
|
||||
handleHooksRequest: async () => false,
|
||||
resolvedAuth,
|
||||
});
|
||||
const wss = new WebSocketServer({ maxPayload: 1024, noServer: true });
|
||||
wss.on("connection", (socket) => socket.close());
|
||||
attachGatewayUpgradeHandler({
|
||||
httpServer,
|
||||
wss,
|
||||
clients,
|
||||
preauthConnectionBudget: createPreauthConnectionBudget(1),
|
||||
resolvedAuth,
|
||||
});
|
||||
await new Promise<void>((resolve) => {
|
||||
httpServer.listen(0, "127.0.0.1", resolve);
|
||||
});
|
||||
const address = httpServer.address();
|
||||
const port = typeof address === "object" && address ? address.port : 0;
|
||||
|
||||
try {
|
||||
await expect(requestUpgradeRejection(port, WORKER_PUBLIC_INGRESS_PATH)).resolves.toEqual({
|
||||
status: 503,
|
||||
body: "Worker websocket ingress unavailable",
|
||||
});
|
||||
} finally {
|
||||
wss.close();
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
httpServer.close((error) => (error ? reject(error) : resolve()));
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
it("rejects worker websocket upgrades after suspension is prepared", async () => {
|
||||
const httpServer = http.createServer();
|
||||
const wss = new WebSocketServer({ maxPayload: 1024, noServer: true });
|
||||
|
||||
@@ -49,6 +49,7 @@ import { resolveSharedGatewaySessionGeneration } from "./ws-shared-generation.js
|
||||
import {
|
||||
GATEWAY_WS_CONNECTION_KIND_PROPERTY,
|
||||
GATEWAY_WS_PREAUTH_BUDGET_PROPERTY,
|
||||
GATEWAY_WS_WORKER_INGRESS_PROPERTY,
|
||||
} from "./ws-types.js";
|
||||
|
||||
async function waitForLazyMessageHandler() {
|
||||
@@ -105,7 +106,7 @@ describe("attachGatewayWsConnectionHandler", () => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("keeps worker sockets off the legacy challenge, plugin surface, and gateway budget", async () => {
|
||||
it("keeps loopback worker sockets off the legacy challenge, plugin surface, and gateway budget", async () => {
|
||||
const socket = createGatewayWsTestSocket();
|
||||
const previous = {
|
||||
socket: { terminate: vi.fn() },
|
||||
@@ -153,6 +154,39 @@ describe("attachGatewayWsConnectionHandler", () => {
|
||||
expect(gatewayBudget.release).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("uses the main budget and auth limiter for public worker sockets", async () => {
|
||||
const socket = createGatewayWsTestSocket();
|
||||
const gatewayBudget = { release: vi.fn() };
|
||||
const rateLimiter = { check: vi.fn() };
|
||||
Object.assign(socket, {
|
||||
[GATEWAY_WS_CONNECTION_KIND_PROPERTY]: "worker",
|
||||
[GATEWAY_WS_WORKER_INGRESS_PROPERTY]: "public",
|
||||
__openclawPreauthBudgetKey: "203.0.113.10",
|
||||
});
|
||||
|
||||
await connectTestWs({
|
||||
socket,
|
||||
options: {
|
||||
preauthConnectionBudget: gatewayBudget as never,
|
||||
rateLimiter: rateLimiter as never,
|
||||
},
|
||||
});
|
||||
|
||||
const handler = firstAttachedWorkerHandlerParams() as {
|
||||
ingress: string;
|
||||
rateLimiter: unknown;
|
||||
rateLimitClientIp: string;
|
||||
setClient(client: never): boolean;
|
||||
};
|
||||
expect(handler).toMatchObject({
|
||||
ingress: "public",
|
||||
rateLimiter,
|
||||
rateLimitClientIp: "203.0.113.10",
|
||||
});
|
||||
expect(handler.setClient({ socket } as never)).toBe(true);
|
||||
expect(gatewayBudget.release).toHaveBeenCalledWith("203.0.113.10");
|
||||
});
|
||||
|
||||
it("threads current auth getters into the handshake handler instead of a stale snapshot", async () => {
|
||||
const initialAuth = createResolvedGatewayTokenAuth("token-before");
|
||||
let currentAuth = initialAuth;
|
||||
|
||||
@@ -51,6 +51,7 @@ import { resolveSharedGatewaySessionGeneration } from "./ws-shared-generation.js
|
||||
import {
|
||||
GATEWAY_WS_CONNECTION_KIND_PROPERTY,
|
||||
GATEWAY_WS_PREAUTH_BUDGET_PROPERTY,
|
||||
GATEWAY_WS_WORKER_INGRESS_PROPERTY,
|
||||
WS_HANDSHAKE_PHASES,
|
||||
type GatewayIngressWebSocket,
|
||||
type GatewayWsClient,
|
||||
@@ -274,11 +275,11 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
let closed = false;
|
||||
const openedAt = Date.now();
|
||||
const connId = randomUUID();
|
||||
const connectionKind =
|
||||
(socket as GatewayIngressWebSocket)[GATEWAY_WS_CONNECTION_KIND_PROPERTY] ?? "gateway";
|
||||
const ingressSocket = socket as GatewayIngressWebSocket;
|
||||
const connectionKind = ingressSocket[GATEWAY_WS_CONNECTION_KIND_PROPERTY] ?? "gateway";
|
||||
const workerIngress = ingressSocket[GATEWAY_WS_WORKER_INGRESS_PROPERTY] ?? "loopback";
|
||||
const connectionPreauthBudget =
|
||||
(socket as GatewayIngressWebSocket)[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY] ??
|
||||
preauthConnectionBudget;
|
||||
ingressSocket[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY] ?? preauthConnectionBudget;
|
||||
const { remoteAddr, remotePort, localAddr, localPort, endpoint } = resolveSocketAddress(socket);
|
||||
const preauthBudgetKey = (
|
||||
socket as WebSocket & {
|
||||
@@ -681,6 +682,9 @@ export function attachGatewayWsConnectionHandler(params: AttachGatewayWsConnecti
|
||||
connId,
|
||||
service: workerConnectionService,
|
||||
isStartupPending,
|
||||
ingress: workerIngress,
|
||||
rateLimiter: workerIngress === "public" ? rateLimiter : undefined,
|
||||
rateLimitClientIp: workerIngress === "public" ? preauthBudgetKey : undefined,
|
||||
send,
|
||||
close,
|
||||
isClosed: () => closed,
|
||||
|
||||
@@ -26,6 +26,7 @@ import {
|
||||
tryBeginGatewaySuspendAdmission,
|
||||
} from "../../../process/gateway-work-admission.js";
|
||||
import { createDeferredCore } from "../../../shared/deferred.js";
|
||||
import type { AuthRateLimiter } from "../../auth-rate-limit.js";
|
||||
import type { WorkerConnectionIdentity } from "../../worker-environments/connection-identity.js";
|
||||
import { createGatewayWsTestSocket } from "../ws-connection.test-helpers.js";
|
||||
import type { GatewayWsClient } from "../ws-types.js";
|
||||
@@ -136,12 +137,27 @@ function createLogger() {
|
||||
return { warn: vi.fn() };
|
||||
}
|
||||
|
||||
function createRateLimiter(overrides: Partial<AuthRateLimiter> = {}): AuthRateLimiter {
|
||||
return {
|
||||
check: vi.fn(() => ({ allowed: true, remaining: 10, retryAfterMs: 0 })),
|
||||
recordFailure: vi.fn(),
|
||||
recordFailureAndDelay: vi.fn(async () => {}),
|
||||
reset: vi.fn(),
|
||||
size: vi.fn(() => 0),
|
||||
prune: vi.fn(),
|
||||
dispose: vi.fn(),
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
function attachHarness(
|
||||
options: {
|
||||
admissionFailure?: WorkerAdmissionFailureReason;
|
||||
commitFailure?: WorkerTranscriptCommitErrorReason;
|
||||
identity?: WorkerConnectionIdentity;
|
||||
liveFailure?: WorkerLiveEventErrorDetails;
|
||||
ingress?: "loopback" | "public";
|
||||
rateLimiter?: AuthRateLimiter;
|
||||
onInferenceLaunch?: (sink: InferenceSink) => void;
|
||||
onSessionTool?: (signal: AbortSignal | undefined) => Promise<WorkerSessionToolResult>;
|
||||
validationFailure?: ReturnType<WorkerConnectionService["validateWorkerConnection"]>;
|
||||
@@ -201,11 +217,15 @@ function attachHarness(
|
||||
});
|
||||
const logGateway = createLogger();
|
||||
const logWsControl = createLogger();
|
||||
const setCloseCause = vi.fn();
|
||||
const setLastFrameMeta = vi.fn();
|
||||
const cleanup = attachWorkerWsMessageHandler({
|
||||
socket: socket as unknown as WebSocket,
|
||||
connId: "worker-connection",
|
||||
service,
|
||||
ingress: options.ingress,
|
||||
rateLimiter: options.rateLimiter,
|
||||
rateLimitClientIp: options.rateLimiter ? "203.0.113.10" : undefined,
|
||||
send: (frame) => responses.push(frame),
|
||||
close,
|
||||
isClosed: () => false,
|
||||
@@ -214,7 +234,7 @@ function attachHarness(
|
||||
setClient,
|
||||
setHandshakeState: vi.fn(),
|
||||
advanceHandshakePhase: vi.fn(),
|
||||
setCloseCause: vi.fn(),
|
||||
setCloseCause,
|
||||
setLastFrameMeta,
|
||||
logGateway,
|
||||
logWsControl,
|
||||
@@ -230,6 +250,7 @@ function attachHarness(
|
||||
responses,
|
||||
service,
|
||||
setClient,
|
||||
setCloseCause,
|
||||
setLastFrameMeta,
|
||||
sendRequest: (method: string, params: unknown, id = "request-1") =>
|
||||
send({ type: "req", id, method, params }),
|
||||
@@ -276,6 +297,85 @@ describe("dedicated worker websocket protocol", () => {
|
||||
expect(harness.setClient).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it.each(["invalid-credential", "environment-mismatch"] as const)(
|
||||
"projects public %s failures to one opaque reason",
|
||||
async (internalReason) => {
|
||||
const recordFailureAndDelay = vi.fn(async () => {});
|
||||
const rateLimiter = createRateLimiter({ recordFailureAndDelay });
|
||||
const harness = attachHarness({
|
||||
admissionFailure: internalReason,
|
||||
ingress: "public",
|
||||
rateLimiter,
|
||||
});
|
||||
harness.sendConnect();
|
||||
|
||||
await waitForWorkerProtocol(() =>
|
||||
expect(harness.close).toHaveBeenCalledWith(1008, "admission-rejected"),
|
||||
);
|
||||
expect(harness.responses[0]).toMatchObject({
|
||||
ok: false,
|
||||
error: { details: { reason: "admission-rejected" } },
|
||||
});
|
||||
expect(harness.logWsControl.warn).toHaveBeenCalledWith(
|
||||
`worker admission rejected reason=${internalReason}`,
|
||||
);
|
||||
expect(harness.setCloseCause).toHaveBeenCalledWith(internalReason);
|
||||
expect(recordFailureAndDelay).toHaveBeenCalledWith("203.0.113.10", "worker-admission");
|
||||
},
|
||||
);
|
||||
|
||||
it("rejects rate-limited public admission before credential verification", async () => {
|
||||
const rateLimiter = createRateLimiter({
|
||||
check: vi.fn(() => ({ allowed: false, remaining: 0, retryAfterMs: 12_000 })),
|
||||
});
|
||||
const harness = attachHarness({ ingress: "public", rateLimiter });
|
||||
harness.sendConnect();
|
||||
|
||||
await waitForWorkerProtocol(() =>
|
||||
expect(harness.close).toHaveBeenCalledWith(1008, "admission-rejected"),
|
||||
);
|
||||
expect(harness.responses[0]).toMatchObject({
|
||||
ok: false,
|
||||
error: {
|
||||
details: { reason: "admission-rejected" },
|
||||
retryable: true,
|
||||
retryAfterMs: 12_000,
|
||||
},
|
||||
});
|
||||
expect(harness.service.admitWorker).not.toHaveBeenCalled();
|
||||
expect(harness.setCloseCause).toHaveBeenCalledWith("rate-limited");
|
||||
});
|
||||
|
||||
it("resets public credential failures after successful admission", async () => {
|
||||
const reset = vi.fn();
|
||||
const rateLimiter = createRateLimiter({ reset });
|
||||
const harness = attachHarness({ ingress: "public", rateLimiter });
|
||||
await admit(harness);
|
||||
|
||||
expect(reset).toHaveBeenCalledWith("203.0.113.10", "worker-admission");
|
||||
});
|
||||
|
||||
it("keeps public ownership failures opaque without charging the credential budget", async () => {
|
||||
const reset = vi.fn();
|
||||
const recordFailureAndDelay = vi.fn(async () => {});
|
||||
const rateLimiter = createRateLimiter({ reset, recordFailureAndDelay });
|
||||
const harness = attachHarness({
|
||||
ingress: "public",
|
||||
rateLimiter,
|
||||
validationFailure: "credential-replaced",
|
||||
});
|
||||
harness.sendConnect();
|
||||
|
||||
await waitForWorkerProtocol(() =>
|
||||
expect(harness.close).toHaveBeenCalledWith(1008, "admission-rejected"),
|
||||
);
|
||||
expect(harness.logWsControl.warn).toHaveBeenCalledWith(
|
||||
"worker admission rejected reason=credential-replaced",
|
||||
);
|
||||
expect(reset).toHaveBeenCalledOnce();
|
||||
expect(recordFailureAndDelay).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it.each([
|
||||
["node.event", { event: "agent.request", payloadJSON: '{"requestId":"r-1"}' }],
|
||||
["health", {}],
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
import {
|
||||
ErrorCodes,
|
||||
type WorkerErrorShape,
|
||||
type WorkerHelloOk,
|
||||
type WorkerLiveEventErrorDetails,
|
||||
type WorkerLiveEventErrorShape,
|
||||
type WorkerProtocolCloseReason,
|
||||
type WorkerTranscriptCommitErrorReason,
|
||||
type WorkerTranscriptCommitErrorShape,
|
||||
WORKER_HEARTBEAT_INTERVAL_MS,
|
||||
WORKER_PROTOCOL_MAX_PAYLOAD_BYTES,
|
||||
} from "../../../../packages/gateway-protocol/src/index.js";
|
||||
import {
|
||||
type WorkerInferenceErrorReason,
|
||||
type WorkerInferenceErrorShape,
|
||||
WORKER_INFERENCE_PROTOCOL_FEATURE,
|
||||
WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES,
|
||||
} from "../../../../packages/gateway-protocol/src/schema/worker-inference.js";
|
||||
import type { WorkerConnectionIdentity } from "../../worker-environments/connection-identity.js";
|
||||
|
||||
export function workerProtocolError(
|
||||
reason: WorkerProtocolCloseReason,
|
||||
options: {
|
||||
code?: WorkerErrorShape["code"];
|
||||
message?: string;
|
||||
retryable?: boolean;
|
||||
retryAfterMs?: number;
|
||||
} = {},
|
||||
): WorkerErrorShape {
|
||||
return {
|
||||
code: options.code ?? ErrorCodes.INVALID_REQUEST,
|
||||
message: options.message ?? "worker protocol request rejected",
|
||||
details: { reason },
|
||||
...(options.retryable === undefined ? {} : { retryable: options.retryable }),
|
||||
...(options.retryAfterMs === undefined ? {} : { retryAfterMs: options.retryAfterMs }),
|
||||
};
|
||||
}
|
||||
|
||||
export function workerMaxPayload(identity: WorkerConnectionIdentity): number {
|
||||
return identity.protocolFeatures.includes(WORKER_INFERENCE_PROTOCOL_FEATURE)
|
||||
? WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES
|
||||
: WORKER_PROTOCOL_MAX_PAYLOAD_BYTES;
|
||||
}
|
||||
|
||||
export function buildWorkerHello(identity: WorkerConnectionIdentity): WorkerHelloOk {
|
||||
return {
|
||||
type: "worker-hello-ok",
|
||||
environmentId: identity.environmentId,
|
||||
sessionId: identity.sessionId,
|
||||
ownerEpoch: identity.ownerEpoch,
|
||||
rpcSetVersion: identity.rpcSetVersion,
|
||||
protocolFeatures: [...identity.protocolFeatures],
|
||||
credentialExpiresAtMs: identity.credentialExpiresAtMs,
|
||||
policy: {
|
||||
heartbeatIntervalMs: WORKER_HEARTBEAT_INTERVAL_MS,
|
||||
maxPayload: workerMaxPayload(identity),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function workerTranscriptCommitError(
|
||||
reason: WorkerTranscriptCommitErrorReason,
|
||||
): WorkerTranscriptCommitErrorShape {
|
||||
return {
|
||||
code: ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker transcript commit rejected",
|
||||
details: { reason },
|
||||
};
|
||||
}
|
||||
|
||||
export function workerLiveEventError(
|
||||
details: WorkerLiveEventErrorDetails,
|
||||
): WorkerLiveEventErrorShape {
|
||||
return {
|
||||
code: ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker live event rejected",
|
||||
details,
|
||||
};
|
||||
}
|
||||
|
||||
export function workerInferenceError(
|
||||
reason: WorkerInferenceErrorReason,
|
||||
): WorkerInferenceErrorShape {
|
||||
return {
|
||||
code: reason === "provider-error" ? ErrorCodes.UNAVAILABLE : ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker inference request rejected",
|
||||
details: { reason },
|
||||
};
|
||||
}
|
||||
@@ -7,7 +7,6 @@ import {
|
||||
type WorkerConnectParams,
|
||||
type WorkerErrorShape,
|
||||
type WorkerHeartbeatResult,
|
||||
type WorkerHelloOk,
|
||||
type WorkerLiveEventErrorDetails,
|
||||
type WorkerLiveEventErrorShape,
|
||||
type WorkerLiveEventParams,
|
||||
@@ -20,7 +19,6 @@ import {
|
||||
type WorkerTranscriptCommitErrorShape,
|
||||
type WorkerTranscriptCommitParams,
|
||||
type WorkerTranscriptCommitResult,
|
||||
WORKER_HEARTBEAT_INTERVAL_MS,
|
||||
WORKER_LIVE_EVENT_PROTOCOL_FEATURE,
|
||||
WORKER_SESSION_TOOLS_PROTOCOL_FEATURE,
|
||||
WORKER_PROTOCOL_MAX_FRAME_ID_LENGTH,
|
||||
@@ -47,7 +45,6 @@ import {
|
||||
type WorkerInferenceTerminalFrame,
|
||||
WORKER_INFERENCE_METHODS,
|
||||
WORKER_INFERENCE_PROTOCOL_FEATURE,
|
||||
WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES,
|
||||
validateWorkerInferenceCancelParams,
|
||||
validateWorkerInferenceStartParams,
|
||||
} from "../../../../packages/gateway-protocol/src/schema/worker-inference.js";
|
||||
@@ -57,9 +54,21 @@ import {
|
||||
runWithGatewayIndependentRootWorkContinuation,
|
||||
tryBeginGatewayRootWorkAdmission,
|
||||
} from "../../../process/gateway-work-admission.js";
|
||||
import {
|
||||
AUTH_RATE_LIMIT_SCOPE_WORKER_ADMISSION,
|
||||
type AuthRateLimiter,
|
||||
} from "../../auth-rate-limit.js";
|
||||
import type { WorkerConnectionIdentity } from "../../worker-environments/connection-identity.js";
|
||||
import { MAX_RUNNING_WORKER_SESSION_TOOL_OPERATIONS } from "../../worker-environments/placement-session-tool-operations.js";
|
||||
import type { GatewayWsClient, WsHandshakePhase } from "../ws-types.js";
|
||||
import type { GatewayWorkerIngress, GatewayWsClient, WsHandshakePhase } from "../ws-types.js";
|
||||
import {
|
||||
buildWorkerHello,
|
||||
workerInferenceError,
|
||||
workerLiveEventError,
|
||||
workerMaxPayload,
|
||||
workerProtocolError,
|
||||
workerTranscriptCommitError,
|
||||
} from "./worker-connection-frames.js";
|
||||
|
||||
type WorkerServiceResult<TResult, TFailure> =
|
||||
| { ok: true; result: TResult }
|
||||
@@ -128,11 +137,22 @@ type WorkerLogger = { warn(message: string): void };
|
||||
const MAX_QUEUED_WORKER_FRAMES = 16;
|
||||
const MAX_QUEUED_WORKER_BYTES = 32 * 1024 * 1024;
|
||||
|
||||
function isWorkerCredentialFailure(reason: WorkerProtocolCloseReason): boolean {
|
||||
return (
|
||||
reason === "invalid-credential" ||
|
||||
reason === "environment-mismatch" ||
|
||||
reason === "credential-expired"
|
||||
);
|
||||
}
|
||||
|
||||
type WorkerWsMessageHandlerParams = {
|
||||
socket: WebSocket;
|
||||
connId: string;
|
||||
service?: WorkerConnectionService;
|
||||
isStartupPending?: () => boolean;
|
||||
ingress?: GatewayWorkerIngress;
|
||||
rateLimiter?: AuthRateLimiter;
|
||||
rateLimitClientIp?: string;
|
||||
send(frame: unknown): void;
|
||||
close(code?: number, reason?: string): void;
|
||||
isClosed(): boolean;
|
||||
@@ -147,46 +167,6 @@ type WorkerWsMessageHandlerParams = {
|
||||
logWsControl: WorkerLogger;
|
||||
};
|
||||
|
||||
function workerProtocolError(
|
||||
reason: WorkerProtocolCloseReason,
|
||||
options: {
|
||||
code?: WorkerErrorShape["code"];
|
||||
message?: string;
|
||||
retryable?: boolean;
|
||||
retryAfterMs?: number;
|
||||
} = {},
|
||||
): WorkerErrorShape {
|
||||
return {
|
||||
code: options.code ?? ErrorCodes.INVALID_REQUEST,
|
||||
message: options.message ?? "worker protocol request rejected",
|
||||
details: { reason },
|
||||
...(options.retryable === undefined ? {} : { retryable: options.retryable }),
|
||||
...(options.retryAfterMs === undefined ? {} : { retryAfterMs: options.retryAfterMs }),
|
||||
};
|
||||
}
|
||||
|
||||
function workerMaxPayload(identity: WorkerConnectionIdentity): number {
|
||||
return identity.protocolFeatures.includes(WORKER_INFERENCE_PROTOCOL_FEATURE)
|
||||
? WORKER_PROTOCOL_MAX_INFERENCE_PAYLOAD_BYTES
|
||||
: WORKER_PROTOCOL_MAX_PAYLOAD_BYTES;
|
||||
}
|
||||
|
||||
function buildWorkerHello(identity: WorkerConnectionIdentity): WorkerHelloOk {
|
||||
return {
|
||||
type: "worker-hello-ok",
|
||||
environmentId: identity.environmentId,
|
||||
sessionId: identity.sessionId,
|
||||
ownerEpoch: identity.ownerEpoch,
|
||||
rpcSetVersion: identity.rpcSetVersion,
|
||||
protocolFeatures: [...identity.protocolFeatures],
|
||||
credentialExpiresAtMs: identity.credentialExpiresAtMs,
|
||||
policy: {
|
||||
heartbeatIntervalMs: WORKER_HEARTBEAT_INTERVAL_MS,
|
||||
maxPayload: workerMaxPayload(identity),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function rejectWorkerRequest(params: {
|
||||
reason: WorkerProtocolCloseReason;
|
||||
respond: WorkerRespond;
|
||||
@@ -198,32 +178,6 @@ function rejectWorkerRequest(params: {
|
||||
queueMicrotask(() => params.close(1008, params.reason));
|
||||
}
|
||||
|
||||
function workerTranscriptCommitError(
|
||||
reason: WorkerTranscriptCommitErrorReason,
|
||||
): WorkerTranscriptCommitErrorShape {
|
||||
return {
|
||||
code: ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker transcript commit rejected",
|
||||
details: { reason },
|
||||
};
|
||||
}
|
||||
|
||||
function workerLiveEventError(details: WorkerLiveEventErrorDetails): WorkerLiveEventErrorShape {
|
||||
return {
|
||||
code: ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker live event rejected",
|
||||
details,
|
||||
};
|
||||
}
|
||||
|
||||
function workerInferenceError(reason: WorkerInferenceErrorReason): WorkerInferenceErrorShape {
|
||||
return {
|
||||
code: reason === "provider-error" ? ErrorCodes.UNAVAILABLE : ErrorCodes.INVALID_REQUEST,
|
||||
message: "worker inference request rejected",
|
||||
details: { reason },
|
||||
};
|
||||
}
|
||||
|
||||
function setSocketMaxPayload(socket: WebSocket, maxPayload: number): void {
|
||||
const receiver = (socket as { _receiver?: { _maxPayload?: number } })["_receiver"];
|
||||
if (receiver) {
|
||||
@@ -441,16 +395,31 @@ export function attachWorkerWsMessageHandler(params: WorkerWsMessageHandlerParam
|
||||
params.send({ type: "res", id, ok: false, error });
|
||||
queueMicrotask(() => closeWorker(code, reason));
|
||||
};
|
||||
const rejectAdmission = (
|
||||
id: string,
|
||||
reason: WorkerProtocolCloseReason,
|
||||
error = workerProtocolError(reason, { message: "worker admission rejected" }),
|
||||
code = 1008,
|
||||
) => {
|
||||
const rejectAdmission = (rejection: {
|
||||
id: string;
|
||||
reason: WorkerProtocolCloseReason;
|
||||
internalReason?: string;
|
||||
error?: WorkerErrorShape;
|
||||
code?: number;
|
||||
}) => {
|
||||
const internalReason = rejection.internalReason ?? rejection.reason;
|
||||
params.setHandshakeState("failed");
|
||||
params.setCloseCause(reason);
|
||||
params.logWsControl.warn(`worker admission rejected reason=${reason}`);
|
||||
sendError(id, reason, error, code);
|
||||
params.setCloseCause(internalReason);
|
||||
params.logWsControl.warn(`worker admission rejected reason=${internalReason}`);
|
||||
sendError(
|
||||
rejection.id,
|
||||
rejection.reason,
|
||||
rejection.error ??
|
||||
workerProtocolError(rejection.reason, { message: "worker admission rejected" }),
|
||||
rejection.code ?? 1008,
|
||||
);
|
||||
};
|
||||
const rejectVerifiedAdmission = (id: string, internalReason: WorkerProtocolCloseReason) => {
|
||||
rejectAdmission({
|
||||
id,
|
||||
reason: params.ingress === "public" ? "admission-rejected" : internalReason,
|
||||
internalReason,
|
||||
});
|
||||
};
|
||||
|
||||
const handleConnect = async (
|
||||
@@ -459,33 +428,58 @@ export function attachWorkerWsMessageHandler(params: WorkerWsMessageHandlerParam
|
||||
admissionOpen: boolean,
|
||||
) => {
|
||||
if (!admissionOpen || params.isStartupPending?.()) {
|
||||
rejectAdmission(
|
||||
rejectAdmission({
|
||||
id,
|
||||
"gateway-unavailable",
|
||||
workerProtocolError("gateway-unavailable", {
|
||||
reason: "gateway-unavailable",
|
||||
error: workerProtocolError("gateway-unavailable", {
|
||||
code: ErrorCodes.UNAVAILABLE,
|
||||
message: "worker gateway unavailable",
|
||||
retryable: true,
|
||||
retryAfterMs: GATEWAY_STARTUP_RETRY_AFTER_MS,
|
||||
}),
|
||||
1013,
|
||||
);
|
||||
code: 1013,
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (connect.minProtocol > PROTOCOL_VERSION || connect.maxProtocol < PROTOCOL_VERSION) {
|
||||
rejectAdmission(id, "protocol-mismatch");
|
||||
rejectAdmission({ id, reason: "protocol-mismatch" });
|
||||
return;
|
||||
}
|
||||
const rateLimit = params.rateLimiter?.check(
|
||||
params.rateLimitClientIp,
|
||||
AUTH_RATE_LIMIT_SCOPE_WORKER_ADMISSION,
|
||||
);
|
||||
if (rateLimit && !rateLimit.allowed) {
|
||||
rejectAdmission({
|
||||
id,
|
||||
reason: "admission-rejected",
|
||||
internalReason: "rate-limited",
|
||||
error: workerProtocolError("admission-rejected", {
|
||||
code: ErrorCodes.UNAVAILABLE,
|
||||
message: "worker admission rejected",
|
||||
retryable: true,
|
||||
retryAfterMs: rateLimit.retryAfterMs,
|
||||
}),
|
||||
});
|
||||
return;
|
||||
}
|
||||
const admission =
|
||||
(await params.service?.admitWorker(connect.admission)) ??
|
||||
({ ok: false, reason: "environment-unavailable" } as const);
|
||||
if (!admission.ok) {
|
||||
rejectAdmission(id, admission.reason);
|
||||
if (isWorkerCredentialFailure(admission.reason)) {
|
||||
await params.rateLimiter?.recordFailureAndDelay(
|
||||
params.rateLimitClientIp,
|
||||
AUTH_RATE_LIMIT_SCOPE_WORKER_ADMISSION,
|
||||
);
|
||||
}
|
||||
rejectVerifiedAdmission(id, admission.reason);
|
||||
return;
|
||||
}
|
||||
params.rateLimiter?.reset(params.rateLimitClientIp, AUTH_RATE_LIMIT_SCOPE_WORKER_ADMISSION);
|
||||
const ownershipFailure = params.service?.validateWorkerConnection(admission.identity);
|
||||
if (ownershipFailure) {
|
||||
rejectAdmission(id, ownershipFailure);
|
||||
rejectVerifiedAdmission(id, ownershipFailure);
|
||||
return;
|
||||
}
|
||||
const client: GatewayWsClient = {
|
||||
|
||||
@@ -7,12 +7,15 @@ import type { WorkerConnectionIdentity } from "../worker-environments/connection
|
||||
|
||||
export const GATEWAY_WS_CONNECTION_KIND_PROPERTY = "__openclawConnectionKind";
|
||||
export const GATEWAY_WS_PREAUTH_BUDGET_PROPERTY = "__openclawPreauthBudget";
|
||||
export const GATEWAY_WS_WORKER_INGRESS_PROPERTY = "__openclawWorkerIngress";
|
||||
type GatewayWsConnectionKind = "gateway" | "worker";
|
||||
export type GatewayWorkerIngress = "loopback" | "public";
|
||||
export type GatewayIngressWebSocket = WebSocket & {
|
||||
[GATEWAY_WS_CONNECTION_KIND_PROPERTY]?: GatewayWsConnectionKind;
|
||||
[GATEWAY_WS_PREAUTH_BUDGET_PROPERTY]?: {
|
||||
release(clientIp: string | undefined): void;
|
||||
};
|
||||
[GATEWAY_WS_WORKER_INGRESS_PROPERTY]?: GatewayWorkerIngress;
|
||||
__openclawPreauthBudgetClaimed?: boolean;
|
||||
__openclawPreauthBudgetKey?: string;
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user