mirror of
https://github.com/lexmount/moli.git
synced 2026-09-28 08:01:37 +00:00
fix(cdp): preserve binary uploads through Fetch request pauses
Resume Window fetch and XHR uploads from their retained bytes rather than lossy display text. Track explicit body overrides across chained sessions and update the uploaded-body snapshot when modified. Expose the original request body at the Fetch pause using existing Network listener and collector visibility. Add five end-to-end regression cases covering pass-through, chained interception, text overrides, and pause-time capture.
This commit is contained in:
@@ -1823,6 +1823,7 @@ impl PendingSubresourceFetchRequest {
|
||||
}
|
||||
if let Some(body) = body {
|
||||
chain.body = Some(body);
|
||||
chain.body_overridden = true;
|
||||
}
|
||||
if let Some(headers) = headers {
|
||||
chain.headers = headers.headers().clone();
|
||||
@@ -1845,7 +1846,7 @@ impl PendingSubresourceFetchRequest {
|
||||
(
|
||||
Some(chain.url.clone()),
|
||||
Some(chain.method.clone()),
|
||||
Some(chain.body.clone()),
|
||||
chain.body_overridden.then(|| chain.body.clone()),
|
||||
chain.header_override.clone(),
|
||||
)
|
||||
}
|
||||
@@ -1869,6 +1870,8 @@ pub struct PendingSubresourceFetchRequestStageChain {
|
||||
pub method: String,
|
||||
pub headers: moli_fetch::RequestHeaders,
|
||||
pub body: Option<String>,
|
||||
/// A display snapshot is not an override of the renderer's binary body.
|
||||
pub body_overridden: bool,
|
||||
pub request_cookie_report: Option<StoredCookieQueryReport>,
|
||||
pub remaining_sessions: Vec<PendingSubresourceFetchRequestStage>,
|
||||
}
|
||||
|
||||
@@ -2430,6 +2430,7 @@ mod tests {
|
||||
method: "GET".to_owned(),
|
||||
headers: Vec::new().into(),
|
||||
body: None,
|
||||
body_overridden: false,
|
||||
request_cookie_report: None,
|
||||
remaining_sessions: vec![
|
||||
PendingSubresourceFetchRequestStage {
|
||||
@@ -2601,6 +2602,7 @@ mod tests {
|
||||
method: "GET".to_owned(),
|
||||
headers: Vec::new().into(),
|
||||
body: None,
|
||||
body_overridden: false,
|
||||
request_cookie_report: None,
|
||||
remaining_sessions: vec![
|
||||
PendingSubresourceFetchRequestStage {
|
||||
|
||||
@@ -2099,6 +2099,7 @@ mod protocol_neutral_tests {
|
||||
method: "GET".to_owned(),
|
||||
headers: vec![("x-old".to_owned(), "1".to_owned())].into(),
|
||||
body: None,
|
||||
body_overridden: false,
|
||||
request_cookie_report: None,
|
||||
remaining_sessions: vec![PendingSubresourceFetchRequestStage {
|
||||
session_id: Some("SID-attached".to_owned()),
|
||||
@@ -2383,6 +2384,7 @@ mod protocol_neutral_tests {
|
||||
method: "GET".to_owned(),
|
||||
headers: Vec::new().into(),
|
||||
body: None,
|
||||
body_overridden: false,
|
||||
request_cookie_report: None,
|
||||
remaining_sessions: vec![PendingSubresourceFetchRequestStage {
|
||||
session_id: Some("BIDI-SID".to_owned()),
|
||||
|
||||
@@ -319,6 +319,7 @@ async fn prepare_subresource_fetch_pause_sources_async(
|
||||
method: info.method.clone(),
|
||||
headers: info.request_headers.clone(),
|
||||
body: info.request_body.clone(),
|
||||
body_overridden: false,
|
||||
request_cookie_report: info.request_cookie_report.clone(),
|
||||
remaining_sessions,
|
||||
}));
|
||||
@@ -363,6 +364,17 @@ pub(crate) fn emit_subresource_fetch_pause_outputs(
|
||||
) {
|
||||
for output in outputs {
|
||||
let network_request_id = output.network_output().network_request_id().to_owned();
|
||||
if !network_session_ids.is_empty() {
|
||||
// A request is already observable at the Fetch pause. Its body
|
||||
// must be readable before the intercepted transport is resumed.
|
||||
network::record_subresource_request_body(
|
||||
conn,
|
||||
owner,
|
||||
&network_request_id,
|
||||
output.network_output().request_body_bytes(),
|
||||
network_session_ids,
|
||||
);
|
||||
}
|
||||
let network_events = network_session_ids
|
||||
.iter()
|
||||
.flat_map(|session_id| {
|
||||
|
||||
@@ -0,0 +1,181 @@
|
||||
use super::*;
|
||||
use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64_STANDARD};
|
||||
|
||||
const BODY: &[u8] = &[0, 255, 195, 40, 128, 65];
|
||||
|
||||
async fn binary_request_round_trip(
|
||||
xhr: bool,
|
||||
chained: bool,
|
||||
replacement: Option<&str>,
|
||||
inspect_paused_body: bool,
|
||||
) {
|
||||
let (sender, mut received) = tokio::sync::mpsc::unbounded_channel();
|
||||
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let addr = listener.local_addr().unwrap();
|
||||
let server = tokio::spawn(async move {
|
||||
axum::serve(
|
||||
listener,
|
||||
Router::new()
|
||||
.route(
|
||||
"/page",
|
||||
get(|| async { "<!doctype html><title>binary request</title>" }),
|
||||
)
|
||||
.route(
|
||||
"/echo",
|
||||
any(move |body: axum::body::Bytes| {
|
||||
let sender = sender.clone();
|
||||
async move {
|
||||
sender.send(body.to_vec()).unwrap();
|
||||
axum::Json(body.to_vec())
|
||||
}
|
||||
}),
|
||||
),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
});
|
||||
let mut ctx = TestContext::new();
|
||||
let page_url = format!("http://{addr}/page");
|
||||
let api_url = format!("http://{addr}/echo");
|
||||
with_loaded_http_document(&mut ctx, &page_url, "SID-1", "TID-1").await;
|
||||
let mut sessions = vec!["SID-1"];
|
||||
if chained {
|
||||
assert!(
|
||||
ctx.conn
|
||||
.browser_context
|
||||
.as_mut()
|
||||
.unwrap()
|
||||
.assign_attached_session_to_target("TID-1", "SID-2".to_owned())
|
||||
);
|
||||
sessions.push("SID-2");
|
||||
}
|
||||
ctx.sent.clear();
|
||||
let mut id = 1;
|
||||
for session in &sessions {
|
||||
for method in ["Network.enable", "Fetch.enable"] {
|
||||
ctx.process_async(json!({"id":id,"method":method,"sessionId":session,
|
||||
"params":{"patterns":[{"urlPattern":"*/echo","requestStage":"Request"}]}}))
|
||||
.await;
|
||||
ctx.expect_result(id, json!({}), Some(session));
|
||||
id += 1;
|
||||
}
|
||||
}
|
||||
let send = if xhr {
|
||||
"new Promise((resolve,reject)=>{const x=new XMLHttpRequest();x.open('POST','/echo');x.onload=()=>resolve(JSON.parse(x.responseText));x.onerror=reject;x.send(bytes);})"
|
||||
} else {
|
||||
"fetch('/echo',{method:'POST',body:bytes}).then(r=>r.json())"
|
||||
};
|
||||
ctx.process_async(json!({"id":id,"method":"Runtime.evaluate","sessionId":"SID-1",
|
||||
"params":{"expression":format!(
|
||||
"globalThis.binaryResult=null; const bytes=new Uint8Array({}); {send}.then(v=>globalThis.binaryResult=v); 'scheduled'",
|
||||
json!(BODY)),"returnByValue":true}})).await;
|
||||
let response = take_response_by_id(&mut ctx, id);
|
||||
assert_eq!(response["result"]["result"]["value"], "scheduled");
|
||||
id += 1;
|
||||
let mut network_id = Value::Null;
|
||||
for (index, session) in sessions.iter().enumerate() {
|
||||
let paused = runtime_fetch::wait_for_request_paused_on_session(
|
||||
&mut ctx,
|
||||
session,
|
||||
&api_url,
|
||||
None,
|
||||
"binary request pause",
|
||||
)
|
||||
.await;
|
||||
network_id = paused["params"]["networkId"].clone();
|
||||
assert!(
|
||||
received.try_recv().is_err(),
|
||||
"paused request must not reach the server"
|
||||
);
|
||||
if inspect_paused_body {
|
||||
ctx.process_async(json!({"id":id,"method":"Network.getRequestPostData",
|
||||
"sessionId":session,"params":{"requestId":network_id}}))
|
||||
.await;
|
||||
ctx.expect_result(
|
||||
id,
|
||||
json!({"postData":BASE64_STANDARD.encode(BODY),"base64Encoded":true}),
|
||||
Some(session),
|
||||
);
|
||||
id += 1;
|
||||
}
|
||||
let mut params = json!({"requestId":paused["params"]["requestId"]});
|
||||
if index == 0
|
||||
&& let Some(replacement) = replacement
|
||||
{
|
||||
params["postData"] = json!(BASE64_STANDARD.encode(replacement));
|
||||
}
|
||||
ctx.process_async(json!({"id":id,"method":"Fetch.continueRequest",
|
||||
"sessionId":session,"params":params}))
|
||||
.await;
|
||||
ctx.expect_result(id, json!({}), Some(session));
|
||||
id += 1;
|
||||
}
|
||||
wait_until_messages(
|
||||
&mut ctx,
|
||||
Some("SID-1"),
|
||||
"binary upload completion",
|
||||
|messages| {
|
||||
messages.iter().any(|m| {
|
||||
m["method"] == "Network.loadingFinished" && m["params"]["requestId"] == network_id
|
||||
})
|
||||
},
|
||||
)
|
||||
.await;
|
||||
let actual = received
|
||||
.try_recv()
|
||||
.expect("server must have consumed the body before completion");
|
||||
let expected = replacement.map(str::as_bytes).unwrap_or(BODY);
|
||||
assert_eq!(
|
||||
actual, expected,
|
||||
"Fetch.continueRequest must preserve the original bytes unless overridden"
|
||||
);
|
||||
assert!(
|
||||
received.try_recv().is_err(),
|
||||
"one continuation must send one POST"
|
||||
);
|
||||
ctx.process_async(
|
||||
json!({"id":id,"method":"Runtime.evaluate","sessionId":"SID-1",
|
||||
"params":{"expression":"globalThis.binaryResult","returnByValue":true}}),
|
||||
)
|
||||
.await;
|
||||
let response = take_response_by_id(&mut ctx, id);
|
||||
assert_eq!(response["result"]["result"]["value"], json!(expected));
|
||||
id += 1;
|
||||
ctx.process_async(
|
||||
json!({"id":id,"method":"Network.getRequestPostData","sessionId":"SID-1",
|
||||
"params":{"requestId":network_id}}),
|
||||
)
|
||||
.await;
|
||||
let body = replacement.map_or_else(|| BASE64_STANDARD.encode(BODY), str::to_owned);
|
||||
ctx.expect_result(
|
||||
id,
|
||||
json!({"postData":body,"base64Encoded":replacement.is_none()}),
|
||||
Some("SID-1"),
|
||||
);
|
||||
server.abort();
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn fetch_noop_preserves_binary_request_body() {
|
||||
binary_request_round_trip(false, false, None, false).await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn xhr_noop_preserves_binary_request_body() {
|
||||
binary_request_round_trip(true, false, None, false).await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn chained_fetch_noop_preserves_binary_request_body() {
|
||||
binary_request_round_trip(false, true, None, false).await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn chained_fetch_text_override_replaces_binary_request_body() {
|
||||
binary_request_round_trip(false, true, Some("replacement"), false).await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread")]
|
||||
async fn paused_fetch_exposes_binary_request_body() {
|
||||
binary_request_round_trip(false, false, None, true).await;
|
||||
}
|
||||
@@ -319,6 +319,7 @@ async fn spawn_digest_proxy(
|
||||
}
|
||||
|
||||
mod basics;
|
||||
mod binary_request;
|
||||
mod command_correlation;
|
||||
mod header_bytes;
|
||||
mod navigation_auth;
|
||||
|
||||
@@ -135,7 +135,7 @@ async fn wait_for_request_paused(ctx: &mut TestContext, url: &str, description:
|
||||
wait_for_request_paused_on_session(ctx, "SID-1", url, None, description).await
|
||||
}
|
||||
|
||||
async fn wait_for_request_paused_on_session(
|
||||
pub(super) async fn wait_for_request_paused_on_session(
|
||||
ctx: &mut TestContext,
|
||||
session_id: &str,
|
||||
url: &str,
|
||||
|
||||
@@ -76,7 +76,7 @@ pub(crate) use agent::{
|
||||
};
|
||||
pub(crate) use backlog::{
|
||||
NetworkBacklogProjectionContext, emit_pending_network_backlog_activity_background_events,
|
||||
emit_prepared_renderer_network_live_background_events,
|
||||
emit_prepared_renderer_network_live_background_events, record_subresource_request_body,
|
||||
};
|
||||
pub(crate) use collectors::NetworkDataCollectorStore;
|
||||
pub(crate) use cookie_context::navigation_cookie_request_context;
|
||||
|
||||
@@ -861,7 +861,7 @@ fn record_subresource_response_body_source(
|
||||
);
|
||||
}
|
||||
|
||||
fn record_subresource_request_body(
|
||||
pub(crate) fn record_subresource_request_body(
|
||||
conn: &mut CdpConnection,
|
||||
owner: &CommandOwnerScope,
|
||||
request_id: &str,
|
||||
|
||||
@@ -67,6 +67,7 @@ pub(crate) struct TargetSubresourceFetchPauseNetworkOutput {
|
||||
method: String,
|
||||
request_headers: Vec<(String, String)>,
|
||||
request_body: Option<String>,
|
||||
request_body_bytes: Option<Vec<u8>>,
|
||||
resource_type: SubresourceResourceType,
|
||||
request_cookie_report: Option<StoredCookieQueryReport>,
|
||||
blocked_intercepts: Vec<DevToolsNetworkInterceptId>,
|
||||
@@ -93,6 +94,7 @@ impl TargetSubresourceFetchPauseNetworkOutput {
|
||||
method: info.method.clone(),
|
||||
request_headers: info.request_headers.to_byte_strings(),
|
||||
request_body: info.request_body.clone(),
|
||||
request_body_bytes: info.request_body_bytes.clone(),
|
||||
resource_type: info.resource_type,
|
||||
request_cookie_report: info.request_cookie_report.clone(),
|
||||
blocked_intercepts: Vec::new(),
|
||||
@@ -149,6 +151,10 @@ impl TargetSubresourceFetchPauseNetworkOutput {
|
||||
self.request_body.as_deref()
|
||||
}
|
||||
|
||||
pub(crate) fn request_body_bytes(&self) -> Option<&[u8]> {
|
||||
self.request_body_bytes.as_deref()
|
||||
}
|
||||
|
||||
pub(crate) fn resource_type(&self) -> SubresourceResourceType {
|
||||
self.resource_type
|
||||
}
|
||||
|
||||
@@ -1138,7 +1138,7 @@ impl ScriptVm {
|
||||
let PendingSubresourceFetchState {
|
||||
redirect_headers,
|
||||
request_origin,
|
||||
info,
|
||||
mut info,
|
||||
load,
|
||||
execution_context,
|
||||
credentials_mode,
|
||||
@@ -1148,6 +1148,11 @@ impl ScriptVm {
|
||||
continuation,
|
||||
deferred_request_started,
|
||||
} = pending;
|
||||
// The text snapshot is only for protocol display. Unless the client
|
||||
// explicitly replaces the body, retain the original binary upload.
|
||||
if let Some(body) = &body {
|
||||
info.request_body_bytes = body.as_ref().map(|body| body.as_bytes().to_vec());
|
||||
}
|
||||
let pending = match continuation {
|
||||
PendingSubresourceContinuation::WebSocket(connection) => {
|
||||
let request_url = url.unwrap_or_else(|| info.url.clone());
|
||||
@@ -1437,10 +1442,10 @@ impl ScriptVm {
|
||||
// up the ambient Page loader here would silently rebind policy/backend
|
||||
// to a newer Document identity.
|
||||
let loader = pending.load.request_client();
|
||||
let mut request = moli_fetch::Request::new(
|
||||
let mut request = moli_fetch::Request::new_bytes(
|
||||
&request_method,
|
||||
request_url.as_str(),
|
||||
request_body.clone(),
|
||||
pending.info.request_body_bytes.clone(),
|
||||
request_headers.clone(),
|
||||
)?
|
||||
.with_redirect_headers(pending.redirect_headers.clone())
|
||||
|
||||
Reference in New Issue
Block a user