fix(navigation): cancel loading without retiring the active document

This commit is contained in:
lanyue-llk
2026-09-22 13:23:43 +08:00
parent 2e9c28529f
commit b43a213dc6
30 changed files with 1051 additions and 48 deletions
+110 -13
View File
@@ -592,16 +592,30 @@ async fn first_nonempty_response_body_chunk(
Some(chunk) if chunk.is_empty() => continue,
Some(chunk) => return Ok(Some(chunk)),
None => {
response
.finish()
.await
.context("failed to read page body from stream")?;
response.finish().await?;
return Ok(None);
}
}
}
}
fn failed_provisional_body_load(
error: anyhow::Error,
context: &str,
) -> anyhow::Result<NavigationLoadOutcome> {
if error
.downcast_ref::<NetworkFetchFailureContext>()
.is_some_and(|failure| {
failure.network_error_text() == moli_fetch::NET_ERR_ABORTED_ERROR_TEXT
})
{
return Ok(NavigationLoadOutcome::network_failure(
moli_fetch::NET_ERR_ABORTED_ERROR_TEXT.to_owned(),
));
}
Err(error.context(context.to_owned()))
}
fn spawn_streaming_body_capture(
mut response: StreamingRawResponse,
initial_chunk: Option<Vec<u8>>,
@@ -959,6 +973,14 @@ impl BackgroundNavigationLoadJob {
network_error_text = failure.network_error_text(),
"main document transport failed before response metadata"
);
// An aborted provisional navigation never commits a replacement
// document. Retain its network failure for the command/event
// completion path, without creating a browser-owned error page.
if failure.network_error_text() == moli_fetch::NET_ERR_ABORTED_ERROR_TEXT {
return Ok(NavigationLoadOutcome::network_failure(
failure.network_error_text().to_owned(),
));
}
return prepare_network_error_page_navigation_with_engine_async(
&mut engine,
self.page_reservation,
@@ -1355,9 +1377,12 @@ async fn build_navigation_from_streaming_raw_response_with_engine_async(
body_progress_source.body_network_progress_for_completed_events(network_events);
let body_progress_source_for_body_finish = body_progress_source.clone();
let mut initial_body_chunk = if response_status_may_use_http_error_page(response_status) {
match first_nonempty_response_body_chunk(&mut response).await? {
Some(chunk) => Some(chunk),
None => {
match first_nonempty_response_body_chunk(&mut response).await {
Err(error) => {
return failed_provisional_body_load(error, "failed to read page body from stream");
}
Ok(Some(chunk)) => Some(chunk),
Ok(None) => {
let body =
CapturedBody::from_string(http_error_page_html(&final_url, response_status));
return prepare_browser_owned_error_page_navigation_with_engine_async(
@@ -1392,10 +1417,9 @@ async fn build_navigation_from_streaming_raw_response_with_engine_async(
.append(&chunk)
.context("failed to capture XML page body")?;
}
response
.finish()
.await
.context("failed to read XML page body from stream")?;
if let Err(error) = response.finish().await {
return failed_provisional_body_load(error, "failed to read XML page body from stream");
}
let captured_body = body_writer
.finish()
.context("failed to finish captured XML page body")?;
@@ -1471,7 +1495,8 @@ async fn build_navigation_from_streaming_raw_response_with_engine_async(
let (body_tx, body_rx) = mpsc::channel(EXTERNAL_RAW_BODY_CHANNEL_CAPACITY);
let (completion_tx, completion_rx) = oneshot::channel();
let raw_body = moli_core::runtime::ExternalRawDocumentBodyStream::new(body_rx, completion_rx);
let raw_body = moli_core::runtime::ExternalRawDocumentBodyStream::new(body_rx, completion_rx)
.with_stop_loading_cancellation(response.cancellation_handle());
let page_storage = load_inputs.page_storage_handles();
let main_document_commit = load_inputs
.main_document_commit_for_final_url(&final_url, None)
@@ -1694,7 +1719,8 @@ impl CdpConnection {
let raw_body = moli_core::runtime::ExternalRawDocumentBodyStream::new(
renderer_body_rx,
renderer_completion_rx,
);
)
.with_stop_loading_cancellation(response.cancellation_handle());
let shared_resource_runtime =
self.shared_resource_runtime_for_navigation_load_inputs(&load_inputs);
let mut engine = self.navigation_engine_handle_for_load_inputs(&load_inputs);
@@ -4019,6 +4045,77 @@ mod tests {
};
use serde_json::json;
#[test]
fn provisional_body_failure_preserves_typed_source_and_context() {
let source = anyhow::Error::new(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"body transport reset",
))
.context("transport completion failed");
let error =
super::failed_provisional_body_load(source, "failed to read XML page body from stream")
.expect_err("transport failure must remain an error");
assert_eq!(
error.downcast_ref::<std::io::Error>().unwrap().kind(),
std::io::ErrorKind::ConnectionReset
);
let message = format!("{error:#}");
assert!(message.contains("failed to read XML page body from stream"));
assert!(message.contains("transport completion failed"));
assert!(message.contains("body transport reset"));
}
#[test]
fn provisional_body_failure_does_not_infer_cancellation_from_error_text() {
let error = super::failed_provisional_body_load(
anyhow::anyhow!(moli_fetch::NET_ERR_ABORTED_ERROR_TEXT),
"failed to read page body from stream",
)
.expect_err("only a typed transport cancellation may become a network outcome");
assert!(format!("{error:#}").contains(moli_fetch::NET_ERR_ABORTED_ERROR_TEXT));
}
#[tokio::test]
async fn response_capture_survives_ordinary_renderer_body_retirement() {
let cancellation = moli_fetch::FetchCancelHandle::new();
let (source_tx, source_rx) = tokio::sync::mpsc::unbounded_channel();
let (source_completion_tx, source_completion_rx) = tokio::sync::oneshot::channel();
let response = moli_fetch::StreamingRawResponse::new(
url::Url::parse("https://capture.test/document").unwrap(),
200,
Vec::new(),
None,
Vec::new(),
false,
Vec::new(),
source_rx,
cancellation.clone(),
source_completion_rx,
);
let (renderer_tx, renderer_rx) = tokio::sync::mpsc::channel(1);
let (renderer_completion_tx, renderer_completion_rx) = tokio::sync::oneshot::channel();
drop(renderer_rx);
drop(renderer_completion_rx);
let capture = super::spawn_streaming_body_capture(
response,
None,
renderer_tx,
renderer_completion_tx,
);
source_tx.send(b"prefix".to_vec()).unwrap();
source_tx
.send(b"captured after parser retirement".to_vec())
.unwrap();
drop(source_tx);
source_completion_tx.send(Ok(())).unwrap();
let body = capture.await.unwrap().unwrap();
assert_eq!(
body.materialize_bytes().unwrap().as_slice(),
b"prefixcaptured after parser retirement"
);
assert!(!cancellation.is_cancelled());
}
#[tokio::test]
async fn streaming_body_failure_preserves_source_for_renderer_and_navigation() {
let (chunks_tx, chunks_rx) = tokio::sync::mpsc::unbounded_channel();
+33
View File
@@ -722,6 +722,14 @@ impl TargetPageSlot {
.is_some_and(|request| request.background_completion_pending)
}
pub(crate) fn cancel_inflight_document_navigation(&self) {
if let Some(request) = self.pending_navigation_request.as_ref() {
// Keep the token installed: the existing completion path owns the
// aborted response and must still settle this exact navigation.
request.cancel();
}
}
pub(crate) fn bind_pending_document_navigation_renderer_page(
&mut self,
token: &DocumentNavigationToken,
@@ -1580,6 +1588,31 @@ mod pending_renderer_page_tests {
);
}
#[test]
fn navigation_cancellation_preserves_completion_token_and_is_target_local() {
let mut slot = TargetPageSlot::default();
let token = slot.start_document_navigation("TID-1".to_owned(), "LOADER-1".to_owned());
let cancellation = slot
.document_navigation_cancellation_handle(&token)
.unwrap();
let mut peer = TargetPageSlot::default();
let peer_token = peer.start_document_navigation("TID-2".to_owned(), "LOADER-2".to_owned());
let peer_cancellation = peer
.document_navigation_cancellation_handle(&peer_token)
.unwrap();
slot.cancel_inflight_document_navigation();
assert!(cancellation.is_cancelled());
assert!(!peer_cancellation.is_cancelled());
assert!(slot.accepts_pending_document_navigation_event(&token));
let next = slot.start_document_navigation("TID-1".to_owned(), "LOADER-next".to_owned());
assert!(
!slot
.document_navigation_cancellation_handle(&next)
.unwrap()
.is_cancelled()
);
}
#[test]
fn navigation_binding_cannot_follow_a_superseding_navigation() {
let mut slot = TargetPageSlot::default();
@@ -574,6 +574,10 @@ impl TargetRuntimeSlot {
self.page_slot.has_inflight_background_navigation()
}
pub(crate) fn cancel_inflight_document_navigation(&self) {
self.page_slot.cancel_inflight_document_navigation();
}
pub(crate) fn accepts_document_body_completion_event(
&self,
token: &DocumentNavigationToken,
@@ -350,17 +350,26 @@ fn materialize_navigation_load_outcome(
materialize_download_navigation_progress(conn, state, *navigation),
),
NavigationLoadOutcome::NetworkFailure(error_text) => {
let document_policy = failed_navigation_document_policy(&error_text);
MaterializedNavigationLoadOutcome::Failed(materialize_failed_navigation_progress(
conn,
state,
error_text,
FailedNavigationDocumentPolicy::InvalidateCommittedDocument,
document_policy,
FailedNavigationResponseMode::CdpErrorTextResult,
))
}
}
}
fn failed_navigation_document_policy(error_text: &str) -> FailedNavigationDocumentPolicy {
if error_text == moli_fetch::NET_ERR_ABORTED_ERROR_TEXT {
FailedNavigationDocumentPolicy::PreserveCommittedDocument
} else {
FailedNavigationDocumentPolicy::InvalidateCommittedDocument
}
}
pub(crate) fn materialize_navigation_load_result(
conn: &mut CdpConnection,
state: &NavigationDispatchState,
@@ -18,6 +18,24 @@ use super::gate::{
};
use super::*;
#[test]
fn aborted_navigation_preserves_document_without_changing_other_failure_policies() {
assert_eq!(
failed_navigation_document_policy(moli_fetch::NET_ERR_ABORTED_ERROR_TEXT),
FailedNavigationDocumentPolicy::PreserveCommittedDocument
);
for error in [
"net::ERR_CONNECTION_REFUSED",
"net::ERR_TIMED_OUT",
"net::ERR_HTTP_RESPONSE_CODE_FAILURE",
] {
assert_eq!(
failed_navigation_document_policy(error),
FailedNavigationDocumentPolicy::InvalidateCommittedDocument
);
}
}
fn completed_events() -> CompletedMainDocumentNetworkEvents {
CompletedMainDocumentNetworkEvents::new(
"GET".to_owned(),
@@ -395,11 +395,13 @@ pub(super) async fn complete_stop_loading_command_dispatch(
owner: &CommandOwnerScope,
) -> PageCommandTaskStep {
let mut out = Vec::new();
if let Ok(slot) = conn.runtime_session_owner_slot_mut_for_owner(owner)
&& let Some(page) = slot.loaded_page_mut()
&& let Err(error) = page.stop_document_lifecycle_async().await
{
tracing::debug!(%error, "failed to stop renderer document lifecycle");
if let Ok(slot) = conn.runtime_session_owner_slot_mut_for_owner(owner) {
slot.cancel_inflight_document_navigation();
if let Some(page) = slot.loaded_page_mut()
&& let Err(error) = page.stop_document_lifecycle_async().await
{
tracing::debug!(%error, "failed to stop renderer document lifecycle");
}
}
let (
pending_navigations,
@@ -5464,6 +5464,373 @@ async fn stop_loading_aborts_paused_request_stage_navigation() {
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_cancels_inflight_unpaused_navigation_transport() {
let mut ctx = TestContext::new();
load_bc_with_session(&mut ctx, "BID-1", "TID-1", "SID-1", "about:blank");
let token = ctx
.conn
.browser_context
.as_mut()
.unwrap()
.start_document_navigation_for_active_target("LOADER-inflight-stop".to_owned())
.unwrap();
let cancellation = ctx
.conn
.document_navigation_cancellation_handle(&token)
.unwrap();
ctx.conn.arm_background_navigation_completion(&token, None);
assert!(!cancellation.is_cancelled());
ctx.process_async(json!({
"id": 901,
"method": "Page.stopLoading",
"sessionId": "SID-1"
}))
.await;
ctx.expect_result(901, json!({}), Some("SID-1"));
assert!(
cancellation.is_cancelled(),
"stopLoading must cancel an ordinary HTTP navigation, not only Fetch-paused requests"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_before_response_preserves_document_and_allows_next_navigation() {
assert_stopped_provisional_navigation_preserves_document(None).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_during_xml_body_preserves_document_and_allows_next_navigation() {
assert_stopped_provisional_navigation_preserves_document(Some((200, "application/xml"))).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_during_http_error_body_preserves_document_and_allows_next_navigation() {
assert_stopped_provisional_navigation_preserves_document(Some((500, "text/html"))).await;
}
async fn assert_stopped_provisional_navigation_preserves_document(
response_head: Option<(u16, &'static str)>,
) {
let request_received = std::sync::Arc::new(tokio::sync::Notify::new());
let release_response = std::sync::Arc::new(tokio::sync::Notify::new());
let handler_received = request_received.clone();
let handler_release = release_response.clone();
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let app = axum::Router::new()
.route(
"/held",
axum::routing::get(move || {
let received = handler_received.clone();
let release = handler_release.clone();
async move {
received.notify_one();
let (status, content_type) = if let Some(head) = response_head {
head
} else {
release.notified().await;
return axum::http::Response::builder()
.header("content-type", "text/html")
.body(axum::body::Body::from("<body>cancelled response</body>"))
.unwrap();
};
let body = futures_util::stream::once(async move {
release.notified().await;
Ok::<_, std::convert::Infallible>(axum::body::Bytes::from_static(
b"<body>cancelled response</body>",
))
});
axum::http::Response::builder()
.status(status)
.header("content-type", content_type)
.body(axum::body::Body::from_stream(body))
.unwrap()
}
}),
)
.route(
"/ready",
axum::routing::get(|| async {
axum::response::Html("<body>subsequent document</body>")
}),
);
axum::serve(listener, app).await.unwrap();
});
let mut ctx = TestContext::new();
load_bc_with_session(&mut ctx, "BID-1", "TID-1", "SID-1", "about:blank");
let initial_url = "data:text/html,<body>previous committed document</body>";
ctx.install_navigation_fixture_for_session_owner(initial_url, Some("SID-1"))
.await;
wait_until_renderer_document_load(&mut ctx, Some("SID-1"), "TID-1", LOADER_ID).await;
for (id, method) in [
(910, "Page.enable"),
(911, "DOM.enable"),
(912, "Network.enable"),
(913, "Runtime.enable"),
] {
ctx.process_async(json!({
"id": id,
"method": method,
"sessionId": "SID-1"
}))
.await;
assert_eq!(take_response_by_id(&mut ctx, id)["result"], json!({}));
}
let previous_html = loaded_page_html_for_test(&mut ctx).await;
ctx.sent.clear();
ctx.enable_background_navigation_scheduler_for_test();
tokio::task::LocalSet::new()
.run_until(async {
ctx.process_async(json!({
"id": 914,
"method": "Page.navigate",
"sessionId": "SID-1",
"params": { "url": format!("http://{addr}/held") }
}))
.await;
tokio::time::timeout(
std::time::Duration::from_secs(5),
request_received.notified(),
)
.await
.expect("the ordinary HTTP navigation should reach the held response");
if let Some((status, _)) = response_head {
wait_until_scheduler_message(
&mut ctx,
"held document response headers",
|message| {
message["method"] == json!("Network.responseReceived")
&& message["params"]["type"] == json!("Document")
&& message["params"]["response"]["status"] == json!(status)
},
)
.await;
} else {
assert!(ctx.sent.iter().all(|message| message["id"] != json!(914)));
}
ctx.process_and_wait_for_response_async(json!({
"id": 915,
"method": "Page.stopLoading",
"sessionId": "SID-1"
}))
.await;
assert_eq!(take_response_by_id(&mut ctx, 915)["result"], json!({}));
wait_until_scheduler_message(&mut ctx, "cancelled navigation response", |message| {
message["id"] == json!(914)
})
.await;
let navigation = take_response_by_id(&mut ctx, 914);
assert_eq!(navigation["result"]["frameId"], json!("TID-1"));
if response_head.is_some_and(|(status, _)| status == 200) {
// Successful response headers acknowledge navigation before
// XML buffering completes; cancellation must not send a
// second reply or replace the still-committed old Document.
assert!(navigation["result"].get("errorText").is_none());
} else {
assert_eq!(navigation["result"]["errorText"], json!("net::ERR_ABORTED"));
}
wait_until_scheduler_message(
&mut ctx,
"cancelled navigation network event",
|message| {
message["method"] == json!("Network.loadingFailed")
&& message["params"]["errorText"] == json!("net::ERR_ABORTED")
},
)
.await;
let failure = ctx
.sent
.iter()
.find(|message| message["method"] == json!("Network.loadingFailed"))
.expect("cancelled request failure");
assert_eq!(failure["params"]["canceled"], json!(true));
assert_eq!(failure["params"]["type"], json!("Document"));
assert!(ctx.sent.iter().all(|message| message["id"] != json!(914)));
assert_eq!(loaded_page_html_for_test(&mut ctx).await, previous_html);
assert_eq!(
ctx.conn.browser_context.as_ref().unwrap().target_url(),
initial_url
);
for method in [
"Page.frameNavigated",
"DOM.documentUpdated",
"Runtime.executionContextsCleared",
"Runtime.executionContextCreated",
"Page.domContentEventFired",
"Page.loadEventFired",
] {
assert!(
ctx.sent
.iter()
.all(|message| message["method"] != json!(method)),
"cancelled pre-response navigation must not replace the document: {method}"
);
}
// The response remains held until cancellation has completed. No
// server delay or timing race can masquerade as transport abort.
release_response.notify_one();
ctx.sent.clear();
ctx.process_and_wait_for_response_async(json!({
"id": 916,
"method": "Page.navigate",
"sessionId": "SID-1",
"params": { "url": format!("http://{addr}/ready") }
}))
.await;
let navigation = take_response_by_id(&mut ctx, 916);
assert!(navigation["result"].get("errorText").is_none());
wait_until_frame_stopped_loading(&mut ctx, "TID-1").await;
let html = loaded_page_html_for_test(&mut ctx).await;
assert!(html.contains("subsequent document"));
assert!(!html.contains("cancelled response"));
assert_eq!(
ctx.conn.browser_context.as_ref().unwrap().target_url(),
format!("http://{addr}/ready")
);
})
.await;
server.abort();
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_after_commit_cancels_transport_without_replacing_partial_document() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let prefix_parsed = std::sync::Arc::new(tokio::sync::Notify::new());
let transport_closed = std::sync::Arc::new(tokio::sync::Notify::new());
let server_parsed = prefix_parsed.clone();
let server_closed = transport_closed.clone();
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
loop {
let (mut socket, _) = listener.accept().await.unwrap();
let parsed = server_parsed.clone();
let closed = server_closed.clone();
tokio::spawn(async move {
let mut request = Vec::new();
let mut byte = [0_u8; 1];
while !request.ends_with(b"\r\n\r\n") {
if socket.read(&mut byte).await.unwrap() == 0 {
return;
}
request.push(byte[0]);
}
if request.starts_with(b"GET /stream ") {
let prefix = b"<!doctype html><body><main id='prefix'>committed prefix</main><script>fetch('/parsed')</script>";
let tail = b"<main id='tail'>unreceived tail</main></body>";
let headers = format!(
"HTTP/1.1 200 OK\r\nContent-Type: text/html\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
prefix.len() + tail.len()
);
socket.write_all(headers.as_bytes()).await.unwrap();
socket.write_all(prefix).await.unwrap();
socket.flush().await.unwrap();
// The server never supplies EOF or the tail. Completion
// here can only come from the client's transport close.
let read = socket.read(&mut byte).await;
assert!(
matches!(read, Ok(0)) || read.is_err(),
"client must close the held transfer: {read:?}"
);
closed.notify_one();
} else if request.starts_with(b"GET /parsed ") {
socket
.write_all(b"HTTP/1.1 204 No Content\r\nConnection: close\r\n\r\n")
.await
.unwrap();
parsed.notify_one();
} else {
socket.write_all(b"HTTP/1.1 200 OK\r\nContent-Type: text/html\r\nContent-Length: 26\r\nConnection: close\r\n\r\n<body>complete page</body>").await.unwrap();
}
});
}
});
let mut ctx = TestContext::new();
load_bc_with_session(&mut ctx, "BID-1", "TID-1", "SID-1", "about:blank");
ctx.enable_page_events_for_test(Some("SID-1"));
ctx.enable_background_navigation_scheduler_for_test();
tokio::task::LocalSet::new()
.run_until(async {
ctx.process_and_wait_for_response_async(json!({
"id": 920, "method": "Page.navigate", "sessionId": "SID-1",
"params": { "url": format!("http://{addr}/stream") }
}))
.await;
assert!(
take_response_by_id(&mut ctx, 920)["result"]
.get("errorText")
.is_none()
);
wait_until_scheduler_message(&mut ctx, "streaming document commit", |message| {
message["method"] == json!("Page.frameNavigated")
})
.await;
tokio::time::timeout(std::time::Duration::from_secs(5), prefix_parsed.notified())
.await
.expect("committed prefix script must execute before stopping");
assert!(
loaded_page_html_for_test(&mut ctx)
.await
.contains("committed prefix")
);
ctx.sent.clear();
ctx.process_and_wait_for_response_async(json!({
"id": 921, "method": "Page.stopLoading", "sessionId": "SID-1"
}))
.await;
assert_eq!(take_response_by_id(&mut ctx, 921)["result"], json!({}));
tokio::time::timeout(
std::time::Duration::from_secs(5),
transport_closed.notified(),
)
.await
.expect("explicit document stop must cancel the committed response transport");
let html = loaded_page_html_for_test(&mut ctx).await;
assert!(html.contains("committed prefix"));
assert!(!html.contains("unreceived tail"));
assert_eq!(
ctx.conn.browser_context.as_ref().unwrap().target_url(),
format!("http://{addr}/stream")
);
assert!(
ctx.sent
.iter()
.all(|message| message["method"] != json!("Page.frameNavigated"))
);
ctx.sent.clear();
ctx.process_and_wait_for_response_async(json!({
"id": 922, "method": "Page.navigate", "sessionId": "SID-1",
"params": { "url": format!("http://{addr}/complete") }
}))
.await;
assert!(
take_response_by_id(&mut ctx, 922)["result"]
.get("errorText")
.is_none()
);
wait_until_frame_stopped_loading(&mut ctx, "TID-1").await;
let completed_html = loaded_page_html_for_test(&mut ctx).await;
assert!(completed_html.contains("complete page"));
ctx.process_and_wait_for_response_async(json!({
"id": 923, "method": "Page.stopLoading", "sessionId": "SID-1"
}))
.await;
assert_eq!(take_response_by_id(&mut ctx, 923)["result"], json!({}));
assert_eq!(loaded_page_html_for_test(&mut ctx).await, completed_html);
})
.await;
server.abort();
}
#[tokio::test(flavor = "multi_thread")]
async fn stop_loading_without_browser_context_returns_empty_result() {
let mut ctx = TestContext::new();
@@ -1826,7 +1826,7 @@ pub(super) fn finish_runtime_mutation_effects(
owner = ?binding.owner(),
element = ?binding.element(),
load_delay_token = ?binding.load_delay_token(),
settled,
?settled,
"settled invalidated connected-style lease at mutation commit"
);
}
+2 -1
View File
@@ -78,7 +78,8 @@ pub(crate) use lifecycle_tasks::{
MainDocumentImageLoadDelayBinding, MainDocumentInteractiveLifecycleAction,
MainDocumentMediaLoadDelayBinding, MainDocumentScriptLoadDelayKind,
MainDocumentScriptLoadDelayLease, MainDocumentScriptLoadDelayRelease,
MainDocumentStyleLoadEventBinding, StylesheetSubresourceLoadDelayBinding,
MainDocumentStyleLoadEventBinding, MainDocumentStyleLoadEventSettlement,
StylesheetSubresourceLoadDelayBinding,
};
pub(crate) use load_delivery_tasks::{
FrameDocumentLoadDeliveryAction, FrameDocumentLoadDeliveryAdmission,
@@ -94,6 +94,11 @@ impl DocumentLifecycleBlockers {
self.window_load.clear();
}
pub(super) fn cancel_for_stop(&mut self) {
self.parser_deferred_scripts.clear();
self.window_load.cancel();
}
#[cfg(test)]
pub(super) fn len(&self) -> usize {
self.parser_deferred_scripts.len() + self.window_load.len()
@@ -81,6 +81,13 @@ impl MainDocumentScriptLoadDelayLease {
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum MainDocumentStyleLoadEventSettlement {
NotOwned,
Released,
CancelledAfterStop,
}
/// Exact main-document ownership for a connected `<style>`/`<link>` load or
/// error event.
///
@@ -12,6 +12,7 @@ use super::records::{DocumentLoadDelayReason, DocumentLoadDelayTokenId};
#[derive(Debug, Default)]
pub(super) struct DocumentLoadGate {
active: BTreeMap<DocumentLoadDelayTokenId, DocumentLoadDelayReason>,
cancelled: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -19,6 +20,9 @@ pub(super) enum DocumentLoadGateRelease {
NotOwned,
StillBlocked,
BecameUnblocked,
/// The exact lease was outstanding when loading stopped. This consumes
/// its terminal ownership but never signals a new load transition.
CancelledAfterStop,
}
impl DocumentLoadGateRelease {
@@ -33,7 +37,10 @@ impl DocumentLoadGate {
token: DocumentLoadDelayTokenId,
reason: DocumentLoadDelayReason,
) -> bool {
if !reason.blocks_window_load_directly() || self.active.contains_key(&token) {
if self.cancelled
|| !reason.blocks_window_load_directly()
|| self.active.contains_key(&token)
{
return false;
}
self.active.insert(token, reason);
@@ -49,7 +56,9 @@ impl DocumentLoadGate {
return DocumentLoadGateRelease::NotOwned;
}
self.active.remove(&token);
if self.active.is_empty() {
if self.cancelled {
DocumentLoadGateRelease::CancelledAfterStop
} else if self.active.is_empty() {
DocumentLoadGateRelease::BecameUnblocked
} else {
DocumentLoadGateRelease::StillBlocked
@@ -61,23 +70,34 @@ impl DocumentLoadGate {
token: DocumentLoadDelayTokenId,
reason: DocumentLoadDelayReason,
) -> bool {
self.active.get(&token) == Some(&reason)
!self.cancelled && self.active.get(&token) == Some(&reason)
}
pub(super) fn owns_any(&self, token: DocumentLoadDelayTokenId) -> bool {
// Cancelled leases still reserve their identity until their terminal
// arrives; unlike `owns`, this is a namespace collision check.
self.active.contains_key(&token)
}
pub(super) fn has_reason(&self, reason: DocumentLoadDelayReason) -> bool {
self.active.values().any(|candidate| *candidate == reason)
!self.cancelled && self.active.values().any(|candidate| *candidate == reason)
}
pub(super) fn is_blocked(&self) -> bool {
!self.active.is_empty()
!self.cancelled && !self.active.is_empty()
}
pub(super) fn clear(&mut self) {
self.active.clear();
self.cancelled = false;
}
/// Stop blocking lifecycle progress without forgetting which exact leases
/// still have a terminal owner. Reuse the existing ledger: a late terminal
/// consumes its cancelled lease once, while a duplicate remains NotOwned.
/// Document retirement drops any terminals that never arrive.
pub(super) fn cancel(&mut self) {
self.cancelled = true;
}
pub(super) fn release_all_document_script_delays(&mut self) -> usize {
@@ -96,6 +116,67 @@ impl DocumentLoadGate {
mod tests {
use super::*;
#[test]
fn stop_cancels_only_outstanding_exact_leases_without_reopening_load() {
let mut gate = DocumentLoadGate::default();
let settled = DocumentLoadDelayTokenId(1);
let cancelled = DocumentLoadDelayTokenId(2);
let other = DocumentLoadDelayTokenId(3);
assert!(gate.acquire(settled, DocumentLoadDelayReason::StyleLoadEvent));
assert!(gate.acquire(cancelled, DocumentLoadDelayReason::StyleLoadEvent));
assert!(gate.acquire(other, DocumentLoadDelayReason::AsyncClassicScript));
assert_eq!(
gate.release(settled, DocumentLoadDelayReason::StyleLoadEvent),
DocumentLoadGateRelease::StillBlocked
);
gate.cancel();
gate.cancel();
assert!(!gate.is_blocked());
assert!(!gate.has_reason(DocumentLoadDelayReason::StyleLoadEvent));
assert!(!gate.owns(cancelled, DocumentLoadDelayReason::StyleLoadEvent));
assert!(!gate.acquire(
DocumentLoadDelayTokenId(4),
DocumentLoadDelayReason::StyleLoadEvent
));
assert_eq!(
gate.release(settled, DocumentLoadDelayReason::StyleLoadEvent),
DocumentLoadGateRelease::NotOwned
);
assert_eq!(
gate.release(cancelled, DocumentLoadDelayReason::Image),
DocumentLoadGateRelease::NotOwned
);
assert_eq!(
gate.release(cancelled, DocumentLoadDelayReason::StyleLoadEvent),
DocumentLoadGateRelease::CancelledAfterStop
);
assert_eq!(
gate.release(cancelled, DocumentLoadDelayReason::StyleLoadEvent),
DocumentLoadGateRelease::NotOwned
);
assert_eq!(
gate.release(other, DocumentLoadDelayReason::AsyncClassicScript),
DocumentLoadGateRelease::CancelledAfterStop
);
assert_eq!(gate.len(), 0);
assert!(!gate.is_blocked());
}
#[test]
fn retirement_discards_cancelled_leases() {
let mut gate = DocumentLoadGate::default();
let token = DocumentLoadDelayTokenId(1);
assert!(gate.acquire(token, DocumentLoadDelayReason::StyleLoadEvent));
gate.cancel();
gate.clear();
assert_eq!(gate.len(), 0);
assert_eq!(
gate.release(token, DocumentLoadDelayReason::StyleLoadEvent),
DocumentLoadGateRelease::NotOwned
);
}
#[test]
fn release_reports_only_the_last_exact_token_as_unblocking() {
let mut gate = DocumentLoadGate::default();
@@ -760,7 +760,7 @@ impl DocumentLifecycleRecord {
let previous_readiness = self.readiness?;
let ready_state_changed = previous_readiness != DocumentReadinessState::Complete;
self.blockers.clear_for_retirement();
self.blockers.cancel_for_stop();
self.incomplete_child_frames.clear();
self.parsing_delay_token = None;
self.interactive_transition_token = None;
@@ -10,7 +10,8 @@ use super::lifecycle_tasks::{
MainDocumentImageLoadDelayBinding, MainDocumentInteractiveLifecycleAction,
MainDocumentMediaLoadDelayBinding, MainDocumentScriptLoadDelayKind,
MainDocumentScriptLoadDelayLease, MainDocumentScriptLoadDelayRelease,
MainDocumentStyleLoadEventBinding, StylesheetSubresourceLoadDelayBinding,
MainDocumentStyleLoadEventBinding, MainDocumentStyleLoadEventSettlement,
StylesheetSubresourceLoadDelayBinding,
};
use super::load_event_gate::DocumentLoadGateRelease;
use super::module_clients::{
@@ -2675,6 +2676,9 @@ impl FrameOwnerStore {
DocumentLoadGateRelease::BecameUnblocked => {
MainDocumentScriptLoadDelayRelease::BecameUnblocked
}
DocumentLoadGateRelease::CancelledAfterStop => {
MainDocumentScriptLoadDelayRelease::AlreadyUnblocked
}
}
}
@@ -2723,18 +2727,29 @@ impl FrameOwnerStore {
pub(crate) fn settle_main_style_load_event(
&mut self,
binding: MainDocumentStyleLoadEventBinding,
) -> bool {
) -> MainDocumentStyleLoadEventSettlement {
if !self.main_style_load_event_is_current(binding) {
return false;
return MainDocumentStyleLoadEventSettlement::NotOwned;
}
let Some(token) = binding.load_delay_token() else {
return true;
return MainDocumentStyleLoadEventSettlement::Released;
};
self.release_document_load_delay(
binding.owner(),
token,
DocumentLoadDelayReason::StyleLoadEvent,
)
let document = self
.documents
.get_mut(&binding.owner().document_id)
.expect("current style event owner must have a Document");
match document
.lifecycle_progress
.release_window_load_delay(token, DocumentLoadDelayReason::StyleLoadEvent)
{
DocumentLoadGateRelease::NotOwned => MainDocumentStyleLoadEventSettlement::NotOwned,
DocumentLoadGateRelease::StillBlocked | DocumentLoadGateRelease::BecameUnblocked => {
MainDocumentStyleLoadEventSettlement::Released
}
DocumentLoadGateRelease::CancelledAfterStop => {
MainDocumentStyleLoadEventSettlement::CancelledAfterStop
}
}
}
pub(crate) fn accept_current_main_stylesheet_subresource_load_delay(
@@ -698,9 +698,13 @@ fn main_style_load_event_binding_delays_complete_until_event_settlement() {
"the exact style event token must block complete"
);
assert!(store.settle_main_style_load_event(binding));
assert!(
!store.settle_main_style_load_event(binding),
assert_eq!(
store.settle_main_style_load_event(binding),
MainDocumentStyleLoadEventSettlement::Released
);
assert_eq!(
store.settle_main_style_load_event(binding),
MainDocumentStyleLoadEventSettlement::NotOwned,
"one style event binding must settle exactly once"
);
assert_eq!(
@@ -742,6 +746,57 @@ fn main_modulepreload_owner_never_allocates_a_load_event_delay() {
);
}
#[test]
fn stopped_style_lease_settlement_is_exact_and_does_not_reopen_lifecycle() {
let mut store = FrameOwnerStore::default();
store.ensure_main_frame(
handle(1),
url("https://example.test/"),
url("https://example.test/"),
"https://example.test".to_owned(),
policy_container(),
policy_context(),
None,
);
let snapshot = store.current_main_owner_snapshot().expect("main owner");
let owner = FrameDocumentTaskOwner::new(
snapshot.scheduler_lane_id,
snapshot.local_window_id,
snapshot.document_id,
);
let binding = store
.accept_current_main_style_load_event(owner, handle(8))
.expect("style lease");
assert!(binding.load_delay_token().is_some());
assert_eq!(store.stop_current_main_document_loading(owner), Some(true));
assert_eq!(store.stop_current_main_document_loading(owner), Some(false));
assert!(store.main_style_load_event_is_current(binding));
assert_eq!(
store.settle_main_style_load_event(binding),
MainDocumentStyleLoadEventSettlement::CancelledAfterStop
);
assert_eq!(
store.settle_main_style_load_event(binding),
MainDocumentStyleLoadEventSettlement::NotOwned
);
assert_eq!(
store.current_main_document_complete_transition_is_ready(owner),
Some(false)
);
let after_stop = store
.accept_current_main_style_load_event(owner, handle(9))
.expect("post-stop style event remains usable");
assert!(after_stop.load_delay_token().is_none());
assert_eq!(
store.settle_main_style_load_event(after_stop),
MainDocumentStyleLoadEventSettlement::Released
);
assert_eq!(
store.current_main_document_complete_transition_is_ready(owner),
Some(false)
);
}
#[test]
fn main_style_load_event_binding_cannot_settle_replacement_document() {
let mut store = FrameOwnerStore::default();
@@ -774,8 +829,9 @@ fn main_style_load_event_binding_cannot_settle_replacement_document() {
.expect("main replacement");
assert!(!store.main_style_load_event_is_current(stale));
assert!(
!store.settle_main_style_load_event(stale),
assert_eq!(
store.settle_main_style_load_event(stale),
MainDocumentStyleLoadEventSettlement::NotOwned,
"stale style completion must not mutate replacement lifecycle"
);
assert_eq!(
@@ -225,6 +225,7 @@ pub(crate) struct ParserResumePermit {
pub(crate) enum ParserStopReason {
DocumentReplacement,
MainResourceLoadFailure,
Stopped,
OwnerDropped,
}
@@ -1055,7 +1055,7 @@ impl JsContextHost {
owner = ?binding.owner(),
element = ?binding.element(),
load_delay_token = ?binding.load_delay_token(),
settled,
?settled,
"settled invalidated connected-style lease at CSSOM commit"
);
}
@@ -6,6 +6,7 @@ use crate::frame_owner_model::{
MainDocumentLoadCompletionState, MainDocumentMediaLoadDelayBinding,
MainDocumentScriptLoadDelayKind, MainDocumentScriptLoadDelayLease,
MainDocumentScriptLoadDelayRelease, MainDocumentStyleLoadEventBinding,
MainDocumentStyleLoadEventSettlement,
};
use crate::{
document_runtime::{
@@ -210,7 +211,7 @@ impl JsContextHost {
pub(crate) fn settle_main_style_load_event(
&mut self,
binding: MainDocumentStyleLoadEventBinding,
) -> bool {
) -> MainDocumentStyleLoadEventSettlement {
self.frame_owner_store.settle_main_style_load_event(binding)
}
@@ -327,6 +327,9 @@ pub(crate) enum PageConnectedStyleLoadDelayEffect {
/// The event task released the exact load-delay token captured when the
/// connected style/link operation was admitted.
ReleasedExactBinding,
/// Loading stopped while this exact lease was outstanding. Its terminal
/// consumed the cancelled lease without reopening the Window-load gate.
ExactBindingCancelledByStop,
/// Event dispatch synchronously replaced the binding's Document.
///
/// Document replacement owns retirement of the old ledger. The selected
@@ -714,6 +714,25 @@ impl LivePageEntry {
.ok_or_else(|| anyhow!("renderer page has no pending phase-one navigation to resume"))
}
pub(super) fn stop_pending_main_document_loading(&mut self) {
let Some(pending) = self.pending_phase_one_navigation.take() else {
return;
};
let (residence, mut metadata) = pending.into_parts();
let browser_context_runtime = residence
.page_vm()
.runtime_hooks
.browser_context_runtime
.clone();
let page_vm = residence.into_stopped_page_vm();
metadata.reject(
None,
&browser_context_runtime,
"Navigation stopped".to_owned(),
);
self.install_resumed_phase_one_page_vm(page_vm);
}
pub(super) fn reject_pending_phase_one_navigation_in_place(&mut self, message: &str) {
let Some(mut pending) = self.pending_phase_one_navigation.take() else {
return;
@@ -2563,6 +2563,9 @@ impl RendererOwnerLocalStore {
// silently writes into the asynchronous Page journal and lets its
// response overtake the resulting protocol fact.
let command_turn_output_scope = entry.page_vm_mut().begin_command_turn_output_scope()?;
if matches!(&command, RendererPageCommand::StopDocumentLifecycle) {
entry.stop_pending_main_document_loading();
}
let replacement_lifecycle_snapshot = entry
.page_vm()
.document_replacement_lifecycle_action_snapshot();
@@ -42,6 +42,106 @@ fn take_next_link_element_event_task_for_test(
Some(task)
}
#[tokio::test(flavor = "current_thread")]
async fn connected_style_event_queued_before_stop_settles_without_reopening_load() {
run_page_vm_async_test(async move {
let loader = crate::network::ResourceRequestClient::new(&FetchConfig::default())?;
let (mut page_vm, _resource_source, _wake_rx) =
page_vm_with_bound_task_sources_and_owner_wake(
&loader,
Url::parse("https://example.com/style-stop")?,
);
let owner = page_vm.vm().current_main_document_task_owner();
page_vm.vm_mut().eval(
r#"
globalThis.__styleStopEvents = [];
window.addEventListener('load', () => __styleStopEvents.push('window-load'));
const style = document.createElement('style');
style.textContent = 'body { color: maroon; }';
style.addEventListener('load', () => __styleStopEvents.push('style-load'));
document.head.append(style);
'queued'
"#,
)?;
page_vm
.vm_mut()
.prime_document_lifecycle_processing_and_record_stylesheet_network_results();
// Hold the exact already-admitted task across the command, just as a
// terminal queued before stop can be selected after the stop turn.
let task = take_next_style_element_event_task_for_test(&mut page_vm)
.expect("connected style task must be ready before stopping");
assert_ne!(
page_vm.vm_mut().eval("document.readyState")?,
"complete",
"the admitted style task must own a pre-complete load-delay lease"
);
page_vm.dispatch_renderer_page_command(RendererPageCommand::StopDocumentLifecycle)?;
let body = page_vm.apply_selected_page_connected_style_event_turn(task)?;
assert!(matches!(
body.action.target_effect,
PageConnectedStyleEventTargetEffect::DispatchedToCurrentOwner { .. }
));
assert_eq!(page_vm.vm().current_main_document_task_owner(), owner);
assert_eq!(
page_vm
.vm_mut()
.eval("document.readyState + '|' + __styleStopEvents.join(',')")?,
"complete|style-load"
);
assert!(page_vm.document_lifecycle.current_snapshot().load.is_none());
Ok::<_, anyhow::Error>(())
})
.await
.expect("stopping must not invalidate a live style task's exact settlement");
}
#[tokio::test(flavor = "current_thread")]
async fn stopped_document_style_event_cannot_settle_replacement_document() {
run_page_vm_async_test(async move {
let loader = crate::network::ResourceRequestClient::new(&FetchConfig::default())?;
let (mut page_vm, _resource_source, _wake_rx) =
page_vm_with_bound_task_sources_and_owner_wake(
&loader,
Url::parse("https://example.com/style-stop-replace")?,
);
let owner = page_vm.vm().current_main_document_task_owner();
page_vm.vm_mut().eval(
r#"
globalThis.__retiredStyleEvent = false;
const style = document.createElement('style');
style.textContent = 'body { color: maroon; }';
style.onload = () => { __retiredStyleEvent = true; };
document.head.append(style);
'queued'
"#,
)?;
page_vm
.vm_mut()
.prime_document_lifecycle_processing_and_record_stylesheet_network_results();
let task = take_next_style_element_event_task_for_test(&mut page_vm)
.expect("style event must be admitted before stop");
page_vm.dispatch_renderer_page_command(RendererPageCommand::StopDocumentLifecycle)?;
page_vm
.vm_mut()
.eval("document.open(); document.write('<main>replacement</main>'); 'replaced'")?;
assert_ne!(page_vm.vm().current_main_document_task_owner(), owner);
let body = page_vm.apply_selected_page_connected_style_event_turn(task)?;
assert_eq!(
body.action.target_effect,
PageConnectedStyleEventTargetEffect::DiscardedStaleOwner
);
assert_eq!(
page_vm
.vm_mut()
.eval("document.readyState + '|' + __retiredStyleEvent")?,
"loading|false"
);
Ok::<_, anyhow::Error>(())
})
.await
.expect("stopped style event must not affect a replacement Document");
}
async fn wait_for_stylesheet_source(
wake_rx: &mut tokio::sync::mpsc::UnboundedReceiver<crate::page_task_queue::RendererOwnerWake>,
expected: RendererOwnerWakeSource,
@@ -69,6 +69,14 @@ pub(in crate::runtime) enum PendingPhaseOneResidence {
}
impl PendingPhaseOneResidence {
pub(in crate::runtime) fn into_stopped_page_vm(self) -> PageVm {
match self {
Self::ParserBlockingSourceLoad { runtime, .. }
| Self::ClosedInputPageWork { runtime, .. } => runtime.into_stopped_page_vm(),
Self::OpenStreaming(continuation) => continuation.into_stopped_page_vm(),
}
}
pub(in crate::runtime) fn parser_blocking_source_load(
runtime: Box<ConcurrentParseTimeRuntime>,
started: Instant,
@@ -199,6 +199,12 @@ impl ConcurrentParseTimeRuntime {
self.page_vm
}
pub(super) fn into_stopped_page_vm(mut self) -> PageVm {
self.state.parser_session.stop(ParserStopReason::Stopped);
drop(self.retire_main_parser_continuation());
self.page_vm
}
pub(super) fn new_parser_owner(
loader: ResourceRequestClient,
stage: PageVmInitStage,
@@ -491,6 +491,7 @@ pub struct ExternalRawDocumentBodyStream {
body_chunks: mpsc::Receiver<Vec<u8>>,
completion: Option<oneshot::Receiver<Result<()>>>,
page_creation_progress: Option<crate::runtime::RendererPageCreationProgress>,
stop_loading_cancellation: Option<moli_fetch::FetchCancelHandle>,
}
impl ExternalRawDocumentBodyStream {
@@ -517,6 +518,7 @@ impl ExternalRawDocumentBodyStream {
body_chunks,
completion: Some(completion),
page_creation_progress: Some(page_creation_progress),
stop_loading_cancellation: None,
},
)
}
@@ -529,9 +531,21 @@ impl ExternalRawDocumentBodyStream {
body_chunks,
completion: Some(completion),
page_creation_progress: None,
stop_loading_cancellation: None,
}
}
/// Shares this exact transfer's cancellation authority with an explicit
/// document stop. Ordinary parser retirement does not cancel the external
/// producer, which may still own browser-side response capture.
pub fn with_stop_loading_cancellation(
mut self,
cancellation: moli_fetch::FetchCancelHandle,
) -> Self {
self.stop_loading_cancellation = Some(cancellation);
self
}
pub fn from_bytes(body: Vec<u8>) -> Self {
let (completion_tx, completion_rx) = oneshot::channel();
let (body_tx, body_stream) = Self::channel(completion_rx);
@@ -556,6 +570,13 @@ pub(super) enum RawDocumentBodySource {
}
impl RawDocumentBodySource {
pub(super) fn stop_loading_cancellation(&self) -> Option<moli_fetch::FetchCancelHandle> {
match self {
Self::FetchResponse(response) => Some(response.cancellation_handle()),
Self::External(source) => source.stop_loading_cancellation.clone(),
}
}
fn fetch_response(response: Box<StreamingRawResponse>) -> Self {
Self::FetchResponse(response)
}
@@ -41,6 +41,7 @@ impl StreamingDocumentInputSender {
/// parked in the owner-local Page slot.
pub(super) struct StreamingDocumentInputSource {
rx: mpsc::Receiver<StreamingDocumentInputEvent>,
stop_loading_cancellation: Option<moli_fetch::FetchCancelHandle>,
}
impl StreamingDocumentInputSource {
@@ -49,6 +50,7 @@ impl StreamingDocumentInputSource {
parser_continuation: RendererPageMainParserContinuationProducer,
task_runner: crate::network::RendererResourceTaskRunner,
) -> Self {
let stop_loading_cancellation = raw_body.stop_loading_cancellation();
let (tx, rx) = mpsc::channel(STREAMING_DOCUMENT_INPUT_BUFFERED_EVENTS);
let sender = StreamingDocumentInputSender {
tx,
@@ -82,7 +84,20 @@ impl StreamingDocumentInputSource {
.send(StreamingDocumentInputEvent::Finished(terminal))
.await;
});
Self { rx }
Self {
rx,
stop_loading_cancellation,
}
}
pub(super) fn stop_loading(self) {
if let Some(cancellation) = self.stop_loading_cancellation.as_ref()
&& !cancellation.response_completion_is_committed()
{
cancellation.cancel();
}
// Dropping the receiver retires queued bytes and wakes a bridge parked
// on input. This is separate from the transport's terminal result.
}
pub(super) fn has_ready_input(&mut self) -> bool {
@@ -275,6 +290,45 @@ mod tests {
));
}
#[tokio::test]
async fn external_body_transfer_is_cancelled_only_by_explicit_stop() {
for explicit_stop in [false, true] {
let cancellation = FetchCancelHandle::new();
let (completion_tx, completion_rx) = tokio::sync::oneshot::channel();
let (body_tx, raw_body) = ExternalRawDocumentBodyStream::channel(completion_rx);
let raw_body = raw_body.with_stop_loading_cancellation(cancellation.clone());
let (_networking, continuation, mut wake_rx) = continuation_fixture(95);
let mut source = StreamingDocumentInputSource::bridge(
RawDocumentBodySource::External(raw_body),
continuation,
crate::network::RendererResourceTaskRunner::from_current_tokio().unwrap(),
);
body_tx
.send(b"queued unparsed tail".to_vec())
.await
.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(1), wake_rx.recv())
.await
.expect("queued body must wake its parser owner")
.expect("parser wake channel must stay open");
assert!(source.has_ready_input());
if explicit_stop {
source.stop_loading();
} else {
drop(source);
}
tokio::time::timeout(std::time::Duration::from_secs(1), body_tx.closed())
.await
.expect("retired parser must release external input and queued tail");
assert_eq!(
cancellation.is_cancelled(),
explicit_stop,
"ordinary external parser retirement must leave browser response capture running"
);
assert!(completion_tx.send(Ok(())).is_err());
}
}
#[tokio::test]
async fn dropping_input_source_cancels_fetch_waiting_for_another_body_chunk() {
let (body_tx, completion_tx, cancel_handle, raw_body) = pending_fetch_body();
@@ -28,6 +28,12 @@ pub(in crate::runtime) struct PendingStreamingPhaseOneContinuation {
}
impl PendingStreamingPhaseOneContinuation {
pub(in crate::runtime) fn into_stopped_page_vm(self) -> PageVm {
let Self { runtime, input, .. } = self;
input.stop_loading();
runtime.into_stopped_page_vm()
}
pub(super) fn publish_pending_page_creation_phase(&self) {
self.runtime.publish_pending_page_creation_phase();
}
@@ -21,6 +21,77 @@ struct OpenStreamingPage {
activity_wake_rx: RendererExternalActivityTestReceiver,
}
#[tokio::test(flavor = "multi_thread")]
async fn explicit_stop_retires_open_parser_and_preserves_committed_partial_document() {
let cancellation = moli_fetch::FetchCancelHandle::new();
let mut stream = OpenStreamingPage::create_with_cancellation(
"https://example.test",
"<!doctype html><body><main id='prefix'>committed prefix</main>",
Some(cancellation.clone()),
)
.await;
let (prefix, _) = tokio::time::timeout(Duration::from_secs(5), stream.page.run_async_command(
RendererPageCommand::EvaluateExpression {
expression: "new Promise(resolve => { const check = () => document.getElementById('prefix') ? resolve(document.getElementById('prefix').textContent) : setTimeout(check, 0); check(); })".to_owned(),
await_promise: true,
},
)).await.unwrap().unwrap();
assert_eq!(
renderer_json_value(prefix),
Some(serde_json::json!("committed prefix"))
);
let original_agent = stream.page.devtools_agent_token();
stream
.page
.run_async_command(RendererPageCommand::StopDocumentLifecycle)
.await
.unwrap();
assert!(cancellation.is_cancelled());
tokio::time::timeout(Duration::from_secs(2), stream.body_tx.closed())
.await
.expect("stop must retire the parser's input receiver");
assert!(
stream
.body_tx
.send(b"<main id='tail'>must not parse</main>".to_vec())
.await
.is_err()
);
let (state, _) = stream.page.run_async_command(RendererPageCommand::EvaluateExpression {
expression: "[document.getElementById('prefix').textContent, !!document.getElementById('tail'), document.readyState].join('|')".to_owned(),
await_promise: false,
}).await.unwrap();
assert_eq!(
renderer_json_value(state),
Some(serde_json::json!("committed prefix|false|complete"))
);
assert_eq!(stream.page.devtools_agent_token(), original_agent);
// Repeated stops and new ordinary Page tasks must not destroy the owner.
stream
.page
.run_async_command(RendererPageCommand::StopDocumentLifecycle)
.await
.unwrap();
let (timer, _) = tokio::time::timeout(
Duration::from_secs(2),
stream
.page
.run_async_command(RendererPageCommand::EvaluateExpression {
expression: "new Promise(resolve => setTimeout(() => resolve('still live'), 0))"
.to_owned(),
await_promise: true,
}),
)
.await
.unwrap()
.unwrap();
assert_eq!(
renderer_json_value(timer),
Some(serde_json::json!("still live"))
);
stream.page.close_async().await.unwrap();
}
#[tokio::test(flavor = "multi_thread")]
async fn owner_loop_executes_async_script_while_main_document_stream_remains_open() {
assert_open_stream_work_executes_before_eof(ASYNC_SCRIPT_HTML, "async script").await;
@@ -160,6 +231,14 @@ async fn assert_open_stream_work_executes_before_eof(html: &str, work_label: &st
impl OpenStreamingPage {
async fn create(base_url: &str, html: &str) -> Self {
Self::create_with_cancellation(base_url, html, None).await
}
async fn create_with_cancellation(
base_url: &str,
html: &str,
cancellation: Option<moli_fetch::FetchCancelHandle>,
) -> Self {
let runtime_owner = JsRuntime::initialize();
let runtime = runtime_owner.handle();
let (activity_wake_tx, activity_wake_rx) = renderer_external_activity_test_channel();
@@ -170,6 +249,10 @@ impl OpenStreamingPage {
let page_url = url::Url::parse(&format!("{base_url}/page")).expect("page url");
let (completion_tx, completion_rx) = oneshot::channel();
let (body_tx, raw_body) = ExternalRawDocumentBodyStream::channel(completion_rx);
let raw_body = match cancellation {
Some(cancellation) => raw_body.with_stop_loading_cancellation(cancellation),
None => raw_body,
};
body_tx
.send(html.as_bytes().to_vec())
.await
@@ -1,4 +1,5 @@
use super::ScriptVm;
use crate::frame_owner_model::MainDocumentStyleLoadEventSettlement;
use crate::{
document_runtime::StylesheetImportCompletionAuthority,
page_task_queue::{
@@ -115,7 +116,7 @@ impl ScriptVm {
.borrow_mut()
.settle_main_style_load_event(binding);
assert!(
settled,
settled != MainDocumentStyleLoadEventSettlement::NotOwned,
"a current connected style event must release its exact load-delay binding"
);
tracing::debug!(
@@ -124,7 +125,9 @@ impl ScriptVm {
load_delay_token = ?binding.load_delay_token(),
"settled main connected style load inside selected event body"
);
if had_load_delay_token {
if settled == MainDocumentStyleLoadEventSettlement::CancelledAfterStop {
PageConnectedStyleLoadDelayEffect::ExactBindingCancelledByStop
} else if had_load_delay_token {
PageConnectedStyleLoadDelayEffect::ReleasedExactBinding
} else {
PageConnectedStyleLoadDelayEffect::NoBindingRequired
@@ -550,11 +550,11 @@ impl DocumentRuntime {
owner = ?binding.owner(),
element = ?binding.element(),
load_delay_token = ?binding.load_delay_token(),
settled,
?settled,
reason,
"settled main connected style load before event posting"
);
settled
settled != crate::frame_owner_model::MainDocumentStyleLoadEventSettlement::NotOwned
}
fn settle_connected_style_load_admission(