Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 22 additions & 1 deletion crates/execution/src/javascript.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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();
Expand Down
132 changes: 118 additions & 14 deletions packages/core/tests/network-http-request.test.ts
Original file line number Diff line number Diff line change
@@ -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();
Expand Down Expand Up @@ -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,",
Expand All @@ -71,20 +72,20 @@ 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) => {',
" responseBody += chunk;",
" });",
' 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");
Expand Down Expand Up @@ -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));",
"})();",
Expand Down Expand Up @@ -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 = [];",
Expand All @@ -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<void>((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<void>((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<never>((_, 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<never>((_, 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<void>((resolve, reject) => {
upstreamWebSocket.close((error) => (error ? reject(error) : resolve()));
});
await new Promise<void>((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: {
Expand All @@ -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());',
Expand Down Expand Up @@ -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<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
await new Promise<void>((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");
Expand Down Expand Up @@ -363,7 +465,9 @@ describe("guest http.request transport", () => {
}
response.end();
});
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
await new Promise<void>((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");
Expand Down
Loading