fix(xhr): dispatch errors when streamed response bodies fail

Complete active page XHRs with the request-error steps when a response
stream fails after delivering partial data. Reset the response and dispatch
error/loadend instead of leaving the request in LOADING.

Check the request identity for stream and body-materialization failures so
queued errors cannot overwrite an aborted or reopened XHR. Exercise real
truncated Content-Length and malformed chunked responses, including abort,
open, and resend before a queued failure is delivered.

Validation: cargo fmt --all; full workspace clippy with warnings denied;
cargo nextest run --no-fail-fast (18,153 passed, 13 skipped).
Focused upstream WPT: 83/88 cases and 133/143 subtests pass, up from 82/88
and 132/143, with no regressions. Page malformed-body error now passes.
This commit is contained in:
ldm0
2026-09-14 18:42:34 +08:00
parent 4e375fc8a9
commit 37eef5dce0
7 changed files with 229 additions and 8 deletions
+4 -3
View File
@@ -181,9 +181,10 @@ pub(crate) use self::xhr::{
XHR_METHOD_SLOT, XHR_OPEN_GENERATION_SLOT, XHR_SEND_FLAG_SLOT, XHR_TIMEOUT_SLOT,
XHR_TIMEOUT_START_MS_SLOT, XHR_TIMEOUT_TIMER_SLOT, XHR_URL_SLOT, XHR_WITH_CREDENTIALS_SLOT,
apply_xhr_failure, apply_xhr_response, apply_xhr_response_body_source,
apply_xhr_response_body_source_with_status_text, apply_xhr_streaming_response_body_source,
apply_xhr_streaming_response_chunk, apply_xhr_streaming_response_head, apply_xhr_timeout,
apply_xhr_upload_event, capture_xhr_upload_listener_flag, dispatch_xhr_loadstart,
apply_xhr_response_body_source_with_status_text, apply_xhr_streaming_failure,
apply_xhr_streaming_response_body_source, apply_xhr_streaming_response_chunk,
apply_xhr_streaming_response_head, apply_xhr_timeout, apply_xhr_upload_event,
capture_xhr_upload_listener_flag, dispatch_xhr_loadstart,
finalize_xml_http_request_event_target_realm_bindings, finish_xhr_abort,
install_progress_event_template_bindings, install_window_xml_http_request_template_bindings,
install_xml_http_request_bindings, install_xml_http_request_event_target_bindings,
+4 -3
View File
@@ -57,9 +57,10 @@ pub(crate) use self::bindings::{
};
pub(crate) use self::delivery::{
apply_xhr_failure, apply_xhr_response, apply_xhr_response_body_source,
apply_xhr_response_body_source_with_status_text, apply_xhr_streaming_response_body_source,
apply_xhr_streaming_response_chunk, apply_xhr_streaming_response_head, apply_xhr_timeout,
finish_xhr_abort, throw_synchronous_xhr_failure,
apply_xhr_response_body_source_with_status_text, apply_xhr_streaming_failure,
apply_xhr_streaming_response_body_source, apply_xhr_streaming_response_chunk,
apply_xhr_streaming_response_head, apply_xhr_timeout, finish_xhr_abort,
throw_synchronous_xhr_failure,
};
pub(crate) use self::encoding::xhr_response_text_decoder;
pub(crate) use self::instance_state::{
@@ -8,7 +8,9 @@ mod timeout;
use super::*;
pub(crate) use self::abort::{apply_xhr_abort, finish_xhr_abort};
pub(crate) use self::failure::{apply_xhr_failure, throw_synchronous_xhr_failure};
pub(crate) use self::failure::{
apply_xhr_failure, apply_xhr_streaming_failure, throw_synchronous_xhr_failure,
};
pub(super) use self::pending::{queue_xhr_failure_delivery, queue_xhr_response_delivery};
pub(in crate::network_host::xhr) use self::progress::clear_xhr_progress_throttle;
pub(super) use self::response::apply_xhr_response_pending_body;
@@ -3,6 +3,16 @@ use super::super::events::{
};
use super::super::*;
pub(crate) fn apply_xhr_streaming_failure(
scope: &mut v8::PinScope<'_, '_>,
xhr: v8::Local<'_, v8::Object>,
internal_id: u64,
) {
if super::progress::xhr_stream_is_current(scope, xhr, internal_id) {
apply_xhr_failure(scope, xhr);
}
}
pub(crate) fn apply_xhr_failure(scope: &mut v8::PinScope<'_, '_>, xhr: v8::Local<'_, v8::Object>) {
super::cancel_xhr_timeout(scope, xhr);
super::clear_xhr_progress_throttle(scope, xhr);
@@ -5734,7 +5734,9 @@ impl ScriptVm {
streaming.pending.continuation
{
let xhr = v8::Local::new(scope, &xhr);
crate::network_host::apply_xhr_failure(scope, xhr);
crate::network_host::apply_xhr_streaming_failure(
scope, xhr, internal_id,
);
}
context_host
.borrow_mut()
@@ -5946,6 +5948,14 @@ impl ScriptVm {
trace_fields,
error_started,
);
if let PendingSubresourceContinuation::Xhr { xhr, .. } =
&streaming.pending.continuation
{
let xhr = v8::Local::new(scope, xhr);
crate::network_host::apply_xhr_streaming_failure(
scope, xhr, internal_id,
);
}
}
}
context_host
@@ -8,6 +8,7 @@ mod forms;
mod misc;
mod query_realms;
mod shadow_dom;
mod streaming_failure;
mod style_invalidation;
mod upload_preflight;
mod upload_transport;
@@ -0,0 +1,196 @@
use super::*;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
#[tokio::test(flavor = "multi_thread")]
async fn xhr_streaming_body_failure_finishes_only_the_current_request() {
for chunked in [false, true] {
for action in ["error", "abort", "reopen", "resend"] {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let origin = format!("http://{}", listener.local_addr().unwrap());
let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let mut head = Vec::new();
while !head.ends_with(b"\r\n\r\n") {
head.push(socket.read_u8().await.unwrap());
}
assert!(head.starts_with(b"GET /broken HTTP/1.1\r\n"));
let framing = if chunked {
"Transfer-Encoding: chunked"
} else {
"Content-Length: 14"
};
socket.write_all(format!(
"HTTP/1.1 200 OK\r\n{framing}\r\nContent-Type: text/plain\r\nConnection: close\r\n\r\n"
).as_bytes()).await.unwrap();
socket
.write_all(if chunked {
b"7\r\npartial\r\n"
} else {
b"partial"
})
.await
.unwrap();
// Expose LOADING and partial response data before the failure.
release_rx.await.unwrap();
if chunked {
let _ = socket.write_all(b"garbage").await;
}
drop(socket);
if action == "resend" {
let (mut socket, _) = listener.accept().await.unwrap();
let mut head = Vec::new();
while !head.ends_with(b"\r\n\r\n") {
head.push(socket.read_u8().await.unwrap());
}
assert!(head.starts_with(b"GET /replacement HTTP/1.1\r\n"));
socket
.write_all(
b"HTTP/1.1 200 OK\r\nContent-Length: 3\r\nConnection: close\r\n\r\nnew",
)
.await
.unwrap();
}
});
let mut config = moli_fetch::FetchConfig::default();
config.set_http_no_proxy(Some("*".to_owned()));
let loader = ResourceRequestClient::new(&config).unwrap();
let (mut vm, mut completions) =
new_storage_test_vm_with_loader_and_resource_completion_queue(
&format!("{origin}/page"),
&loader,
);
vm.eval(r#"
globalThis.events = [];
globalThis.partial = null;
globalThis.xhr = new XMLHttpRequest();
globalThis.snapshot = () => ({
state: xhr.readyState, status: xhr.status, text: xhr.responseText,
responseURL: xhr.responseURL, headers: xhr.getAllResponseHeaders(),
contentType: xhr.getResponseHeader("content-type")
});
for (const type of ["error", "abort", "load", "loadend"])
xhr.addEventListener(type, e => events.push([type, e.loaded, e.total, e.lengthComputable]));
xhr.onprogress = () => { if (!partial) partial = snapshot(); };
xhr.open("GET", "/broken");
xhr.send();
"#).unwrap();
let mut old_id = None;
tokio::time::timeout(Duration::from_secs(5), async {
while vm.eval("partial !== null").unwrap() != "true" {
assert!(completions.wait_for_arrival_without_timeout().await);
while let Some(event) = completions.pop_next_async_subresource_event() {
if let crate::types::AsyncSubresourceFetchEvent::StreamingStarted(started) =
&event
{
old_id = Some(started.internal_id);
}
let activity = vm
.complete_async_subresource_fetch_event_body(event)
.unwrap();
vm.finish_async_subresource_body_checkpoint_for_test(activity)
.unwrap();
}
}
})
.await
.expect("partial XHR response should be observable");
assert_eq!(vm.eval("JSON.stringify([partial.state, partial.status, partial.text, partial.contentType])").unwrap(), "[3,200,\"partial\",\"text/plain\"]");
let old_id = old_id.expect("stream should start before delivering bytes");
release_tx.send(()).unwrap();
let failure =
tokio::time::timeout(Duration::from_secs(5), async {
loop {
assert!(completions.wait_for_arrival_without_timeout().await);
while let Some(event) = completions.pop_next_async_subresource_event() {
if let crate::types::AsyncSubresourceFetchEvent::StreamingFinished(
finished,
) = &event
&& finished.internal_id == old_id
{
assert!(finished.result.is_err());
return event;
}
let activity = vm
.complete_async_subresource_fetch_event_body(event)
.unwrap();
vm.finish_async_subresource_body_checkpoint_for_test(activity)
.unwrap();
}
}
})
.await
.expect("native transport should queue the body failure");
// The network has finished, but its terminal task has not entered
// JS yet. Reusing the XHR must invalidate that queued delivery.
match action {
"abort" => {
vm.eval("xhr.abort()").unwrap();
}
"reopen" => {
vm.eval("xhr.open('GET', '/replacement')").unwrap();
}
"resend" => {
vm.eval("xhr.open('GET', '/replacement'); xhr.send()")
.unwrap();
}
_ => {}
}
let activity = vm
.complete_async_subresource_fetch_event_body(failure)
.unwrap();
vm.finish_async_subresource_body_checkpoint_for_test(activity)
.unwrap();
tokio::time::timeout(Duration::from_secs(5), async {
while action == "resend" && vm.eval("xhr.readyState === 4").unwrap() != "true" {
assert!(completions.wait_for_arrival_without_timeout().await);
while let Some(event) = completions.pop_next_async_subresource_event() {
let activity = vm
.complete_async_subresource_fetch_event_body(event)
.unwrap();
vm.finish_async_subresource_body_checkpoint_for_test(activity)
.unwrap();
}
}
})
.await
.unwrap_or_else(|_| panic!("stream finish stalled: {chunked}/{action}"));
let expected_events = match action {
"error" => "[[\"error\",0,0,false],[\"loadend\",0,0,false]]",
"abort" => "[[\"abort\",0,0,false],[\"loadend\",0,0,false]]",
"reopen" => "[]",
"resend" => "[[\"load\",3,3,true],[\"loadend\",3,3,true]]",
_ => unreachable!(),
};
assert_eq!(
vm.eval("JSON.stringify(events)").unwrap(),
expected_events,
"{chunked}/{action}"
);
if action == "resend" {
assert_eq!(
vm.eval("JSON.stringify([xhr.status, xhr.responseText])")
.unwrap(),
"[200,\"new\"]"
);
} else {
let state = match action {
"error" => 4,
"reopen" => 1,
_ => 0,
};
assert_eq!(
vm.eval("JSON.stringify(snapshot())").unwrap(),
format!(
r#"{{"state":{state},"status":0,"text":"","responseURL":"","headers":"","contentType":null}}"#
)
);
}
tokio::time::timeout(Duration::from_secs(3), server)
.await
.unwrap()
.unwrap();
}
}
}