diff --git a/crates/native-sidecar/src/execution/mod.rs b/crates/native-sidecar/src/execution/mod.rs index d9234ebb4d..05e7f25352 100644 --- a/crates/native-sidecar/src/execution/mod.rs +++ b/crates/native-sidecar/src/execution/mod.rs @@ -292,6 +292,7 @@ use url::Url; const DEFAULT_KERNEL_STDIN_READ_MAX_BYTES: usize = 64 * 1024; const DEFAULT_KERNEL_STDIN_READ_TIMEOUT_MS: u64 = 100; +const JAVASCRIPT_NET_CLOSE_SENTINEL: &str = "__agentos_net_close__"; const JAVASCRIPT_NET_TIMEOUT_SENTINEL: &str = "__agentos_net_timeout__"; const PYTHON_PYODIDE_GUEST_ROOT: &str = "/__agentos_pyodide"; const PYTHON_PYODIDE_CACHE_GUEST_ROOT: &str = "/__agentos_pyodide_cache"; diff --git a/crates/native-sidecar/src/execution/network/tcp.rs b/crates/native-sidecar/src/execution/network/tcp.rs index b321a5b16b..e503eac025 100644 --- a/crates/native-sidecar/src/execution/network/tcp.rs +++ b/crates/native-sidecar/src/execution/network/tcp.rs @@ -609,7 +609,14 @@ impl ActiveTcpSocket { .fetch_add(1, Ordering::Relaxed); } self.saw_remote_end.store(true, Ordering::SeqCst); - Ok(Some(JavascriptTcpSocketEvent::End)) + if kernel + .socket_get(socket_id) + .is_some_and(|record| record.peer_socket_id().is_none()) + { + Ok(Some(JavascriptTcpSocketEvent::Close { had_error: false })) + } else { + Ok(Some(JavascriptTcpSocketEvent::End)) + } } Err(error) if error.code() == "EAGAIN" => { if trace_enabled { @@ -634,7 +641,16 @@ impl ActiveTcpSocket { } if revents.intersects(POLLHUP) { self.saw_remote_end.store(true, Ordering::SeqCst); - return Ok(Some(JavascriptTcpSocketEvent::End)); + return Ok(Some( + if kernel + .socket_get(socket_id) + .is_some_and(|record| record.peer_socket_id().is_none()) + { + JavascriptTcpSocketEvent::Close { had_error: false } + } else { + JavascriptTcpSocketEvent::End + }, + )); } if revents.intersects(POLLERR) { return Ok(Some(JavascriptTcpSocketEvent::Error { @@ -2996,8 +3012,9 @@ pub(in crate::execution) fn javascript_net_read_value( Some(JavascriptTcpSocketEvent::Data { bytes, .. }) => Ok(Value::String( base64::engine::general_purpose::STANDARD.encode(bytes), )), - Some(JavascriptTcpSocketEvent::End | JavascriptTcpSocketEvent::Close { .. }) => { - Ok(Value::Null) + Some(JavascriptTcpSocketEvent::End) => Ok(Value::Null), + Some(JavascriptTcpSocketEvent::Close { .. }) => { + Ok(Value::String(String::from(JAVASCRIPT_NET_CLOSE_SENTINEL))) } Some(JavascriptTcpSocketEvent::Error { code, message }) => { let detail = code.unwrap_or_else(|| String::from("socket read")); diff --git a/crates/vfs/src/engine/types.rs b/crates/vfs/src/engine/types.rs index 6292c36390..fe50ca349a 100644 --- a/crates/vfs/src/engine/types.rs +++ b/crates/vfs/src/engine/types.rs @@ -36,7 +36,7 @@ pub(crate) fn decode_unwritten_extents( )); } let mut ranges = Vec::with_capacity(encoded.len() / 16); - for chunk in encoded.chunks_exact(16) { + for chunk in encoded.as_chunks::<16>().0 { let start = u64::from_le_bytes(chunk[..8].try_into().expect("eight-byte extent start")); let end = u64::from_le_bytes(chunk[8..].try_into().expect("eight-byte extent end")); if start >= end diff --git a/examples/quickstart/sandbox/index.ts b/examples/quickstart/sandbox/index.ts index cb93c6d767..6bcb32a46e 100644 --- a/examples/quickstart/sandbox/index.ts +++ b/examples/quickstart/sandbox/index.ts @@ -39,10 +39,10 @@ try { const runCommandResult = await vm.process.exec( "agentos-sandbox run-command --command echo --args 'hello from Docker sandbox'", ); - console.log("Sandbox command:", runCommandResult.stdout.trim()); + console.log("Sandbox command:", (runCommandResult.stdout ?? "").trim()); const processList = await vm.process.exec("agentos-sandbox list-processes"); - console.log("Sandbox processes:", processList.stdout.trim()); + console.log("Sandbox processes:", (processList.stdout ?? "").trim()); const ANTHROPIC_API_KEY = process.env.ANTHROPIC_API_KEY; if (ANTHROPIC_API_KEY) { diff --git a/examples/workflows/server.ts b/examples/workflows/server.ts index 5fd75deae0..fb15231579 100644 --- a/examples/workflows/server.ts +++ b/examples/workflows/server.ts @@ -67,7 +67,7 @@ async function runTests( ): Promise { const agent = step.client().vm.getOrCreate("bug-fixer"); const tests = await agent.process.exec("cd /home/agentos/repo && npm test"); - return tests.exitCode; + return tests.exitCode ?? 1; } // docs:end basic diff --git a/packages/agentos-apps/tests/apps.test.ts b/packages/agentos-apps/tests/apps.test.ts index 7fbf1c8446..b62c956885 100644 --- a/packages/agentos-apps/tests/apps.test.ts +++ b/packages/agentos-apps/tests/apps.test.ts @@ -1202,7 +1202,7 @@ describe("replica artifact lifecycle", () => { vi.stubEnv("RIVET_TOKEN", "host-management-token"); const definitions = createAppsActors(); const replicaDefinition = definitions.agentOSAppsReplica; - const actions = replicaDefinition.config.actions as Record< + const actions = replicaDefinition.config.actions as unknown as Record< string, (...args: any[]) => any >; diff --git a/packages/build-tools/bridge-src/builtins/net.ts b/packages/build-tools/bridge-src/builtins/net.ts index 6fb740ce97..234a7529b6 100644 --- a/packages/build-tools/bridge-src/builtins/net.ts +++ b/packages/build-tools/bridge-src/builtins/net.ts @@ -1282,10 +1282,11 @@ function createAcceptedClientHandle(socketId, info) { }; } -// Must match JAVASCRIPT_NET_TIMEOUT_SENTINEL in crates/native-sidecar/src/execution/mod.rs. +// Must match the sentinels in crates/native-sidecar/src/execution/mod.rs. // A mismatched sentinel is NOT a soft failure: every no-data poll response then // falls through to base64 decoding and injects the decoded sentinel bytes into // the socket stream as phantom data. +var NET_BRIDGE_CLOSE_SENTINEL = "__agentos_net_close__"; var NET_BRIDGE_TIMEOUT_SENTINEL = "__agentos_net_timeout__"; function isNetBridgeTraceEnabled() { @@ -2897,6 +2898,13 @@ var NetSocket = class _NetSocket extends CanonicalDuplex { countNetBridgeMetric("readWaitsForWake"); return; } + if (chunk === NET_BRIDGE_CLOSE_SENTINEL) { + countNetBridgeMetric("readCloseEvents"); + this._pendingBridgeWake = false; + this._pendingBridgeWakeRetries = 0; + this.destroy(); + return; + } if (chunk === null) { if (firstPumpRun && !firstPumpResultRecorded) { firstPumpResultRecorded = true; @@ -3830,6 +3838,7 @@ export { isValidIPv6Zone, isValidTcpPort, maxNetBridgeMetric, + NET_BRIDGE_CLOSE_SENTINEL, NET_BRIDGE_MAX_RAW_WRITE_BYTES, NET_BRIDGE_TIMEOUT_SENTINEL, NET_SERVER_HANDLE_PREFIX, diff --git a/packages/core/src/agent-os.ts b/packages/core/src/agent-os.ts index 1ff9dec103..2c1a340d3e 100644 --- a/packages/core/src/agent-os.ts +++ b/packages/core/src/agent-os.ts @@ -325,6 +325,14 @@ export interface HttpResponse { body: Uint8Array; } +function headersToRecord(headers: Headers): Record { + const result: Record = {}; + headers.forEach((value, name) => { + result[name] = value; + }); + return result; +} + export interface ProcessOutput { pid: number; stream: "stdout" | "stderr"; @@ -5088,9 +5096,7 @@ export class AgentOs { port, method: request.method, path: `${url.pathname}${url.search}`, - headersJson: JSON.stringify( - Object.fromEntries(request.headers.entries()), - ), + headersJson: JSON.stringify(headersToRecord(request.headers)), ...(request.method !== "GET" && request.method !== "HEAD" ? { bodyBase64: Buffer.from(await request.arrayBuffer()).toString( @@ -5131,9 +5137,7 @@ export class AgentOs { port, method: request.method, path: `${url.pathname}${url.search}`, - headersJson: JSON.stringify( - Object.fromEntries(request.headers.entries()), - ), + headersJson: JSON.stringify(headersToRecord(request.headers)), ...(request.method !== "GET" && request.method !== "HEAD" ? { bodyBase64: Buffer.from(await request.arrayBuffer()).toString( diff --git a/packages/core/src/cron/timer-driver.ts b/packages/core/src/cron/timer-driver.ts index 62f2acaefe..a109d41dab 100644 --- a/packages/core/src/cron/timer-driver.ts +++ b/packages/core/src/cron/timer-driver.ts @@ -1,3 +1,5 @@ +/// + import type { LongTimeout } from "long-timeout"; import { clearTimeout as clearLongTimeout, diff --git a/packages/core/src/long-timeout.d.ts b/packages/core/src/long-timeout.d.ts new file mode 100644 index 0000000000..daa38748b6 --- /dev/null +++ b/packages/core/src/long-timeout.d.ts @@ -0,0 +1,11 @@ +declare module "long-timeout" { + export interface LongTimeout {} + + export function setTimeout( + callback: (...args: unknown[]) => void, + delay: number, + ...args: unknown[] + ): LongTimeout; + + export function clearTimeout(timeout: LongTimeout): void; +} diff --git a/packages/core/src/sandbox.ts b/packages/core/src/sandbox.ts index d85fd7dbbb..81446d745a 100644 --- a/packages/core/src/sandbox.ts +++ b/packages/core/src/sandbox.ts @@ -143,7 +143,11 @@ function normalizeHeaders( } if (headers instanceof Headers) { - return Object.fromEntries(headers.entries()); + const normalized: Record = {}; + headers.forEach((value, name) => { + normalized[name] = value; + }); + return normalized; } if (Array.isArray(headers)) { diff --git a/packages/core/tests/network-http-request.test.ts b/packages/core/tests/network-http-request.test.ts index 8d96d56d1f..49471d3e54 100644 --- a/packages/core/tests/network-http-request.test.ts +++ b/packages/core/tests/network-http-request.test.ts @@ -12,7 +12,7 @@ async function runSpawnedProcess( ): Promise<{ exitCode: number; stdout: string; stderr: string }> { const stdoutChunks: string[] = []; const stderrChunks: string[] = []; - const { pid } = vm.spawn(command, args, { + const { pid } = await vm.spawn(command, args, { onStdout: (chunk) => { stdoutChunks.push(textDecoder.decode(chunk)); }, @@ -22,8 +22,9 @@ async function runSpawnedProcess( }, }); + const exit = await vm.process.wait(pid); return { - exitCode: await vm.waitProcess(pid), + exitCode: exit.exitCode ?? -1, stdout: stdoutChunks.join(""), stderr: stderrChunks.join(""), }; @@ -248,7 +249,7 @@ describe("guest http.request transport", () => { }); }); - test("streams a guest response after the handler opens an outbound websocket", async () => { + test("reclaims cancelled guest response streams while an outbound websocket stays open", async () => { const upstream = createServer(); const upstreamWebSocket = new WebSocketServer({ server: upstream }); upstreamWebSocket.on("connection", (socket) => { @@ -282,7 +283,7 @@ describe("guest http.request transport", () => { }); const script = [ 'const http = require("node:http");', - "const server = http.createServer(async (_request, response) => {", + "void (async () => {", ` 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) => {", @@ -290,12 +291,13 @@ describe("guest http.request transport", () => { " 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"));', + " const server = http.createServer((_request, response) => {", + ' 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"));', + "})().catch((error) => { console.error(error); process.exitCode = 1; });", ].join("\n"); const child = await vm.spawn("node", ["-e", script], { onStdout: (chunk) => { @@ -320,7 +322,11 @@ describe("guest http.request transport", () => { ), ]); - for (let requestIndex = 0; requestIndex < 10; requestIndex++) { + const requestCount = Number.parseInt( + process.env.AGENTOS_VM_FETCH_STREAM_REQUESTS ?? "300", + 10, + ); + for (let requestIndex = 0; requestIndex < requestCount; requestIndex++) { const head = await Promise.race([ vm.fetchStreamStart( 3000, @@ -328,7 +334,12 @@ describe("guest http.request transport", () => { ), new Promise((_, reject) => setTimeout( - () => reject(new Error("stream response head timed out")), + () => + reject( + new Error( + `stream response head timed out at request ${requestIndex + 1}/${requestCount}`, + ), + ), 5_000, ), ), diff --git a/scripts/check-layout.mjs b/scripts/check-layout.mjs index ed56dcb00b..7290151338 100644 --- a/scripts/check-layout.mjs +++ b/scripts/check-layout.mjs @@ -49,6 +49,7 @@ const allowedTestHomes = [ /^software\/[^/]+\/test\/.+\.test\.ts$/, /^toolchain\/conformance\/.+\.test\.ts$/, /^packages\/[^/]+\/tests\/.+\.test\.ts$/, + /^benchmarks\/[^/]+\/src\/.+\.test\.ts$/, /^experiments\/[^/]+\/.+\.test\.ts$/, /^scripts\/.+\.test\.ts$/, ]; diff --git a/scripts/check-layout.test.mjs b/scripts/check-layout.test.mjs index 1fc8bb8a33..48c6f799c6 100644 --- a/scripts/check-layout.test.mjs +++ b/scripts/check-layout.test.mjs @@ -8,7 +8,7 @@ import test from "node:test"; const script = join(dirname(fileURLToPath(import.meta.url)), "check-layout.mjs"); -test("allows experiment tests and ignores nested Claude worktrees", () => { +test("allows benchmark and experiment tests and ignores nested Claude worktrees", () => { const root = mkdtempSync(join(tmpdir(), "agentos-layout-")); try { const nestedTest = join( @@ -20,6 +20,9 @@ test("allows experiment tests and ignores nested Claude worktrees", () => { const experimentTest = join(root, "experiments/gigacode/gate.test.ts"); mkdirSync(dirname(experimentTest), { recursive: true }); writeFileSync(experimentTest, "export {};\n"); + const benchmarkTest = join(root, "benchmarks/apps/src/load.test.ts"); + mkdirSync(dirname(benchmarkTest), { recursive: true }); + writeFileSync(benchmarkTest, "export {};\n"); const bin = join(root, "bin"); mkdirSync(bin);