mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
perf(secrets): reuse upstream TLS trust and connections (#156357)
Parse the immutable proxy CA bundle once and give each process grant an origin-pooled HTTPS agent. Revoke active and idle connections with their grant. A 1,000-request local TLS rig reduced SecureContext creation and handshakes from 1,000 to 1, and main-thread CPU from 5.048 to 0.383 ms/request with identical forwarded payloads. All 99 focused proxy tests pass on Blacksmith Testbox.
This commit is contained in:
parent
c8d3f81045
commit
668eef7522
5 changed files with 226 additions and 47 deletions
|
|
@ -696,6 +696,60 @@ it("rejects oversized direct bridge responses", async () => {
|
|||
});
|
||||
});
|
||||
|
||||
it("binds direct bridge tokens to the relay they were issued for", async () => {
|
||||
await withOpenClawTestState({ label: "relay-token-binding" }, async () => {
|
||||
const first = registerOwnedNativeHookRelay({
|
||||
provider: "codex",
|
||||
relayId: "codex-first-bridge-session",
|
||||
sessionId: "session-1",
|
||||
runId: "run-1",
|
||||
allowedEvents: ["pre_tool_use"],
|
||||
});
|
||||
const second = registerOwnedNativeHookRelay({
|
||||
provider: "codex",
|
||||
relayId: "codex-second-bridge-session",
|
||||
sessionId: "session-2",
|
||||
runId: "run-2",
|
||||
allowedEvents: ["pre_tool_use"],
|
||||
});
|
||||
try {
|
||||
await Promise.all([first.ready, second.ready]);
|
||||
const firstRecord = await store.readNativeHookRelayBridgeRecord({ relayId: first.relayId });
|
||||
if (!firstRecord) {
|
||||
throw new Error("test bridge registration unavailable");
|
||||
}
|
||||
await store.writeNativeHookRelayBridgeRecord({
|
||||
record: { ...firstRecord, relayId: second.relayId, expiresAtMs: Date.now() + 10_000 },
|
||||
});
|
||||
// Cold locator startup must not consume this token-binding fixture's caller deadline.
|
||||
const clock = vi.spyOn(Date, "now").mockReturnValue(Date.now());
|
||||
try {
|
||||
await expect(
|
||||
invokeNativeHookRelayBridge({
|
||||
provider: "codex",
|
||||
relayId: second.relayId,
|
||||
generation: second.generation,
|
||||
event: "pre_tool_use",
|
||||
timeoutMs: 500,
|
||||
rawPayload: {
|
||||
hook_event_name: "PreToolUse",
|
||||
tool_name: "Bash",
|
||||
tool_input: { command: "pnpm test" },
|
||||
},
|
||||
}),
|
||||
).rejects.toThrow("native hook relay bridge target mismatch");
|
||||
expect(testing.getNativeHookRelayInvocationsForTests()).toStrictEqual([]);
|
||||
} finally {
|
||||
clock.mockRestore();
|
||||
}
|
||||
} finally {
|
||||
first.unregister();
|
||||
second.unregister();
|
||||
await Promise.all([first.drain(), second.drain()]);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
it("does not start transport when locator lookup consumes the caller deadline", async () => {
|
||||
await withOpenClawTestState({ label: "relay-lookup-deadline" }, async () => {
|
||||
const relay = registerOwnedNativeHookRelay({
|
||||
|
|
|
|||
|
|
@ -2189,49 +2189,6 @@ describe("native hook relay registry", () => {
|
|||
).rejects.toThrow("native hook relay bridge not found");
|
||||
});
|
||||
|
||||
it("binds direct bridge tokens to the relay they were issued for", async () => {
|
||||
const first = registerNativeHookRelay({
|
||||
provider: "codex",
|
||||
relayId: "codex-first-bridge-session",
|
||||
sessionId: "session-1",
|
||||
runId: "run-1",
|
||||
allowedEvents: ["pre_tool_use"],
|
||||
});
|
||||
const second = registerNativeHookRelay({
|
||||
provider: "codex",
|
||||
relayId: "codex-second-bridge-session",
|
||||
sessionId: "session-2",
|
||||
runId: "run-2",
|
||||
allowedEvents: ["pre_tool_use"],
|
||||
});
|
||||
|
||||
const firstRecord = await waitForNativeHookRelayBridgeRecord(first.relayId);
|
||||
await waitForNativeHookRelayBridgeRecord(second.relayId);
|
||||
await nativeHookRelayStore.writeNativeHookRelayBridgeRecord({
|
||||
record: {
|
||||
...firstRecord,
|
||||
relayId: second.relayId,
|
||||
expiresAtMs: Date.now() + 10_000,
|
||||
},
|
||||
});
|
||||
|
||||
await expect(
|
||||
invokeNativeHookRelayBridge({
|
||||
provider: "codex",
|
||||
relayId: second.relayId,
|
||||
generation: second.generation,
|
||||
event: "pre_tool_use",
|
||||
timeoutMs: 500,
|
||||
rawPayload: {
|
||||
hook_event_name: "PreToolUse",
|
||||
tool_name: "Bash",
|
||||
tool_input: { command: "pnpm test" },
|
||||
},
|
||||
}),
|
||||
).rejects.toThrow("native hook relay bridge target mismatch");
|
||||
expect(testing.getNativeHookRelayInvocationsForTests()).toStrictEqual([]);
|
||||
});
|
||||
|
||||
it("accepts an allowed Codex invocation and preserves raw payload", async () => {
|
||||
const relay = registerNativeHookRelay({
|
||||
provider: "codex",
|
||||
|
|
|
|||
|
|
@ -189,6 +189,13 @@ function sendSecretEgressRequest(
|
|||
},
|
||||
),
|
||||
);
|
||||
upstream.once("socket", () => {
|
||||
// An agent can queue prepared credentials; recheck before Node flushes them.
|
||||
if (!forward.isActive()) {
|
||||
upstream.destroy();
|
||||
forward.response.destroy();
|
||||
}
|
||||
});
|
||||
const bodyTransform = forward.ownResource(
|
||||
createSecretEgressBodyTransform({
|
||||
isActive: forward.isActive,
|
||||
|
|
|
|||
|
|
@ -188,6 +188,154 @@ afterEach(async () => {
|
|||
});
|
||||
|
||||
describe("secret egress registration lifecycle", () => {
|
||||
it.each(["keep-alive", "close"])(
|
||||
"rejects queued credentials after shared revocation before socket assignment (%s)",
|
||||
async (connection) => {
|
||||
const authority = new Int32Array(new SharedArrayBuffer(4));
|
||||
grant = proxy.registerProcess(
|
||||
[{ name: "SERVICE_API_KEY", sentinel, allowedHosts: ["localhost"] }],
|
||||
() => Atomics.load(authority, 0) === 0,
|
||||
);
|
||||
const queued = createDeferredCore<{ request: http.ClientRequest; agent: https.Agent }>();
|
||||
const { request: requestUpstream } = await vi.importActual<typeof https>("node:https");
|
||||
vi.mocked(https.request).mockImplementation((...args) => {
|
||||
const agent = (args[0] as https.RequestOptions).agent;
|
||||
if (!(agent instanceof https.Agent)) {
|
||||
throw new Error("Expected the grant's HTTPS agent");
|
||||
}
|
||||
agent.maxSockets = 1;
|
||||
agent.maxTotalSockets = 1;
|
||||
const request = requestUpstream(...args);
|
||||
if (request.path === "/queued") {
|
||||
queued.resolve({ request, agent });
|
||||
}
|
||||
return request;
|
||||
});
|
||||
const held = createDeferredCore<ServerResponse>();
|
||||
const paths: string[] = [];
|
||||
origin.removeAllListeners("request");
|
||||
origin.on("request", (request, response) => {
|
||||
paths.push(request.url!);
|
||||
expect(request.headers.authorization).toBe(`Bearer ${value}`);
|
||||
request.resume();
|
||||
request.once("end", () => {
|
||||
if (request.url === "/held") {
|
||||
held.resolve(response);
|
||||
} else {
|
||||
response.end("ok");
|
||||
}
|
||||
});
|
||||
});
|
||||
const first = await openTlsTunnel();
|
||||
first.write(
|
||||
`GET /held HTTP/1.1\r\nHost: localhost:${originPort}\r\nConnection: ${connection}\r\nAuthorization: Bearer ${sentinel}\r\nContent-Length: 0\r\n\r\n`,
|
||||
);
|
||||
const heldResponse = await held.promise;
|
||||
const second = await openTlsTunnel();
|
||||
second.write(
|
||||
`GET /queued HTTP/1.1\r\nHost: localhost:${originPort}\r\nAuthorization: Bearer ${sentinel}\r\nContent-Length: 0\r\n\r\n`,
|
||||
);
|
||||
const { request: pending, agent } = await queued.promise;
|
||||
expect(pending.socket).toBeNull();
|
||||
expect(pending.writableEnded).toBe(true);
|
||||
const pendingClosed = new Promise<void>((resolve) => {
|
||||
pending.once("close", resolve);
|
||||
});
|
||||
const secondClosed = onClose(second);
|
||||
expect(Object.values(agent.requests).flat()).toContain(pending);
|
||||
// Revoke after response guards finish, before Node dispatches the socket queue.
|
||||
// The Worker cleanup message deliberately has not called grant.revoke().
|
||||
agent.prependOnceListener("free", () => Atomics.store(authority, 0, 1));
|
||||
heldResponse.end("ok");
|
||||
await pendingClosed;
|
||||
expect(paths).toEqual(["/held"]);
|
||||
await secondClosed;
|
||||
await sendCredential(await openTlsTunnel(register().env));
|
||||
expect(paths).toEqual(["/held", "/"]);
|
||||
},
|
||||
);
|
||||
|
||||
it("reuses upstream TLS only within a live grant and releases idle connections on revocation", async () => {
|
||||
const peers: Socket[] = [];
|
||||
origin.removeAllListeners("request");
|
||||
origin.on("request", (request, response) => {
|
||||
peers.push(request.socket);
|
||||
expect(request.headers.authorization).toBe(`Bearer ${value}`);
|
||||
expect(request.headers["proxy-authorization"]).toBeUndefined();
|
||||
request.resume();
|
||||
response.writeHead(200, { "Content-Length": 2 });
|
||||
response.end("ok");
|
||||
});
|
||||
const send = async (processGrant: SecretEgressProcessGrant, connection = "keep-alive") => {
|
||||
const url = new URL(processGrant.env.HTTPS_PROXY!);
|
||||
return new Promise<void>((resolve, reject) => {
|
||||
const request = httpRequest(
|
||||
{
|
||||
hostname: url.hostname,
|
||||
port: url.port,
|
||||
path: `https://localhost:${originPort}/`,
|
||||
agent: false,
|
||||
headers: {
|
||||
Connection: connection,
|
||||
Authorization: `Bearer ${sentinel}`,
|
||||
"Proxy-Authorization": `Basic ${Buffer.from(`openclaw:${url.password}`).toString("base64")}`,
|
||||
},
|
||||
},
|
||||
(response) => {
|
||||
let body = "";
|
||||
response.on("data", (chunk: Buffer) => (body += chunk.toString()));
|
||||
response.once("end", () => {
|
||||
expect(response.statusCode).toBe(200);
|
||||
expect(body).toBe("ok");
|
||||
resolve();
|
||||
});
|
||||
},
|
||||
);
|
||||
request.once("error", reject);
|
||||
request.end();
|
||||
});
|
||||
};
|
||||
const connect = vi.spyOn(tls, "connect");
|
||||
const sibling = register();
|
||||
await send(grant);
|
||||
await send(grant);
|
||||
await send(sibling);
|
||||
expect(peers[1]).toBe(peers[0]);
|
||||
expect(peers[2]).not.toBe(peers[0]);
|
||||
const revokedPeerClosed = onClose(peers[0]!);
|
||||
grant.revoke();
|
||||
await revokedPeerClosed;
|
||||
await send(sibling);
|
||||
expect(peers[3]).toBe(peers[2]);
|
||||
// Even a peer-requested reconnect must reuse parsed trust, without retaining
|
||||
// request credentials in TLS options or an agent's persistent options.
|
||||
const siblingPeerClosed = onClose(peers[2]!);
|
||||
await send(sibling, "close");
|
||||
await siblingPeerClosed;
|
||||
await send(sibling);
|
||||
expect(peers[5]).not.toBe(peers[2]);
|
||||
const tlsOptions = connect.mock.calls.map(([options]) => options as tls.ConnectionOptions);
|
||||
expect(tlsOptions).toHaveLength(3);
|
||||
expect(tlsOptions[0]?.secureContext).toBeDefined();
|
||||
for (const option of tlsOptions) {
|
||||
expect(option.secureContext).toBe(tlsOptions[0]?.secureContext);
|
||||
expect(option.ca).toBeUndefined();
|
||||
expect(option.key).toBeUndefined();
|
||||
expect(option.cert).toBeUndefined();
|
||||
}
|
||||
for (const [options] of vi.mocked(https.request).mock.calls) {
|
||||
const agent = (options as https.RequestOptions).agent;
|
||||
expect(agent).toBeInstanceOf(https.Agent);
|
||||
if (agent instanceof https.Agent) {
|
||||
expect(JSON.stringify(agent.options)).not.toContain(value);
|
||||
expect(agent.options).not.toHaveProperty("headers");
|
||||
}
|
||||
}
|
||||
const stoppedPeerClosed = onClose(peers[5]!);
|
||||
await proxy.stop();
|
||||
await stoppedPeerClosed;
|
||||
});
|
||||
|
||||
it.each(["header", "body", "url"] as const)(
|
||||
"keeps an in-flight %s credential bound to its process snapshot",
|
||||
async (location) => {
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ import { Agent as HttpsAgent } from "node:https";
|
|||
import net, { type Socket } from "node:net";
|
||||
import path from "node:path";
|
||||
import { Readable, type Duplex, type Writable } from "node:stream";
|
||||
import { createServer as createTlsServer, rootCertificates } from "node:tls";
|
||||
import { createSecureContext, createServer as createTlsServer, rootCertificates } from "node:tls";
|
||||
import { URL } from "node:url";
|
||||
import { normalizeExactAllowedHost as normalizeHostname } from "../exact-hostname.js";
|
||||
import {
|
||||
|
|
@ -78,6 +78,7 @@ type RegisteredProcess = {
|
|||
resolveSentinel: (sentinel: string) => string | undefined;
|
||||
resources: Set<Readable | Writable>;
|
||||
tlsServers: Map<string, SecretEgressTlsContext>;
|
||||
upstreamTlsAgent: HttpsAgent;
|
||||
};
|
||||
|
||||
function parseConnectTarget(rawTarget: string | undefined): ConnectTarget {
|
||||
|
|
@ -243,7 +244,9 @@ export async function startSecretEgressProxyServer(params: {
|
|||
const { caPem } = certificates;
|
||||
const trustBundlePath = path.join(params.caDir, "trust-bundle.pem");
|
||||
fs.writeFileSync(trustBundlePath, `${rootCertificates.join("\n")}\n${caPem}`, { mode: 0o644 });
|
||||
const upstreamTlsAgent = new HttpsAgent({
|
||||
// The CA set is immutable for this proxy lifetime. A new proxy owns new trust;
|
||||
// leaf renewal does not change it. Share only parsed CAs across process grants.
|
||||
const upstreamSecureContext = createSecureContext({
|
||||
ca: [...rootCertificates, caPem],
|
||||
});
|
||||
const bypassHosts = new Set((params.bypassHosts ?? []).map(normalizeHostname));
|
||||
|
|
@ -288,6 +291,7 @@ export async function startSecretEgressProxyServer(params: {
|
|||
const revokeRegistration = (registered: RegisteredProcess) => {
|
||||
registrations.delete(registered);
|
||||
registered.sentinelBindings.clear();
|
||||
registered.upstreamTlsAgent.destroy();
|
||||
for (const resource of registered.resources) {
|
||||
resource.destroy();
|
||||
}
|
||||
|
|
@ -427,7 +431,7 @@ export async function startSecretEgressProxyServer(params: {
|
|||
substituted: swappedUrl.substituted || swappedHeaders.substituted,
|
||||
};
|
||||
},
|
||||
upstreamTlsAgent,
|
||||
upstreamTlsAgent: forward.registered.upstreamTlsAgent,
|
||||
isActive: forward.registered.isActive,
|
||||
ownResource: (resource) => ownResource(forward.registered, resource),
|
||||
releaseResponse: () => {
|
||||
|
|
@ -644,6 +648,16 @@ export async function startSecretEgressProxyServer(params: {
|
|||
resolveSentinel: params.resolveSentinel ?? resolveSecretSentinel,
|
||||
resources: new Set(),
|
||||
tlsServers: new Map(),
|
||||
// Node pools by origin. Grant ownership keeps connections and TLS sessions
|
||||
// isolated and lets revocation destroy idle sockets as well as live work.
|
||||
upstreamTlsAgent: new HttpsAgent({
|
||||
secureContext: upstreamSecureContext,
|
||||
keepAlive: true,
|
||||
maxSockets: 64,
|
||||
maxTotalSockets: 64,
|
||||
maxFreeSockets: 4,
|
||||
timeout: 30_000,
|
||||
}),
|
||||
};
|
||||
registrations.add(registered);
|
||||
// Basic is deliberately used because curl and Go net/http derive it from
|
||||
|
|
@ -674,7 +688,6 @@ export async function startSecretEgressProxyServer(params: {
|
|||
for (const registered of registrations.values()) {
|
||||
revokeRegistration(registered);
|
||||
}
|
||||
upstreamTlsAgent.destroy();
|
||||
for (const socket of sockets) {
|
||||
socket.destroy();
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue