fix(xhr): queue early worker request failures

Route URL and network policy rejections through the pending worker XHR
completion queue so send() returns before error delivery. Immediate abort()
and open() can cancel the queued failure. Validate request identity before
applying a completion so old failures cannot overwrite a new send.

Cover six rejection causes, listeners installed after send, microtask abort,
reopen and resend, synchronous NetworkError, and duplicate failure records.

Validation: cargo fmt --all; full workspace clippy with warnings denied;
cargo nextest run --no-fail-fast (18,152 passed, 13 skipped).
Focused upstream WPT: 77/86 cases and 127/141 subtests pass, up from 76/86
and 126/141, with no regressions. Worker abort-with-error now passes.
This commit is contained in:
ldm0
2026-09-13 17:07:51 +08:00
parent 77962500c8
commit 2d8390f335
3 changed files with 291 additions and 92 deletions
+48 -92
View File
@@ -329,90 +329,42 @@ pub(crate) fn try_worker_xhr_send_callback<'s>(
crate::content_security_policy::ContentSecurityPolicyResourceKind::WorkerConnect,
)
};
if let Some(violation) = csp_violation {
let request_error = if let Some(violation) = csp_violation {
dispatch_worker_content_security_policy_violation_event_for_state(
scope, &state, &violation,
);
let message = worker_content_security_policy_error_message(&violation, "xhr");
record_worker_subresource_failure(
&state.borrow(),
prepared.document_url,
prepared.resolved_url,
prepared.method,
prepared.request_headers,
request_body_text(&prepared.send_body),
SubresourceResourceType::Xhr,
message,
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
Some(worker_content_security_policy_error_message(
&violation, "xhr",
))
} else if let Err(error) = url_policy {
Some(error.to_string())
} else if should_request_be_blocked_due_to_bad_port(&prepared.resolved_url) {
Some(format!(
"xhr: blocked bad port for `{}`",
prepared.resolved_url
))
} else if worker_url_blocked(&blocked_url_patterns, &prepared.resolved_url) {
Some(BLOCKED_BY_CLIENT_ERROR_TEXT.to_owned())
} else if network_offline {
Some("Network emulation offline".to_owned())
} else {
None
};
if xhr_state_bool_property(scope, xhr, XHR_ABORTED_SLOT).unwrap_or(false)
|| worker_xhr_open_generation_changed(scope, xhr, open_generation)
{
return true;
}
if let Err(error) = url_policy {
record_worker_subresource_failure(
&state.borrow(),
prepared.document_url,
prepared.resolved_url,
prepared.method,
prepared.request_headers,
request_body_text(&prepared.send_body),
SubresourceResourceType::Xhr,
error.to_string(),
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
return true;
}
if should_request_be_blocked_due_to_bad_port(&prepared.resolved_url) {
record_worker_subresource_failure(
&state.borrow(),
prepared.document_url,
prepared.resolved_url.clone(),
prepared.method,
prepared.request_headers,
request_body_text(&prepared.send_body),
SubresourceResourceType::Xhr,
format!("xhr: blocked bad port for `{}`", prepared.resolved_url),
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
return true;
}
if worker_url_blocked(&blocked_url_patterns, &prepared.resolved_url) {
record_worker_subresource_failure(
&state.borrow(),
prepared.document_url,
prepared.resolved_url,
prepared.method,
prepared.request_headers,
request_body_text(&prepared.send_body),
SubresourceResourceType::Xhr,
BLOCKED_BY_CLIENT_ERROR_TEXT.to_owned(),
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
return true;
}
if network_offline {
record_worker_subresource_failure(
&state.borrow(),
prepared.document_url,
prepared.resolved_url,
prepared.method,
prepared.request_headers,
request_body_text(&prepared.send_body),
SubresourceResourceType::Xhr,
"Network emulation offline".to_owned(),
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
return true;
}
let local_response = local_url_response_with_blob_entry(
&prepared.resolved_url,
&prepared.method,
prepared.blob_url_entry.as_ref(),
);
// A rejected fetch still completes in a networking task. Keep the pending
// XHR and its load lease until delivery so abort/open can cancel that task.
let local_response = request_error.map(Err).or_else(|| {
local_url_response_with_blob_entry(
&prepared.resolved_url,
&prepared.method,
prepared.blob_url_entry.as_ref(),
)
});
if !async_request && let Some(result) = local_response {
match result {
Ok(response) => apply_xhr_response(scope, xhr, response),
@@ -452,7 +404,7 @@ pub(crate) fn try_worker_xhr_send_callback<'s>(
SubresourceResourceType::Xhr,
"Synchronous XMLHttpRequest interception is not supported".to_owned(),
);
apply_worker_xhr_request_failure(scope, xhr, async_request, &request_url);
throw_synchronous_xhr_failure(scope, xhr, &request_url, "NetworkError");
return true;
}
@@ -801,19 +753,6 @@ fn send_synchronous_worker_xhr(
}
}
fn apply_worker_xhr_request_failure(
scope: &mut v8::PinScope<'_, '_>,
xhr: v8::Local<'_, v8::Object>,
async_request: bool,
request_url: &str,
) {
if async_request {
apply_xhr_failure(scope, xhr);
} else {
throw_synchronous_xhr_failure(scope, xhr, request_url, "NetworkError");
}
}
struct SynchronousWorkerXhrTimeout {
wait_delay: Duration,
configured_timeout: Duration,
@@ -968,6 +907,23 @@ pub(in crate::worker) fn drain_worker_xhr_completion(
state: &Rc<RefCell<WorkerGlobalState>>,
completion: WorkerXhrCompletion,
) {
let xhr = {
let state = state.borrow();
let Some(pending) = state.pending_xhrs.get(&completion.xhr_id) else {
return;
};
v8::Local::new(scope, &pending.xhr)
};
if xhr_state_number_property(scope, xhr, XHR_ACTIVE_INTERNAL_ID_SLOT)
!= Some(f64::from(completion.xhr_id))
|| !xhr_state_bool_property(scope, xhr, XHR_SEND_FLAG_SLOT).unwrap_or(false)
|| xhr_state_bool_property(scope, xhr, XHR_ABORTED_SLOT).unwrap_or(false)
{
if let Some(pending) = state.borrow_mut().pending_xhrs.remove(&completion.xhr_id) {
pending.load.cancel();
}
return;
}
let parent_tx = state.borrow().parent_tx.clone();
if let Some(network_request_headers) = completion.network_request_headers.as_ref()
&& let Some(record) = state
@@ -1526,3 +1526,4 @@ mod network;
mod postmessage;
mod tls;
mod trusted_types_reporting;
mod xhr_failure;
@@ -0,0 +1,242 @@
use super::*;
const REJECTIONS: [(&str, &str, &str); 6] = [
(
"scheme",
"invalid-protocol://example.test/xhr",
"not supported",
),
("file", "file:///moli-policy-must-not-open", "not supported"),
("port", "http://example.test:25/xhr", "blocked bad port"),
(
"blocked",
"http://example.test/xhr",
"net::ERR_BLOCKED_BY_CLIENT",
),
(
"offline",
"http://example.test/xhr",
"Network emulation offline",
),
("csp", "http://example.test/xhr", "Content Security Policy"),
];
fn rejected_xhr_worker(case: &str, script: String) -> WorkerTestHandle {
let policy = WorkerNetworkPolicy {
blocked_url_patterns: if case == "blocked" {
vec!["http://example.test/*".to_owned()]
} else {
vec![]
},
network_offline: case == "offline",
..WorkerNetworkPolicy::default()
};
spawn_test_worker_with_options(
WorkerSpawnOptions::new(script, "http://example.test/worker.js".to_owned())
.with_network_policy(policy)
.with_content_security_policies(if case == "csp" {
vec!["connect-src 'none'".to_owned()]
} else {
vec![]
}),
)
}
async fn collect_rejected_xhr_output(
mut handle: WorkerTestHandle,
) -> (serde_json::Value, Vec<SubresourceNetworkRecord>) {
let mut posts = vec![];
let mut records = vec![];
while let Some(message) = timeout(TIMEOUT, handle.recv())
.await
.expect("worker XHR should finish and close")
{
match message {
WorkerToParentMessage::Post(payload) => {
posts.push(serde_json::from_str(&stringify_payload(&payload)).unwrap());
}
WorkerToParentMessage::SubresourceNetwork(record) => records.push(record),
other => panic!("unexpected worker XHR message: {other:?}"),
}
}
assert_eq!(posts.len(), 1, "worker should post exactly one result");
(posts.remove(0), records)
}
#[tokio::test]
async fn worker_xhr_early_failure_runs_after_send_and_microtasks() {
ensure_v8();
for (case, url, error_text) in REJECTIONS {
let script = r#"
const events = [];
const xhr = new XMLHttpRequest();
xhr.onreadystatechange = () => events.push("state:" + xhr.readyState);
for (const [target, prefix] of [[xhr, "xhr"], [xhr.upload, "upload"]]) {
for (const type of ["loadstart", "progress", "error", "abort", "load", "loadend"]) {
target.addEventListener(type, event => events.push(
`${prefix}:${type}:${event.loaded}:${event.total}:${event.lengthComputable}`));
}
}
xhr.onloadend = () => {
postMessage({ events, readyState: xhr.readyState, status: xhr.status });
close();
};
xhr.open("POST", URL);
xhr.send("payload");
events.push("returned:" + xhr.readyState);
queueMicrotask(() => {
events.push("microtask");
xhr.onerror = () => events.push("late-error-listener");
});
"#
.replace("URL", &serde_json::to_string(url).unwrap());
let (result, records) =
collect_rejected_xhr_output(rejected_xhr_worker(case, script)).await;
assert_eq!(
result,
serde_json::json!({
"events": [
"state:1", "xhr:loadstart:0:0:false", "upload:loadstart:0:7:true",
"returned:1", "microtask", "state:4", "upload:error:0:0:false",
"upload:loadend:0:0:false", "xhr:error:0:0:false", "late-error-listener",
"xhr:loadend:0:0:false"
],
"readyState": 4, "status": 0
}),
"{case}"
);
assert_eq!(records.len(), 1, "{case}: exactly one failure record");
assert_eq!(records[0].url().as_str(), url);
assert_eq!(records[0].resource_type(), SubresourceResourceType::Xhr);
assert!(
matches!(
records[0].outcome(),
SubresourceNetworkOutcome::Failure { error_text: actual } if actual.contains(error_text)
),
"{case}: {:?}",
records[0].outcome()
);
}
}
#[tokio::test]
async fn worker_xhr_early_failure_can_abort_or_reopen_before_delivery() {
ensure_v8();
for (case, url, _) in REJECTIONS {
for action in ["abort", "microtask-abort", "reopen"] {
let script = r#"
const events = [];
const xhr = new XMLHttpRequest();
for (const target of [xhr, xhr.upload]) {
for (const type of ["error", "abort", "load", "loadend"])
target.addEventListener(type, () => events.push(
`${target === xhr ? 'xhr' : 'upload'}:${type}`));
}
xhr.open("POST", URL);
xhr.send("payload");
const returnedState = xhr.readyState;
function cancel() {
if (ACTION === "reopen") xhr.open("GET", "data:text/plain,reopened");
else xhr.abort();
}
if (ACTION === "microtask-abort") queueMicrotask(cancel);
else cancel();
// A second XHR completes through the same queue, after the
// canceled failure, so the observation cannot race a timer.
const barrier = new XMLHttpRequest();
barrier.onloadend = () => {
postMessage({ returnedState, readyState: xhr.readyState, status: xhr.status, events });
close();
};
barrier.open("GET", URL);
barrier.send();
"#
.replace("URL", &serde_json::to_string(url).unwrap())
.replace("ACTION", &serde_json::to_string(action).unwrap());
let (result, records) =
collect_rejected_xhr_output(rejected_xhr_worker(case, script)).await;
let events = if action == "reopen" {
vec![]
} else {
vec!["upload:abort", "upload:loadend", "xhr:abort", "xhr:loadend"]
};
assert_eq!(
result,
serde_json::json!({
"returnedState": 1, "readyState": if action == "reopen" { 1 } else { 0 },
"status": 0, "events": events
}),
"{case}/{action}"
);
assert_eq!(
records.len(),
1,
"{case}/{action}: only the barrier may report failure"
);
}
}
}
#[tokio::test]
async fn worker_xhr_early_failure_does_not_overwrite_a_new_send() {
ensure_v8();
let script = r#"
const xhr = new XMLHttpRequest();
const events = [];
for (const type of ["error", "abort", "load", "loadend"])
xhr.addEventListener(type, () => events.push(type));
xhr.onloadend = () => {
postMessage({ events, status: xhr.status, text: xhr.responseText });
close();
};
xhr.open("GET", "invalid-protocol://example.test/old");
xhr.send();
xhr.open("GET", "data:text/plain,new-request");
xhr.send();
"#;
let (result, records) =
collect_rejected_xhr_output(rejected_xhr_worker("scheme", script.to_owned())).await;
assert_eq!(
result,
serde_json::json!({
"events": ["load", "loadend"], "status": 200, "text": "new-request"
})
);
assert!(
records.is_empty(),
"superseded request must not report a late failure"
);
}
#[tokio::test]
async fn worker_xhr_early_failure_still_throws_for_synchronous_requests() {
ensure_v8();
for (case, url, _) in REJECTIONS {
let script = r#"
const xhr = new XMLHttpRequest();
const events = [];
xhr.onreadystatechange = () => events.push("state:" + xhr.readyState);
for (const target of [xhr, xhr.upload])
for (const type of ["loadstart", "progress", "error", "abort", "load", "loadend"])
target.addEventListener(type, () => events.push(type));
xhr.open("POST", URL, false);
let exception = null;
try { xhr.send("payload"); }
catch (error) { exception = { name: error.name, domException: error instanceof DOMException }; }
postMessage({ exception, events, readyState: xhr.readyState, status: xhr.status });
close();
"#
.replace("URL", &serde_json::to_string(url).unwrap());
let (result, records) =
collect_rejected_xhr_output(rejected_xhr_worker(case, script)).await;
assert_eq!(
result,
serde_json::json!({
"exception": { "name": "NetworkError", "domException": true },
"events": ["state:1"], "readyState": 4, "status": 0
}),
"{case}"
);
assert_eq!(records.len(), 1, "{case}: exactly one synchronous failure");
}
}