fix(fetch): reject locked bodies before consumption and cloning

This commit is contained in:
ldm0
2026-10-04 12:31:25 +08:00
parent 7106a4e4b6
commit d74f65b06f
11 changed files with 128 additions and 8 deletions
+47
View File
@@ -0,0 +1,47 @@
async function runLockedBodyProbe(url) {
const errors = [];
const check = (value, label) => { if (!value) errors.push(label); };
const throwsTypeError = (callback, label) => {
try { callback(); errors.push(label + ': accepted'); }
catch (error) { check(error instanceof TypeError, label + ': wrong error ' + error); }
};
const rejectsTypeError = async (callback, label) => {
let promise;
try { promise = callback(); }
catch (error) { errors.push(label + ': synchronous throw ' + error); return; }
try { await promise; errors.push(label + ': fulfilled'); }
catch (error) { check(error instanceof TypeError, label + ': wrong rejection ' + error); }
};
const methods = ['text', 'json', 'arrayBuffer', 'bytes', 'blob', 'formData'];
const streamFrom = value => new ReadableStream({
start(controller) { controller.enqueue(new TextEncoder().encode(value)); controller.close(); }
});
const makeBody = async (source, method = 'text') => {
const text = method === 'formData' ? 'name=value' : '{"value":1}';
const headers = {'Content-Type': method === 'formData'
? 'application/x-www-form-urlencoded' : 'application/json'};
if (source === 'request') return new Request(url, {method: 'POST', body: text, headers});
if (source === 'response') return new Response(text, {headers});
if (source === 'stream') return new Response(streamFrom(text), {headers});
return fetch(url);
};
const sources = ['request', 'response', 'stream', 'fetch'];
for (const source of sources) {
for (const method of methods) {
const body = await makeBody(source, method);
const reader = body.body.getReader();
Object.defineProperty(body.body, 'locked', {get() {
throw new Error('public locked getter consulted');
}});
check(!body.bodyUsed, source + ': locking alone does not disturb');
throwsTypeError(() => body.clone(), source + ': locked clone');
await rejectsTypeError(() => body[method](), source + ': locked ' + method);
check(!body.bodyUsed, source + ': failed ' + method + ' must not disturb');
reader.releaseLock();
try { await body.text(); }
catch (error) { errors.push(source + ': released unused body cannot be consumed ' + error); }
}
}
return {errors};
}
+13 -1
View File
@@ -11,6 +11,7 @@
<script>
(async () => {
const response = new Response('hello world');
const textResponse = response.clone();
const reader = response.body.getReader();
const first = await reader.read();
document.body.setAttribute(
@@ -18,7 +19,18 @@
new TextDecoder().decode(first.value),
);
document.body.setAttribute('data-response-text', await response.text());
reader.releaseLock();
let consumedBodyRejected = false;
try {
await response.text();
} catch (error) {
if (!(error instanceof TypeError)) throw error;
consumedBodyRejected = true;
}
if (!consumedBodyRejected || !response.bodyUsed) {
throw new Error('Consumed Response body must remain unusable');
}
document.body.setAttribute('data-response-text', await textResponse.text());
const input = new ReadableStream({
start(controller) {
+3
View File
@@ -1,3 +1,6 @@
#[path = "web_apis/locked_body.rs"]
mod locked_body;
#[path = "web_apis/body_utf8.rs"]
mod body_utf8;
+22
View File
@@ -0,0 +1,22 @@
use super::*;
#[tokio::test(flavor = "multi_thread")]
async fn locked_bodies_reject_consumption_and_clone_without_becoming_disturbed() -> Result<()> {
let server = FixtureServer::spawn().await?;
let browser = Browser::new(AppConfig::default())?;
let source = format!(
"{}\nrunLockedBodyProbe({}).then(finish, error => finish({{error: String(error)}}));",
include_str!("../fixtures/runtime/locked_body.js"),
serde_json::to_string(&server.url("/compat/child-dynamic-markup-document"))?,
);
for target in ["window", "child", "worker"] {
let observed = tokio::time::timeout(
Duration::from_secs(20),
super::pipe_disturbed::run_probe(&browser, &server, target, &source),
)
.await??;
assert_eq!(observed, serde_json::json!({"errors": []}), "{target}");
}
server.shutdown().await;
Ok(())
}
+2 -1
View File
@@ -426,7 +426,8 @@ pub use self::storage_buckets::{
};
pub(crate) use self::stream_adapter::{
cancel_readable_stream, close_stream, enqueue_byte_chunk, error_stream,
readable_stream_disturbed, readable_stream_has_pipe_owner, require_internal_stream_value,
readable_stream_disturbed, readable_stream_has_pipe_owner, readable_stream_locked,
require_internal_stream_value,
};
pub(crate) use self::streams::{
ReadableStreamClonePayload, TransformStreamClonePayload, WritableStreamClonePayload,
@@ -107,8 +107,8 @@ pub(in crate::context_bootstrap) use readable_state::{
};
pub(crate) use readable_state::{close_stream, enqueue_chunk, error_stream};
pub(super) use readable_state::{
readable_stream_closed, readable_stream_error, readable_stream_locked,
reject_pending_read_requests, remove_pending_closed_promise, writable_stream_locked,
readable_stream_closed, readable_stream_error, reject_pending_read_requests,
remove_pending_closed_promise, writable_stream_locked,
};
pub(in crate::context_bootstrap::stream_adapter) use transform_finish::transform_stream_readable_cancel_callback;
pub(super) use utils::{
@@ -140,6 +140,7 @@ pub(in crate::context_bootstrap::stream_adapter) use read_request::{
error_read_request, fulfill_read_request, new_internal_read_request,
};
pub(in crate::context_bootstrap::stream_adapter) use readable_state::enqueue_pending_closed_promise;
pub(crate) use readable_state::readable_stream_locked;
pub(in crate::context_bootstrap::stream_adapter) use readable_state::{
enqueue_pending_read, finish_readable_stream_close,
readable_stream_controller_algorithm_object, readable_stream_controller_algorithm_value,
@@ -568,7 +568,7 @@ fn reject_readable_stream_closed_promises<'s>(
);
}
pub(in crate::context_bootstrap) fn readable_stream_locked<'s>(
pub(crate) fn readable_stream_locked<'s>(
scope: &mut v8::PinScope<'s, '_>,
stream: v8::Local<'s, v8::Object>,
) -> bool {
+3 -1
View File
@@ -37,7 +37,9 @@ pub(crate) use self::async_fetch::{
};
pub(crate) use self::beacon::{navigator_send_beacon_callback, send_link_audit_ping};
pub(super) use self::bindings::install_window_network_bindings;
pub(in crate::network_host) use self::body::{PreparedBodyInit, body_init};
pub(in crate::network_host) use self::body::{
PreparedBodyInit, body_init, body_stream_object, readable_body_stream_unusable,
};
pub(crate) use self::body::{append_default_body_content_type, has_header};
#[cfg(test)]
pub(crate) use self::body_source::pending_network_body_source_buffered_len_for_test;
+21
View File
@@ -6,6 +6,27 @@ pub(in crate::network_host) const URL_SEARCH_PARAMS_CONTENT_TYPE: &str =
"application/x-www-form-urlencoded;charset=UTF-8";
pub(in crate::network_host) const TEXT_CONTENT_TYPE: &str = "text/plain;charset=UTF-8";
pub(in crate::network_host) fn body_stream_object<'s>(
scope: &mut v8::PinScope<'s, '_>,
object: v8::Local<'s, v8::Object>,
) -> Option<v8::Local<'s, v8::Object>> {
if is_branded_request_object(scope, object) {
request_slot_object(scope, object, REQUEST_BODY_SLOT)
} else if is_branded_response_object(scope, object) {
response_slot_object(scope, object, RESPONSE_BODY_SLOT)
} else {
None
}
}
pub(in crate::network_host) fn readable_body_stream_unusable<'s>(
scope: &mut v8::PinScope<'s, '_>,
stream: v8::Local<'s, v8::Object>,
) -> bool {
crate::context_bootstrap::readable_stream_locked(scope, stream)
|| crate::context_bootstrap::readable_stream_disturbed(scope, stream)
}
#[derive(Debug, Clone)]
pub(in crate::network_host) struct PreparedBodyInit {
pub(in crate::network_host) bytes: Vec<u8>,
@@ -800,7 +800,10 @@ fn request_clone_callback<'s>(
let Some(this) = require_request_receiver(scope, args.this()) else {
return;
};
if body_already_used(scope, this, REQUEST_BODY_USED_SLOT) {
if body_already_used(scope, this, REQUEST_BODY_USED_SLOT)
|| body_stream_object(scope, this)
.is_some_and(|stream| readable_body_stream_unusable(scope, stream))
{
throw_type_error(
scope,
"Failed to execute 'clone' on 'Request': body stream already used",
@@ -831,7 +834,10 @@ fn response_clone_callback<'s>(
let Some(this) = require_response_receiver(scope, args.this()) else {
return;
};
if body_already_used(scope, this, RESPONSE_BODY_USED_SLOT) {
if body_already_used(scope, this, RESPONSE_BODY_USED_SLOT)
|| body_stream_object(scope, this)
.is_some_and(|stream| readable_body_stream_unusable(scope, stream))
{
throw_type_error(
scope,
"Failed to execute 'clone' on 'Response': body stream already used",
@@ -208,6 +208,11 @@ fn begin_body_consumption<'s>(
object: v8::Local<'s, v8::Object>,
receiver: BodyReceiver,
) -> bool {
if body_stream_object(scope, object)
.is_some_and(|stream| readable_body_stream_unusable(scope, stream))
{
return false;
}
match receiver {
BodyReceiver::Response => {
if response_slot_bool(scope, object, RESPONSE_BODY_USED_SLOT) {