diff --git a/changelog.d/10888-create-connection-upgrade.md b/changelog.d/10888-create-connection-upgrade.md new file mode 100644 index 0000000000..6051c7f864 --- /dev/null +++ b/changelog.d/10888-create-connection-upgrade.md @@ -0,0 +1,7 @@ +### Fixed + +- Fixed HTTP upgrade requests that use a custom `createConnection` socket. + Perry now preserves `Connection: Upgrade`, emits the request's `upgrade` + event for a `101` response, and hands the same live socket back to callers. + This allows source-compiled `ws` clients to complete their handshake instead + of receiving `Unexpected server response: 426` (#10888). diff --git a/crates/perry-ext-http/src/client_connect_override.rs b/crates/perry-ext-http/src/client_connect_override.rs index bbe5bbf7e1..21c5e517aa 100644 --- a/crates/perry-ext-http/src/client_connect_override.rs +++ b/crates/perry-ext-http/src/client_connect_override.rs @@ -49,10 +49,10 @@ pub(crate) fn request_create_connection_socket( } /// Serialize an HTTP/1.1 request (request line + headers + body) into the -/// bytes to write onto a socket. Forces `Connection: close` (the raw socket -/// path reads until EOF), drops any caller-supplied `Connection`/`Host` -/// header (we set `Host` from the URL), and adds `Content-Length` when a -/// body is present and the caller didn't. +/// bytes to write onto a socket. Ordinary responses force `Connection: close` +/// because this path reads until EOF. Upgrade requests preserve the caller's +/// `Connection: Upgrade` header so a `101` can hand the live socket back to +/// JavaScript. `Host` is always derived from the URL. fn serialize_http_request( method: &str, path: &str, @@ -60,13 +60,16 @@ fn serialize_http_request( headers: &HashMap, body: &[u8], ) -> Vec { + let wants_upgrade = crate::client_upgrade::wants_upgrade(headers); let mut req = format!("{} {} HTTP/1.1\r\nHost: {}\r\n", method, path, host_header); let mut has_content_length = false; for (k, v) in headers { if k.eq_ignore_ascii_case("content-length") { has_content_length = true; } - if k.eq_ignore_ascii_case("connection") || k.eq_ignore_ascii_case("host") { + if k.eq_ignore_ascii_case("host") + || (k.eq_ignore_ascii_case("connection") && !wants_upgrade) + { continue; } req.push_str(k); @@ -74,7 +77,9 @@ fn serialize_http_request( req.push_str(v); req.push_str("\r\n"); } - req.push_str("Connection: close\r\n"); + if !wants_upgrade { + req.push_str("Connection: close\r\n"); + } if !body.is_empty() && !has_content_length { req.push_str(&format!("Content-Length: {}\r\n", body.len())); } @@ -86,11 +91,12 @@ fn serialize_http_request( /// #2154 — run an HTTP exchange over a socket that a `createConnection` /// override (Agent-level or, since #10469, request-level) produced -/// (`socket_id`), instead of through reqwest. Writes the serialized -/// request, reads the response until the peer closes (we force -/// `Connection: close`), parses it with [`parse_http_response`], and pushes -/// the same `Response` / `Error` event the reqwest path produces — so the -/// IncomingMessage surface is identical. +/// (`socket_id`), instead of through reqwest. Ordinary responses force +/// `Connection: close` and read to EOF. A `101` response to an upgrade request +/// detaches the still-live socket from the raw reader and pushes `Upgrade` with +/// any bytes following the header block. Other responses are parsed with +/// [`parse_http_response`] and produce the same `Response` / `Error` events as +/// the reqwest path. /// /// The socket I/O goes through perry-ffi's raw-net vtable (published by /// perry-ext-net), so this crate needs no link edge to perry-ext-net. If no @@ -129,6 +135,7 @@ pub(crate) fn dispatch_request_over_socket( path.push_str(q); } let req_bytes = serialize_http_request(&method, &path, &host_header, &headers, &body); + let wants_upgrade = crate::client_upgrade::wants_upgrade(&headers); let deadline = std::time::Duration::from_millis(timeout_ms.unwrap_or(30_000)); spawn_blocking(move || { @@ -176,6 +183,39 @@ pub(crate) fn dispatch_request_over_socket( let n = (vtable.poll_read)(socket_id, chunk.as_mut_ptr(), chunk.len()); if n > 0 { raw.extend_from_slice(&chunk[..n as usize]); + if wants_upgrade { + if let Some(header_end) = raw.windows(4).position(|w| w == b"\r\n\r\n") { + let header_end = header_end + 4; + let status = std::str::from_utf8(&raw[..header_end]) + .ok() + .and_then(|head| head.lines().next()) + .and_then(|line| line.split_whitespace().nth(1)) + .and_then(|code| code.parse::().ok()); + if status == Some(101) { + match parse_http_response(&raw[..header_end]) { + Ok(parsed) => { + (vtable.detach)(socket_id); + push_event(PendingHttpEvent::Upgrade { + request_handle, + status: parsed.status, + status_message: parsed.status_message, + headers: parsed.headers, + socket_handle: socket_id, + head: raw[header_end..].to_vec(), + }); + } + Err(error_message) => { + (vtable.close)(socket_id); + push_event(PendingHttpEvent::Error { + request_handle, + error_message, + }); + } + } + return; + } + } + } } else if n == 0 { break; // clean EOF — peer closed after the response } else { @@ -209,3 +249,27 @@ pub(crate) fn dispatch_request_over_socket( std::mem::forget(jh); }); } + +#[cfg(test)] +mod tests { + use super::serialize_http_request; + use std::collections::HashMap; + + #[test] + fn websocket_upgrade_keeps_connection_header() { + let headers = HashMap::from([ + ("Connection".to_string(), "Upgrade".to_string()), + ("Upgrade".to_string(), "websocket".to_string()), + ]); + let request = String::from_utf8(serialize_http_request( + "GET", + "/socket", + "localhost:1234", + &headers, + &[], + )) + .unwrap(); + assert!(request.contains("Connection: Upgrade\r\n"), "{request}"); + assert!(!request.contains("Connection: close\r\n"), "{request}"); + } +} diff --git a/crates/perry-ext-net/src/raw_bridge.rs b/crates/perry-ext-net/src/raw_bridge.rs index 9a51c88d76..6a111d489b 100644 --- a/crates/perry-ext-net/src/raw_bridge.rs +++ b/crates/perry-ext-net/src/raw_bridge.rs @@ -139,6 +139,18 @@ extern "C" fn perry_net_raw_poll_read(socket_id: i64, out: *mut u8, max: usize) n as isize } +/// Return an attached socket to ordinary JS event delivery while keeping the +/// underlying connection alive. HTTP uses this after parsing a `101` response +/// so the request's `upgrade` listener receives the same live `net.Socket` +/// supplied by `createConnection`. +extern "C" fn perry_net_raw_detach(socket_id: i64) { + if let Ok(mut sockets) = statics::sockets().lock() { + if let Some(socket) = sockets.get_mut(&socket_id) { + socket.raw = None; + } + } +} + /// Tear down socket `socket_id` (equivalent to `socket.destroy()`) and /// unregister it. `perry-ext-http` calls this once the HTTP exchange is done /// draining. Because raw mode suppresses the JS `Close` event (the path that @@ -172,6 +184,7 @@ pub(crate) fn register() { attach: perry_net_raw_attach, write: perry_net_raw_write, poll_read: perry_net_raw_poll_read, + detach: perry_net_raw_detach, close: perry_net_raw_close, }); } diff --git a/crates/perry-ffi/src/raw_net.rs b/crates/perry-ffi/src/raw_net.rs index bf4364349b..3d914b6396 100644 --- a/crates/perry-ffi/src/raw_net.rs +++ b/crates/perry-ffi/src/raw_net.rs @@ -45,6 +45,10 @@ pub struct RawNetVtable { /// is drained and the peer closed, or `-1` when no bytes are /// currently available but the socket is still open ("would block"). pub poll_read: extern "C" fn(socket_id: i64, out: *mut u8, max: usize) -> isize, + /// Return a raw-consumer socket to normal JS event delivery without + /// closing it. Used when an HTTP 101 hands the caller-supplied socket to + /// the request's `upgrade` listener. + pub detach: extern "C" fn(socket_id: i64), /// Tear the socket down (sends the equivalent of `socket.destroy()`). pub close: extern "C" fn(socket_id: i64), } diff --git a/crates/perry-runtime/src/object/field_get_set/get_field_by_name.rs b/crates/perry-runtime/src/object/field_get_set/get_field_by_name.rs index 7478b27e48..f8be43bad5 100644 --- a/crates/perry-runtime/src/object/field_get_set/get_field_by_name.rs +++ b/crates/perry-runtime/src/object/field_get_set/get_field_by_name.rs @@ -1265,6 +1265,37 @@ pub(crate) fn get_field_by_name_past_inherited_cache( } return JSValue::from_bits(value.to_bits()); } + // A capture-carrying class declaration is materialized as a + // heap class object (`ClassExprFresh`). References that were + // lowered before the declaration's runtime binding existed + // (the class constructor itself and earlier helper closures) + // still carry the template `ClassRef`. Runtime additions such + // as `Object.defineProperty(C, "OPEN", { value: 1 })` live on + // the materialized object, so consulting only the template + // tables makes `C.OPEN` undefined in those bodies even though + // the same expression at the declaration site reads `1`. + // + // `CLASS_OBJECT_VALUES` is already the runtime identity used + // by `instance.constructor`. Read that same current + // evaluation first, preserving its own-property and pinned + // static-parent semantics; a miss continues through the + // ordinary ClassRef registry path below. + if !is_prototype_ref { + if let Some(class_object) = + super::super::class_registry::class_object_value_for_cid(class_id) + { + let class_object = JSValue::from_bits(class_object.to_bits()); + if class_object.is_pointer() { + let class_object = class_object.as_pointer::(); + if !class_object.is_null() && class_object as usize != obj as usize { + let value = js_object_get_field_by_name(class_object, key); + if !value.is_undefined() { + return value; + } + } + } + } + } // Instance (prototype) methods must only resolve when reading // off the prototype ref (`C.prototype.m`), NOT off the class ref // itself (`C.m`). In JS a class object does not expose its diff --git a/crates/perry/tests/issue_10888_create_connection_upgrade.rs b/crates/perry/tests/issue_10888_create_connection_upgrade.rs new file mode 100644 index 0000000000..03b6de67fc --- /dev/null +++ b/crates/perry/tests/issue_10888_create_connection_upgrade.rs @@ -0,0 +1,166 @@ +//! Regression for #10888: an upgrade request that supplies `createConnection` +//! must preserve `Connection: Upgrade` and return that same live socket from +//! the request's `upgrade` event. + +use std::path::PathBuf; +use std::process::{Command, Stdio}; +use std::time::{Duration, Instant}; + +const SOURCE: &str = r#" +import { createServer, request } from "node:http"; +import { connect } from "node:net"; + +const watchdog = setTimeout(() => { + console.log("timeout"); + process.exit(1); +}, 5000); + +const server = createServer((_req: any, res: any) => { + console.log("ordinary request"); + res.statusCode = 426; + res.end(); +}); +server.on("upgrade", (_req: any, socket: any) => { + console.log("server upgrade"); + socket.write( + "HTTP/1.1 101 Switching Protocols\r\n" + + "Connection: Upgrade\r\n" + + "Upgrade: probe\r\n\r\n", + ); + setTimeout(() => socket.write("hello"), 20); +}); +server.listen(0, "127.0.0.1", () => { + const port = server.address().port; + const req = request({ + host: "127.0.0.1", + port, + headers: { Connection: "Upgrade", Upgrade: "probe" }, + createConnection: (options: any) => connect(options.port, options.host), + }); + req.on("upgrade", (res: any, socket: any, head: any) => { + console.log("client upgrade", res.statusCode, res.headers.upgrade, head.length); + socket.on("data", (data: any) => { + console.log("client data", data.toString()); + clearTimeout(watchdog); + socket.destroy(); + server.close(() => process.exit(0)); + }); + }); + req.on("response", (res: any) => console.log("response", res.statusCode)); + req.on("error", (error: any) => console.log("request error", error.message)); + req.end(); +}); +"#; + +const CAPTURED_CLASS_STATIC_SOURCE: &str = r#" +function make() { + const marker = "captured"; + function state() { + return SocketState.CONNECTING; + } + class SocketState { + constructor() { + this.marker = marker; + this.state = SocketState.CONNECTING; + } + } + Object.defineProperty(SocketState, "CONNECTING", { value: 0 }); + return { SocketState, state }; +} + +const { SocketState, state } = make(); +const socket = new SocketState(); +console.log(socket.marker, socket.state, state(), SocketState.CONNECTING); +"#; + +#[test] +fn custom_connection_preserves_upgrade_and_hands_back_the_live_socket() { + let dir = tempfile::tempdir().expect("tempdir"); + let entry = dir.path().join("main.ts"); + let binary = dir.path().join("main"); + std::fs::write(&entry, SOURCE).unwrap(); + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../.."); + let compile = Command::new(env!("CARGO_BIN_EXE_perry")) + .args([ + "compile", + entry.to_str().unwrap(), + "-o", + binary.to_str().unwrap(), + "--no-cache", + ]) + .env("PERRY_WORKSPACE_ROOT", root) + .output() + .expect("compile"); + assert!( + compile.status.success(), + "compile failed: {}", + String::from_utf8_lossy(&compile.stderr) + ); + + let mut child = Command::new(binary) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .unwrap(); + let deadline = Instant::now() + Duration::from_secs(10); + while child.try_wait().unwrap().is_none() { + if Instant::now() >= deadline { + child.kill().unwrap(); + let output = child.wait_with_output().unwrap(); + panic!( + "compiled fixture hung: {}\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + } + std::thread::sleep(Duration::from_millis(20)); + } + let output = child.wait_with_output().unwrap(); + assert!( + output.status.success(), + "compiled fixture failed: {}\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + assert_eq!( + String::from_utf8(output.stdout).unwrap(), + "server upgrade\nclient upgrade 101 probe 0\nclient data hello\n" + ); +} + +#[test] +fn captured_class_self_reads_runtime_static_properties() { + let dir = tempfile::tempdir().expect("tempdir"); + let entry = dir.path().join("main.ts"); + let binary = dir.path().join("main"); + std::fs::write(&entry, CAPTURED_CLASS_STATIC_SOURCE).unwrap(); + let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../.."); + let compile = Command::new(env!("CARGO_BIN_EXE_perry")) + .args([ + "compile", + entry.to_str().unwrap(), + "-o", + binary.to_str().unwrap(), + "--no-cache", + ]) + .env("PERRY_WORKSPACE_ROOT", root) + .output() + .expect("compile"); + assert!( + compile.status.success(), + "compile failed: {}", + String::from_utf8_lossy(&compile.stderr) + ); + + let output = Command::new(binary).output().expect("run compiled fixture"); + assert!( + output.status.success(), + "compiled fixture failed: {}\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + assert_eq!( + String::from_utf8(output.stdout).unwrap(), + "captured 0 0 0\n" + ); +}