mirror of
https://github.com/lexmount/moli.git
synced 2026-09-28 08:01:37 +00:00
fix(xhr): drive upload events from network transmission
Observe transmitted request bytes in the native transport and queue upload progress and completion in the owning page or worker realm. Cover empty bodies, resumed requests, redirects, and partial uploads without fabricating completion during send(). Mark uploads complete before final callbacks and reject stale events after abort or reopen. Preserve upload error events and update local smoke tests to collect the full final upload task. Validation: cargo fmt --all; full workspace clippy with warnings denied; cargo nextest run --no-fail-fast (18,148 passed, 13 skipped). Focused upstream WPT: 67/73 passing versus 65/73, 114/125 subtests versus 112/125, with no regressions.
This commit is contained in:
@@ -25,6 +25,7 @@ mod runtime;
|
||||
mod streaming_response;
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
mod upload;
|
||||
mod url_pattern;
|
||||
|
||||
#[cfg(any(test, feature = "test-support"))]
|
||||
@@ -85,4 +86,5 @@ pub use runtime::{
|
||||
FetchRuntimeIdentity, FetchRuntimeJoinReport, FetchRuntimeJoinStatus, FetchRuntimePanicReport,
|
||||
};
|
||||
pub use streaming_response::{StreamingHtmlResponse, StreamingRawResponse};
|
||||
pub use upload::{UploadEvent, UploadObserver};
|
||||
pub use url_pattern::url_pattern_matches;
|
||||
|
||||
@@ -54,6 +54,7 @@ pub struct Request {
|
||||
timeout_policy: RequestTimeoutPolicy,
|
||||
network_observation_recorder: Option<NetworkObservationRecorder>,
|
||||
browser_identity: Option<std::sync::Arc<moli_browser_profile::BrowserIdentityProfile>>,
|
||||
upload_observer: Option<crate::UploadObserver>,
|
||||
}
|
||||
|
||||
/// The lifetime of request headers supplied by an interception command.
|
||||
@@ -453,6 +454,7 @@ impl Request {
|
||||
timeout_policy: RequestTimeoutPolicy::default(),
|
||||
network_observation_recorder: None,
|
||||
browser_identity: None,
|
||||
upload_observer: None,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -486,6 +488,7 @@ impl Request {
|
||||
timeout_policy: RequestTimeoutPolicy::default(),
|
||||
network_observation_recorder: None,
|
||||
browser_identity: None,
|
||||
upload_observer: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -552,6 +555,7 @@ impl Request {
|
||||
timeout_policy: RequestTimeoutPolicy::default(),
|
||||
network_observation_recorder: None,
|
||||
browser_identity: None,
|
||||
upload_observer: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -715,6 +719,15 @@ impl Request {
|
||||
self.use_page_network_policy
|
||||
}
|
||||
|
||||
pub fn with_upload_observer(mut self, observer: crate::UploadObserver) -> Self {
|
||||
self.upload_observer = Some(observer);
|
||||
self
|
||||
}
|
||||
|
||||
pub fn upload_observer(&self) -> Option<&crate::UploadObserver> {
|
||||
self.upload_observer.as_ref()
|
||||
}
|
||||
|
||||
pub(crate) fn with_network_observation_recorder(
|
||||
mut self,
|
||||
recorder: NetworkObservationRecorder,
|
||||
|
||||
@@ -2669,6 +2669,7 @@ struct ActiveRawStreamingTransferContext {
|
||||
pub(crate) struct FetchTransferHandler {
|
||||
response: FetchResponseCollector,
|
||||
network_observation_recorder: Option<NetworkObservationRecorder>,
|
||||
upload_observer: Option<crate::UploadObserver>,
|
||||
proxy_connect_response_recorder: ProxyConnectResponseRecorder,
|
||||
}
|
||||
|
||||
@@ -2695,6 +2696,7 @@ impl FetchTransferHandler {
|
||||
Self {
|
||||
response,
|
||||
network_observation_recorder: None,
|
||||
upload_observer: None,
|
||||
proxy_connect_response_recorder: ProxyConnectResponseRecorder::default(),
|
||||
}
|
||||
}
|
||||
@@ -2745,7 +2747,9 @@ impl FetchTransferHandler {
|
||||
&mut self,
|
||||
network_observation_recorder: Option<NetworkObservationRecorder>,
|
||||
capture_proxy_connect_response: bool,
|
||||
upload_observer: Option<crate::UploadObserver>,
|
||||
) {
|
||||
self.upload_observer = upload_observer;
|
||||
self.network_observation_recorder = network_observation_recorder;
|
||||
self.proxy_connect_response_recorder
|
||||
.begin_transfer(capture_proxy_connect_response);
|
||||
@@ -2788,7 +2792,7 @@ impl Handler for FetchTransferHandler {
|
||||
}
|
||||
|
||||
fn progress(&mut self, dltotal: f64, dlnow: f64, ultotal: f64, ulnow: f64) -> bool {
|
||||
match &mut self.response {
|
||||
let keep_going = match &mut self.response {
|
||||
FetchResponseCollector::Buffered(collector) => {
|
||||
collector.progress(dltotal, dlnow, ultotal, ulnow)
|
||||
}
|
||||
@@ -2798,7 +2802,11 @@ impl Handler for FetchTransferHandler {
|
||||
FetchResponseCollector::StreamingRaw(collector) => {
|
||||
collector.progress(dltotal, dlnow, ultotal, ulnow)
|
||||
}
|
||||
};
|
||||
if keep_going && let Some(observer) = &self.upload_observer {
|
||||
observer.bytes_sent(ulnow as u64);
|
||||
}
|
||||
keep_going
|
||||
}
|
||||
|
||||
fn debug(&mut self, kind: InfoType, data: &[u8]) {
|
||||
@@ -2807,6 +2815,9 @@ impl Handler for FetchTransferHandler {
|
||||
let is_proxy_connect = self
|
||||
.proxy_connect_response_recorder
|
||||
.record_outgoing_header_block(data);
|
||||
if !is_proxy_connect && let Some(observer) = &self.upload_observer {
|
||||
observer.request_headers_sent();
|
||||
}
|
||||
if !is_proxy_connect
|
||||
&& let Some(recorder) = self.network_observation_recorder.as_ref()
|
||||
{
|
||||
@@ -2828,12 +2839,17 @@ fn configure_network_observation(
|
||||
capture_proxy_connect_response: bool,
|
||||
) -> Result<()> {
|
||||
let recorder = request.network_observation_recorder().cloned();
|
||||
let verbose = recorder.is_some() || capture_proxy_connect_response;
|
||||
let upload_observer = request
|
||||
.body
|
||||
.as_ref()
|
||||
.and(request.upload_observer())
|
||||
.cloned();
|
||||
let verbose = recorder.is_some() || capture_proxy_connect_response || upload_observer.is_some();
|
||||
if let Some(recorder) = recorder.as_ref() {
|
||||
recorder.set_current_request_cookie_report(request_cookie_report.cloned());
|
||||
}
|
||||
easy.get_mut()
|
||||
.begin_transfer(recorder, capture_proxy_connect_response);
|
||||
.begin_transfer(recorder, capture_proxy_connect_response, upload_observer);
|
||||
easy.verbose(verbose)
|
||||
.context("failed to configure curl network observation")
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ mod mixed_transport;
|
||||
mod request_security;
|
||||
mod support;
|
||||
mod tls_credentials;
|
||||
mod upload;
|
||||
mod websocket_transport;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
|
||||
@@ -0,0 +1,208 @@
|
||||
use super::*;
|
||||
use crate::{UploadEvent, UploadObserver};
|
||||
|
||||
async fn upload_response(
|
||||
client: &FetchClient,
|
||||
mut request: Request,
|
||||
transport: &str,
|
||||
cancel: FetchCancelHandle,
|
||||
) -> Result<u16> {
|
||||
match transport {
|
||||
"buffered" => {
|
||||
request.follow_redirects = false;
|
||||
Ok(client.fetch_with_cancel(request, cancel).await?.status)
|
||||
}
|
||||
"html" => {
|
||||
let mut response = client.fetch_html_stream(request).await?;
|
||||
while response.next_chunk().await.is_some() {}
|
||||
response.finish().await?;
|
||||
Ok(response.status)
|
||||
}
|
||||
"raw" => {
|
||||
let mut response = client.fetch_raw_stream_with_cancel(request, cancel).await?;
|
||||
while response.next_chunk().await.is_some() {}
|
||||
response.finish().await?;
|
||||
Ok(response.status)
|
||||
}
|
||||
_ => unreachable!(),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upload_observer_finishes_before_response_and_distinguishes_empty_from_absent() -> Result<()>
|
||||
{
|
||||
for transport in ["buffered", "html", "raw"] {
|
||||
for body in [None, Some(Vec::new()), Some(b"\x00body\xff".to_vec())] {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await?;
|
||||
let url = format!("http://{}/upload", listener.local_addr()?);
|
||||
let expected = body.clone().unwrap_or_default();
|
||||
let total = expected.len() as u64;
|
||||
let (received_tx, received_rx) = oneshot::channel();
|
||||
let (release_tx, release_rx) = oneshot::channel();
|
||||
let server = tokio::spawn(async move {
|
||||
let (mut stream, _) = listener.accept().await.unwrap();
|
||||
let head = read_http_request_head(&mut stream).await.unwrap();
|
||||
assert!(head.starts_with("POST /upload HTTP/1.1\r\n"));
|
||||
let mut bytes = vec![0; expected.len()];
|
||||
stream.read_exact(&mut bytes).await.unwrap();
|
||||
assert_eq!(bytes, expected);
|
||||
received_tx.send(()).unwrap();
|
||||
release_rx.await.unwrap();
|
||||
stream
|
||||
.write_all(
|
||||
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
let (event_tx, mut events) = mpsc::unbounded_channel();
|
||||
let observer = UploadObserver::new(total, move |event| {
|
||||
let _ = event_tx.send(event);
|
||||
});
|
||||
let request = Request::new_bytes("POST", &url, body.clone(), Vec::new())?
|
||||
.with_upload_observer(observer);
|
||||
let client =
|
||||
FetchClient::new(&FetchConfig::default(), new_shared_browser_cookie_store());
|
||||
let fetch = upload_response(&client, request, transport, FetchCancelHandle::new());
|
||||
let observe = async {
|
||||
received_rx.await.unwrap();
|
||||
if body.is_some() {
|
||||
loop {
|
||||
let event = tokio::time::timeout(Duration::from_secs(3), events.recv())
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
if let UploadEvent::Complete {
|
||||
loaded,
|
||||
total: observed_total,
|
||||
} = event
|
||||
{
|
||||
assert_eq!((loaded, observed_total), (total, total), "{transport}");
|
||||
break;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
assert!(events.try_recv().is_err());
|
||||
}
|
||||
release_tx.send(()).unwrap();
|
||||
};
|
||||
let (result, ()) = tokio::join!(fetch, observe);
|
||||
assert_eq!(result?, 200);
|
||||
assert!(
|
||||
events.try_recv().is_err(),
|
||||
"duplicate upload completion: {transport}"
|
||||
);
|
||||
server.await?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upload_observer_does_not_infer_completion_from_an_early_response() -> Result<()> {
|
||||
for transport in ["buffered", "html", "raw"] {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await?;
|
||||
let url = format!("http://{}/reject", listener.local_addr()?);
|
||||
let server = tokio::spawn(async move {
|
||||
let (mut stream, _) = listener.accept().await.unwrap();
|
||||
let head = read_http_request_head(&mut stream).await.unwrap();
|
||||
assert!(
|
||||
head.to_ascii_lowercase()
|
||||
.contains("expect: 100-continue\r\n")
|
||||
);
|
||||
stream.write_all(b"HTTP/1.1 413 Content Too Large\r\nContent-Length: 0\r\nConnection: close\r\n\r\n").await.unwrap();
|
||||
});
|
||||
let (event_tx, mut events) = mpsc::unbounded_channel();
|
||||
let body = vec![b'x'; 1024 * 1024];
|
||||
let observer = UploadObserver::new(body.len() as u64, move |event| {
|
||||
let _ = event_tx.send(event);
|
||||
});
|
||||
let request = Request::new_bytes(
|
||||
"POST",
|
||||
&url,
|
||||
Some(body),
|
||||
vec![("Expect".into(), "100-continue".into())],
|
||||
)?
|
||||
.with_upload_observer(observer);
|
||||
let client = FetchClient::new(&FetchConfig::default(), new_shared_browser_cookie_store());
|
||||
assert_eq!(
|
||||
upload_response(&client, request, transport, FetchCancelHandle::new()).await?,
|
||||
413
|
||||
);
|
||||
assert!(
|
||||
events.try_recv().is_err(),
|
||||
"unsent body reported as uploaded: {transport}"
|
||||
);
|
||||
server.await?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upload_observer_reports_partial_bytes_and_stops_after_cancellation() -> Result<()> {
|
||||
for transport in ["buffered", "raw"] {
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await?;
|
||||
let url = format!("http://{}/partial", listener.local_addr()?);
|
||||
let server = tokio::spawn(async move {
|
||||
let (mut stream, _) = listener.accept().await.unwrap();
|
||||
let _ = read_http_request_head(&mut stream).await.unwrap();
|
||||
// Hold the receive window closed long enough to expose a partial
|
||||
// upload, then let transmission advance to the next progress task.
|
||||
tokio::time::sleep(Duration::from_millis(150)).await;
|
||||
let mut received = 0;
|
||||
let mut chunk = [0; 65536];
|
||||
while let Ok(count) = stream.read(&mut chunk).await {
|
||||
if count == 0 {
|
||||
break;
|
||||
}
|
||||
received += count;
|
||||
}
|
||||
received
|
||||
});
|
||||
let cancel = FetchCancelHandle::new();
|
||||
let cancel_upload = cancel.clone();
|
||||
let (event_tx, mut events) = mpsc::unbounded_channel();
|
||||
let total = 16 * 1024 * 1024;
|
||||
let observer = UploadObserver::new(total, move |event| {
|
||||
if matches!(event, UploadEvent::Progress { loaded, total } if loaded > 0 && loaded < total)
|
||||
{
|
||||
cancel_upload.cancel();
|
||||
}
|
||||
let _ = event_tx.send(event);
|
||||
});
|
||||
let request = Request::new_bytes(
|
||||
"POST",
|
||||
&url,
|
||||
Some(vec![b'x'; total as usize]),
|
||||
vec![("Expect".into(), String::new())],
|
||||
)?
|
||||
.with_upload_observer(observer);
|
||||
let client = FetchClient::new(&FetchConfig::default(), new_shared_browser_cookie_store());
|
||||
let result = tokio::time::timeout(
|
||||
Duration::from_secs(5),
|
||||
upload_response(&client, request, transport, cancel),
|
||||
)
|
||||
.await?;
|
||||
assert!(
|
||||
result.is_err(),
|
||||
"partial cancellation must abort {transport}"
|
||||
);
|
||||
let mut partial = false;
|
||||
while let Ok(event) = events.try_recv() {
|
||||
match event {
|
||||
UploadEvent::Progress {
|
||||
loaded,
|
||||
total: observed_total,
|
||||
} => {
|
||||
assert_eq!(observed_total, total);
|
||||
assert!(loaded > 0 && loaded < total);
|
||||
partial = true;
|
||||
}
|
||||
UploadEvent::Complete { .. } => panic!("cancelled upload completed: {transport}"),
|
||||
}
|
||||
}
|
||||
assert!(partial, "missing partial upload: {transport}");
|
||||
assert!(tokio::time::timeout(Duration::from_secs(3), server).await?? < total as usize);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
use parking_lot::Mutex;
|
||||
use std::{
|
||||
fmt,
|
||||
sync::Arc,
|
||||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
/// Facts about transmission of a request body, delivered on the fetch runtime
|
||||
/// thread. Consumers must enqueue work rather than enter a script engine here.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub enum UploadEvent {
|
||||
Progress { loaded: u64, total: u64 },
|
||||
Complete { loaded: u64, total: u64 },
|
||||
}
|
||||
|
||||
/// A single body upload, shared by redirect and authentication retries. Byte
|
||||
/// counts never go backwards and completion is reported only once.
|
||||
#[derive(Clone)]
|
||||
pub struct UploadObserver(Arc<UploadObserverInner>);
|
||||
|
||||
struct UploadObserverInner {
|
||||
total: u64,
|
||||
callback: Box<dyn Fn(UploadEvent) + Send + Sync>,
|
||||
state: Mutex<UploadState>,
|
||||
}
|
||||
|
||||
struct UploadState {
|
||||
loaded: u64,
|
||||
last_progress: Instant,
|
||||
complete: bool,
|
||||
}
|
||||
|
||||
impl fmt::Debug for UploadObserver {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
f.debug_struct("UploadObserver")
|
||||
.field("total", &self.0.total)
|
||||
.finish_non_exhaustive()
|
||||
}
|
||||
}
|
||||
|
||||
impl UploadObserver {
|
||||
pub fn new(total: u64, callback: impl Fn(UploadEvent) + Send + Sync + 'static) -> Self {
|
||||
Self(Arc::new(UploadObserverInner {
|
||||
total,
|
||||
callback: Box::new(callback),
|
||||
state: Mutex::new(UploadState {
|
||||
loaded: 0,
|
||||
last_progress: Instant::now(),
|
||||
complete: false,
|
||||
}),
|
||||
}))
|
||||
}
|
||||
|
||||
pub(crate) fn request_headers_sent(&self) {
|
||||
// There is no positive curl upload count for a present, empty body.
|
||||
// Its request headers, unlike a response, establish that it was sent.
|
||||
if self.0.total == 0 {
|
||||
self.observe(0);
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn bytes_sent(&self, loaded: u64) {
|
||||
if loaded > 0 {
|
||||
self.observe(loaded);
|
||||
}
|
||||
}
|
||||
|
||||
fn observe(&self, loaded: u64) {
|
||||
let total = self.0.total;
|
||||
let (progress, completion) = {
|
||||
let mut state = self.0.state.lock();
|
||||
if state.complete {
|
||||
return;
|
||||
}
|
||||
let loaded = loaded.min(total);
|
||||
let progress = if loaded > state.loaded {
|
||||
state.loaded = loaded;
|
||||
let now = Instant::now();
|
||||
if now.duration_since(state.last_progress) >= Duration::from_millis(50) {
|
||||
state.last_progress = now;
|
||||
Some(UploadEvent::Progress { loaded, total })
|
||||
} else {
|
||||
None
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let completion = (loaded == total).then(|| {
|
||||
state.complete = true;
|
||||
UploadEvent::Complete { loaded, total }
|
||||
});
|
||||
(progress, completion)
|
||||
};
|
||||
// Keep the byte-progress and end-of-body tasks distinct. In particular,
|
||||
// a consumer can abort while processing the former before the latter.
|
||||
for event in progress.into_iter().chain(completion) {
|
||||
(self.0.callback)(event);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1592,6 +1592,9 @@ impl JsContextHost {
|
||||
use crate::types::AsyncSubresourceFetchEventTarget;
|
||||
|
||||
match target {
|
||||
AsyncSubresourceFetchEventTarget::Upload { internal_id } => {
|
||||
self.pending_xhr_for_upload(internal_id).is_some()
|
||||
}
|
||||
AsyncSubresourceFetchEventTarget::Completion { internal_id } => {
|
||||
self.pending_subresource_fetches.contains_key(&internal_id)
|
||||
|| self.running_subresource_fetches.contains_key(&internal_id)
|
||||
@@ -1615,6 +1618,58 @@ impl JsContextHost {
|
||||
}
|
||||
}
|
||||
|
||||
fn pending_xhr_for_upload(&self, internal_id: u64) -> Option<&PendingSubresourceFetchState> {
|
||||
let pending = self
|
||||
.pending_subresource_fetches
|
||||
.get(&internal_id)
|
||||
.or_else(|| {
|
||||
self.running_subresource_fetches
|
||||
.get(&internal_id)
|
||||
.map(|state| &state.pending)
|
||||
})
|
||||
.or_else(|| {
|
||||
self.streaming_subresource_fetches
|
||||
.get(&internal_id)
|
||||
.map(|state| &state.pending)
|
||||
})?;
|
||||
pending.continuation.is_window_xhr().then_some(pending)
|
||||
}
|
||||
|
||||
pub(crate) fn xhr_upload_delivery<'s>(
|
||||
&self,
|
||||
scope: &mut v8::PinScope<'s, '_, ()>,
|
||||
internal_id: u64,
|
||||
) -> Option<(
|
||||
v8::Local<'s, v8::Object>,
|
||||
crate::types::PendingSubresourceExecutionContext,
|
||||
)> {
|
||||
let pending = self.pending_xhr_for_upload(internal_id)?;
|
||||
if let Some(target) = pending.execution_context.window_request_target()
|
||||
&& !self
|
||||
.window_execution_context_owner_is_current(target.owner(), target.dispatch_scope())
|
||||
{
|
||||
return None;
|
||||
}
|
||||
let PendingSubresourceContinuation::Xhr { xhr, .. } = &pending.continuation else {
|
||||
return None;
|
||||
};
|
||||
use crate::types::PendingSubresourceExecutionContext;
|
||||
let execution = match &pending.execution_context {
|
||||
PendingSubresourceExecutionContext::Window(binding) => {
|
||||
PendingSubresourceExecutionContext::Window(binding.clone())
|
||||
}
|
||||
PendingSubresourceExecutionContext::Adapter {
|
||||
dispatch_scope,
|
||||
context,
|
||||
} => PendingSubresourceExecutionContext::Adapter {
|
||||
dispatch_scope: *dispatch_scope,
|
||||
context: context.clone(),
|
||||
},
|
||||
_ => return None,
|
||||
};
|
||||
Some((v8::Local::new(scope, xhr), execution))
|
||||
}
|
||||
|
||||
pub(crate) fn take_pending_subresource_fetch(
|
||||
&mut self,
|
||||
internal_id: u64,
|
||||
|
||||
@@ -36,7 +36,7 @@ pub(crate) use self::async_fetch::{
|
||||
fetch_browser_subresource_with_preflight_and_network_metadata,
|
||||
fetch_browser_subresource_with_preflight_headers,
|
||||
fetch_browser_subresource_with_preflight_headers_and_network_metadata, fetch_cors_script_text,
|
||||
spawn_async_subresource_fetch,
|
||||
observe_async_xhr_upload, spawn_async_subresource_fetch,
|
||||
};
|
||||
pub(crate) use self::beacon::{navigator_send_beacon_callback, send_link_audit_ping};
|
||||
pub(super) use self::bindings::install_window_network_bindings;
|
||||
@@ -187,7 +187,7 @@ pub(crate) use self::xhr::{
|
||||
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,
|
||||
capture_xhr_upload_listener_flag, dispatch_xhr_loadstart, dispatch_xhr_upload_complete,
|
||||
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,
|
||||
|
||||
@@ -553,6 +553,28 @@ async fn run_cors_preflight_if_needed(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) fn observe_async_xhr_upload(
|
||||
request: Request,
|
||||
completion_tx: &RendererResourceCompletionSender,
|
||||
internal_id: u64,
|
||||
) -> Request {
|
||||
let mut request = request;
|
||||
if request.browser_request_metadata() == Some(BrowserRequestMetadata::Xhr)
|
||||
&& request.upload_observer().is_none()
|
||||
&& let Some(body) = request.body.as_ref()
|
||||
{
|
||||
let upload_tx = completion_tx.clone();
|
||||
let observer = moli_fetch::UploadObserver::new(body.len() as u64, move |event| {
|
||||
let _ = upload_tx.send_async_subresource_event(AsyncSubresourceFetchEvent::Upload {
|
||||
internal_id,
|
||||
event,
|
||||
});
|
||||
});
|
||||
request = request.with_upload_observer(observer);
|
||||
}
|
||||
request
|
||||
}
|
||||
|
||||
pub(crate) fn spawn_async_subresource_fetch(
|
||||
task_runner: crate::network::RendererResourceTaskRunner,
|
||||
completion_tx: RendererResourceCompletionSender,
|
||||
@@ -572,6 +594,7 @@ pub(crate) fn spawn_async_subresource_fetch(
|
||||
Some(BrowserRequestMetadata::Image)
|
||||
)
|
||||
.then(|| loader.parkable_image_manager(&task_runner));
|
||||
let request = observe_async_xhr_upload(request, &completion_tx, internal_id);
|
||||
task_runner.spawn(async move {
|
||||
let preflight_observer =
|
||||
CorsPreflightNetworkObserver::new(completion_tx.clone(), network_context);
|
||||
@@ -1736,6 +1759,7 @@ mod tests {
|
||||
head_sent_rx
|
||||
.await
|
||||
.expect("server should publish the response head and first bytes");
|
||||
expect_upload_completion(&mut queue, 42).await?;
|
||||
let body_source_id = match tokio::time::timeout(
|
||||
Duration::from_secs(2),
|
||||
next_async_subresource_event(&mut queue),
|
||||
@@ -1899,6 +1923,7 @@ mod tests {
|
||||
head_sent_rx
|
||||
.await
|
||||
.expect("server should publish the redirected final response head");
|
||||
expect_upload_completion(&mut queue, 43).await?;
|
||||
match next_async_subresource_event(&mut queue).await? {
|
||||
AsyncSubresourceFetchEvent::ObservedNetworkRecord(record) => {
|
||||
assert_eq!(record.document_url(), &document_url);
|
||||
@@ -1948,4 +1973,27 @@ mod tests {
|
||||
server.await?;
|
||||
Ok(())
|
||||
}
|
||||
async fn expect_upload_completion(
|
||||
queue: &mut RendererResourceCompletionTestHarness,
|
||||
expected_id: u64,
|
||||
) -> anyhow::Result<()> {
|
||||
loop {
|
||||
match next_async_subresource_event(queue).await? {
|
||||
AsyncSubresourceFetchEvent::Upload { internal_id, event } => {
|
||||
assert_eq!(internal_id, expected_id);
|
||||
match event {
|
||||
moli_fetch::UploadEvent::Progress { loaded, total } => {
|
||||
assert_eq!(total, 7);
|
||||
assert!(loaded > 0 && loaded <= total);
|
||||
}
|
||||
moli_fetch::UploadEvent::Complete { loaded, total } => {
|
||||
assert_eq!((loaded, total), (7, 7));
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
}
|
||||
other => anyhow::bail!("expected upload before response headers: {other:?}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ mod header_surface;
|
||||
mod instance_state;
|
||||
mod response_type;
|
||||
mod send;
|
||||
mod upload;
|
||||
|
||||
use super::*;
|
||||
use crate::context_bootstrap::{
|
||||
@@ -36,6 +37,7 @@ pub(crate) const XHR_ABORTED_SLOT: &str = "__lmXhrAborted";
|
||||
pub(crate) const XHR_ASYNC_SLOT: &str = "__lmXhrAsync";
|
||||
pub(crate) const XHR_SEND_FLAG_SLOT: &str = "__lmXhrSendFlag";
|
||||
pub(crate) const XHR_UPLOAD_IN_PROGRESS_SLOT: &str = "__lmXhrUploadInProgress";
|
||||
const XHR_UPLOAD_LOADED_SLOT: &str = "__lmXhrUploadLoaded";
|
||||
const XHR_UPLOAD_LISTENER_SLOT: &str = "__lmXhrUploadListener";
|
||||
const XHR_PENDING_KIND_SLOT: &str = "__lmXhrPendingKind";
|
||||
const XHR_PENDING_STATUS_SLOT: &str = "__lmXhrPendingStatus";
|
||||
@@ -78,8 +80,9 @@ use self::response_type::{XmlHttpRequestResponseType, xhr_response_type};
|
||||
pub(crate) use self::send::prepare_xhr_send_body;
|
||||
pub(crate) use self::send::{
|
||||
PreparedXhrSendBody, capture_xhr_upload_listener_flag, dispatch_xhr_loadstart,
|
||||
dispatch_xhr_upload_complete, convert_xhr_send_body_from_args, xhr_author_request_headers,
|
||||
convert_xhr_send_body_from_args, xhr_author_request_headers,
|
||||
};
|
||||
pub(crate) use self::upload::apply_xhr_upload_event;
|
||||
|
||||
pub(crate) fn install_progress_event_template_bindings<'s>(
|
||||
scope: &mut v8::PinScope<'s, '_, ()>,
|
||||
|
||||
@@ -13,7 +13,7 @@ pub(crate) fn finish_xhr_abort(scope: &mut v8::PinScope<'_, '_>, xhr: v8::Local<
|
||||
set_xhr_state_number(scope, xhr, XHR_READY_STATE_SLOT, 4.0);
|
||||
super::reset_xhr_response_for_request_error(scope, xhr);
|
||||
super::super::events::xhr_fire_readystatechange(scope, xhr, 4);
|
||||
super::super::send::dispatch_xhr_upload_abort_if_in_progress(scope, xhr);
|
||||
super::super::upload::dispatch_xhr_upload_error_if_in_progress(scope, xhr, "abort");
|
||||
xhr_dispatch_progress_event(scope, xhr, "abort", 0.0, 0.0);
|
||||
xhr_dispatch_progress_event(scope, xhr, "loadend", 0.0, 0.0);
|
||||
}
|
||||
@@ -34,7 +34,7 @@ pub(crate) fn apply_xhr_abort(scope: &mut v8::PinScope<'_, '_>, xhr: v8::Local<'
|
||||
set_xhr_state_number(scope, xhr, XHR_ACTIVE_INTERNAL_ID_SLOT, 0.0);
|
||||
set_xhr_state_number(scope, xhr, XHR_READY_STATE_SLOT, 0.0);
|
||||
super::reset_xhr_response_for_request_error(scope, xhr);
|
||||
super::super::send::dispatch_xhr_upload_abort_if_in_progress(scope, xhr);
|
||||
super::super::upload::dispatch_xhr_upload_error_if_in_progress(scope, xhr, "abort");
|
||||
xhr_dispatch_progress_event(scope, xhr, "abort", 0.0, 0.0);
|
||||
xhr_dispatch_progress_event(scope, xhr, "loadend", 0.0, 0.0);
|
||||
}
|
||||
|
||||
@@ -21,6 +21,7 @@ pub(crate) fn apply_xhr_failure(scope: &mut v8::PinScope<'_, '_>, xhr: v8::Local
|
||||
if xhr_is_aborted(scope, xhr) {
|
||||
return;
|
||||
}
|
||||
super::super::upload::dispatch_xhr_upload_error_if_in_progress(scope, xhr, "error");
|
||||
xhr_dispatch_progress_event(scope, xhr, "error", 0.0, 0.0);
|
||||
if scope.is_execution_terminating() {
|
||||
return;
|
||||
|
||||
@@ -203,6 +203,7 @@ pub(crate) fn apply_xhr_timeout(scope: &mut v8::PinScope<'_, '_>, xhr: v8::Local
|
||||
if xhr_is_aborted(scope, xhr) {
|
||||
return;
|
||||
}
|
||||
super::super::upload::dispatch_xhr_upload_error_if_in_progress(scope, xhr, "timeout");
|
||||
xhr_dispatch_progress_event(scope, xhr, "timeout", 0.0, 0.0);
|
||||
xhr_dispatch_progress_event(scope, xhr, "loadend", 0.0, 0.0);
|
||||
}
|
||||
|
||||
@@ -5,6 +5,6 @@ mod state;
|
||||
pub(super) use self::dispatch::xhr_fire_readystatechange;
|
||||
pub(crate) use self::dispatch::{
|
||||
xhr_dispatch_progress_event, xhr_dispatch_progress_event_with_length_computable,
|
||||
xhr_dispatch_upload_progress_event,
|
||||
xhr_dispatch_upload_progress_event, xhr_dispatch_upload_progress_events,
|
||||
};
|
||||
pub(super) use self::state::{xhr_is_aborted, xhr_is_async};
|
||||
|
||||
@@ -139,6 +139,16 @@ pub(crate) fn xhr_dispatch_upload_progress_event(
|
||||
event_type: &str,
|
||||
loaded: f64,
|
||||
total: f64,
|
||||
) {
|
||||
xhr_dispatch_upload_progress_events(scope, xhr, &[event_type], loaded, total);
|
||||
}
|
||||
|
||||
pub(crate) fn xhr_dispatch_upload_progress_events(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
xhr: v8::Local<'_, v8::Object>,
|
||||
event_types: &[&str],
|
||||
loaded: f64,
|
||||
total: f64,
|
||||
) {
|
||||
if !xhr_state_bool_property(scope, xhr, XHR_UPLOAD_LISTENER_SLOT).unwrap_or(false) {
|
||||
return;
|
||||
@@ -146,13 +156,26 @@ pub(crate) fn xhr_dispatch_upload_progress_event(
|
||||
let Some(upload) = xhr_upload_object(scope, xhr) else {
|
||||
return;
|
||||
};
|
||||
if !xhr_has_event_observers(scope, upload, event_type) {
|
||||
return;
|
||||
// The listener flag is tested once on entry to the event algorithm, even
|
||||
// when an earlier event reopens the XHR and resets its request state.
|
||||
for event_type in event_types {
|
||||
if scope.is_execution_terminating() {
|
||||
return;
|
||||
}
|
||||
if !xhr_has_event_observers(scope, upload, event_type) {
|
||||
continue;
|
||||
}
|
||||
let event = super::progress::make_progress_event(
|
||||
scope,
|
||||
event_type,
|
||||
upload,
|
||||
total > 0.0,
|
||||
loaded,
|
||||
total,
|
||||
);
|
||||
let handler_name = format!("on{event_type}");
|
||||
xhr_invoke_handler(scope, upload, &handler_name, event);
|
||||
}
|
||||
let event =
|
||||
super::progress::make_progress_event(scope, event_type, upload, total > 0.0, loaded, total);
|
||||
let handler_name = format!("on{event_type}");
|
||||
xhr_invoke_handler(scope, upload, &handler_name, event);
|
||||
}
|
||||
|
||||
fn local_object_in_scope<'s>(
|
||||
|
||||
@@ -81,6 +81,8 @@ struct XmlHttpRequestStateDeclaration {
|
||||
send_flag: (),
|
||||
#[webapi(slot = XHR_UPLOAD_IN_PROGRESS_SLOT, init = false)]
|
||||
upload_in_progress: (),
|
||||
#[webapi(slot = XHR_UPLOAD_LOADED_SLOT, init = 0.0)]
|
||||
upload_loaded: (),
|
||||
#[webapi(slot = XHR_UPLOAD_LISTENER_SLOT, init = false)]
|
||||
upload_listener: (),
|
||||
#[webapi(slot = XHR_ACTIVE_INTERNAL_ID_SLOT, init = 0)]
|
||||
|
||||
@@ -108,7 +108,6 @@ pub(super) fn xhr_send_callback<'s>(
|
||||
if !dispatch_xhr_loadstart(scope, xhr, prepared.send_body.as_deref()) {
|
||||
return;
|
||||
}
|
||||
dispatch_xhr_upload_complete(scope, xhr, prepared.send_body.as_deref());
|
||||
if xhr_is_aborted(scope, xhr) || xhr_open_generation_changed(scope, xhr, open_generation) {
|
||||
return;
|
||||
}
|
||||
@@ -289,6 +288,7 @@ pub(crate) fn dispatch_xhr_loadstart(
|
||||
xhr_state_number_property(scope, xhr, XHR_OPEN_GENERATION_SLOT).unwrap_or(0.0);
|
||||
// The XHR loadstart listener can abort before upload.loadstart runs.
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, send_body.is_some());
|
||||
set_xhr_state_number(scope, xhr, XHR_UPLOAD_LOADED_SLOT, 0.0);
|
||||
xhr_dispatch_progress_event(scope, xhr, "loadstart", 0.0, 0.0);
|
||||
if xhr_is_aborted(scope, xhr) || xhr_open_generation_changed(scope, xhr, open_generation) {
|
||||
return false;
|
||||
@@ -299,37 +299,6 @@ pub(crate) fn dispatch_xhr_loadstart(
|
||||
!xhr_is_aborted(scope, xhr) && !xhr_open_generation_changed(scope, xhr, open_generation)
|
||||
}
|
||||
|
||||
pub(crate) fn dispatch_xhr_upload_complete(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
xhr: v8::Local<'_, v8::Object>,
|
||||
send_body: Option<&[u8]>,
|
||||
) {
|
||||
let Some(send_body) = send_body else {
|
||||
return;
|
||||
};
|
||||
let total = send_body.len() as f64;
|
||||
for event_type in ["progress", "load", "loadend"] {
|
||||
if xhr_is_aborted(scope, xhr) {
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, false);
|
||||
return;
|
||||
}
|
||||
xhr_dispatch_upload_progress_event(scope, xhr, event_type, total, total);
|
||||
}
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, false);
|
||||
}
|
||||
|
||||
pub(crate) fn dispatch_xhr_upload_abort_if_in_progress(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
xhr: v8::Local<'_, v8::Object>,
|
||||
) {
|
||||
if !xhr_state_bool_property(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT).unwrap_or(false) {
|
||||
return;
|
||||
}
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, false);
|
||||
xhr_dispatch_upload_progress_event(scope, xhr, "abort", 0.0, 0.0);
|
||||
xhr_dispatch_upload_progress_event(scope, xhr, "loadend", 0.0, 0.0);
|
||||
}
|
||||
|
||||
fn send_synchronous_network_xhr(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
host: &mut JsContextHost,
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
use super::events::xhr_dispatch_upload_progress_events;
|
||||
use super::*;
|
||||
|
||||
pub(crate) fn apply_xhr_upload_event(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
xhr: v8::Local<'_, v8::Object>,
|
||||
internal_id: u64,
|
||||
event: moli_fetch::UploadEvent,
|
||||
) -> bool {
|
||||
let is_current = |scope: &mut v8::PinScope<'_, '_>| {
|
||||
xhr_state_number_property(scope, xhr, XHR_ACTIVE_INTERNAL_ID_SLOT)
|
||||
== Some(internal_id as f64)
|
||||
&& xhr_state_bool_property(scope, xhr, XHR_SEND_FLAG_SLOT).unwrap_or(false)
|
||||
&& !events::xhr_is_aborted(scope, xhr)
|
||||
};
|
||||
if !is_current(scope) {
|
||||
return false;
|
||||
}
|
||||
if !xhr_state_bool_property(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT).unwrap_or(false) {
|
||||
return true;
|
||||
}
|
||||
match event {
|
||||
moli_fetch::UploadEvent::Progress { loaded, total } => {
|
||||
let previous =
|
||||
xhr_state_number_property(scope, xhr, XHR_UPLOAD_LOADED_SLOT).unwrap_or(0.0);
|
||||
if loaded as f64 > previous {
|
||||
set_xhr_state_number(scope, xhr, XHR_UPLOAD_LOADED_SLOT, loaded as f64);
|
||||
xhr_dispatch_upload_progress_events(
|
||||
scope,
|
||||
xhr,
|
||||
&["progress"],
|
||||
loaded as f64,
|
||||
total as f64,
|
||||
);
|
||||
}
|
||||
}
|
||||
moli_fetch::UploadEvent::Complete { loaded, total } => {
|
||||
// The end-of-body algorithm marks the upload complete before
|
||||
// firing events. No state writes follow the reentrant callbacks.
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, false);
|
||||
xhr_dispatch_upload_progress_events(
|
||||
scope,
|
||||
xhr,
|
||||
&["progress", "load", "loadend"],
|
||||
loaded as f64,
|
||||
total as f64,
|
||||
);
|
||||
}
|
||||
}
|
||||
is_current(scope)
|
||||
}
|
||||
|
||||
pub(super) fn dispatch_xhr_upload_error_if_in_progress(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
xhr: v8::Local<'_, v8::Object>,
|
||||
event: &str,
|
||||
) {
|
||||
if !xhr_state_bool_property(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT).unwrap_or(false) {
|
||||
return;
|
||||
}
|
||||
set_xhr_state_bool(scope, xhr, XHR_UPLOAD_IN_PROGRESS_SLOT, false);
|
||||
xhr_dispatch_upload_progress_events(scope, xhr, &[event, "loadend"], 0.0, 0.0);
|
||||
}
|
||||
@@ -4173,6 +4173,8 @@ impl ScriptVm {
|
||||
host.begin_active_subresource_request();
|
||||
host.record_running_subresource_fetch(state);
|
||||
}
|
||||
let request =
|
||||
crate::network_host::observe_async_xhr_upload(request, &completion_tx, internal_id);
|
||||
task_runner.spawn(async move {
|
||||
let result = if let Some(result) = local_response {
|
||||
result.map(crate::protocol_types::NavigationResponse::from)
|
||||
@@ -4511,6 +4513,9 @@ impl ScriptVm {
|
||||
let trace_fields = async_subresource_trace_fields_for_event(&event);
|
||||
trace_async_subresource_stage("async_subresource_event_start", trace_fields, trace_started);
|
||||
let result = match event {
|
||||
AsyncSubresourceFetchEvent::Upload { internal_id, event } => {
|
||||
self.apply_async_xhr_upload_event(internal_id, event)
|
||||
}
|
||||
AsyncSubresourceFetchEvent::Completion(completion) => {
|
||||
self.complete_async_subresource_fetch_body(*completion)
|
||||
}
|
||||
@@ -4539,6 +4544,52 @@ impl ScriptVm {
|
||||
result
|
||||
}
|
||||
|
||||
fn apply_async_xhr_upload_event(
|
||||
&mut self,
|
||||
internal_id: u64,
|
||||
event: moli_fetch::UploadEvent,
|
||||
) -> Result<AsyncSubresourceFetchBodyActivity> {
|
||||
let context_host = self._context_host.clone();
|
||||
self.renderer_document_isolate
|
||||
.with_entered_renderer_document_isolate(|isolate| {
|
||||
let scope = pin!(v8::HandleScope::new(isolate));
|
||||
let scope = &mut scope.init();
|
||||
let Some((xhr, execution)) = context_host
|
||||
.borrow()
|
||||
.xhr_upload_delivery(scope, internal_id)
|
||||
else {
|
||||
return Ok(AsyncSubresourceFetchBodyActivity::NoWindowRealmEntered);
|
||||
};
|
||||
let Some(context) = execution.context_global() else {
|
||||
return Ok(AsyncSubresourceFetchBodyActivity::NoWindowRealmEntered);
|
||||
};
|
||||
let context = v8::Local::new(scope, context);
|
||||
let scope = &mut v8::ContextScope::new(scope, context);
|
||||
if execution.window_realm_binding().is_some_and(|binding| {
|
||||
crate::native_bridge::current_runtime_observable_context_token(scope)
|
||||
!= Some(binding.realm_token())
|
||||
|| !binding.is_current(&context_host.borrow())
|
||||
}) {
|
||||
let _ = context_host
|
||||
.borrow_mut()
|
||||
.abort_subresource_fetch(internal_id);
|
||||
return Ok(AsyncSubresourceFetchBodyActivity::NoWindowRealmEntered);
|
||||
}
|
||||
let dispatch_scope = execution.dispatch_scope();
|
||||
let previous =
|
||||
enter_subresource_owner_async_scope(&context_host, scope, dispatch_scope);
|
||||
let current =
|
||||
crate::network_host::apply_xhr_upload_event(scope, xhr, internal_id, event);
|
||||
defer_subresource_owner_async_scope(&context_host, scope, dispatch_scope, previous);
|
||||
if !current {
|
||||
let _ = context_host
|
||||
.borrow_mut()
|
||||
.abort_subresource_fetch(internal_id);
|
||||
}
|
||||
Ok(AsyncSubresourceFetchBodyActivity::WindowRealmEntered)
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn async_subresource_fetch_event_target_is_current(
|
||||
&self,
|
||||
target: crate::types::AsyncSubresourceFetchEventTarget,
|
||||
@@ -6472,6 +6523,11 @@ fn async_subresource_trace_fields_for_event(
|
||||
event: &AsyncSubresourceFetchEvent,
|
||||
) -> AsyncSubresourceTraceFields {
|
||||
match event {
|
||||
AsyncSubresourceFetchEvent::Upload { internal_id, .. } => AsyncSubresourceTraceFields {
|
||||
event_kind: Some("upload"),
|
||||
internal_id: Some(*internal_id),
|
||||
..AsyncSubresourceTraceFields::default()
|
||||
},
|
||||
AsyncSubresourceFetchEvent::Completion(completion) => AsyncSubresourceTraceFields {
|
||||
event_kind: Some("completion"),
|
||||
internal_id: Some(completion.internal_id),
|
||||
|
||||
@@ -16,6 +16,7 @@ mod send_body;
|
||||
mod shadow_dom;
|
||||
mod style_invalidation;
|
||||
mod upload_preflight;
|
||||
mod upload_transport;
|
||||
mod xhr;
|
||||
|
||||
mod response_type;
|
||||
|
||||
@@ -0,0 +1,312 @@
|
||||
use super::*;
|
||||
use std::time::Duration;
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
async fn read_upload_request(socket: &mut tokio::net::TcpStream) -> (String, Vec<u8>) {
|
||||
let mut head = Vec::new();
|
||||
while !head.ends_with(b"\r\n\r\n") {
|
||||
head.push(socket.read_u8().await.unwrap());
|
||||
}
|
||||
let head = String::from_utf8(head).unwrap();
|
||||
let length = head
|
||||
.lines()
|
||||
.find_map(|line| {
|
||||
let (name, value) = line.split_once(':')?;
|
||||
name.eq_ignore_ascii_case("content-length")
|
||||
.then(|| value.trim().parse().unwrap())
|
||||
})
|
||||
.unwrap_or(0);
|
||||
let mut body = vec![0; length];
|
||||
socket.read_exact(&mut body).await.unwrap();
|
||||
(head, body)
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn xhr_upload_completes_before_response_and_preserves_reentrant_event_steps() {
|
||||
for mode in [
|
||||
"normal",
|
||||
"empty",
|
||||
"absent",
|
||||
"intercept",
|
||||
"redirect",
|
||||
"abort-complete",
|
||||
"reopen-complete",
|
||||
] {
|
||||
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 (head, body) = read_upload_request(&mut socket).await;
|
||||
assert!(head.starts_with("POST /upload HTTP/1.1\r\n"));
|
||||
let expected = if matches!(mode, "empty" | "absent") {
|
||||
b"".as_slice()
|
||||
} else {
|
||||
b"payload".as_slice()
|
||||
};
|
||||
assert_eq!(body, expected);
|
||||
// No response bytes can be available until JS observes upload end.
|
||||
release_rx.await.unwrap();
|
||||
if mode == "redirect" {
|
||||
socket.write_all(b"HTTP/1.1 307 Temporary Redirect\r\nLocation: /final\r\nContent-Length: 0\r\nConnection: close\r\n\r\n").await.unwrap();
|
||||
} else {
|
||||
let _ = socket
|
||||
.write_all(
|
||||
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
|
||||
)
|
||||
.await;
|
||||
}
|
||||
drop(socket);
|
||||
if matches!(mode, "redirect" | "reopen-complete") {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let (head, body) = read_upload_request(&mut socket).await;
|
||||
if mode == "redirect" {
|
||||
assert!(head.starts_with("POST /final HTTP/1.1\r\n"));
|
||||
assert_eq!(body, b"payload");
|
||||
} else {
|
||||
assert!(head.starts_with("GET /replacement HTTP/1.1\r\n"));
|
||||
assert!(body.is_empty());
|
||||
}
|
||||
socket
|
||||
.write_all(
|
||||
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
|
||||
)
|
||||
.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,
|
||||
);
|
||||
if mode == "intercept" {
|
||||
vm.set_fetch_subresource_interception(
|
||||
true,
|
||||
Some(crate::types::SubresourceResourceType::Xhr),
|
||||
);
|
||||
}
|
||||
vm.eval(&format!(r#"
|
||||
globalThis.result = null;
|
||||
globalThis.events = [];
|
||||
const xhr = new XMLHttpRequest();
|
||||
const mode = {mode:?};
|
||||
for (const type of ["loadstart", "progress", "load", "loadend", "abort", "error"]) {{
|
||||
xhr.upload.addEventListener(type, e => events.push([type, e.loaded, e.total, e.lengthComputable]));
|
||||
}}
|
||||
xhr.upload.onload = () => {{
|
||||
if (mode === "abort-complete") xhr.abort();
|
||||
if (mode === "reopen-complete") {{ xhr.open("GET", "/replacement"); xhr.send(); }}
|
||||
}};
|
||||
xhr.onloadend = () => {{ result = xhr.status; }};
|
||||
xhr.open("POST", "/upload");
|
||||
if (mode === "absent") xhr.send();
|
||||
else xhr.send(mode === "empty" ? "" : "payload");
|
||||
globalThis.initialEvents = events.map(e => e[0]);
|
||||
"#)).unwrap();
|
||||
assert_eq!(
|
||||
vm.eval("JSON.stringify(initialEvents)").unwrap(),
|
||||
if mode == "absent" {
|
||||
"[]"
|
||||
} else {
|
||||
"[\"loadstart\"]"
|
||||
},
|
||||
"{mode}: upload completed inside send()"
|
||||
);
|
||||
if mode == "intercept" {
|
||||
let pending = vm.take_pending_subresource_fetch_infos();
|
||||
assert_eq!(pending.len(), 1);
|
||||
vm.set_fetch_subresource_interception(false, None);
|
||||
vm.continue_pending_subresource_fetch(
|
||||
pending[0].internal_id,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
let mut release = Some(release_tx);
|
||||
tokio::time::timeout(Duration::from_secs(10), async {
|
||||
loop {
|
||||
if (mode == "absent"
|
||||
|| vm.eval("events.some(e => e[0] === 'loadend')").unwrap() == "true")
|
||||
&& let Some(release) = release.take()
|
||||
{
|
||||
release.send(()).unwrap();
|
||||
}
|
||||
if vm.eval("result !== null").unwrap() == "true" && release.is_none() {
|
||||
break;
|
||||
}
|
||||
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!("upload stalled: {mode}"));
|
||||
let has_body = mode != "absent";
|
||||
let total = if mode == "empty" { 0 } else { 7 };
|
||||
assert_eq!(
|
||||
vm.eval("String(result)").unwrap(),
|
||||
if mode == "abort-complete" { "0" } else { "200" },
|
||||
"{mode}"
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval("events.filter(e => e[0] === 'load').length")
|
||||
.unwrap(),
|
||||
if has_body { "1" } else { "0" },
|
||||
"{mode}"
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval("events.filter(e => e[0] === 'loadend').length")
|
||||
.unwrap(),
|
||||
if has_body { "1" } else { "0" },
|
||||
"{mode}"
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval("events.some(e => e[0] === 'abort' || e[0] === 'error')")
|
||||
.unwrap(),
|
||||
"false",
|
||||
"{mode}"
|
||||
);
|
||||
if has_body {
|
||||
assert_eq!(
|
||||
vm.eval("JSON.stringify(events.find(e => e[0] === 'loadend').slice(1))")
|
||||
.unwrap(),
|
||||
format!("[{total},{total},{}]", total > 0),
|
||||
"{mode}"
|
||||
);
|
||||
}
|
||||
tokio::time::timeout(Duration::from_secs(3), server)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn xhr_partial_upload_abort_and_reopen_discard_old_network_events() {
|
||||
for reopen in [false, true] {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let origin = format!("http://{}", listener.local_addr().unwrap());
|
||||
let total = 16 * 1024 * 1024;
|
||||
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"POST /upload HTTP/1.1\r\n"));
|
||||
tokio::time::sleep(Duration::from_millis(150)).await;
|
||||
let mut bytes = 0;
|
||||
let mut chunk = [0; 65536];
|
||||
while let Ok(count) = socket.read(&mut chunk).await {
|
||||
if count == 0 {
|
||||
break;
|
||||
}
|
||||
bytes += count;
|
||||
}
|
||||
if reopen {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let (head, body) = read_upload_request(&mut socket).await;
|
||||
assert!(head.starts_with("GET /replacement HTTP/1.1\r\n"));
|
||||
assert!(body.is_empty());
|
||||
socket
|
||||
.write_all(
|
||||
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
bytes
|
||||
});
|
||||
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(&format!(r#"
|
||||
globalThis.result = null;
|
||||
globalThis.partial = null;
|
||||
globalThis.events = [];
|
||||
const xhr = new XMLHttpRequest();
|
||||
for (const type of ["loadstart", "progress", "load", "loadend", "abort", "error"]) {{
|
||||
xhr.upload.addEventListener(type, e => events.push([type, e.loaded, e.total, e.lengthComputable]));
|
||||
}}
|
||||
xhr.upload.onprogress = e => {{
|
||||
if (!partial && e.loaded > 0 && e.loaded < e.total) {{
|
||||
partial = [e.loaded, e.total];
|
||||
if ({reopen}) {{ xhr.open("GET", "/replacement"); xhr.send(); }}
|
||||
else xhr.abort();
|
||||
}}
|
||||
}};
|
||||
xhr.onloadend = () => {{ result = xhr.status; }};
|
||||
xhr.open("POST", "/upload");
|
||||
xhr.send("x".repeat({total}));
|
||||
"#)).unwrap();
|
||||
assert_eq!(
|
||||
vm.eval("JSON.stringify(events.map(e => e[0]))").unwrap(),
|
||||
"[\"loadstart\"]"
|
||||
);
|
||||
tokio::time::timeout(Duration::from_secs(10), async {
|
||||
while vm.eval("result === null").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
|
||||
.expect("partial upload should be cancelled");
|
||||
assert_eq!(
|
||||
vm.eval("String(result)").unwrap(),
|
||||
if reopen { "200" } else { "0" }
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval(&format!(
|
||||
"partial[0] > 0 && partial[0] < {total} && partial[1] === {total}"
|
||||
))
|
||||
.unwrap(),
|
||||
"true"
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval("events.some(e => e[0] === 'load' || e[0] === 'error')")
|
||||
.unwrap(),
|
||||
"false"
|
||||
);
|
||||
assert_eq!(
|
||||
vm.eval("JSON.stringify(events.filter(e => e[0] === 'abort' || e[0] === 'loadend'))")
|
||||
.unwrap(),
|
||||
if reopen {
|
||||
"[]"
|
||||
} else {
|
||||
"[[\"abort\",0,0,false],[\"loadend\",0,0,false]]"
|
||||
}
|
||||
);
|
||||
assert!(
|
||||
tokio::time::timeout(Duration::from_secs(3), server)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
< total as usize
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -807,6 +807,9 @@ pub(super) struct AsyncSubresourceStreamingFinished {
|
||||
/// another kind of resident.
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
pub(crate) enum AsyncSubresourceFetchEventTarget {
|
||||
Upload {
|
||||
internal_id: u64,
|
||||
},
|
||||
Completion {
|
||||
internal_id: u64,
|
||||
},
|
||||
@@ -828,6 +831,10 @@ pub(crate) enum AsyncSubresourceFetchEventTarget {
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(super) enum AsyncSubresourceFetchEvent {
|
||||
Upload {
|
||||
internal_id: u64,
|
||||
event: moli_fetch::UploadEvent,
|
||||
},
|
||||
Completion(Box<AsyncSubresourceFetchCompletion>),
|
||||
ObservedNetworkRecord(Box<SubresourceNetworkRecord>),
|
||||
StreamingStarted(Box<AsyncSubresourceStreamingStarted>),
|
||||
@@ -838,6 +845,9 @@ pub(super) enum AsyncSubresourceFetchEvent {
|
||||
impl AsyncSubresourceFetchEvent {
|
||||
pub(crate) fn target(&self) -> AsyncSubresourceFetchEventTarget {
|
||||
match self {
|
||||
Self::Upload { internal_id, .. } => AsyncSubresourceFetchEventTarget::Upload {
|
||||
internal_id: *internal_id,
|
||||
},
|
||||
Self::Completion(completion) => AsyncSubresourceFetchEventTarget::Completion {
|
||||
internal_id: completion.internal_id,
|
||||
},
|
||||
|
||||
@@ -511,7 +511,7 @@ fn spawn_worker_fetch_service_worker(
|
||||
|
||||
pub(in crate::worker) fn spawn_worker_xhr_network(
|
||||
load: ResourceLoadLease,
|
||||
completion_tx: mpsc::UnboundedSender<WorkerXhrCompletion>,
|
||||
completion_tx: mpsc::UnboundedSender<WorkerXhrEvent>,
|
||||
xhr_id: u32,
|
||||
cancel_handle: FetchCancelHandle,
|
||||
document_url: Url,
|
||||
@@ -546,6 +546,13 @@ pub(in crate::worker) fn spawn_worker_xhr_network(
|
||||
.with_network_partition_key(network_partition_key.clone())
|
||||
.with_browser_request_metadata(BrowserRequestMetadata::Xhr)
|
||||
.with_use_cors_preflight(use_cors_preflight);
|
||||
if let Some(body) = request.body.as_ref() {
|
||||
let upload_tx = completion_tx.clone();
|
||||
let observer = moli_fetch::UploadObserver::new(body.len() as u64, move |event| {
|
||||
let _ = upload_tx.send(WorkerXhrEvent::Upload { xhr_id, event });
|
||||
});
|
||||
request = request.with_upload_observer(observer);
|
||||
}
|
||||
if let Some(referrer_policy) = referrer_policy {
|
||||
request =
|
||||
request.with_script_fetch_metadata(moli_fetch::ScriptFetchRequestMetadata {
|
||||
@@ -621,11 +628,11 @@ pub(in crate::worker) fn spawn_worker_xhr_network(
|
||||
}
|
||||
Err(error) => (Err(format!("xhr: failed to build request: {error}")), None),
|
||||
};
|
||||
let _ = completion_tx.send(WorkerXhrCompletion {
|
||||
let _ = completion_tx.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers,
|
||||
result,
|
||||
});
|
||||
})));
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1094,11 +1101,11 @@ pub(in crate::worker) fn fail_pending_worker_xhr(
|
||||
}
|
||||
completion_tx
|
||||
};
|
||||
let _ = completion_tx.send(WorkerXhrCompletion {
|
||||
let _ = completion_tx.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers: None,
|
||||
result: Err(error_text),
|
||||
});
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn fail_pending_worker_xhr_auth(
|
||||
@@ -1118,11 +1125,11 @@ pub(in crate::worker) fn fail_pending_worker_xhr_auth(
|
||||
}
|
||||
completion_tx
|
||||
};
|
||||
let _ = completion_tx.send(WorkerXhrCompletion {
|
||||
let _ = completion_tx.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers: None,
|
||||
result: Err(error_text),
|
||||
});
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn fulfill_pending_worker_xhr(
|
||||
@@ -1158,11 +1165,13 @@ pub(in crate::worker) fn fulfill_pending_worker_xhr(
|
||||
worker_response_from_body(request.url, response_code, response_headers, response_body);
|
||||
(completion_tx, response, xhr_id)
|
||||
};
|
||||
let _ = completion.0.send(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Materialized(Box::new(completion.1))),
|
||||
});
|
||||
let _ = completion
|
||||
.0
|
||||
.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Materialized(Box::new(completion.1))),
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn continue_pending_worker_xhr_response(
|
||||
@@ -1197,14 +1206,16 @@ pub(in crate::worker) fn continue_pending_worker_xhr_response(
|
||||
}
|
||||
(completion_tx, response, xhr_id)
|
||||
};
|
||||
let _ = completion.0.send(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Streamed {
|
||||
head: Box::new(completion.1.head),
|
||||
body: completion.1.body,
|
||||
}),
|
||||
});
|
||||
let _ = completion
|
||||
.0
|
||||
.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Streamed {
|
||||
head: Box::new(completion.1.head),
|
||||
body: completion.1.body,
|
||||
}),
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn fail_pending_worker_xhr_response(
|
||||
@@ -1227,11 +1238,11 @@ pub(in crate::worker) fn fail_pending_worker_xhr_response(
|
||||
pending.paused_response = None;
|
||||
completion_tx
|
||||
};
|
||||
let _ = completion_tx.send(WorkerXhrCompletion {
|
||||
let _ = completion_tx.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers: None,
|
||||
result: Err(error_text),
|
||||
});
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn fulfill_pending_worker_xhr_response(
|
||||
@@ -1258,11 +1269,13 @@ pub(in crate::worker) fn fulfill_pending_worker_xhr_response(
|
||||
worker_response_from_body(request.url, response_code, response_headers, response_body);
|
||||
(completion_tx, response, xhr_id)
|
||||
};
|
||||
let _ = completion.0.send(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Materialized(Box::new(completion.1))),
|
||||
});
|
||||
let _ = completion
|
||||
.0
|
||||
.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id: completion.2,
|
||||
network_request_headers: None,
|
||||
result: Ok(WorkerXhrResponse::Materialized(Box::new(completion.1))),
|
||||
})));
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn record_worker_websocket_subresource_failure(
|
||||
|
||||
@@ -70,14 +70,14 @@ use crate::network_host::{
|
||||
XHR_ACTIVE_INTERNAL_ID_SLOT, XHR_ASYNC_SLOT, 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, append_default_body_content_type, apply_xhr_failure,
|
||||
apply_xhr_response, apply_xhr_response_body_source, apply_xhr_timeout,
|
||||
apply_xhr_response, apply_xhr_response_body_source, apply_xhr_timeout, apply_xhr_upload_event,
|
||||
browser_request_needs_manual_preflight_redirects,
|
||||
build_fetch_response_object_from_body_source_for_request_mode,
|
||||
build_fetch_response_object_from_stream_for_request_mode,
|
||||
build_fetch_response_object_from_subresource_body_for_request_mode,
|
||||
capture_xhr_upload_listener_flag, close_pending_network_body_stream, dispatch_xhr_loadstart,
|
||||
dispatch_xhr_upload_complete, enqueue_pending_network_body_chunk,
|
||||
error_pending_network_body_stream_with_reason, extract_subresource_auth_challenge,
|
||||
enqueue_pending_network_body_chunk, error_pending_network_body_stream_with_reason,
|
||||
extract_subresource_auth_challenge,
|
||||
fetch_browser_subresource_raw_stream_with_preflight_headers_and_network_metadata,
|
||||
fetch_browser_subresource_with_preflight_headers_and_network_metadata,
|
||||
filter_cors_exposed_response_headers, filter_headers_for_guard, is_cors_policy_failure_message,
|
||||
@@ -1239,6 +1239,14 @@ pub(super) struct PausedWorkerSubresourceResponse {
|
||||
pub(super) body: SubresourceResponseBody,
|
||||
}
|
||||
|
||||
pub(super) enum WorkerXhrEvent {
|
||||
Upload {
|
||||
xhr_id: u32,
|
||||
event: moli_fetch::UploadEvent,
|
||||
},
|
||||
Completion(Box<WorkerXhrCompletion>),
|
||||
}
|
||||
|
||||
pub(super) struct WorkerXhrCompletion {
|
||||
pub(super) xhr_id: u32,
|
||||
pub(super) network_request_headers: Option<Vec<(String, String)>>,
|
||||
@@ -1548,7 +1556,7 @@ pub(crate) struct WorkerGlobalState {
|
||||
/// Fetch id counter.
|
||||
pub(super) next_fetch_id: u32,
|
||||
/// Async XHR completions routed back onto the worker event loop.
|
||||
pub(super) xhr_completion_tx: mpsc::UnboundedSender<WorkerXhrCompletion>,
|
||||
pub(super) xhr_completion_tx: mpsc::UnboundedSender<WorkerXhrEvent>,
|
||||
/// In-flight worker XHR requests keyed by internal id.
|
||||
pub(super) pending_xhrs: HashMap<u32, PendingWorkerXhr>,
|
||||
/// Worker XHR id counter.
|
||||
|
||||
@@ -335,7 +335,6 @@ pub(crate) fn try_worker_xhr_send_callback<'s>(
|
||||
if !dispatch_xhr_loadstart(scope, xhr, prepared.send_body.as_deref()) {
|
||||
return true;
|
||||
}
|
||||
dispatch_xhr_upload_complete(scope, xhr, prepared.send_body.as_deref());
|
||||
if xhr_state_bool_property(scope, xhr, XHR_ABORTED_SLOT).unwrap_or(false)
|
||||
|| worker_xhr_open_generation_changed(scope, xhr, open_generation)
|
||||
{
|
||||
@@ -483,11 +482,16 @@ pub(crate) fn try_worker_xhr_send_callback<'s>(
|
||||
if let Some(result) = local_response {
|
||||
// Local responses and network errors still complete asynchronously, so
|
||||
// abort(), open() and timeout processing use the ordinary pending XHR.
|
||||
let _ = state.borrow().xhr_completion_tx.send(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers: None,
|
||||
result: result.map(|response| WorkerXhrResponse::Materialized(Box::new(response))),
|
||||
});
|
||||
let _ = state
|
||||
.borrow()
|
||||
.xhr_completion_tx
|
||||
.send(WorkerXhrEvent::Completion(Box::new(WorkerXhrCompletion {
|
||||
xhr_id,
|
||||
network_request_headers: None,
|
||||
result: result
|
||||
.map(|response| WorkerXhrResponse::Materialized(Box::new(response)))
|
||||
.map_err(|error| error.to_string()),
|
||||
})));
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -906,6 +910,32 @@ pub(in crate::worker) fn record_worker_xhr_failure(
|
||||
);
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn drain_worker_xhr_event(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
state: &Rc<RefCell<WorkerGlobalState>>,
|
||||
event: WorkerXhrEvent,
|
||||
) {
|
||||
match event {
|
||||
WorkerXhrEvent::Upload { xhr_id, event } => {
|
||||
let xhr = {
|
||||
let state = state.borrow();
|
||||
let Some(pending) = state.pending_xhrs.get(&xhr_id) else {
|
||||
return;
|
||||
};
|
||||
v8::Local::new(scope, &pending.xhr)
|
||||
};
|
||||
if !apply_xhr_upload_event(scope, xhr, u64::from(xhr_id), event)
|
||||
&& let Some(pending) = state.borrow_mut().pending_xhrs.remove(&xhr_id)
|
||||
{
|
||||
pending.load.cancel();
|
||||
}
|
||||
}
|
||||
WorkerXhrEvent::Completion(completion) => {
|
||||
drain_worker_xhr_completion(scope, state, *completion)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::worker) fn drain_worker_xhr_completion(
|
||||
scope: &mut v8::PinScope<'_, '_>,
|
||||
state: &Rc<RefCell<WorkerGlobalState>>,
|
||||
|
||||
@@ -70,7 +70,7 @@ use runtime_inspector::WorkerRuntimeInspector;
|
||||
|
||||
use super::global_scope::{
|
||||
WorkerFetchEvent, WorkerGlobalState, WorkerIsolateTimerQueues, WorkerOpfsCompletion,
|
||||
WorkerWebCryptoCompletion, WorkerXhrCompletion, close_worker_owned_broadcast_channels,
|
||||
WorkerWebCryptoCompletion, WorkerXhrEvent, close_worker_owned_broadcast_channels,
|
||||
close_worker_owned_message_ports, continue_pending_worker_csp_report,
|
||||
continue_pending_worker_fetch, continue_pending_worker_fetch_response,
|
||||
continue_pending_worker_xhr, continue_pending_worker_xhr_response,
|
||||
@@ -86,7 +86,7 @@ use super::global_scope::{
|
||||
drain_service_worker_push_unsubscribe_result, drain_service_worker_show_notification_result,
|
||||
drain_service_worker_sync_get_tags_result, drain_service_worker_sync_registration_result,
|
||||
drain_worker_fetch_completion, drain_worker_opfs_completion, drain_worker_webcrypto_completion,
|
||||
drain_worker_xhr_completion, fail_pending_worker_csp_report, fail_pending_worker_fetch,
|
||||
drain_worker_xhr_event, fail_pending_worker_csp_report, fail_pending_worker_fetch,
|
||||
fail_pending_worker_fetch_auth, fail_pending_worker_fetch_response, fail_pending_worker_xhr,
|
||||
fail_pending_worker_xhr_auth, fail_pending_worker_xhr_response,
|
||||
fulfill_pending_worker_csp_report, fulfill_pending_worker_fetch,
|
||||
@@ -1608,8 +1608,7 @@ async fn worker_main(
|
||||
// Worker global state (accessible from JS callbacks).
|
||||
let (fetch_completion_tx, mut fetch_completion_rx) =
|
||||
mpsc::unbounded_channel::<WorkerFetchEvent>();
|
||||
let (xhr_completion_tx, mut xhr_completion_rx) =
|
||||
mpsc::unbounded_channel::<WorkerXhrCompletion>();
|
||||
let (xhr_completion_tx, mut xhr_completion_rx) = mpsc::unbounded_channel::<WorkerXhrEvent>();
|
||||
let (module_graph_fetch_tx, mut module_graph_fetch_rx) =
|
||||
mpsc::unbounded_channel::<WorkerModuleGraphFetchCompletion>();
|
||||
let (module_evaluation_tx, mut module_evaluation_rx) =
|
||||
@@ -1981,7 +1980,7 @@ async fn worker_main(
|
||||
enum WorkerLoopWake {
|
||||
Message(Option<WorkerMessage>),
|
||||
Fetch(Option<WorkerFetchEvent>),
|
||||
Xhr(Option<WorkerXhrCompletion>),
|
||||
Xhr(Option<WorkerXhrEvent>),
|
||||
ModuleGraphFetch(Option<Box<WorkerModuleGraphFetchCompletion>>),
|
||||
ModuleEvaluation(Option<WorkerModuleEvaluationCompletion>),
|
||||
ModuleRuntime,
|
||||
@@ -3032,7 +3031,7 @@ async fn worker_main(
|
||||
let scope = &mut scope.init();
|
||||
let ctx = v8::Local::new(scope, &context);
|
||||
let scope = &mut v8::ContextScope::new(scope, ctx);
|
||||
drain_worker_xhr_completion(scope, &state, completion);
|
||||
drain_worker_xhr_event(scope, &state, completion);
|
||||
perform_worker_microtask_checkpoint_and_report_pending_promise_rejections(scope);
|
||||
drain_worker_dynamic_module_imports(scope, &state, &module_graph_fetch_tx);
|
||||
}
|
||||
|
||||
@@ -3749,6 +3749,96 @@ async fn worker_xhr_upload_listener_preflight_survives_sync_and_interception() {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_xhr_partial_upload_can_abort_or_reopen_without_stale_completion() {
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
ensure_v8();
|
||||
for reopen in [false, true] {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let origin = format!("http://{}", listener.local_addr().unwrap());
|
||||
let total = 16 * 1024 * 1024;
|
||||
let server = tokio::spawn(async move {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let head = read_http_request_head(&mut socket).await.unwrap();
|
||||
assert!(head.starts_with("POST /upload HTTP/1.1"));
|
||||
tokio::time::sleep(Duration::from_millis(150)).await;
|
||||
let mut bytes = 0;
|
||||
let mut chunk = [0; 65536];
|
||||
while let Ok(count) = socket.read(&mut chunk).await {
|
||||
if count == 0 {
|
||||
break;
|
||||
}
|
||||
bytes += count;
|
||||
}
|
||||
if reopen {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let head = read_http_request_head(&mut socket).await.unwrap();
|
||||
assert!(head.starts_with("GET /replacement HTTP/1.1"));
|
||||
socket
|
||||
.write_all(
|
||||
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
bytes
|
||||
});
|
||||
let mut config = FetchConfig::default();
|
||||
config.set_http_no_proxy(Some("*".to_owned()));
|
||||
let loader = ResourceRequestClient::new(&config).unwrap();
|
||||
let mut handle = spawn_worker_with_request_client(
|
||||
format!(
|
||||
r#"
|
||||
const xhr = new XMLHttpRequest();
|
||||
const events = [];
|
||||
let partial = null;
|
||||
for (const type of ["loadstart", "progress", "load", "loadend", "abort", "error"]) {{
|
||||
xhr.upload.addEventListener(type, e => events.push([type, e.loaded, e.total, e.lengthComputable]));
|
||||
}}
|
||||
xhr.upload.onprogress = e => {{
|
||||
if (!partial && e.loaded > 0 && e.loaded < e.total) {{
|
||||
partial = [e.loaded, e.total];
|
||||
if ({reopen}) {{ xhr.open("GET", "/replacement"); xhr.send(); }}
|
||||
else xhr.abort();
|
||||
}}
|
||||
}};
|
||||
xhr.onloadend = () => {{ postMessage({{status: xhr.status, partial, events}}); close(); }};
|
||||
xhr.open("POST", "/upload");
|
||||
xhr.send("x".repeat({total}));
|
||||
postMessage(events.map(e => e[0]));
|
||||
"#
|
||||
),
|
||||
format!("{origin}/worker.js"),
|
||||
loader,
|
||||
);
|
||||
assert_eq!(recv_post_json(&mut handle).await, "[\"loadstart\"]");
|
||||
let result: serde_json::Value =
|
||||
serde_json::from_str(&recv_post_json(&mut handle).await).unwrap();
|
||||
assert_eq!(result["status"], if reopen { 200 } else { 0 });
|
||||
assert_eq!(result["partial"][1], total);
|
||||
let loaded = result["partial"][0].as_u64().unwrap();
|
||||
assert!(loaded > 0 && loaded < total as u64);
|
||||
let events = result["events"].as_array().unwrap();
|
||||
assert!(
|
||||
!events
|
||||
.iter()
|
||||
.any(|event| event[0] == "load" || event[0] == "error")
|
||||
);
|
||||
let aborts: Vec<_> = events.iter().filter(|event| event[0] == "abort").collect();
|
||||
let ends: Vec<_> = events
|
||||
.iter()
|
||||
.filter(|event| event[0] == "loadend")
|
||||
.collect();
|
||||
assert_eq!(aborts.len(), usize::from(!reopen));
|
||||
assert_eq!(ends.len(), usize::from(!reopen));
|
||||
if !reopen {
|
||||
assert_eq!(*aborts[0], serde_json::json!(["abort", 0, 0, false]));
|
||||
assert_eq!(*ends[0], serde_json::json!(["loadend", 0, 0, false]));
|
||||
}
|
||||
assert!(timeout(TIMEOUT, server).await.unwrap().unwrap() < total as usize);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn worker_fetch_uses_worker_script_base_url_and_resolves_response_text() {
|
||||
ensure_v8();
|
||||
|
||||
@@ -43,15 +43,19 @@ self.onmessage = function (event) {
|
||||
};
|
||||
xhr.onloadend = function () {
|
||||
xhrEvents.push("loadend");
|
||||
postMessage({
|
||||
kind,
|
||||
status: xhr.status,
|
||||
response: xhr.response,
|
||||
uploadEvents,
|
||||
uploadOrder,
|
||||
xhrEvents,
|
||||
});
|
||||
close();
|
||||
// An abort from the final upload progress callback fires XHR loadend
|
||||
// before the upload's remaining events. Collect after that task finishes.
|
||||
setTimeout(() => {
|
||||
postMessage({
|
||||
kind,
|
||||
status: xhr.status,
|
||||
response: xhr.response,
|
||||
uploadEvents,
|
||||
uploadOrder,
|
||||
xhrEvents,
|
||||
});
|
||||
close();
|
||||
}, 0);
|
||||
};
|
||||
xhr.open("POST", "/wpt/runtime/xhr/echo-body");
|
||||
xhr.responseType = "json";
|
||||
|
||||
@@ -210,6 +210,9 @@ promise_test(async function () {
|
||||
|
||||
promise_test(async function () {
|
||||
const result = await run_worker_upload_case({ kind: "abort-from-progress" });
|
||||
// loaded === total does not distinguish a chunk callback from the final
|
||||
// end-of-body callback, which has already marked the upload complete.
|
||||
const interrupted = result.uploadEvents.some(event => event.startsWith("abort:"));
|
||||
|
||||
assert_equals(result.status, 0, "worker XHR abort should reset status");
|
||||
assert_array_equals(
|
||||
@@ -217,10 +220,10 @@ promise_test(async function () {
|
||||
[
|
||||
"loadstart:true:true:true:0:12",
|
||||
"progress:true:true:true:12:12",
|
||||
"abort:true:true:false:0:0",
|
||||
"loadend:true:true:false:0:0",
|
||||
interrupted ? "abort:true:true:false:0:0" : "load:true:true:true:12:12",
|
||||
interrupted ? "loadend:true:true:false:0:0" : "loadend:true:true:true:12:12",
|
||||
],
|
||||
"worker XHR upload abort from progress should stop completion and emit upload abort/loadend",
|
||||
"worker XHR upload abort from progress should respect whether the upload has completed",
|
||||
);
|
||||
assert_array_equals(
|
||||
result.uploadOrder,
|
||||
@@ -231,9 +234,9 @@ promise_test(async function () {
|
||||
"listener-before:progress",
|
||||
"handler:progress",
|
||||
"listener-after:progress",
|
||||
"listener-before:abort",
|
||||
"handler:abort",
|
||||
"listener-after:abort",
|
||||
"listener-before:" + (interrupted ? "abort" : "load"),
|
||||
"handler:" + (interrupted ? "abort" : "load"),
|
||||
"listener-after:" + (interrupted ? "abort" : "load"),
|
||||
"listener-before:loadend",
|
||||
"handler:loadend",
|
||||
"listener-after:loadend",
|
||||
|
||||
@@ -121,7 +121,9 @@ function abort_from_upload_event_case(abortEventType) {
|
||||
};
|
||||
xhr.onloadend = function () {
|
||||
xhrEvents.push("loadend");
|
||||
resolve({ xhr, uploadEvents, xhrEvents });
|
||||
// A final upload progress callback can abort the XHR before the
|
||||
// upload's remaining load/loadend events. Observe the completed task.
|
||||
setTimeout(() => resolve({ xhr, uploadEvents, xhrEvents }), 0);
|
||||
};
|
||||
xhr.open("POST", ECHO_BODY_URL);
|
||||
xhr.responseType = "json";
|
||||
@@ -269,16 +271,19 @@ promise_test(async function () {
|
||||
|
||||
promise_test(async function () {
|
||||
const result = await abort_from_upload_event_case("progress");
|
||||
// A full byte count can be a chunk-progress callback or the end-of-body
|
||||
// callback. Only the former still has an incomplete upload to abort.
|
||||
const interrupted = result.uploadEvents.some(event => event.startsWith("abort:"));
|
||||
|
||||
assert_array_equals(
|
||||
result.uploadEvents,
|
||||
[
|
||||
"loadstart:true:true:true:0:12",
|
||||
"progress:true:true:true:12:12",
|
||||
"abort:true:true:false:0:0",
|
||||
"loadend:true:true:false:0:0",
|
||||
interrupted ? "abort:true:true:false:0:0" : "load:true:true:true:12:12",
|
||||
interrupted ? "loadend:true:true:false:0:0" : "loadend:true:true:true:12:12",
|
||||
],
|
||||
"aborting from upload progress should stop upload completion and emit upload abort/loadend",
|
||||
"aborting from upload progress should respect whether the upload has completed",
|
||||
);
|
||||
assert_array_equals(
|
||||
result.xhrEvents,
|
||||
|
||||
Reference in New Issue
Block a user