From a1b00be35d79ae798c2d52a8d8715f60c11f06da Mon Sep 17 00:00:00 2001 From: Nathan Flurry Date: Mon, 31 Aug 2026 04:26:02 -0700 Subject: [PATCH] fix: support RivetKit Node dependencies in AgentOS --- crates/execution/src/javascript.rs | 23 ++- .../core/tests/network-http-request.test.ts | 132 ++++++++++++++++-- 2 files changed, 140 insertions(+), 15 deletions(-) diff --git a/crates/execution/src/javascript.rs b/crates/execution/src/javascript.rs index 6eb863fa28..1b169046e5 100644 --- a/crates/execution/src/javascript.rs +++ b/crates/execution/src/javascript.rs @@ -6926,7 +6926,20 @@ fn builtin_named_exports(module_name: &str) -> &'static [&'static str] { "validateHeaderName", "validateHeaderValue", ], - "http2" => &["connect", "createServer", "createSecureServer"], + "http2" => &[ + "Http2ServerRequest", + "Http2ServerResponse", + "Http2Session", + "Http2Stream", + "connect", + "constants", + "createServer", + "createSecureServer", + "getDefaultSettings", + "getPackedSettings", + "getUnpackedSettings", + "sensitiveHeaders", + ], "https" => &[ "Agent", "ClientRequest", @@ -7457,6 +7470,14 @@ mod tests { use std::time::{SystemTime, UNIX_EPOCH}; use tempfile::tempdir; + #[test] + fn http2_builtin_exposes_node_compatibility_classes() { + let exports = builtin_named_exports("http2"); + assert!(exports.contains(&"Http2ServerRequest")); + assert!(exports.contains(&"Http2ServerResponse")); + assert!(exports.contains(&"constants")); + } + #[test] fn dispose_context_reclaims_one_shot_metadata_without_reusing_ids() { let mut engine = JavascriptExecutionEngine::default(); diff --git a/packages/core/tests/network-http-request.test.ts b/packages/core/tests/network-http-request.test.ts index e5aaf54b0a..ae4cc9d368 100644 --- a/packages/core/tests/network-http-request.test.ts +++ b/packages/core/tests/network-http-request.test.ts @@ -1,5 +1,6 @@ import { createServer } from "node:http"; import { afterEach, describe, expect, test } from "vitest"; +import { WebSocketServer } from "ws"; import { AgentOs } from "../src/index.js"; const textDecoder = new TextDecoder(); @@ -54,10 +55,10 @@ describe("guest http.request transport", () => { ' socket.on("data", (chunk) => {', " buffered += chunk;", ' if (!buffered.includes("\\r\\n\\r\\n")) return;', - ' socket.end([', + " socket.end([", ' "HTTP/1.1 200 OK",', ' "Content-Type: application/json",', - ' `Content-Length: ${Buffer.byteLength(body)}`,', + " `Content-Length: ${Buffer.byteLength(body)}`,", ' "Connection: close",', ' "",', " body,", @@ -71,7 +72,7 @@ describe("guest http.request transport", () => { " process.exit(1);", " return;", " }", - ' const req = http.get(`http://127.0.0.1:${address.port}/transport-check`, (res) => {', + " const req = http.get(`http://127.0.0.1:${address.port}/transport-check`, (res) => {", ' let responseBody = "";', ' res.setEncoding("utf8");', ' res.on("data", (chunk) => {', @@ -79,12 +80,12 @@ describe("guest http.request transport", () => { " });", ' res.on("end", () => {', " console.log(JSON.stringify({ statusCode: res.statusCode, body: responseBody }));", - ' server.close(() => process.exit(0));', + " server.close(() => process.exit(0));", " });", " });", ' req.on("error", (error) => {', - ' console.error(error?.stack ?? String(error));', - ' server.close(() => process.exit(1));', + " console.error(error?.stack ?? String(error));", + " server.close(() => process.exit(1));", " });", "});", ].join("\n"); @@ -176,11 +177,11 @@ describe("guest http.request transport", () => { 'const http = require("node:http");', "void (async () => {", "const server = http.createServer((request, response) => {", - ' response.end(`self:${request.url}`);', + " response.end(`self:${request.url}`);", "});", 'await new Promise((resolve, reject) => { server.once("error", reject); server.listen(0, "127.0.0.1", resolve); });', "const address = server.address();", - 'const response = await fetch(`http://127.0.0.1:${address.port}/self-fetch`);', + "const response = await fetch(`http://127.0.0.1:${address.port}/self-fetch`);", "console.log(await response.text());", "await new Promise((resolve) => server.close(resolve));", "})();", @@ -215,12 +216,12 @@ describe("guest http.request transport", () => { ' response.write("data: ready\\n\\n");', " return;", " }", - ' response.end(`control:${request.url}`);', + " response.end(`control:${request.url}`);", "});", 'await new Promise((resolve, reject) => { server.once("error", reject); server.listen(0, "127.0.0.1", resolve); });', "const address = server.address();", "const abort = new AbortController();", - 'const events = await fetch(`http://127.0.0.1:${address.port}/events`, { signal: abort.signal });', + "const events = await fetch(`http://127.0.0.1:${address.port}/events`, { signal: abort.signal });", "const reader = events.body.getReader();", "await reader.read();", "const controls = [];", @@ -241,11 +242,110 @@ describe("guest http.request transport", () => { expect(result, result.stderr).toMatchObject({ exitCode: 0, - stdout: "control:/control-0,control:/control-1,control:/control-2,control:/control-3,control:/control-4\n", + stdout: + "control:/control-0,control:/control-1,control:/control-2,control:/control-3,control:/control-4\n", stderr: "", }); }); + test("streams a guest response after the handler opens an outbound websocket", async () => { + const upstream = createServer(); + const upstreamWebSocket = new WebSocketServer({ server: upstream }); + upstreamWebSocket.on("connection", (socket) => { + socket.once("message", () => { + socket.send(Uint8Array.from([1, 2, 3]), { binary: true }); + }); + }); + await new Promise((resolve) => + upstream.listen(0, "127.0.0.1", resolve), + ); + const upstreamAddress = upstream.address(); + if (!upstreamAddress || typeof upstreamAddress === "string") { + throw new Error("missing upstream TCP address"); + } + + try { + vm = await AgentOs.create({ + sidecar: { kind: "shared", pool: "reentrant-websocket-test" }, + loopbackExemptPorts: [upstreamAddress.port], + permissions: { + fs: "allow", + network: "allow", + childProcess: "allow", + }, + }); + + let resolveReady = () => {}; + let stderr = ""; + const ready = new Promise((resolve) => { + resolveReady = resolve; + }); + const script = [ + 'const http = require("node:http");', + "const server = http.createServer(async (_request, response) => {", + ` const socket = new WebSocket("ws://127.0.0.1:${upstreamAddress.port}/", ["rivet", "rivet_token.token"]);`, + ' socket.binaryType = "arraybuffer";', + " const binaryLength = await new Promise((resolve, reject) => {", + " socket.onopen = () => queueMicrotask(() => socket.send(new Uint8Array([4, 5, 6])));", + " socket.onmessage = (event) => resolve(event.data.byteLength);", + " socket.onerror = reject;", + " });", + " socket.close();", + ' response.writeHead(200, { "Content-Type": "text/event-stream" });', + " response.flushHeaders();", + ' response.write(`data: websocket-${binaryLength}\\n\\n`);', + "});", + 'server.listen(3000, "0.0.0.0", () => console.log("READY"));', + ].join("\n"); + const child = await vm.spawn("node", ["-e", script], { + onStdout: (chunk) => { + if (textDecoder.decode(chunk).includes("READY")) resolveReady(); + }, + onStderr: (chunk) => { + stderr += textDecoder.decode(chunk); + }, + }); + await Promise.race([ + ready, + vm.process.wait(child.pid).then((result) => { + throw new Error( + `guest server exited with ${result.exitCode}: ${stderr}`, + ); + }), + new Promise((_, reject) => + setTimeout( + () => reject(new Error("guest server readiness timed out")), + 5_000, + ), + ), + ]); + + const head = await Promise.race([ + vm.fetchStreamStart(3000, new Request("http://guest/events")), + new Promise((_, reject) => + setTimeout( + () => reject(new Error("stream response head timed out")), + 5_000, + ), + ), + ]); + const chunk = await vm.fetchStreamRead(head.streamId); + + expect(head.status, stderr).toBe(200); + expect(textDecoder.decode(chunk.body)).toBe("data: websocket-3\n\n"); + await vm.fetchStreamCancel(head.streamId); + await vm.process.kill(child.pid, "SIGKILL"); + } finally { + upstreamWebSocket.clients.forEach((socket) => socket.terminate()); + await new Promise((resolve, reject) => { + upstreamWebSocket.close((error) => (error ? reject(error) : resolve())); + }); + await new Promise((resolve, reject) => { + upstream.close((error) => (error ? reject(error) : resolve())); + }); + } + }); + test("supports writeFileSync on an fd from a nested guest child", async () => { vm = await AgentOs.create({ permissions: { @@ -265,7 +365,7 @@ describe("guest http.request transport", () => { ' await fs.mkdir("/workspace/nested-fs", { recursive: true });', ' const fd = fsSync.openSync("/workspace/nested-fs/session.jsonl", "wx");', ' try { fsSync.writeFileSync(fd, "first\\n"); fsSync.writeFileSync(fd, "second\\n"); }', - ' finally { fsSync.closeSync(fd); }', + " finally { fsSync.closeSync(fd); }", ' await fs.writeFile("/workspace/nested-fs/result.txt", "nested-ok", "utf8");', ' console.log(await fs.readFile("/workspace/nested-fs/result.txt", "utf8"));', ' console.log((await fs.readFile("/workspace/nested-fs/session.jsonl", "utf8")).trim());', @@ -293,7 +393,9 @@ describe("guest http.request transport", () => { response.write("data: first\n\n"); setTimeout(() => response.end("data: done\n\n"), 100); }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + await new Promise((resolve) => + server.listen(0, "127.0.0.1", resolve), + ); const address = server.address(); if (!address || typeof address === "string") { throw new Error("missing host TCP address"); @@ -363,7 +465,9 @@ describe("guest http.request transport", () => { } response.end(); }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + await new Promise((resolve) => + server.listen(0, "127.0.0.1", resolve), + ); const address = server.address(); if (!address || typeof address === "string") { throw new Error("missing host TCP address");