From 59387f1e903948bcc111e0328ce0008aaa3f9118 Mon Sep 17 00:00:00 2001 From: Pierre Le Fevre Date: Sun, 17 May 2026 12:44:33 +0200 Subject: [PATCH] Implement streaming fetch request bodies --- crates/e2e/tests/phase20_network.rs | 104 ++++++++ crates/js/src/fetch.rs | 353 +++++++++++++++++++++++++++- crates/net/src/client.rs | 62 ++++- crates/net/src/http.rs | 66 ++++++ 4 files changed, 578 insertions(+), 7 deletions(-) diff --git a/crates/e2e/tests/phase20_network.rs b/crates/e2e/tests/phase20_network.rs index 638304a..26c71c6 100644 --- a/crates/e2e/tests/phase20_network.rs +++ b/crates/e2e/tests/phase20_network.rs @@ -50,6 +50,36 @@ fn read_http_request(stream: &mut TcpStream) -> String { String::from_utf8_lossy(&request).into_owned() } +fn read_chunked_request_body(stream: &mut TcpStream) -> Vec { + let mut body = Vec::new(); + loop { + let mut line = Vec::new(); + let mut byte = [0; 1]; + while !line.ends_with(b"\r\n") { + stream.read_exact(&mut byte).expect("read chunk size byte"); + line.push(byte[0]); + } + let line_str = String::from_utf8_lossy(&line); + let size_hex = line_str.trim().split(';').next().unwrap_or(""); + let size = usize::from_str_radix(size_hex, 16).expect("parse chunk size"); + if size == 0 { + let mut trailer_end = [0; 2]; + stream + .read_exact(&mut trailer_end) + .expect("read chunk trailer"); + assert_eq!(&trailer_end, b"\r\n"); + break; + } + let mut chunk = vec![0; size]; + stream.read_exact(&mut chunk).expect("read chunk"); + let mut crlf = [0; 2]; + stream.read_exact(&mut crlf).expect("read chunk crlf"); + assert_eq!(&crlf, b"\r\n"); + body.extend_from_slice(&chunk); + } + body +} + fn http_ok(stream: &mut TcpStream, content_type: &str, body: &str) { write!( stream, @@ -254,6 +284,7 @@ fn fetch_large_body_through_transformstream() { var state = {{}}; state.result = ""; fetch("http://{addr}/large").then(function(resp) {{ + try {{ var networkReader = resp.body.getReader(); var ts = new TransformStream({{ transform: function(chunk, controller) {{ @@ -288,6 +319,11 @@ fn fetch_large_body_through_transformstream() { readTransformed(); return readNetwork(); + }} catch (err) {{ + state.result = "error:" + err; + }} + }}, function(err) {{ + state.result = "fetch-rejected:" + err; }}); "# ); @@ -303,6 +339,74 @@ fn fetch_large_body_through_transformstream() { we_js::fetch::reset_fetch_state(); } +#[test] +fn fetch_posts_readablestream_body_as_chunked_upload() { + we_js::fetch::reset_fetch_state(); + we_js::timers::reset_timers(); + + let listener = TcpListener::bind("127.0.0.1:0").expect("bind fetch server"); + let addr = listener.local_addr().expect("fetch server addr"); + let server = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept fetch"); + let request = read_http_request(&mut stream); + assert!(request.starts_with("POST /upload HTTP/1.1\r\n")); + assert!( + request.contains("Transfer-Encoding: chunked\r\n"), + "request should use chunked upload: {request}" + ); + assert!( + !request.contains("Content-Length:"), + "stream upload should not buffer a Content-Length: {request}" + ); + let body = read_chunked_request_body(&mut stream); + assert_eq!(String::from_utf8_lossy(&body), "hello stream"); + http_ok(&mut stream, "text/plain", "ok"); + }); + + let mut vm = Vm::new(); + we_js::fetch::init_fetch_api(&mut vm); + let source = format!( + r#" + var state = {{}}; + state.result = ""; + var upload = new ReadableStream({{ + start: function(controller) {{ + controller.enqueue("hello "); + controller.enqueue("stream"); + controller.close(); + }} + }}); + fetch("http://{addr}/upload", {{ method: "POST", body: upload }}).then( + function(resp) {{ + state.result = "status:" + resp.status + ":used:" + upload.__we_fetch_body_used; + }}, + function(err) {{ + state.result = "rejected:" + err; + }} + ); + "# + ); + execute(&mut vm, &source); + + let ok = pump_until(&mut vm, Duration::from_secs(2), |vm| { + read_object_property(vm, "state", "result") + .to_js_string(&vm.gc) + .starts_with("status:") + }); + server.join().expect("fetch server join"); + assert!( + ok, + "fetch did not complete: {}", + read_object_property(&vm, "state", "result").to_js_string(&vm.gc) + ); + assert_eq!( + read_object_property(&vm, "state", "result").to_js_string(&vm.gc), + "status:200:used:true" + ); + + we_js::fetch::reset_fetch_state(); +} + fn websocket_server_handshake(stream: &mut TcpStream) { let request = read_http_request(stream); let key = request diff --git a/crates/js/src/fetch.rs b/crates/js/src/fetch.rs index 7c40014..474b02a 100644 --- a/crates/js/src/fetch.rs +++ b/crates/js/src/fetch.rs @@ -5,7 +5,7 @@ use std::cell::RefCell; use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Condvar, Mutex}; use crate::builtins::{ create_promise_object_pub, make_native, reject_promise_internal, resolve_promise_internal, @@ -40,12 +40,14 @@ struct PendingFetch { thread_local! { static FETCH_STATE: RefCell> = const { RefCell::new(Vec::new()) }; static FETCH_BODY_STREAMS: RefCell> = const { RefCell::new(Vec::new()) }; + static FETCH_UPLOAD_STREAMS: RefCell> = const { RefCell::new(Vec::new()) }; /// The origin of the document that owns the current JS execution context. /// Used for Same-Origin Policy enforcement in fetch(). static DOCUMENT_ORIGIN: RefCell> = const { RefCell::new(None) }; } static NEXT_BODY_STREAM_ID: AtomicU64 = AtomicU64::new(1); +static NEXT_UPLOAD_STREAM_ID: AtomicU64 = AtomicU64::new(1); #[derive(Debug)] enum FetchBodyEvent { @@ -63,6 +65,32 @@ struct FetchBodyStream { closed: bool, } +#[derive(Debug)] +enum FetchUploadEvent { + Chunk(Vec), + Close, + Error(String), +} + +struct FetchUploadState { + queue: Mutex>, + ready: Condvar, +} + +impl FetchUploadState { + fn new() -> Self { + Self { + queue: Mutex::new(Vec::new()), + ready: Condvar::new(), + } + } +} + +struct FetchUploadStream { + id: u64, + state: Arc, +} + /// Set the document origin for SOP enforcement in the fetch API. /// /// Should be called before executing scripts so that `fetch()` can check @@ -86,6 +114,7 @@ pub(crate) fn get_document_origin() -> Option { pub fn reset_fetch_state() { FETCH_STATE.with(|s| s.borrow_mut().clear()); FETCH_BODY_STREAMS.with(|s| s.borrow_mut().clear()); + FETCH_UPLOAD_STREAMS.with(|s| s.borrow_mut().clear()); } /// Returns true if there are any in-flight fetches. @@ -163,6 +192,7 @@ pub fn fetch_native(args: &[Value], ctx: &mut NativeContext) -> Result = Vec::new(); let mut body: Option> = None; + let mut upload_stream: Option> = None; let mut cors_mode = "cors".to_string(); let mut credentials_mode = "same-origin".to_string(); @@ -178,6 +208,15 @@ pub fn fetch_native(args: &[Value], ctx: &mut NativeContext) -> Result Result, + upload_stream: Option>, document_origin: Option<&str>, cors_mode: &str, credentials_str: &str, @@ -518,12 +559,23 @@ fn do_fetch_streaming( let mut response_meta: Option = None; let stream_enabled = Arc::new(Mutex::new(false)); let stream_enabled_for_chunks = stream_enabled.clone(); + let mut next_upload_chunk = upload_stream.as_ref().map(|state| { + let state = state.clone(); + move || next_fetch_upload_event(&state) + }); + let request_body = if let Some(next_chunk) = next_upload_chunk.as_mut() { + Some(we_net::client::RequestBody::Chunked( + next_chunk as &mut dyn FnMut() -> we_net::client::Result>>, + )) + } else { + body.map(we_net::client::RequestBody::Bytes) + }; client - .request_streaming( + .request_streaming_with_body( http_method, &url, &req_headers, - body, + request_body, |response| { if is_cross_origin && cors_mode == "cors" { let doc_origin = document_origin.unwrap(); @@ -878,6 +930,138 @@ fn set_response_body_used(gc: &mut Gc, shapes: &mut ShapeTable, this } } +fn fetch_upload_start(_args: &[Value], _ctx: &mut NativeContext) -> Result { + let id = NEXT_UPLOAD_STREAM_ID.fetch_add(1, Ordering::Relaxed); + FETCH_UPLOAD_STREAMS.with(|streams| { + streams.borrow_mut().push(FetchUploadStream { + id, + state: Arc::new(FetchUploadState::new()), + }); + }); + Ok(Value::Number(id as f64)) +} + +fn fetch_upload_chunk(args: &[Value], ctx: &mut NativeContext) -> Result { + let id = args.first().and_then(value_to_stream_id).unwrap_or(0); + let bytes = args + .get(1) + .map(|value| bytes_from_upload_value(ctx, value)) + .unwrap_or_default(); + push_fetch_upload_event(id, FetchUploadEvent::Chunk(bytes)); + Ok(Value::Undefined) +} + +fn fetch_upload_close(args: &[Value], _ctx: &mut NativeContext) -> Result { + let id = args.first().and_then(value_to_stream_id).unwrap_or(0); + push_fetch_upload_event(id, FetchUploadEvent::Close); + Ok(Value::Undefined) +} + +fn fetch_upload_error(args: &[Value], ctx: &mut NativeContext) -> Result { + let id = args.first().and_then(value_to_stream_id).unwrap_or(0); + let message = args + .get(1) + .map(|value| value.to_js_string(ctx.gc)) + .unwrap_or_else(|| "stream upload failed".to_string()); + push_fetch_upload_event(id, FetchUploadEvent::Error(message)); + Ok(Value::Undefined) +} + +fn take_fetch_upload_stream(id: u64) -> Option> { + FETCH_UPLOAD_STREAMS.with(|streams| { + streams + .borrow() + .iter() + .find(|stream| stream.id == id) + .map(|stream| stream.state.clone()) + }) +} + +fn push_fetch_upload_event(id: u64, event: FetchUploadEvent) { + FETCH_UPLOAD_STREAMS.with(|streams| { + if let Some(stream) = streams.borrow().iter().find(|stream| stream.id == id) { + let mut queue = stream.state.queue.lock().unwrap(); + queue.push(event); + stream.state.ready.notify_one(); + } + }); +} + +fn next_fetch_upload_event(state: &FetchUploadState) -> we_net::client::Result>> { + let mut queue = state.queue.lock().unwrap(); + loop { + if !queue.is_empty() { + return match queue.remove(0) { + FetchUploadEvent::Chunk(bytes) => Ok(Some(bytes)), + FetchUploadEvent::Close => Ok(None), + FetchUploadEvent::Error(message) => { + Err(we_net::client::ClientError::InvalidUrl(message)) + } + }; + } + queue = state.ready.wait(queue).unwrap(); + } +} + +fn bytes_from_upload_value(ctx: &NativeContext, value: &Value) -> Vec { + match value { + Value::String(s) => s.as_bytes().to_vec(), + Value::Object(obj_ref) => bytes_from_upload_object(ctx, *obj_ref), + _ => value.to_js_string(ctx.gc).into_bytes(), + } +} + +fn bytes_from_upload_object(ctx: &NativeContext, obj_ref: GcRef) -> Vec { + if let Some(Value::Object(buffer_ref)) = object_prop(ctx, obj_ref, "buffer") { + let offset = object_prop(ctx, obj_ref, "byteOffset") + .map(|v| v.to_number() as usize) + .unwrap_or(0); + let len = object_prop(ctx, obj_ref, "byteLength") + .map(|v| v.to_number() as usize) + .unwrap_or(0); + return bytes_from_array_buffer(ctx, buffer_ref, offset, len); + } + let len = object_prop(ctx, obj_ref, "byteLength") + .or_else(|| object_prop(ctx, obj_ref, "length")) + .map(|v| v.to_number() as usize) + .unwrap_or(0); + bytes_from_array_buffer(ctx, obj_ref, 0, len) +} + +fn bytes_from_array_buffer( + ctx: &NativeContext, + obj_ref: GcRef, + offset: usize, + len: usize, +) -> Vec { + let mut out = Vec::with_capacity(len); + if let Some(Value::Object(bytes_ref)) = object_prop(ctx, obj_ref, "_bytes") { + for i in 0..len { + out.push( + object_prop(ctx, bytes_ref, &(offset + i).to_string()) + .map(|v| v.to_number() as u8) + .unwrap_or(0), + ); + } + return out; + } + for i in 0..len { + out.push( + object_prop(ctx, obj_ref, &(offset + i).to_string()) + .map(|v| v.to_number() as u8) + .unwrap_or(0), + ); + } + out +} + +fn object_prop(ctx: &NativeContext, obj_ref: GcRef, key: &str) -> Option { + match ctx.gc.get(obj_ref) { + Some(HeapObject::Object(data)) => data.get_property(key, ctx.shapes).map(|p| p.value), + _ => None, + } +} + fn fetch_stream_cancel(args: &[Value], _ctx: &mut NativeContext) -> Result { let id = args.first().and_then(value_to_stream_id).unwrap_or(0); FETCH_BODY_STREAMS.with(|streams| { @@ -1292,8 +1476,12 @@ pub fn fetch_text_sync(url: &str) -> Result { pub fn init_fetch_api(vm: &mut Vm) { vm.define_native("__we_fetch_stream_cancel", fetch_stream_cancel); vm.define_native("__we_fetch_reader_read", fetch_reader_read); + vm.define_native("__we_fetch_upload_start", fetch_upload_start); + vm.define_native("__we_fetch_upload_chunk", fetch_upload_chunk); + vm.define_native("__we_fetch_upload_close", fetch_upload_close); + vm.define_native("__we_fetch_upload_error", fetch_upload_error); let fetch_fn = make_native(&mut vm.gc, "fetch", fetch_native); - vm.set_global("fetch", Value::Function(fetch_fn)); + vm.set_global("__we_fetch_native", Value::Function(fetch_fn)); init_fetch_stream_preamble(vm); } @@ -1318,6 +1506,97 @@ fn init_fetch_stream_preamble(vm: &mut Vm) { } const FETCH_STREAM_PREAMBLE: &str = r#" +function __we_fetch_is_stream(value) { + return value !== null && value !== undefined && typeof value.getReader === 'function'; +} + +function __we_fetch_stream_locked(value) { + return !!(value && (value.locked || value._reader)); +} + +function __we_fetch_start_upload(stream) { + if (__we_fetch_stream_locked(stream) || stream.__we_fetch_body_used) { + throw new TypeError("body is locked or used"); + } + var uploadId = __we_fetch_upload_start(); + var reader = stream.getReader(); + stream.__we_fetch_body_used = true; + function pump() { + return reader.read().then( + function(record) { + if (record.done) { + __we_fetch_upload_close(uploadId); + return; + } + __we_fetch_upload_chunk(uploadId, record.value); + return pump(); + }, + function(reason) { + __we_fetch_upload_error(uploadId, reason); + } + ); + } + pump(); + return { __we_upload_stream_id: uploadId }; +} + +function Request(input, init) { + var self = Object.create(Request.prototype); + init = init || {}; + if (input !== null && typeof input === 'object' && input.url !== undefined) { + if (input.bodyUsed || (input.body && __we_fetch_stream_locked(input.body))) { + throw new TypeError("body is locked or used"); + } + self.url = String(input.url); + self.method = init.method !== undefined ? String(init.method) : input.method; + self.headers = init.headers !== undefined ? init.headers : input.headers; + self.body = init.body !== undefined ? init.body : input.body; + } else { + self.url = String(input); + self.method = init.method !== undefined ? String(init.method) : "GET"; + self.headers = init.headers || {}; + self.body = init.body; + } + self.bodyUsed = false; + return self; +} + +function __we_fetch_prepare(input, init) { + init = init || {}; + var url = input; + var opts = {}; + var sourceRequest = null; + if (input !== null && typeof input === 'object' && input.url !== undefined) { + sourceRequest = input; + url = input.url; + opts.method = input.method; + opts.headers = input.headers; + opts.body = input.body; + } + if (init.method !== undefined) opts.method = init.method; + if (init.headers !== undefined) opts.headers = init.headers; + if (init.mode !== undefined) opts.mode = init.mode; + if (init.credentials !== undefined) opts.credentials = init.credentials; + if (init.body !== undefined) opts.body = init.body; + + if (__we_fetch_is_stream(opts.body)) { + if (sourceRequest && sourceRequest.bodyUsed) throw new TypeError("body is locked or used"); + opts.body = __we_fetch_start_upload(opts.body); + if (sourceRequest) sourceRequest.bodyUsed = true; + } else if (sourceRequest && sourceRequest.body !== undefined && init.body === undefined) { + sourceRequest.bodyUsed = true; + } + return { url: url, options: opts }; +} + +function fetch(input, init) { + if (input === undefined && init === undefined) { + throw new TypeError("fetch requires at least one argument"); + } + var prepared = __we_fetch_prepare(input, init); + return __we_fetch_native(prepared.url, prepared.options); +} + function __we_fetch_body_unusable(response) { return response.bodyUsed || (response.body !== null && response.body.locked); } @@ -1943,6 +2222,72 @@ mod tests { ); } + #[test] + fn test_upload_stream_preamble_drains_readablestream_chunks() { + reset_fetch_state(); + crate::timers::reset_timers(); + let source = r#" + var upload = new ReadableStream({ + start: function(controller) { + controller.enqueue("AB"); + controller.enqueue("C"); + controller.close(); + } + }); + __we_fetch_start_upload(upload); + "#; + let program = Parser::parse(source).expect("parse failed"); + let func = compiler::compile(&program).expect("compile failed"); + let mut vm = Vm::new(); + init_fetch_api(&mut vm); + vm.execute(&func).expect("execute failed"); + vm.run_event_loop(100).expect("event loop failed"); + + FETCH_UPLOAD_STREAMS.with(|streams| { + let streams = streams.borrow(); + let stream = streams.last().expect("upload stream"); + let queue = stream.state.queue.lock().unwrap(); + assert_eq!(queue.len(), 3, "queue: {queue:?}"); + match &queue[0] { + FetchUploadEvent::Chunk(bytes) => assert_eq!(bytes, b"AB"), + other => panic!("expected first chunk, got {other:?}"), + } + match &queue[1] { + FetchUploadEvent::Chunk(bytes) => assert_eq!(bytes, b"C"), + other => panic!("expected second chunk, got {other:?}"), + } + assert!(matches!(queue[2], FetchUploadEvent::Close)); + }); + } + + #[test] + fn test_fetch_rejects_reusing_readablestream_body() { + reset_fetch_state(); + crate::timers::reset_timers(); + let source = r#" + var upload = new ReadableStream({ + start: function(controller) { + controller.enqueue("one-shot"); + controller.close(); + } + }); + fetch("http://127.0.0.1:1/upload", { method: "POST", body: upload }); + try { + fetch("http://127.0.0.1:1/upload", { method: "POST", body: upload }); + "reuse accepted"; + } catch (err) { + "reuse rejected"; + } + "#; + let program = Parser::parse(source).expect("parse failed"); + let func = compiler::compile(&program).expect("compile failed"); + let mut vm = Vm::new(); + init_fetch_api(&mut vm); + let result = vm.execute(&func).expect("execute failed"); + assert_eq!(result.to_js_string(&vm.gc), "reuse rejected"); + reset_fetch_state(); + } + #[test] fn test_response_object_properties() { // Test that create_response_object creates the right structure. diff --git a/crates/net/src/client.rs b/crates/net/src/client.rs index 37c47b3..411e63f 100644 --- a/crates/net/src/client.rs +++ b/crates/net/src/client.rs @@ -118,6 +118,12 @@ enum Connection { Tls(TlsStream), } +/// Request body source for HTTP/1.1 uploads. +pub enum RequestBody<'a> { + Bytes(&'a [u8]), + Chunked(&'a mut dyn FnMut() -> Result>>), +} + impl Connection { fn read(&mut self, buf: &mut [u8]) -> Result { match self { @@ -404,6 +410,31 @@ impl HttpClient { url: &Url, headers: &Headers, body: Option<&[u8]>, + on_response: H, + on_chunk: F, + ) -> Result + where + H: FnMut(&HttpResponse) -> Result<()>, + F: FnMut(&[u8]) -> Result<()>, + { + self.request_streaming_with_body( + method, + url, + headers, + body.map(RequestBody::Bytes), + on_response, + on_chunk, + ) + } + + /// Perform a single HTTP/1.1 request with a possibly streaming request + /// body and stream response body bytes to `on_chunk` as they arrive. + pub fn request_streaming_with_body( + &mut self, + method: Method, + url: &Url, + headers: &Headers, + body: Option>, mut on_response: H, mut on_chunk: F, ) -> Result @@ -629,7 +660,7 @@ impl HttpClient { path: &str, host: &str, headers: &Headers, - body: Option<&[u8]>, + body: Option>, url: &Url, on_response: &mut H, on_chunk: &mut F, @@ -640,8 +671,33 @@ impl HttpClient { { conn.set_read_timeout(Some(self.read_timeout))?; - let request_bytes = http::serialize_request(method, path, host, headers, body); - conn.write_all(&request_bytes)?; + match body { + Some(RequestBody::Bytes(bytes)) => { + let request_bytes = + http::serialize_request(method, path, host, headers, Some(bytes)); + conn.write_all(&request_bytes)?; + } + Some(RequestBody::Chunked(next_chunk)) => { + let request_head = + http::serialize_request_head_chunked(method, path, host, headers); + conn.write_all(&request_head)?; + while let Some(chunk) = next_chunk()? { + if chunk.is_empty() { + continue; + } + let size = format!("{:x}\r\n", chunk.len()); + conn.write_all(size.as_bytes())?; + conn.write_all(&chunk)?; + conn.write_all(b"\r\n")?; + conn.flush()?; + } + conn.write_all(b"0\r\n\r\n")?; + } + None => { + let request_bytes = http::serialize_request(method, path, host, headers, None); + conn.write_all(&request_bytes)?; + } + } conn.flush()?; let response = read_response_streaming(&mut conn, on_response, on_chunk)?; diff --git a/crates/net/src/http.rs b/crates/net/src/http.rs index fd2ac27..371d4e4 100644 --- a/crates/net/src/http.rs +++ b/crates/net/src/http.rs @@ -320,6 +320,57 @@ pub fn serialize_request( buf } +/// Serialize an HTTP/1.1 request head for a body that will be written +/// separately with chunked transfer coding. +pub fn serialize_request_head_chunked( + method: Method, + path: &str, + host: &str, + headers: &Headers, +) -> Vec { + let mut buf = Vec::with_capacity(256); + + buf.extend_from_slice(method.as_str().as_bytes()); + buf.push(b' '); + buf.extend_from_slice(path.as_bytes()); + buf.push(b' '); + buf.extend_from_slice(HTTP_VERSION.as_bytes()); + buf.extend_from_slice(CRLF.as_bytes()); + + if !headers.contains("Host") { + buf.extend_from_slice(b"Host: "); + buf.extend_from_slice(host.as_bytes()); + buf.extend_from_slice(CRLF.as_bytes()); + } + if !headers.contains("User-Agent") { + buf.extend_from_slice(b"User-Agent: "); + buf.extend_from_slice(USER_AGENT.as_bytes()); + buf.extend_from_slice(CRLF.as_bytes()); + } + if !headers.contains("Accept") { + buf.extend_from_slice(b"Accept: */*"); + buf.extend_from_slice(CRLF.as_bytes()); + } + if !headers.contains("Connection") { + buf.extend_from_slice(b"Connection: keep-alive"); + buf.extend_from_slice(CRLF.as_bytes()); + } + if !headers.contains("Content-Length") && !headers.contains("Transfer-Encoding") { + buf.extend_from_slice(b"Transfer-Encoding: chunked"); + buf.extend_from_slice(CRLF.as_bytes()); + } + + for (name, value) in headers.iter() { + buf.extend_from_slice(name.as_bytes()); + buf.extend_from_slice(b": "); + buf.extend_from_slice(value.as_bytes()); + buf.extend_from_slice(CRLF.as_bytes()); + } + + buf.extend_from_slice(CRLF.as_bytes()); + buf +} + // --------------------------------------------------------------------------- // HTTP Response // --------------------------------------------------------------------------- @@ -755,6 +806,21 @@ mod tests { assert!(req_str.ends_with("{\"key\": \"value\"}")); } + #[test] + fn serialize_chunked_request_head() { + let mut headers = Headers::new(); + headers.add("Content-Type", "application/octet-stream"); + let req = serialize_request_head_chunked(Method::Post, "/upload", "example.com", &headers); + let req_str = String::from_utf8(req).unwrap(); + + assert!(req_str.starts_with("POST /upload HTTP/1.1\r\n")); + assert!(req_str.contains("Host: example.com\r\n")); + assert!(req_str.contains("Transfer-Encoding: chunked\r\n")); + assert!(req_str.contains("Content-Type: application/octet-stream\r\n")); + assert!(!req_str.contains("Content-Length:")); + assert!(req_str.ends_with("\r\n\r\n")); + } + #[test] fn serialize_request_with_path() { let headers = Headers::new(); -- 2.51.2