fix(network): preserve request metadata and response evidence

This commit is contained in:
lanyue-llk
2026-09-22 13:23:43 +08:00
parent b0f237b095
commit db091f6976
27 changed files with 1904 additions and 142 deletions
+46
View File
@@ -66,6 +66,21 @@ where
self.limits
}
/// Changes budgets, returning oversized and then oldest entries that no
/// longer fit. Retained entries keep their original insertion order.
pub fn set_limits(&mut self, limits: ByteLimits) -> Vec<(K, V)> {
self.limits = limits;
let mut evicted = self.retain_entries(|_, entry| entry.byte_len <= limits.max_entry_bytes);
while self.used_bytes > limits.max_total_bytes {
if let Some(entry) = self.pop_oldest() {
evicted.push(entry);
} else {
break;
}
}
evicted
}
/// Returns the sum of the logical byte charges for retained entries.
pub fn used_bytes(&self) -> usize {
self.used_bytes
@@ -151,12 +166,43 @@ where
}
/// Removes all entries and releases all byte charges.
///
/// This does not return removed entries; use [`Self::retain`] to take ownership.
pub fn clear(&mut self) {
self.entries.clear();
self.insertion_order.clear();
self.used_bytes = 0;
}
/// Visits each value once in insertion order, returning removed entries in
/// that order and releasing their byte charges. Mutating a value does not
/// change its recorded charge; use [`Self::insert`] to replace its charge.
pub fn retain(&mut self, mut keep: impl FnMut(&K, &mut V) -> bool) -> Vec<(K, V)> {
self.retain_entries(|key, entry| keep(key, &mut entry.value))
}
fn retain_entries(
&mut self,
mut keep: impl FnMut(&K, &mut BufferedValue<V>) -> bool,
) -> Vec<(K, V)> {
let mut removed = Vec::new();
self.insertion_order.retain(|key| {
if self
.entries
.get_mut(key)
.is_some_and(|entry| keep(key, entry))
{
return true;
}
if let Some(entry) = self.entries.remove(key) {
self.used_bytes -= entry.byte_len;
removed.push((key.clone(), entry.value));
}
false
});
removed
}
fn pop_oldest(&mut self) -> Option<(K, V)> {
while let Some(key) = self.insertion_order.pop_front() {
let Some(entry) = self.entries.remove(&key) else {
+73
View File
@@ -1,5 +1,44 @@
use super::{BoundedByteBuffer, ByteLimits, InsertOutcome};
#[test]
fn shrinking_entry_limit_keeps_interleaved_survivors_in_fifo_order() {
let mut buffer = BoundedByteBuffer::new(ByteLimits::new(100, 20));
for (key, bytes) in [(0, 10), (1, 1), (2, 10), (3, 1), (4, 10)] {
buffer.insert(key, key, bytes);
}
assert_eq!(
buffer.set_limits(ByteLimits::new(100, 1)),
vec![(0, 0), (2, 2), (4, 4)]
);
assert_eq!(buffer.used_bytes(), 2);
assert_eq!(buffer.set_limits(ByteLimits::new(1, 1)), vec![(1, 1)]);
assert_eq!(buffer.get(&3), Some(&3));
assert_eq!(buffer.len(), 1);
}
#[test]
fn changing_limits_evicts_oversized_then_oldest_without_reordering_survivors() {
let mut buffer = BoundedByteBuffer::new(ByteLimits::new(20, 10));
buffer.insert("a", 1, 3);
buffer.insert("b", 2, 6);
buffer.insert("c", 3, 3);
assert_eq!(
buffer.set_limits(ByteLimits::new(4, 4)),
vec![("b", 2), ("a", 1)]
);
assert_eq!(buffer.used_bytes(), 3);
assert_eq!(buffer.get(&"c"), Some(&3));
assert!(buffer.set_limits(ByteLimits::new(10, 10)).is_empty());
assert_eq!(
buffer.insert("d", 4, 8),
InsertOutcome::Stored {
evicted: vec![("c", 3)]
}
);
assert_eq!(buffer.set_limits(ByteLimits::new(0, 0)), vec![("d", 4)]);
assert!(buffer.is_empty());
}
#[test]
fn accepts_entries_at_exact_limits() {
let mut buffer = BoundedByteBuffer::new(ByteLimits::new(4, 4));
@@ -115,3 +154,37 @@ fn remove_and_clear_return_all_byte_charges() {
assert!(buffer.is_empty());
assert_eq!(buffer.used_bytes(), 0);
}
#[test]
fn retain_visits_once_preserves_fifo_and_releases_only_removed_charges() {
let mut buffer = BoundedByteBuffer::new(ByteLimits::new(10, 10));
for key in 0..4 {
let _ = buffer.insert(key, key * 10, 2);
}
let mut visited = Vec::new();
let removed = buffer.retain(|key, value| {
visited.push(*key);
*value += 1;
key % 2 == 0
});
assert_eq!(visited, vec![0, 1, 2, 3]);
assert_eq!(removed, vec![(1, 11), (3, 31)]);
assert_eq!(buffer.used_bytes(), 4);
assert_eq!(buffer.get(&2), Some(&21));
assert_eq!(
buffer.insert(4, 40, 8),
InsertOutcome::Stored {
evicted: vec![(0, 1)]
}
);
assert!(buffer.retain(|_, _| true).is_empty());
assert_eq!(buffer.used_bytes(), 10);
assert_eq!(buffer.retain(|_, _| false), vec![(2, 21), (4, 40)]);
assert_eq!(buffer.used_bytes(), 0);
assert!(buffer.is_empty());
assert!(
buffer
.retain(|_, _| panic!("empty buffer must not visit"))
.is_empty()
);
}
+2 -1
View File
@@ -68,7 +68,8 @@ pub use request::{
RequestAuthScheme, RequestAuthTarget, RequestCacheMode, RequestCredentialsMode,
RequestHeaderOverride, RequestMode, RequestPriorityHints, RequestRedirectMode,
RequestResourceType, ResourceLoadPriority, ScriptFetchRequestMetadata,
ScriptFetchSchedulerPriority, SubresourceRequestMetadata,
ScriptFetchSchedulerPriority, SubresourceRequestMetadata, is_request_body_header_name,
redirect_status_rewrites_to_get,
};
pub use request_policy::{is_bad_port, should_request_be_blocked_due_to_bad_port};
pub use response::{
+33 -1
View File
@@ -14,6 +14,7 @@ const MAX_OBSERVED_HEADER_BLOCK_BYTES: usize = 256 * 1024;
/// Request headers observed after the transport has finalized an HTTP exchange.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct NetworkRequestObservation {
method: Option<String>,
headers: Vec<(String, String)>,
cookie_report: Option<StoredCookieQueryReport>,
truncated: bool,
@@ -22,6 +23,16 @@ pub struct NetworkRequestObservation {
impl NetworkRequestObservation {
pub fn new(headers: Vec<(String, String)>) -> Self {
Self {
method: None,
headers,
cookie_report: None,
truncated: false,
}
}
pub fn new_with_method(method: impl Into<String>, headers: Vec<(String, String)>) -> Self {
Self {
method: Some(method.into()),
headers,
cookie_report: None,
truncated: false,
@@ -31,12 +42,17 @@ impl NetworkRequestObservation {
fn from_header_block(data: &[u8], cookie_report: Option<StoredCookieQueryReport>) -> Self {
let (data, truncated) = bounded_header_data(data);
Self {
method: parse_request_method(data),
headers: parse_request_header_block(data),
cookie_report,
truncated,
}
}
pub fn method(&self) -> Option<&str> {
self.method.as_deref()
}
pub fn headers(&self) -> &[(String, String)] {
&self.headers
}
@@ -338,6 +354,14 @@ fn parse_request_header_block(data: &[u8]) -> Vec<(String, String)> {
.collect()
}
fn parse_request_method(data: &[u8]) -> Option<String> {
let line = String::from_utf8_lossy(data);
let line = line.lines().next()?.trim_end_matches('\r');
is_http_request_line(line)
.then(|| line.split_whitespace().next().map(str::to_owned))
.flatten()
}
fn is_http_request_line(line: &str) -> bool {
let mut fields = line.split_whitespace();
let (Some(method), Some(target), Some(version), None) =
@@ -648,7 +672,7 @@ mod tests {
fn recorder_preserves_redirect_exchange_order_and_raw_response_status() {
let recorder = NetworkObservationRecorder::default();
recorder.record_request_header_block(
b"GET /start HTTP/1.1\r\nHost: example.test\r\nAccept-Encoding: gzip\r\n\r\n",
b"POST /start HTTP/1.1\r\nHost: example.test\r\nAccept-Encoding: gzip\r\n\r\n",
);
recorder.record_response_header_line(b"HTTP/1.1 302 Found\r\n");
recorder.record_response_header_line(b"Location: /final\r\n");
@@ -662,6 +686,7 @@ mod tests {
let journal = recorder.snapshot();
assert_eq!(journal.exchanges().len(), 2);
assert_eq!(journal.exchanges()[0].request().method(), Some("POST"));
assert_eq!(
journal.exchanges()[0].request().headers(),
[
@@ -676,6 +701,13 @@ mod tests {
.status(),
302
);
assert_eq!(
journal
.final_request_observation()
.expect("final request")
.method(),
Some("GET")
);
assert_eq!(
journal
.final_request_observation()
+4 -2
View File
@@ -1122,12 +1122,14 @@ impl Request {
}
}
fn redirect_status_rewrites_to_get(status: u16, method: &str) -> bool {
/// Whether following this HTTP redirect replaces the request method with GET.
pub fn redirect_status_rewrites_to_get(status: u16, method: &str) -> bool {
status == 303 && !method.eq_ignore_ascii_case("GET") && !method.eq_ignore_ascii_case("HEAD")
|| matches!(status, 301 | 302) && method.eq_ignore_ascii_case("POST")
}
fn is_request_body_header_name(name: &str) -> bool {
/// Whether a header must be removed when a redirect discards the request body.
pub fn is_request_body_header_name(name: &str) -> bool {
matches!(
name.to_ascii_lowercase().as_str(),
"content-encoding"
+8
View File
@@ -3282,6 +3282,14 @@ pub struct ChildFrameDocumentNetworkSnapshot {
pub request_method: String,
#[serde(default)]
pub request_headers: Vec<(String, String)>,
/// Exact in-process transport observations for the actual child-navigation
/// hops. Serialized snapshots retain their historical lightweight shape.
#[serde(skip)]
pub network_observation_journal: moli_fetch::NetworkObservationJournal,
/// In-process browser-facing redirect metadata paired with the transport
/// journal. Serialized snapshots retain their historical lightweight shape.
#[serde(skip)]
pub redirect_chain: Vec<moli_fetch::RedirectInfo>,
pub final_url: String,
pub status: u16,
#[serde(default)]
@@ -134,6 +134,14 @@ impl TargetNetworkListenerOwnerMut<'_> {
}
self.remove_network_observation_cursor(listener_session_id.as_deref());
self.remove_captured_response_body_visibility_for_session(listener_session_id.as_deref());
if !self.is_attached_session()
&& let Some(primary_session_id) = self.target.session_id().map(str::to_owned)
{
// Primary commands share one listener whether routed on the root
// channel or its flattened session. Captured events can carry the
// explicit primary ID, so revoke both wire identities on disable.
self.remove_captured_response_body_visibility_for_session(Some(&primary_session_id));
}
self.clear_network_observation_artifacts_if_unobserved();
true
}
@@ -302,7 +302,7 @@ impl PageTargetHost {
needs_fetch_navigation_request_id: bool,
) -> (String, Option<String>, Option<String>) {
if clear_captured_response_bodies {
self.runtime_slot.clear_captured_response_bodies();
self.runtime_slot.prepare_response_bodies_for_navigation();
}
let mut allocator = self.runtime_slot.request_id_allocator();
let document_loader_id =
@@ -1133,6 +1133,19 @@ impl TargetRuntimeSlot {
self.network_agent.clear_captured_response_bodies();
}
pub(crate) fn configure_durable_response_bodies(
&mut self,
session_id: Option<&str>,
limits: Option<moli_bounded_buffer::ByteLimits>,
) {
self.network_agent
.configure_durable_response_bodies(session_id, limits);
}
pub(crate) fn prepare_response_bodies_for_navigation(&mut self) {
self.network_agent.prepare_response_bodies_for_navigation();
}
pub(crate) fn clear_network_body_artifacts(&mut self) {
self.network_agent.clear_body_artifacts();
}
+13
View File
@@ -44,6 +44,7 @@ mod load_resource;
mod main_document_progress;
mod output;
mod output_queue;
mod redirect_request;
mod response_body;
pub(crate) mod settings;
#[cfg(test)]
@@ -504,6 +505,14 @@ fn start_set_network_domain_enabled_command(
cmd: &Cmd<'_>,
enabled: bool,
) -> NetworkCommandTaskStep {
let durable_limits = if enabled {
match settings::durable_body_limits(cmd) {
Ok(limits) => limits,
Err(plan) => return NetworkCommandTaskStep::Complete(plan),
}
} else {
None
};
let updated = if enabled {
conn.enable_network_listener_for_session_owner(cmd.session_id)
} else {
@@ -516,6 +525,10 @@ fn start_set_network_domain_enabled_command(
));
}
if enabled && let Ok(slot) = conn.runtime_session_owner_slot_mut(cmd.session_id) {
slot.configure_durable_response_bodies(cmd.session_id, durable_limits);
}
let kind = if enabled {
PendingNetworkCommandKind::Enable
} else {
+302 -15
View File
@@ -1,4 +1,4 @@
use std::collections::{HashMap, HashSet};
use std::collections::{HashMap, HashSet, VecDeque};
use moli_bounded_buffer::{BoundedByteBuffer, ByteLimits, InsertOutcome};
use moli_core::page::{ScriptNetworkOutputItem, SubresourceNetworkRequestHandle};
@@ -15,6 +15,8 @@ use super::{
const RESPONSE_BODY_BUFFER_MAX_TOTAL_BYTES: usize = 20_000_000;
const RESPONSE_BODY_BUFFER_MAX_ENTRY_BYTES: usize = 2_000_000;
// Bound zero-byte bodies and payload-free failure/eviction metadata as well.
const DURABLE_RESPONSE_BODY_MAX_ENTRIES: usize = 4096;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RendererSubresourceTeardownDisposition {
@@ -705,6 +707,24 @@ impl TargetNetworkAgentState {
.clear_captured_response_bodies();
}
pub(crate) fn configure_durable_response_bodies(
&mut self,
session_id: Option<&str>,
limits: Option<ByteLimits>,
) {
self.artifacts
.body_artifacts
.captured_response_bodies
.configure_durable(session_id, limits);
}
pub(crate) fn prepare_response_bodies_for_navigation(&mut self) {
self.artifacts
.body_artifacts
.captured_response_bodies
.prepare_navigation();
}
pub(crate) fn clear_body_artifacts(&mut self) {
self.artifacts.body_artifacts.clear_session_scoped();
}
@@ -1336,6 +1356,9 @@ struct CapturedResponseBodyStore {
/// `buffered_bodies`, making the bounded buffer the single body owner.
bodies: HashMap<String, CapturedResponseBody>,
buffered_bodies: BoundedByteBuffer<String, CapturedResponseBody>,
durable_sessions: HashMap<Option<String>, ByteLimits>,
durable_entry_order: VecDeque<String>,
retained_request_ids: HashSet<String>,
}
impl Default for CapturedResponseBodyStore {
@@ -1352,9 +1375,136 @@ impl CapturedResponseBodyStore {
Self {
bodies: HashMap::new(),
buffered_bodies: BoundedByteBuffer::new(limits),
durable_sessions: HashMap::new(),
durable_entry_order: VecDeque::new(),
retained_request_ids: HashSet::new(),
}
}
fn configure_durable(&mut self, session_id: Option<&str>, limits: Option<ByteLimits>) {
let session = session_id.map(str::to_owned);
if let Some(limits) = limits {
self.durable_sessions.insert(session, limits);
} else {
self.durable_sessions.remove(&session);
self.revoke_retained_visibility(session_id);
}
for (id, mut body) in self
.buffered_bodies
.set_limits(self.effective_durable_limits())
{
body.mark_evicted();
self.bodies.insert(id, body);
}
self.synchronize_durable_entries();
}
fn revoke_retained_visibility(&mut self, session_id: Option<&str>) {
// Revocation affects only old-document claims, not ordinary access to
// responses from the current document.
let keep = |request_id: &String, body: &mut CapturedResponseBody| {
!self.retained_request_ids.contains(request_id)
|| body.remove_session_visibility(session_id)
};
self.bodies.retain(keep);
self.buffered_bodies.retain(keep);
self.prune_removed_entry_indexes();
}
fn prune_removed_entry_indexes(&mut self) {
let exists =
|id: &String| self.bodies.contains_key(id) || self.buffered_bodies.contains_key(id);
self.durable_entry_order.retain(exists);
self.retained_request_ids.retain(exists);
}
fn effective_durable_limits(&self) -> ByteLimits {
// One target owns one payload store. Multiple clients cannot multiply
// its budget; each requested budget is an upper bound, not a reserve.
self.durable_sessions.values().fold(
ByteLimits::new(
RESPONSE_BODY_BUFFER_MAX_TOTAL_BYTES,
RESPONSE_BODY_BUFFER_MAX_ENTRY_BYTES,
),
|acc, limits| {
ByteLimits::new(
acc.max_total_bytes.min(limits.max_total_bytes),
acc.max_entry_bytes.min(limits.max_entry_bytes),
)
},
)
}
fn synchronize_durable_entries(&mut self) {
if self.durable_sessions.is_empty() {
self.durable_entry_order.clear();
self.retained_request_ids.clear();
return;
}
let tracked: HashSet<_> = self.durable_entry_order.iter().collect();
let mut untracked = self
.bodies
.keys()
.chain(self.buffered_bodies.iter().map(|(id, _)| id))
.filter(|id| !tracked.contains(id))
.cloned()
.collect::<Vec<_>>();
untracked.sort();
self.durable_entry_order.extend(untracked);
self.trim_durable_entries();
}
fn track_durable_entry(&mut self, request_id: &str) {
if self.durable_sessions.is_empty() {
return;
}
self.durable_entry_order.retain(|id| id != request_id);
self.durable_entry_order.push_back(request_id.to_owned());
self.trim_durable_entries();
}
fn trim_durable_entries(&mut self) {
let mut expired_payloads = HashSet::new();
while self.durable_entry_order.len() > DURABLE_RESPONSE_BODY_MAX_ENTRIES {
if let Some(id) = self.durable_entry_order.pop_front() {
self.bodies.remove(&id);
self.retained_request_ids.remove(&id);
if self.buffered_bodies.contains_key(&id) {
expired_payloads.insert(id);
}
}
}
if !expired_payloads.is_empty() {
self.buffered_bodies
.retain(|id, _| !expired_payloads.contains(id));
}
}
fn prepare_navigation(&mut self) {
if self.durable_sessions.is_empty() {
self.clear();
return;
}
let retain = |body: &mut CapturedResponseBody| {
body.session_ids
.retain(|id| self.durable_sessions.contains_key(id));
// CDP durability does not opt BiDi collectors into retention.
body.collector_ids.clear();
!body.session_ids.is_empty()
&& !matches!(body.state, CapturedResponseBodyState::Pending)
};
self.bodies.retain(|_, body| retain(body));
self.buffered_bodies.retain(|_, body| retain(body));
self.retained_request_ids = self
.bodies
.keys()
.chain(self.buffered_bodies.iter().map(|(id, _)| id))
.cloned()
.collect();
self.durable_entry_order
.retain(|id| self.retained_request_ids.contains(id));
}
#[cfg(test)]
pub(crate) fn insert(
&mut self,
@@ -1380,6 +1530,7 @@ impl CapturedResponseBodyStore {
collection_was_gated: bool,
) {
let byte_len = body.len();
let tracked_request_id = request_id.clone();
let captured = CapturedResponseBody::from_captured_body_with_collector_scope(
body,
session_ids,
@@ -1402,6 +1553,7 @@ impl CapturedResponseBodyStore {
self.bodies.insert(request_id, rejected_body);
}
}
self.track_durable_entry(&tracked_request_id);
}
pub(crate) fn insert_pending_with_collector_scope(
@@ -1416,6 +1568,7 @@ impl CapturedResponseBodyStore {
{
return;
}
let tracked_request_id = request_id.clone();
self.bodies.insert(
request_id,
CapturedResponseBody::pending_with_collector_scope(
@@ -1424,6 +1577,7 @@ impl CapturedResponseBodyStore {
collection_was_gated,
),
);
self.track_durable_entry(&tracked_request_id);
}
pub(crate) fn insert_failed_with_collector_scope(
@@ -1434,6 +1588,7 @@ impl CapturedResponseBodyStore {
collector_ids: impl IntoIterator<Item = String>,
collection_was_gated: bool,
) {
let tracked_request_id = request_id.clone();
self.buffered_bodies.remove(&request_id);
self.bodies.insert(
request_id,
@@ -1444,6 +1599,7 @@ impl CapturedResponseBodyStore {
collection_was_gated,
),
);
self.track_durable_entry(&tracked_request_id);
}
pub(crate) fn get(&self, request_id: &str) -> Option<&CapturedResponseBody> {
@@ -1463,6 +1619,8 @@ impl CapturedResponseBodyStore {
pub(crate) fn clear(&mut self) {
self.bodies.clear();
self.buffered_bodies.clear();
self.durable_entry_order.clear();
self.retained_request_ids.clear();
}
#[cfg(test)]
@@ -1476,23 +1634,13 @@ impl CapturedResponseBodyStore {
}
pub(crate) fn remove_session_visibility(&mut self, session_id: Option<&str>) {
self.configure_durable(session_id, None);
self.bodies
.retain(|_, body| body.remove_session_visibility(session_id));
let buffered_request_ids = self
.buffered_bodies
.iter()
.map(|(request_id, _)| request_id.clone())
.collect::<Vec<_>>();
for request_id in buffered_request_ids {
let retain = self
.buffered_bodies
.get_mut(request_id.as_str())
.is_some_and(|body| body.remove_session_visibility(session_id));
if !retain {
self.buffered_bodies.remove(request_id.as_str());
}
}
self.buffered_bodies
.retain(|_, body| body.remove_session_visibility(session_id));
self.prune_removed_entry_indexes();
}
#[cfg(test)]
@@ -1959,6 +2107,145 @@ mod tests {
);
}
#[test]
fn durable_response_body_store_bounds_payload_metadata_and_zero_byte_entries() {
let mut store = CapturedResponseBodyStore::default();
store.configure_durable(
Some("owner"),
Some(moli_bounded_buffer::ByteLimits::new(5, 3)),
);
store.insert("a".into(), "aaa".into(), [Some("owner".into())]);
store.insert("b".into(), "bbb".into(), [Some("owner".into())]);
assert!(store.get("a").unwrap().body_bytes_limited(10).is_err());
assert_eq!(store.buffered_body_bytes(), 3);
store.insert("oversize".into(), "xxxx".into(), [Some("owner".into())]);
assert!(
store
.get("oversize")
.unwrap()
.body_bytes_limited(10)
.is_err()
);
for n in 0..super::DURABLE_RESPONSE_BODY_MAX_ENTRIES + 10 {
store.insert(format!("empty-{n}"), String::new(), [Some("owner".into())]);
store.insert_failed_with_collector_scope(
format!("failed-{n}"),
"failed".into(),
[Some("owner".into())],
[],
false,
);
}
assert!(
store.bodies.len() + store.buffered_bodies.len()
<= super::DURABLE_RESPONSE_BODY_MAX_ENTRIES
);
assert_eq!(
store.durable_entry_order.len(),
super::DURABLE_RESPONSE_BODY_MAX_ENTRIES
);
assert!(!store.contains_key("a"));
store.prepare_navigation();
store.remove_session_visibility(Some("owner"));
assert!(store.is_empty());
assert!(store.durable_sessions.is_empty());
assert!(store.retained_request_ids.is_empty());
assert!(store.durable_entry_order.is_empty());
}
#[test]
fn durable_response_body_revocation_does_not_evict_other_sessions_below_budget() {
for disable_retention in [false, true] {
let mut store = CapturedResponseBodyStore::default();
for session in ["a", "b"] {
store.configure_durable(
Some(session),
Some(moli_bounded_buffer::ByteLimits::new(100_000, 100)),
);
}
store.insert("keep".into(), "original".into(), [Some("a".into())]);
for n in 1..super::DURABLE_RESPONSE_BODY_MAX_ENTRIES {
store.insert(format!("b-{n}"), "other".into(), [Some("b".into())]);
}
if disable_retention {
store.prepare_navigation();
store.configure_durable(Some("b"), None);
} else {
store.remove_session_visibility(Some("b"));
}
assert!(store.get("b-1").is_none());
store.insert("new".into(), "next".into(), [Some("a".into())]);
assert_eq!(
store
.get("keep")
.expect("unrelated response must remain readable")
.body_bytes_limited(100)
.unwrap(),
b"original"
);
assert_eq!(
store.get("new").unwrap().body_bytes_limited(100).unwrap(),
b"next"
);
store.prepare_navigation();
assert!(store.get("keep").unwrap().is_visible_to_session(Some("a")));
assert!(!store.get("keep").unwrap().is_visible_to_session(Some("b")));
}
}
#[test]
fn durable_response_body_store_navigation_preserves_only_existing_opted_claims() {
let mut store = CapturedResponseBodyStore::default();
store.configure_durable(
Some("a"),
Some(moli_bounded_buffer::ByteLimits::new(100, 50)),
);
store.configure_durable(
Some("b"),
Some(moli_bounded_buffer::ByteLimits::new(80, 40)),
);
store.insert(
"old".into(),
"body".into(),
[Some("a".into()), Some("ordinary".into())],
);
store.insert_pending_with_collector_scope("pending".into(), [Some("a".into())], [], false);
store.prepare_navigation();
let body = store.get("old").unwrap();
assert!(body.is_visible_to_session(Some("a")));
assert!(!body.is_visible_to_session(Some("b")));
assert!(!body.is_visible_to_session(Some("ordinary")));
assert!(!store.contains_key("pending"));
assert_eq!(
store.buffered_bodies.limits(),
moli_bounded_buffer::ByteLimits::new(80, 40)
);
store.insert("current".into(), "current".into(), [Some("a".into())]);
store.configure_durable(Some("a"), None);
assert!(!store.contains_key("old"));
assert!(store.contains_key("current"));
store.prepare_navigation();
assert!(store.is_empty());
let other_target = CapturedResponseBodyStore::default();
assert!(!other_target.contains_key("current"));
}
#[test]
fn durable_response_body_default_navigation_still_clears_and_budget_shrink_releases_payloads() {
let mut store = CapturedResponseBodyStore::default();
store.insert("ordinary".into(), "body".into(), [None]);
store.prepare_navigation();
assert!(store.is_empty());
store.configure_durable(None, Some(moli_bounded_buffer::ByteLimits::new(20, 10)));
store.insert("a".into(), "aaaa".into(), [None]);
store.insert("b".into(), "bbbb".into(), [None]);
store.configure_durable(None, Some(moli_bounded_buffer::ByteLimits::new(5, 3)));
assert_eq!(store.buffered_body_bytes(), 0);
store.clear();
assert!(store.is_empty());
assert!(store.retained_request_ids.is_empty());
}
#[test]
fn response_body_store_marks_oldest_payload_evicted_when_total_budget_fills() {
let mut store =
+26 -6
View File
@@ -24,6 +24,7 @@ use super::output_queue::{
TargetSubresourceRequestNetworkDeliveryOutput, TargetSubresourceResponseNetworkDeliveryOutput,
TargetWebSocketDeliveryRecord, TargetWebSocketLifecycleDeliveryKind,
};
use super::redirect_request::RedirectRequest;
pub(crate) struct NetworkBacklogProjectionContext<'a> {
owner: CommandOwnerScope,
@@ -261,6 +262,7 @@ fn emit_complete_subresource_network_delivery_record(
.then(|| output.request_cookie_report())
.flatten(),
&[],
true,
);
}
}
@@ -286,7 +288,13 @@ fn emit_complete_subresource_network_delivery_record(
response_headers,
response_body_len,
} => {
let mut redirected_request = RedirectRequest::new(
output.method(),
output.request_body(),
output.request_headers(),
);
for redirect in redirect_chain {
redirected_request.follow(redirect.status);
for event_session_id in event_session_ids {
emit_request_will_be_sent(
out,
@@ -297,9 +305,9 @@ fn emit_complete_subresource_network_delivery_record(
timestamp,
record_document_url,
&redirect.to_url,
output.method(),
output.request_body(),
output.request_headers(),
redirected_request.method,
redirected_request.body,
&redirected_request.headers,
resource_type,
output.request_initiator_type(),
Some((
@@ -312,6 +320,7 @@ fn emit_complete_subresource_network_delivery_record(
!redirect.cookie_set_reports.is_empty(),
redirect.request_cookie_report.as_ref(),
&[],
true,
);
emit_redirect_response_received_extra_info(
out,
@@ -453,6 +462,10 @@ fn emit_staged_subresource_request_started(
false,
output.request_cookie_report(),
&[],
// Cookie selection precedes transport header generation. Keep
// it in the main event, but defer ExtraInfo until completion
// supplies the actual headers (or a terminal fallback).
false,
);
}
}
@@ -502,7 +515,13 @@ fn emit_staged_subresource_response_started(
let loader_id = request.loader_id();
let timestamp = base_timestamp + ((output.index() + 1) as f64 * 0.000_001);
let resource_type = request.resource_type().into();
let mut redirected_request = RedirectRequest::new(
request.method(),
request.request_body(),
request.request_headers(),
);
for redirect in output.redirect_chain() {
redirected_request.follow(redirect.status);
for event_session_id in event_session_ids {
emit_request_will_be_sent(
out,
@@ -513,9 +532,9 @@ fn emit_staged_subresource_response_started(
timestamp,
request.document_url(),
&redirect.to_url,
request.method(),
request.request_body(),
request.request_headers(),
redirected_request.method,
redirected_request.body,
&redirected_request.headers,
resource_type,
request.request_initiator_type(),
Some((
@@ -528,6 +547,7 @@ fn emit_staged_subresource_response_started(
!redirect.cookie_set_reports.is_empty(),
redirect.request_cookie_report.as_ref(),
&[],
true,
);
emit_redirect_response_received_extra_info(
out,
+3 -1
View File
@@ -279,6 +279,7 @@ pub(crate) fn emit_request_will_be_sent(
redirect_has_extra_info: bool,
cookie_access_report: Option<&moli_cookie_jar::StoredCookieQueryReport>,
blocked_intercepts: &[DevToolsNetworkInterceptId],
emit_cookie_extra_info: bool,
) {
if redirect_response.is_some_and(|(_, _, _, from_cache, _)| from_cache) {
emit_request_served_from_cache(out, session_id, request_id);
@@ -343,7 +344,7 @@ pub(crate) fn emit_request_will_be_sent(
}),
session_id,
);
if let Some(cookie_access_report) = cookie_access_report {
if emit_cookie_extra_info && let Some(cookie_access_report) = cookie_access_report {
emit_request_will_be_sent_extra_info(
out,
session_id,
@@ -1536,6 +1537,7 @@ mod tests {
false,
None,
&[],
true,
);
assert_eq!(events.len(), 2);
@@ -611,7 +611,7 @@ impl MainDocumentBodyProgressSource {
response_headers,
response_cookie_reports,
network_observation_journal,
redirect_chain.len(),
redirect_chain,
network_extra_info_available,
response_from_cache,
negotiated_http_version,
@@ -704,7 +704,7 @@ impl MainDocumentBodyProgressSource {
};
let final_exchange = navigation_exchange_group(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
redirect_chain.len(),
)
.and_then(|group| group.last());
@@ -721,6 +721,7 @@ impl MainDocumentBodyProgressSource {
redirect_chain,
network_observation_journal,
final_network_extra_info_available,
SubresourceRequestInitiatorType::Other,
);
live_source.send_progress_events(MainDocumentProgressPhase::RequestStarted, events);
}
@@ -759,7 +760,7 @@ impl MainDocumentBodyProgressSource {
if network_extra_info_available {
let (response_status, response_headers) = observed_response_metadata(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
redirect_chain.len(),
response_status,
response_headers,
@@ -944,6 +945,7 @@ impl MainDocumentLiveNetworkProgressSource {
redirect_chain,
network_observation_journal,
final_network_extra_info_available,
SubresourceRequestInitiatorType::Other,
);
self.send_progress_events(MainDocumentProgressPhase::RequestStarted, events);
}
@@ -967,6 +969,7 @@ impl MainDocumentLiveNetworkProgressSource {
redirect_chain,
network_observation_journal,
final_network_extra_info_available,
SubresourceRequestInitiatorType::Other,
);
self.send_progress_events_into_output(
MainDocumentProgressPhase::RequestStarted,
@@ -984,13 +987,15 @@ impl MainDocumentLiveNetworkProgressSource {
redirect_chain: &[RedirectInfo],
network_observation_journal: &NetworkObservationJournal,
final_network_extra_info_available: bool,
request_initiator_type: SubresourceRequestInitiatorType,
) -> Vec<MainDocumentNavigationProgressEvent> {
let mut events = Vec::new();
let initial_network_extra_info_available = !network_observation_journal.is_empty()
|| redirect_chain
.first()
.is_some_and(|redirect| redirect.network_extra_info_available)
|| (redirect_chain.is_empty() && final_network_extra_info_available);
let initial_network_extra_info_available =
navigation_exchange_group(network_observation_journal, redirect_chain, 0).is_some()
|| redirect_chain
.first()
.is_some_and(|redirect| redirect.network_extra_info_available)
|| (redirect_chain.is_empty() && final_network_extra_info_available);
if initial_network_extra_info_available {
let initial_request_cookie_report = if redirect_chain.is_empty() {
final_request_cookie_report
@@ -1003,7 +1008,7 @@ impl MainDocumentLiveNetworkProgressSource {
self.progress_target(),
observed_request_headers(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
0,
&self.initial_request_headers,
),
@@ -1012,15 +1017,27 @@ impl MainDocumentLiveNetworkProgressSource {
}
for (index, redirect) in redirect_chain.iter().enumerate() {
let target = self.progress_target();
let request_method = observed_request_method(
network_observation_journal,
redirect_chain,
index + 1,
final_request_method,
);
let request_headers = observed_request_headers(
network_observation_journal,
redirect_chain,
index + 1,
final_request_headers,
);
let observed_response_available =
navigation_exchange_group(network_observation_journal, redirect_chain.len(), index)
navigation_exchange_group(network_observation_journal, redirect_chain, index)
.and_then(|group| group.last())
.and_then(NetworkExchangeObservation::response)
.is_some();
if observed_response_available && !redirect.network_extra_info_available {
let (status, headers) = observed_response_metadata(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
index,
redirect.status,
&redirect.headers,
@@ -1037,10 +1054,12 @@ impl MainDocumentLiveNetworkProgressSource {
events.push(MainDocumentNavigationProgressEvent::RequestWillBeSent {
target: target.clone(),
url: redirect.to_url.clone(),
method: final_request_method.to_owned(),
request_body: request_body.map(str::to_owned),
request_headers: final_request_headers.to_vec(),
request_initiator_type: SubresourceRequestInitiatorType::Other,
method: request_method.clone(),
request_body: (request_method == final_request_method)
.then(|| request_body.map(str::to_owned))
.flatten(),
request_headers,
request_initiator_type,
redirect_response: Box::new(Some(MainDocumentRedirectResponse {
url: redirect.from_url.clone(),
status: redirect.status,
@@ -1061,7 +1080,7 @@ impl MainDocumentLiveNetworkProgressSource {
if observed_response_available && redirect.network_extra_info_available {
let (status, headers) = observed_response_metadata(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
index,
redirect.status,
&redirect.headers,
@@ -1075,14 +1094,13 @@ impl MainDocumentLiveNetworkProgressSource {
},
);
}
let request_network_extra_info_available = network_observation_journal
.exchanges()
.get(index + 1)
.is_some()
|| redirect_chain
.get(index + 1)
.is_some_and(|next| next.network_extra_info_available)
|| (index + 1 == redirect_chain.len() && final_network_extra_info_available);
let request_network_extra_info_available =
navigation_exchange_group(network_observation_journal, redirect_chain, index + 1)
.is_some()
|| redirect_chain
.get(index + 1)
.is_some_and(|next| next.network_extra_info_available)
|| (index + 1 == redirect_chain.len() && final_network_extra_info_available);
if request_network_extra_info_available {
let cookie_report = redirect.request_cookie_report.clone().or_else(|| {
(index + 1 == redirect_chain.len())
@@ -1093,7 +1111,7 @@ impl MainDocumentLiveNetworkProgressSource {
target,
observed_request_headers(
network_observation_journal,
redirect_chain.len(),
redirect_chain,
index + 1,
final_request_headers,
),
@@ -1111,15 +1129,15 @@ impl MainDocumentLiveNetworkProgressSource {
response_headers: &[(String, String)],
response_cookie_reports: &[StoredCookieSetReport],
network_observation_journal: &NetworkObservationJournal,
redirect_count: usize,
redirect_chain: &[RedirectInfo],
network_extra_info_available: bool,
response_from_cache: bool,
negotiated_http_version: Option<NegotiatedHttpVersion>,
) {
let (extra_info_status, extra_info_headers) = observed_response_metadata(
network_observation_journal,
redirect_count,
redirect_count,
redirect_chain,
redirect_chain.len(),
response_status,
response_headers,
);
@@ -1354,32 +1372,85 @@ pub(crate) fn emit_child_document_navigation_network_background_events(
frame_id: frame_id.to_owned(),
timestamp,
};
let redirect_count = network.redirect_chain.len();
let request_method = observed_request_method(
&network.network_observation_journal,
&network.redirect_chain,
0,
&network.request_method,
);
let request_headers = observed_request_headers(
&network.network_observation_journal,
&network.redirect_chain,
0,
&network.request_headers,
);
let mut output = MainDocumentProgressOutputTarget::background_events(out);
output.emit_event(MainDocumentNavigationProgressEvent::RequestWillBeSent {
target: target.clone(),
url: request_url,
method: network.request_method.clone(),
method: request_method,
request_body: None,
request_headers: network.request_headers.clone(),
request_headers,
request_initiator_type: SubresourceRequestInitiatorType::Parser,
redirect_response: Box::new(None),
redirect_has_extra_info: false,
cookie_access_report: None,
});
let source = MainDocumentLiveNetworkProgressSource {
sender: None,
progress_queue: MainDocumentProgressQueueHandle::from_source(
MainDocumentProgressSource::streaming(),
),
session_ids: session_ids.clone(),
request_id: request_id.to_owned(),
loader_id: loader_id.to_owned(),
frame_id: frame_id.to_owned(),
timestamp,
initial_request_headers: network.request_headers.clone(),
initial_request_cookie_report: None,
};
for event in source.redirect_request_events(
&network.request_method,
None,
&network.request_headers,
None,
&network.redirect_chain,
&network.network_observation_journal,
!network.network_observation_journal.is_empty(),
SubresourceRequestInitiatorType::Parser,
) {
output.emit_event(event);
}
let (extra_info_status, extra_info_headers) = observed_response_metadata(
&network.network_observation_journal,
&network.redirect_chain,
redirect_count,
network.status,
&network.response_headers,
);
let network_extra_info_available = navigation_exchange_group(
&network.network_observation_journal,
&network.redirect_chain,
redirect_count,
)
.and_then(|group| group.last())
.and_then(NetworkExchangeObservation::response)
.is_some();
output.emit_event(MainDocumentNavigationProgressEvent::ResponseReceived {
target: target.clone(),
final_url,
status: network.status,
headers: network.response_headers.clone(),
cookie_set_reports: Vec::new(),
extra_info_status: network.status,
extra_info_headers: network.response_headers.clone(),
network_extra_info_available: false,
emit_extra_info: false,
extra_info_status,
extra_info_headers,
network_extra_info_available,
emit_extra_info: network_extra_info_available,
encoded_data_length: 0,
from_cache: network.from_cache,
negotiated_http_version: None,
has_extra_info: false,
has_extra_info: network_extra_info_available,
});
record_child_document_response_body(
conn,
@@ -1634,7 +1705,12 @@ impl CompletedMainDocumentProgressContext {
cookie_access_report: initial_request_cookie_report.clone(),
});
}
let initial_network_extra_info_available = !events.network_observation_journal.is_empty()
let initial_network_extra_info_available = navigation_exchange_group(
&events.network_observation_journal,
&events.redirect_chain,
0,
)
.is_some()
|| events
.redirect_chain
.first()
@@ -1645,7 +1721,7 @@ impl CompletedMainDocumentProgressContext {
target.clone(),
observed_request_headers(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
0,
&self.request_headers,
),
@@ -1653,9 +1729,21 @@ impl CompletedMainDocumentProgressContext {
));
}
for (index, redirect) in events.redirect_chain.iter().enumerate() {
let request_method = observed_request_method(
&events.network_observation_journal,
&events.redirect_chain,
index + 1,
&events.request_method,
);
let request_headers = observed_request_headers(
&events.network_observation_journal,
&events.redirect_chain,
index + 1,
&events.request_headers,
);
let observed_response_available = navigation_exchange_group(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
index,
)
.and_then(|group| group.last())
@@ -1664,7 +1752,7 @@ impl CompletedMainDocumentProgressContext {
if observed_response_available && !redirect.network_extra_info_available {
let (status, headers) = observed_response_metadata(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
index,
redirect.status,
&redirect.headers,
@@ -1681,9 +1769,11 @@ impl CompletedMainDocumentProgressContext {
progress_events.push(MainDocumentNavigationProgressEvent::RequestWillBeSent {
target: target.clone(),
url: redirect.to_url.clone(),
method: events.request_method.clone(),
request_body: self.request_body.clone(),
request_headers: events.request_headers.clone(),
method: request_method.clone(),
request_body: (request_method == self.request_method)
.then(|| self.request_body.clone())
.flatten(),
request_headers,
request_initiator_type: SubresourceRequestInitiatorType::Other,
redirect_response: Box::new(Some(MainDocumentRedirectResponse {
url: redirect.from_url.clone(),
@@ -1705,7 +1795,7 @@ impl CompletedMainDocumentProgressContext {
if observed_response_available && redirect.network_extra_info_available {
let (status, headers) = observed_response_metadata(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
index,
redirect.status,
&redirect.headers,
@@ -1719,11 +1809,12 @@ impl CompletedMainDocumentProgressContext {
},
);
}
let request_network_extra_info_available = events
.network_observation_journal
.exchanges()
.get(index + 1)
.is_some()
let request_network_extra_info_available = navigation_exchange_group(
&events.network_observation_journal,
&events.redirect_chain,
index + 1,
)
.is_some()
|| events
.redirect_chain
.get(index + 1)
@@ -1740,7 +1831,7 @@ impl CompletedMainDocumentProgressContext {
target.clone(),
observed_request_headers(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
index + 1,
&events.request_headers,
),
@@ -1762,7 +1853,7 @@ impl CompletedMainDocumentProgressContext {
};
let (extra_info_status, extra_info_headers) = observed_response_metadata(
&events.network_observation_journal,
events.redirect_chain.len(),
&events.redirect_chain,
events.redirect_chain.len(),
events.response_status,
&events.response_headers,
@@ -1907,62 +1998,107 @@ fn request_extra_info_event(
}
}
fn navigation_exchange_group(
journal: &NetworkObservationJournal,
redirect_count: usize,
fn navigation_exchange_group<'a>(
journal: &'a NetworkObservationJournal,
redirect_chain: &[impl RedirectTransportBoundary],
hop_index: usize,
) -> Option<&[NetworkExchangeObservation]> {
) -> Option<&'a [NetworkExchangeObservation]> {
// A truncated journal still proves that a transport exchange happened, but
// no longer provides a trustworthy tail for redirect-hop correlation.
if journal.truncated() {
return None;
}
if hop_index > redirect_count {
if hop_index > redirect_chain.len() {
return None;
}
let exchanges = journal.exchanges();
let mut group_start = 0;
let mut current_hop = 0;
for (index, exchange) in exchanges.iter().enumerate() {
let ends_redirect_hop = exchange.response().is_some_and(|response| {
matches!(response.status(), 301 | 302 | 303 | 307 | 308)
|| response.headers().iter().any(|(name, value)| {
name.eq_ignore_ascii_case("critical-ch") && !value.trim().is_empty()
})
});
if !ends_redirect_hop || current_hop >= redirect_count {
for (redirect_index, redirect) in redirect_chain.iter().enumerate() {
let Some(expected_status) = redirect.transport_response_status() else {
if redirect_index == hop_index {
return None;
}
continue;
};
let group_end = exchanges[group_start..]
.iter()
.position(|exchange| {
exchange
.response()
.is_some_and(|response| response.status() == expected_status)
})?
.saturating_add(group_start);
if redirect_index == hop_index {
return Some(&exchanges[group_start..=group_end]);
}
if current_hop == hop_index {
return Some(&exchanges[group_start..=index]);
}
current_hop += 1;
group_start = index.saturating_add(1);
group_start = group_end.saturating_add(1);
}
(hop_index == redirect_chain.len() && group_start < exchanges.len())
.then_some(&exchanges[group_start..])
}
trait RedirectTransportBoundary {
fn transport_response_status(&self) -> Option<u16>;
}
impl RedirectTransportBoundary for RedirectInfo {
fn transport_response_status(&self) -> Option<u16> {
self.response_extra_info
.as_ref()
.map(|response| response.status)
.or_else(|| {
(self.source == moli_fetch::RedirectSource::Network && !self.from_cache)
.then_some(self.status)
})
}
}
impl RedirectTransportBoundary for NavigationRedirect {
fn transport_response_status(&self) -> Option<u16> {
self.response_extra_info
.as_ref()
.map(|response| response.status)
.or_else(|| {
(self.source == moli_fetch::RedirectSource::Network && !self.from_cache)
.then_some(self.status)
})
}
(current_hop == hop_index && group_start < exchanges.len()).then_some(&exchanges[group_start..])
}
fn observed_request_headers(
journal: &NetworkObservationJournal,
redirect_count: usize,
redirect_chain: &[impl RedirectTransportBoundary],
hop_index: usize,
fallback: &[(String, String)],
) -> Vec<(String, String)> {
navigation_exchange_group(journal, redirect_count, hop_index)
navigation_exchange_group(journal, redirect_chain, hop_index)
.and_then(|group| group.first())
.map(|exchange| exchange.request().headers().to_vec())
.unwrap_or_else(|| fallback.to_vec())
}
fn observed_request_method(
journal: &NetworkObservationJournal,
redirect_chain: &[impl RedirectTransportBoundary],
hop_index: usize,
fallback: &str,
) -> String {
navigation_exchange_group(journal, redirect_chain, hop_index)
.and_then(|group| group.first())
.and_then(|exchange| exchange.request().method())
.unwrap_or(fallback)
.to_owned()
}
fn observed_response_metadata(
journal: &NetworkObservationJournal,
redirect_count: usize,
redirect_chain: &[impl RedirectTransportBoundary],
hop_index: usize,
fallback_status: u16,
fallback_headers: &[(String, String)],
) -> (u16, Vec<(String, String)>) {
navigation_exchange_group(journal, redirect_count, hop_index)
navigation_exchange_group(journal, redirect_chain, hop_index)
.and_then(|group| group.last())
.and_then(NetworkExchangeObservation::response)
.map(|response| (response.status(), response.headers().to_vec()))
@@ -4,7 +4,8 @@ use moli_cookie_jar::{
};
use moli_fetch::{
NegotiatedHttpVersion, NetworkExchangeObservation, NetworkObservationJournal,
NetworkRequestObservation, NetworkResponseObservation, RedirectInfo,
NetworkRequestExtraInfo, NetworkRequestObservation, NetworkResponseExtraInfo,
NetworkResponseObservation, RedirectInfo,
};
use serde_json::{Value, json};
use tokio::sync::mpsc::unbounded_channel;
@@ -48,6 +49,22 @@ fn observation_journal(
)
}
fn response_extra_info(
request_headers: Vec<(String, String)>,
status: u16,
response_headers: Vec<(String, String)>,
) -> NetworkResponseExtraInfo {
NetworkResponseExtraInfo {
request_extra_info: NetworkRequestExtraInfo {
headers: request_headers,
cookie_report: StoredCookieQueryReport::default(),
},
status,
headers: response_headers,
cookie_set_reports: Vec::new(),
}
}
fn completed_progress_context() -> CompletedMainDocumentProgressContext {
CompletedMainDocumentProgressContext::new(
vec![Some("SID-1".to_owned())],
@@ -419,25 +436,36 @@ fn live_progress_source_serializes_through_progress_emissions() {
from_cache: false,
negotiated_http_version: None,
};
let journal = observation_journal(vec![
(
vec![("Host".to_owned(), "example.test".to_owned())],
302,
vec![("Location".to_owned(), final_url.to_string())],
let journal = NetworkObservationJournal::from_exchanges(vec![
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"POST",
vec![("Host".to_owned(), "example.test".to_owned())],
),
Some(NetworkResponseObservation::new(
302,
vec![("Location".to_owned(), final_url.to_string())],
)),
),
(
vec![("Host".to_owned(), "example.test".to_owned())],
200,
vec![("Content-Type".to_owned(), "text/html".to_owned())],
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"GET",
vec![("Host".to_owned(), "example.test".to_owned())],
),
Some(NetworkResponseObservation::new(
200,
vec![("Content-Type".to_owned(), "text/html".to_owned())],
)),
),
]);
let redirects = vec![redirect];
source.emit_redirect_requests(
"GET",
None,
"POST",
Some("account=1"),
&[("Accept".to_owned(), "text/html".to_owned())],
None,
&[redirect],
&redirects,
&journal,
true,
);
@@ -447,7 +475,7 @@ fn live_progress_source_serializes_through_progress_emissions() {
&[("Content-Type".to_owned(), "text/html".to_owned())],
&[],
&journal,
1,
&redirects,
true,
false,
None,
@@ -475,6 +503,11 @@ fn live_progress_source_serializes_through_progress_emissions() {
json!("Network.requestWillBeSent")
);
assert_eq!(redirect_request["params"]["requestId"], json!("REQ-1"));
assert_eq!(
redirect_request["params"]["request"]["method"],
json!("GET")
);
assert!(redirect_request["params"]["request"]["postData"].is_null());
assert_eq!(
redirect_request["params"]["redirectResponse"]["status"],
json!(302)
@@ -546,6 +579,135 @@ fn live_progress_source_serializes_through_progress_emissions() {
assert_eq!(finished["params"]["encodedDataLength"], json!(17));
}
#[test]
fn redirect_request_preserves_post_body_when_transport_preserves_method() {
let (source, mut receiver) = live_progress_source();
let final_url = Url::parse("http://example.test/final").unwrap();
let redirect = RedirectInfo {
source: moli_fetch::RedirectSource::Network,
from_url: Url::parse("http://example.test/start").unwrap(),
to_url: final_url.clone(),
status: 307,
headers: vec![("Location".to_owned(), final_url.to_string())],
network_extra_info_available: true,
request_extra_info: None,
response_extra_info: None,
redirect_has_extra_info: true,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
from_cache: false,
negotiated_http_version: None,
};
let journal = NetworkObservationJournal::from_exchanges(vec![
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method("POST", Vec::new()),
Some(NetworkResponseObservation::new(307, Vec::new())),
),
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method("POST", Vec::new()),
Some(NetworkResponseObservation::new(200, Vec::new())),
),
]);
source.emit_redirect_requests(
"POST",
Some("account=1"),
&[],
None,
&[redirect],
&journal,
true,
);
let _initial_request_extra = receiver
.try_recv()
.expect("initial request extra info event");
let (redirect_request, _) = receiver
.try_recv()
.expect("redirect request event")
.into_parts();
assert_eq!(
redirect_request["params"]["request"]["method"],
json!("POST")
);
assert_eq!(
redirect_request["params"]["request"]["postData"],
json!("account=1")
);
}
#[test]
fn observed_request_headers_do_not_restore_a_referrer_removed_by_transport_policy() {
let journal = NetworkObservationJournal::from_exchanges(vec![NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"GET",
vec![("Host".to_owned(), "example.test".to_owned())],
),
Some(NetworkResponseObservation::new(200, Vec::new())),
)]);
let headers = super::observed_request_headers(
&journal,
&[] as &[RedirectInfo],
0,
&[("Referer".to_owned(), "https://stale.test/".to_owned())],
);
assert_eq!(
headers,
vec![("Host".to_owned(), "example.test".to_owned())]
);
}
#[test]
fn truncated_observation_journal_falls_back_instead_of_guessing_redirect_alignment() {
let final_url = Url::parse("http://example.test/final").unwrap();
let redirect = RedirectInfo {
source: moli_fetch::RedirectSource::Network,
from_url: Url::parse("http://example.test/start").unwrap(),
to_url: final_url.clone(),
status: 302,
headers: vec![("Location".to_owned(), final_url.to_string())],
network_extra_info_available: true,
request_extra_info: None,
response_extra_info: None,
redirect_has_extra_info: true,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
from_cache: false,
negotiated_http_version: None,
};
let journal = NetworkObservationJournal::from_exchanges(
(0..33)
.map(|index| {
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"GET",
vec![("X-Wire-Hop".to_owned(), index.to_string())],
),
Some(NetworkResponseObservation::new(200, Vec::new())),
)
})
.collect(),
);
assert!(journal.truncated());
let fallback_headers = vec![("X-Logical-Hop".to_owned(), "final".to_owned())];
assert_eq!(
super::observed_request_headers(
&journal,
std::slice::from_ref(&redirect),
1,
&fallback_headers
),
fallback_headers
);
assert_eq!(
super::observed_request_method(&journal, &[redirect], 1, "POST"),
"POST"
);
}
#[test]
fn completed_body_http_response_emits_correlated_empty_cookie_extra_info() {
let final_url = Url::parse("http://example.test/final").unwrap();
@@ -773,19 +935,46 @@ fn completed_body_http_redirect_emits_correlated_no_cookie_extra_info() {
from_cache: false,
negotiated_http_version: None,
}];
events = events.with_network_observation_journal(observation_journal(vec![
(
vec![("X-Request-Hop".to_owned(), "initial".to_owned())],
302,
vec![("X-Response-Hop".to_owned(), "redirect".to_owned())],
),
(
vec![("X-Request-Hop".to_owned(), "final".to_owned())],
200,
vec![("X-Response-Hop".to_owned(), "final".to_owned())],
),
]));
let batches = completed_progress_context().event_batches(&events, &final_url, 17);
events.request_method = "POST".to_owned();
events =
events.with_network_observation_journal(NetworkObservationJournal::from_exchanges(vec![
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"POST",
vec![("X-Request-Hop".to_owned(), "initial".to_owned())],
),
Some(NetworkResponseObservation::new(
302,
vec![("X-Response-Hop".to_owned(), "redirect".to_owned())],
)),
),
NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"GET",
vec![("X-Request-Hop".to_owned(), "final".to_owned())],
),
Some(NetworkResponseObservation::new(
200,
vec![("X-Response-Hop".to_owned(), "final".to_owned())],
)),
),
]));
let context = CompletedMainDocumentProgressContext::new(
vec![Some("SID-1".to_owned())],
Some("REQ-1".to_owned()),
false,
Url::parse("http://example.test/start").unwrap(),
"POST".to_owned(),
Some("account=1".to_owned()),
vec![(
"Content-Type".to_owned(),
"application/x-www-form-urlencoded".to_owned(),
)],
"LOADER-1".to_owned(),
"FRAME-1".to_owned(),
12.5,
);
let batches = context.event_batches(&events, &final_url, 17);
let mut drain =
MainDocumentProgressDrain::from_source(MainDocumentProgressSource::completed_body(batches));
@@ -826,6 +1015,11 @@ fn completed_body_http_redirect_emits_correlated_no_cookie_extra_info() {
redirected_request["params"]["redirectHasExtraInfo"],
json!(true)
);
assert_eq!(
redirected_request["params"]["request"]["method"],
json!("GET")
);
assert!(redirected_request["params"]["request"]["postData"].is_null());
let request_extras = out
.iter()
.filter(|message| message["method"] == json!("Network.requestWillBeSentExtraInfo"))
@@ -931,7 +1125,14 @@ fn critical_client_hint_restart_keeps_discarded_response_extra_info_separate_fro
headers: vec![("Location".to_owned(), navigation_url.to_string())],
network_extra_info_available: false,
request_extra_info: None,
response_extra_info: None,
response_extra_info: Some(response_extra_info(
vec![("Sec-CH-UA".to_owned(), "\"Chromium\";v=\"145\"".to_owned())],
403,
vec![
("Accept-CH".to_owned(), "Sec-CH-UA-Arch".to_owned()),
("Critical-CH".to_owned(), "Sec-CH-UA-Arch".to_owned()),
],
)),
redirect_has_extra_info: false,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
@@ -995,6 +1196,69 @@ fn critical_client_hint_restart_keeps_discarded_response_extra_info_separate_fro
assert_eq!(relevant[6]["params"]["hasExtraInfo"], json!(true));
}
#[test]
fn transportless_https_upgrade_attributes_wire_headers_to_upgraded_request() {
let initial_url = Url::parse("http://example.test/start").unwrap();
let final_url = Url::parse("https://example.test/start").unwrap();
let journal = NetworkObservationJournal::from_exchanges(vec![NetworkExchangeObservation::new(
NetworkRequestObservation::new_with_method(
"GET",
vec![("X-Wire-Upgraded".to_owned(), "yes".to_owned())],
),
Some(NetworkResponseObservation::new(
200,
vec![("Content-Type".to_owned(), "text/html".to_owned())],
)),
)]);
let mut events = completed_events().with_network_observation_journal(journal);
events.network_extra_info_available = true;
events.redirect_chain = vec![moli_core::page::NavigationRedirect {
source: moli_fetch::RedirectSource::Internal,
from_url: initial_url,
to_url: final_url.clone(),
status: 307,
headers: vec![("Location".to_owned(), final_url.to_string())],
network_extra_info_available: false,
request_extra_info: None,
response_extra_info: None,
redirect_has_extra_info: false,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
from_cache: false,
negotiated_http_version: Some(NegotiatedHttpVersion::Http11),
}];
let batches = completed_progress_context().event_batches(&events, &final_url, 17);
let mut drain =
MainDocumentProgressDrain::from_source(MainDocumentProgressSource::completed_body(batches));
drain.mark_output_visible_until(MainDocumentProgressOutputBoundary::BodyFinishedVisible);
let out = drain_into_protocol_messages(&mut drain);
let redirected_request = out
.iter()
.filter(|message| message["method"] == json!("Network.requestWillBeSent"))
.nth(1)
.expect("upgraded request");
assert_eq!(
redirected_request["params"]["request"]["headers"]["X-Wire-Upgraded"],
json!("yes"),
"the sole wire exchange belongs to the post-upgrade HTTPS request"
);
let request_extras = out
.iter()
.filter(|message| message["method"] == json!("Network.requestWillBeSentExtraInfo"))
.collect::<Vec<_>>();
assert_eq!(
request_extras.len(),
1,
"a transportless internal hop must not receive wire ExtraInfo"
);
assert_eq!(
request_extras[0]["params"]["headers"]["X-Wire-Upgraded"],
json!("yes")
);
}
#[test]
fn completed_body_uses_negotiated_protocol_for_redirect_and_final_response() {
let final_url = Url::parse("https://example.test/final").unwrap();
@@ -2976,17 +2976,18 @@ impl TargetNetworkOutputQueue {
let Some(request) = self.staged_subresource_requests.get(&handle) else {
return;
};
if request.request_cookie_report().is_some() {
return;
}
let request_headers = network_request_headers
.unwrap_or_else(|| request.request_headers())
.to_vec();
let Some(request_cookie_report) = request_cookie_report.cloned().or_else(|| {
network_request_headers
.is_some()
.then(StoredCookieQueryReport::default)
}) else {
let Some(request_cookie_report) = request_cookie_report
.cloned()
.or_else(|| request.request_cookie_report().cloned())
.or_else(|| {
network_request_headers
.is_some()
.then(StoredCookieQueryReport::default)
})
else {
return;
};
self.append_subresource_request_extra_info(
@@ -3686,6 +3687,60 @@ mod tests {
));
}
#[test]
fn staged_cookie_report_does_not_suppress_observed_transport_headers() {
let handle = SubresourceNetworkRequestHandle::new(81);
let document_url = Url::parse("https://example.com/page").unwrap();
let request_url = Url::parse("https://example.com/api").unwrap();
let request = SubresourceRequestStarted::new(
handle,
Some("FRAME-1".to_owned()),
document_url,
request_url.clone(),
"POST".to_owned(),
vec![("Content-Type".to_owned(), "text/plain".to_owned())].into(),
Some("body".to_owned()),
SubresourceResourceType::Fetch,
SubresourceRequestInitiatorType::Script,
Some(StoredCookieQueryReport::default()),
);
let headers = vec![
("Referer".to_owned(), "https://example.com/page".to_owned()),
("User-Agent".to_owned(), "Moli/Test".to_owned()),
];
let record = subresource_record(SubresourceResourceType::Fetch, request_url.as_str())
.with_request_handle(handle)
.with_network_request_headers(Some(headers.clone()));
let mut queue = TargetNetworkOutputQueue::default();
append_concrete_items_for_test(
&mut queue,
&[
ScriptNetworkOutputItem::SubresourceRequestStarted(Box::new(request)),
ScriptNetworkOutputItem::SubresourceNetworkRecord(Box::new(record)),
],
"LOADER-1",
);
let activity = PendingSubresourceNetworkActivity::from_sessions(vec![
PendingSubresourceNetworkActivitySession::new(None, 0),
]);
let mut request_ids = StableSubresourceHandleRequestIds::default();
let snapshot = pending_delivery_snapshot(&queue, activity, None, &mut request_ids).unwrap();
let outputs = subresource_outputs(&snapshot);
let extras = outputs
.iter()
.filter_map(|output| match output {
TargetSubresourceNetworkDeliveryOutput::RequestExtraInfo(output) => Some(output),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
extras.len(),
1,
"cookie metadata is not transport-header evidence"
);
assert_eq!(extras[0].output().request_headers(), headers.as_slice());
}
#[test]
fn staged_request_followed_by_redirect_complete_record_preserves_redirect_chain() {
let handle = SubresourceNetworkRequestHandle::new(7);
@@ -0,0 +1,101 @@
/// Project the request metadata for successive followed redirect hops.
///
/// This shares FetchRequest::apply_redirect_status rules: the emitted CDP request
/// must describe the new hop, not repeat the original POST after it became GET.
pub(super) struct RedirectRequest<'a> {
pub(super) method: &'a str,
pub(super) body: Option<&'a str>,
pub(super) headers: Vec<(String, String)>,
}
impl<'a> RedirectRequest<'a> {
pub(super) fn new(
method: &'a str,
body: Option<&'a str>,
headers: &[(String, String)],
) -> Self {
Self {
method,
body,
headers: headers.to_vec(),
}
}
pub(super) fn follow(&mut self, status: u16) {
if moli_fetch::redirect_status_rewrites_to_get(status, self.method) {
self.method = "GET";
self.body = None;
self.headers
.retain(|(name, _)| !moli_fetch::is_request_body_header_name(name));
}
}
}
#[cfg(test)]
mod tests {
use super::RedirectRequest;
#[test]
fn redirect_status_method_matrix_matches_fetch_semantics() {
// Explicit HTTP expectations, independent of the implementation's
// classification formula. None means the redirect discards the body.
for (status, method, expected_method, expected_body) in [
(301, "GET", "GET", Some("value")),
(301, "HEAD", "HEAD", Some("value")),
(301, "POST", "GET", None),
(301, "PUT", "PUT", Some("value")),
(301, "DELETE", "DELETE", Some("value")),
(301, "post", "GET", None),
(302, "GET", "GET", Some("value")),
(302, "HEAD", "HEAD", Some("value")),
(302, "POST", "GET", None),
(302, "PUT", "PUT", Some("value")),
(302, "DELETE", "DELETE", Some("value")),
(302, "post", "GET", None),
(303, "GET", "GET", Some("value")),
(303, "HEAD", "HEAD", Some("value")),
(303, "POST", "GET", None),
(303, "PUT", "GET", None),
(303, "DELETE", "GET", None),
(303, "post", "GET", None),
(307, "GET", "GET", Some("value")),
(307, "HEAD", "HEAD", Some("value")),
(307, "POST", "POST", Some("value")),
(307, "PUT", "PUT", Some("value")),
(307, "DELETE", "DELETE", Some("value")),
(307, "post", "post", Some("value")),
(308, "GET", "GET", Some("value")),
(308, "HEAD", "HEAD", Some("value")),
(308, "POST", "POST", Some("value")),
(308, "PUT", "PUT", Some("value")),
(308, "DELETE", "DELETE", Some("value")),
(308, "post", "post", Some("value")),
] {
let headers = vec![
("Content-Type".into(), "text/plain".into()),
("CONTENT-LENGTH".into(), "5".into()),
("Accept".into(), "*/*".into()),
];
let mut request = RedirectRequest::new(method, Some("value"), &headers);
request.follow(status);
assert_eq!(request.method, expected_method, "{status} {method}");
assert_eq!(request.body, expected_body, "{status} {method}");
let expected_headers = if expected_body.is_none() {
vec![("Accept".into(), "*/*".into())]
} else {
headers.clone()
};
assert_eq!(request.headers, expected_headers, "{status} {method}");
assert_eq!(headers.len(), 3, "original metadata stays unchanged");
}
}
#[test]
fn redirected_get_does_not_recover_original_post_on_later_307() {
let mut request = RedirectRequest::new("POST", Some("value"), &[]);
request.follow(302);
request.follow(307);
assert_eq!(request.method, "GET");
assert_eq!(request.body, None);
}
}
@@ -13,6 +13,35 @@ struct SetBlockedUrlsParams {
urls: Vec<String>,
}
#[derive(Default, Deserialize)]
#[serde(rename_all = "camelCase")]
struct DurableBodyParams {
#[serde(default)]
enable_durable_messages: bool,
max_total_buffer_size: Option<usize>,
max_resource_buffer_size: Option<usize>,
}
pub(super) fn durable_body_limits(
cmd: &Cmd<'_>,
) -> Result<Option<moli_bounded_buffer::ByteLimits>, CommandOutputPlan> {
let params = cmd
.get_params::<DurableBodyParams>()
.map_err(|_| CommandOutputPlan::error(-32602, "InvalidParams"))?
.unwrap_or_default();
if !params.enable_durable_messages {
return Ok(None);
}
let Some(total) = params.max_total_buffer_size.filter(|size| *size > 0) else {
return Err(CommandOutputPlan::error(-32602, "InvalidParams"));
};
let resource = params.max_resource_buffer_size.unwrap_or(2_000_000);
if resource == 0 {
return Err(CommandOutputPlan::error(-32602, "InvalidParams"));
}
Ok(Some(moli_bounded_buffer::ByteLimits::new(total, resource)))
}
pub(super) fn enabled_command_output_plan(
conn: &mut CdpConnection,
session_id: Option<&str>,
@@ -145,6 +145,7 @@ mod cache;
mod cookies;
mod load_resource;
mod navigation;
mod redirect_wire;
mod response_body;
mod runtime;
mod service_worker;
@@ -0,0 +1,190 @@
use super::*;
#[derive(Clone, Debug)]
struct WireRequest {
path: String,
method: String,
headers: HeaderMap,
body: String,
}
fn event_header<'a>(headers: &'a serde_json::Value, name: &str) -> Option<&'a str> {
headers.as_object()?.iter().find_map(|(key, value)| {
key.eq_ignore_ascii_case(name)
.then(|| value.as_str())
.flatten()
})
}
#[tokio::test(flavor = "multi_thread")]
async fn form_post_redirect_chain_cdp_metadata_matches_wire_requests() {
async fn capture(
State(seen): State<Arc<Mutex<Vec<WireRequest>>>>,
method: axum::http::Method,
uri: axum::http::Uri,
headers: HeaderMap,
body: String,
) -> axum::response::Response {
let path = uri.path().to_owned();
seen.lock().push(WireRequest {
path: path.clone(),
method: method.to_string(),
headers,
body,
});
match path.as_str() {
"/start" => (StatusCode::FOUND, [(LOCATION, "/middle")]).into_response(),
"/middle" => (StatusCode::TEMPORARY_REDIRECT, [(LOCATION, "/final")]).into_response(),
"/final" => axum::response::Html("<!doctype html><p>finished</p>").into_response(),
_ => unreachable!("only the three redirect routes use this handler"),
}
}
let seen = Arc::new(Mutex::new(Vec::<WireRequest>::new()));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = Router::new()
.route(
"/page",
get(|| async {
axum::response::Html(
"<!doctype html><form id='f' method='POST' action='/start'><input name='account' value='alice'></form>",
)
}),
)
// Accept any method so an incorrect redirect reaches an assertion,
// rather than being hidden behind the router's 405 response.
.route("/start", axum::routing::any(capture))
.route("/middle", axum::routing::any(capture))
.route("/final", axum::routing::any(capture))
.with_state(seen.clone());
let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let page_url = format!("http://{addr}/page");
let mut ctx = TestContext::new();
let mut bc = BrowserContext::new("BID-1".into());
bc.set_active_target_id("TID-1".to_owned());
bc.attach_active_session("SID-1".to_owned());
ctx.conn.install_browser_context_fixture_for_test(bc);
ctx.install_navigation_fixture_for_session_owner(&page_url, Some("SID-1"))
.await;
ctx.sent.clear();
ctx.process_async(json!({
"id": 79_300, "method": "Network.enable", "sessionId": "SID-1"
}))
.await;
ctx.expect_result(79_300, json!({}), Some("SID-1"));
ctx.process_async(json!({
"id": 79_301, "method": "Runtime.evaluate", "sessionId": "SID-1",
"params": {"expression": "document.getElementById('f').submit(); 'scheduled'"}
}))
.await;
let reply = ctx.take_response_by_id(79_301);
assert!(reply.get("error").is_none(), "{reply}");
assert!(reply["result"].get("exceptionDetails").is_none(), "{reply}");
let final_url = format!("http://{addr}/final");
wait_until_messages(
&mut ctx,
"SID-1",
"completed form redirect chain",
|messages| {
let id = messages.iter().find_map(|message| {
(message["method"] == json!("Network.requestWillBeSent")
&& message["params"]["request"]["url"] == json!(final_url))
.then(|| message["params"]["requestId"].as_str())
.flatten()
});
id.is_some_and(|id| {
messages.iter().any(|message| {
message["method"] == json!("Network.loadingFinished")
&& message["params"]["requestId"] == json!(id)
})
})
},
)
.await;
server.abort();
let wire = seen.lock().clone();
assert_eq!(
wire.len(),
3,
"exactly one request per redirect hop: {wire:?}"
);
// HTTP expectations are explicit, independent of both the server records
// and the protocol implementation, so two equally wrong views cannot pass.
for (request, (path, method, body)) in wire.iter().zip([
("/start", "POST", "account=alice"),
("/middle", "GET", ""),
("/final", "GET", ""),
]) {
assert_eq!(request.path, path);
assert_eq!(request.method, method);
assert_eq!(request.body, body);
assert_eq!(request.headers.get("referer").unwrap(), page_url.as_str());
}
assert_eq!(
wire[0].headers.get("content-type").unwrap(),
"application/x-www-form-urlencoded"
);
for request in &wire[1..] {
assert!(!request.headers.contains_key("content-type"));
assert!(!request.headers.contains_key("content-length"));
}
let requests = ctx
.sent
.iter()
.filter(|message| {
message["method"] == json!("Network.requestWillBeSent")
&& wire.iter().any(|request| {
message["params"]["request"]["url"]
== json!(format!("http://{addr}{}", request.path))
})
})
.collect::<Vec<_>>();
assert_eq!(requests.len(), 3);
let request_id = &requests[0]["params"]["requestId"];
for (index, (event, actual)) in requests.iter().zip(&wire).enumerate() {
assert_eq!(&event["params"]["requestId"], request_id);
assert_eq!(event["params"]["request"]["method"], actual.method);
assert_eq!(
event["params"]["request"]["url"],
format!("http://{addr}{}", actual.path)
);
// Ordinary request events describe the request before all transport
// headers are finalized. Compare wire headers through ExtraInfo below;
// Chromium can retain Content-Type here after a redirect removed it.
if index == 0 {
assert_eq!(event["params"]["request"]["postData"], "account=alice");
} else {
assert!(event["params"]["request"].get("postData").is_none());
assert_eq!(
event["params"]["redirectResponse"]["status"],
[302, 307][index - 1]
);
}
}
let extra = ctx
.sent
.iter()
.filter(|message| {
message["method"] == json!("Network.requestWillBeSentExtraInfo")
&& &message["params"]["requestId"] == request_id
})
.collect::<Vec<_>>();
assert_eq!(extra.len(), 3, "one wire header observation per hop");
for (event, actual) in extra.iter().zip(&wire) {
for name in ["referer", "content-type", "content-length"] {
assert_eq!(
event_header(&event["params"]["headers"], name),
actual
.headers
.get(name)
.map(|value| value.to_str().unwrap()),
"{} {name}",
actual.path
);
}
}
}
@@ -15,6 +15,168 @@ use crate::domains::network::{
use super::*;
#[tokio::test(flavor = "multi_thread")]
async fn durable_response_body_survives_navigation_only_for_opted_session() {
let mut ctx = TestContext::new();
let mut bc = BrowserContext::new("BID-durable".into());
bc.set_active_target_id("TID-durable".to_owned());
bc.attach_active_session("SID-durable".to_owned());
assert!(bc.assign_attached_session_to_target("TID-durable", "SID-ordinary".to_owned()));
ctx.conn.install_browser_context_fixture_for_test(bc);
for (id, session, params) in [
(
79_100,
"SID-durable",
json!({"enableDurableMessages":true,"maxTotalBufferSize":1024,"maxResourceBufferSize":128}),
),
(79_101, "SID-ordinary", json!({})),
] {
ctx.process_async(
json!({"id":id,"method":"Network.enable","sessionId":session,"params":params}),
)
.await;
ctx.expect_result(id, json!({}), Some(session));
}
let bc = ctx.conn.browser_context.as_mut().unwrap();
bc.record_captured_response_body(
"REQ-durable".into(),
"retained".into(),
[Some("SID-durable".into()), Some("SID-ordinary".into())],
);
bc.active_page_target_mut()
.prepare_document_navigation_request_ids(
&mut crate::conn::ConnectionNetworkRequestIdAllocator::default(),
true,
true,
false,
);
ctx.process_async(json!({"id":79_102,"method":"Network.getResponseBody","sessionId":"SID-durable","params":{"requestId":"REQ-durable"}})).await;
ctx.expect_result(
79_102,
json!({"body":"retained","base64Encoded":false}),
Some("SID-durable"),
);
ctx.process_async(json!({"id":79_103,"method":"Network.getResponseBody","sessionId":"SID-ordinary","params":{"requestId":"REQ-durable"}})).await;
ctx.expect_error(79_103, -32000, "No resource with given identifier found");
}
#[tokio::test(flavor = "multi_thread")]
async fn durable_response_body_requires_explicit_positive_total_budget() {
let mut ctx = TestContext::new();
let mut bc = BrowserContext::new("BID-durable".into());
bc.set_active_target_id("TID-durable".to_owned());
bc.attach_active_session("SID-durable".to_owned());
ctx.conn.install_browser_context_fixture_for_test(bc);
for (index, params) in [
json!({"enableDurableMessages":true}),
json!({"enableDurableMessages":true,"maxTotalBufferSize":0}),
json!({"enableDurableMessages":true,"maxTotalBufferSize":-1}),
json!({"enableDurableMessages":true,"maxTotalBufferSize":100,"maxResourceBufferSize":0}),
json!({"enableDurableMessages":"true","maxTotalBufferSize":100}),
]
.into_iter()
.enumerate()
{
let id = 79_104 + index as u64;
ctx.process_async(
json!({"id":id,"method":"Network.enable","sessionId":"SID-durable","params":params}),
)
.await;
ctx.expect_error(id, -32602, "InvalidParams");
}
assert!(
!ctx.conn
.browser_context
.as_ref()
.unwrap()
.active_page_target()
.runtime_slot
.has_network_event_listeners()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn durable_response_body_disable_detach_close_and_other_target_do_not_leak() {
let mut ctx = TestContext::new();
let mut bc = BrowserContext::new("BID-durable".into());
bc.set_active_target_id("TID-durable".to_owned());
bc.attach_active_session("SID-owner".to_owned());
assert!(bc.assign_attached_session_to_target("TID-durable", "SID-detach".to_owned()));
bc.insert_page_target_host(PageTargetHost::new(
"TID-other".into(),
Some("SID-other".into()),
TargetIdentityState::about_blank(),
TargetPageSlot::empty_for_test_fixture(),
));
ctx.conn.install_browser_context_fixture_for_test(bc);
let params =
json!({"enableDurableMessages":true,"maxTotalBufferSize":1024,"maxResourceBufferSize":128});
for (index, session) in ["SID-owner", "SID-detach", "SID-other"]
.into_iter()
.enumerate()
{
let id = 79_200 + index as u64;
ctx.process_async(
json!({"id":id,"method":"Network.enable","sessionId":session,"params":params}),
)
.await;
ctx.expect_result(id, json!({}), Some(session));
}
let bc = ctx.conn.browser_context.as_mut().unwrap();
bc.page_target_mut("TID-durable")
.unwrap()
.runtime_slot
.record_captured_response_body(
"REQ-old".into(),
"retained".into(),
[Some("SID-owner".into()), Some("SID-detach".into())],
);
bc.page_target_mut("TID-durable")
.unwrap()
.prepare_document_navigation_request_ids(
&mut crate::conn::ConnectionNetworkRequestIdAllocator::default(),
true,
true,
false,
);
ctx.process_async(json!({"id":79_203,"method":"Network.getResponseBody","sessionId":"SID-other","params":{"requestId":"REQ-old"}})).await;
ctx.expect_error(79_203, -32000, "No resource with given identifier found");
ctx.process_async(json!({"id":79_204,"method":"Network.disable","sessionId":"SID-owner"}))
.await;
ctx.expect_result(79_204, json!({}), Some("SID-owner"));
ctx.process_async(
json!({"id":79_205,"method":"Network.enable","sessionId":"SID-owner","params":params}),
)
.await;
ctx.expect_result(79_205, json!({}), Some("SID-owner"));
ctx.process_async(json!({"id":79_206,"method":"Network.getResponseBody","sessionId":"SID-owner","params":{"requestId":"REQ-old"}})).await;
ctx.expect_error(79_206, -32000, "No resource with given identifier found");
ctx.process_async(
json!({"id":79_207,"method":"Target.detachFromTarget","params":{"sessionId":"SID-detach"}}),
)
.await;
ctx.expect_result(79_207, json!({}), None);
assert!(
ctx.conn
.browser_context
.as_ref()
.unwrap()
.page_target("TID-durable")
.unwrap()
.runtime_slot
.captured_response_body("REQ-old")
.is_none()
);
ctx.process_async(
json!({"id":79_208,"method":"Target.closeTarget","params":{"targetId":"TID-durable"}}),
)
.await;
ctx.expect_result(79_208, json!({"success":true}), None);
let bc = ctx.conn.browser_context.as_ref().unwrap();
assert!(bc.page_target("TID-durable").is_none());
assert!(bc.page_target("TID-other").is_some());
}
fn bidi_network_context(session_id: &str) -> crate::devtools_runtime::DevToolsCommandContext {
crate::devtools_runtime::DevToolsCommandContext {
protocol: crate::devtools_runtime::DevToolsProtocol::WebDriverBidi,
+110
View File
@@ -4121,6 +4121,8 @@ mod producer_tests {
request_url: "https://example.test/child".to_owned(),
request_method: "GET".to_owned(),
request_headers: vec![("Accept".to_owned(), "text/html".to_owned())],
network_observation_journal: Default::default(),
redirect_chain: Vec::new(),
final_url: "https://example.test/child".to_owned(),
status: 200,
response_headers: vec![("Content-Type".to_owned(), "text/html".to_owned())],
@@ -4259,6 +4261,8 @@ mod producer_tests {
request_url: "https://example.test/retired-child".to_owned(),
request_method: "GET".to_owned(),
request_headers: Vec::new(),
network_observation_journal: Default::default(),
redirect_chain: Vec::new(),
final_url: "https://example.test/retired-child".to_owned(),
status: 200,
response_headers: vec![("Content-Type".to_owned(), "text/html".to_owned())],
@@ -4367,6 +4371,8 @@ mod producer_tests {
request_url: "https://example.test/legacy-child".to_owned(),
request_method: "GET".to_owned(),
request_headers: Vec::new(),
network_observation_journal: Default::default(),
redirect_chain: Vec::new(),
final_url: "https://example.test/legacy-child".to_owned(),
status: 200,
response_headers: vec![("Content-Type".to_owned(), "text/html".to_owned())],
@@ -4408,6 +4414,110 @@ mod producer_tests {
);
}
#[tokio::test(flavor = "multi_thread")]
async fn child_document_network_reports_actual_redirect_methods_and_request_headers() {
let mut conn = CdpConnection::default();
let mut bc = BrowserContext::new("BID-1".into());
bc.set_active_target_id("TID-1");
bc.attach_active_session("SID-1");
conn.install_browser_context_fixture_for_test(bc);
assert!(conn.enable_network_listener_for_session_owner(Some("SID-1")));
let final_url = url::Url::parse("https://example.test/done").unwrap();
let snapshot = ChildFrameDocumentNetworkSnapshot {
request_url: "https://example.test/profile".to_owned(),
request_method: "POST".to_owned(),
request_headers: Vec::new(),
network_observation_journal: moli_fetch::NetworkObservationJournal::from_exchanges(
vec![
moli_fetch::NetworkExchangeObservation::new(
moli_fetch::NetworkRequestObservation::new_with_method(
"POST",
vec![(
"Referer".to_owned(),
"https://example.test/source".to_owned(),
)],
),
Some(moli_fetch::NetworkResponseObservation::new(
302,
vec![("Location".to_owned(), final_url.to_string())],
)),
),
moli_fetch::NetworkExchangeObservation::new(
moli_fetch::NetworkRequestObservation::new_with_method(
"GET",
vec![(
"Referer".to_owned(),
"https://example.test/source".to_owned(),
)],
),
Some(moli_fetch::NetworkResponseObservation::new(
200,
vec![("Content-Type".to_owned(), "text/html".to_owned())],
)),
),
],
),
redirect_chain: vec![moli_fetch::RedirectInfo {
source: moli_fetch::RedirectSource::Network,
from_url: url::Url::parse("https://example.test/profile").unwrap(),
to_url: final_url.clone(),
status: 302,
headers: vec![("Location".to_owned(), final_url.to_string())],
network_extra_info_available: true,
request_extra_info: None,
response_extra_info: None,
redirect_has_extra_info: true,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
from_cache: false,
negotiated_http_version: None,
}],
final_url: final_url.to_string(),
status: 200,
response_headers: vec![("Content-Type".to_owned(), "text/html".to_owned())],
encoded_data_length: 0,
response_body: None,
from_cache: false,
};
let mut background_events = Vec::new();
crate::domains::network::emit_child_document_navigation_network_background_events(
&mut conn,
&mut background_events,
Some("SID-1"),
"CHILD-FRAME-REDIRECT",
"LID-CHILD-REDIRECT",
"LID-CHILD-REDIRECT",
12.5,
&snapshot,
);
let messages = protocol_messages_from_background_events(background_events);
let requests = messages
.iter()
.filter(|message| message["method"] == json!("Network.requestWillBeSent"))
.collect::<Vec<_>>();
assert_eq!(requests.len(), 2, "redirect hop must remain observable");
assert_eq!(requests[0]["params"]["request"]["method"], json!("POST"));
assert_eq!(requests[1]["params"]["request"]["method"], json!("GET"));
assert_eq!(
requests[1]["params"]["redirectResponse"]["status"],
json!(302)
);
assert_eq!(
requests[0]["params"]["request"]["headers"]["Referer"],
json!("https://example.test/source")
);
let request_extra = messages
.iter()
.filter(|message| message["method"] == json!("Network.requestWillBeSentExtraInfo"))
.collect::<Vec<_>>();
assert_eq!(request_extra.len(), 2);
assert!(request_extra.iter().all(|message| {
message["params"]["headers"]["Referer"] == json!("https://example.test/source")
}));
}
#[tokio::test(flavor = "multi_thread")]
async fn child_frame_activity_drain_preserves_prepared_attachment_only_token() {
let mut conn = CdpConnection::default();
@@ -765,6 +765,7 @@ fn register_dedicated_worker_target(
false,
None,
&[],
true,
);
if !events.is_empty() {
outputs.push(WorkerTargetLifecycleOutput::DedicatedWorkerEvents {
@@ -3369,6 +3370,7 @@ fn emit_service_worker_fetch_diagnostic_events(
false,
None,
&[],
true,
);
tag_service_worker_fetch_diagnostic_event(out.get_mut(request_event_index), diagnostic);
@@ -229,20 +229,24 @@ impl JsContextHost {
request_url,
request_method,
request_headers,
moli_fetch::NetworkObservationJournal::default(),
head,
body,
&parent_character_set,
);
}
let response = task_resource_loader
.fetch_raw(request)
.fetch_raw_with_network_metadata(request)
.await
.map_err(|error| error.to_string())?;
let (response, network_observation_journal) =
response.into_parts_with_observation_journal();
let (head, body) = response.into_body();
child_document_load_outcome_from_response(
request_url,
request_method,
request_headers,
network_observation_journal,
head,
body,
&parent_character_set,
@@ -737,7 +741,8 @@ fn child_document_load_outcome_from_response(
request_url: String,
request_method: String,
request_headers: Vec<(String, String)>,
head: moli_fetch::ResponseHead,
network_observation_journal: moli_fetch::NetworkObservationJournal,
mut head: moli_fetch::ResponseHead,
body: moli_fetch::ResponseBody,
parent_character_set: &str,
) -> Result<ChildDocumentLoadOutcome, String> {
@@ -763,6 +768,7 @@ fn child_document_load_outcome_from_response(
};
let policy_container =
DocumentPolicyContainer::from_navigation_response_headers(&head.headers, &head.final_url);
let redirect_chain = std::mem::take(&mut head.redirect_chain);
Ok(ChildDocumentLoadOutcome::Loaded(Box::new(
LoadedChildDocument {
final_url: head.final_url.clone(),
@@ -774,6 +780,8 @@ fn child_document_load_outcome_from_response(
request_url,
request_method,
request_headers,
network_observation_journal,
redirect_chain,
final_url: head.final_url.as_str().to_owned(),
status: head.status,
response_headers: head.headers.clone(),
@@ -832,20 +840,45 @@ mod tests {
use super::*;
#[test]
fn loaded_child_document_retains_exact_network_response_body() {
fn loaded_child_document_retains_exact_network_response_and_transport_metadata() {
let body_bytes = b"<!doctype html><p>child network body \xff</p>".to_vec();
let journal = moli_fetch::NetworkObservationJournal::from_exchanges(vec![
moli_fetch::NetworkExchangeObservation::new(
moli_fetch::NetworkRequestObservation::new_with_method(
"GET",
vec![("Referer".to_owned(), "https://example.test/".to_owned())],
),
Some(moli_fetch::NetworkResponseObservation::new(200, Vec::new())),
),
]);
let redirect = moli_fetch::RedirectInfo {
source: moli_fetch::RedirectSource::Network,
from_url: url::Url::parse("https://example.test/start").unwrap(),
to_url: url::Url::parse("https://example.test/child").unwrap(),
status: 302,
headers: Vec::new(),
network_extra_info_available: true,
request_extra_info: None,
response_extra_info: None,
redirect_has_extra_info: true,
request_cookie_report: None,
cookie_set_reports: Vec::new(),
from_cache: false,
negotiated_http_version: None,
};
let outcome = child_document_load_outcome_from_response(
"https://example.test/child".to_owned(),
"GET".to_owned(),
Vec::new(),
journal.clone(),
moli_fetch::ResponseHead {
final_url: url::Url::parse("https://example.test/child").unwrap(),
status: 200,
headers: vec![("Content-Type".to_owned(), "text/html".to_owned())],
request_cookie_report: None,
cookie_set_reports: Vec::new(),
redirected: false,
redirect_chain: Vec::new(),
redirected: true,
redirect_chain: vec![redirect.clone()],
from_cache: false,
negotiated_http_version: None,
},
@@ -861,6 +894,8 @@ mod tests {
.as_ref()
.expect("loaded child document should retain Network metadata");
assert_eq!(network.encoded_data_length, body_bytes.len());
assert_eq!(network.network_observation_journal, journal);
assert_eq!(network.redirect_chain, vec![redirect]);
assert_eq!(
network
.response_body
@@ -244,6 +244,43 @@ impl NavigationResourceLoader {
}
}
pub async fn fetch_raw_with_network_metadata(
&self,
request: Request,
) -> Result<NetworkFetchResult<RawResponse>> {
self.begin_fetch()?;
match self
.request_client
.fetch_raw_stream_with_cancel_and_network_metadata(
request.with_page_network_policy(),
self.inner.cancel.clone(),
)
.await
{
Ok(response) => {
let (response, observation_journal) =
response.into_parts_with_observation_journal();
match response.into_materialized_raw_response().await {
Ok(response) => {
self.finish_response_ready()?;
Ok(NetworkFetchResult::with_observation_journal(
response,
observation_journal,
))
}
Err(error) => {
self.finish_failed();
Err(error)
}
}
}
Err(error) => {
self.finish_failed();
Err(error)
}
}
}
pub async fn fetch_raw_stream(&self, request: Request) -> Result<StreamingRawResponse> {
self.fetch_raw_stream_with_network_metadata(request)
.await
@@ -732,6 +732,8 @@ pub(super) fn stale_loaded_completion(
request_url: request_url.to_owned(),
request_method: "GET".to_owned(),
request_headers: Vec::new(),
network_observation_journal: Default::default(),
redirect_chain: Vec::new(),
final_url: request_url.to_owned(),
status: 200,
response_headers: Vec::new(),
@@ -254,6 +254,45 @@ fn child_frame_document_network_charge(
.sum::<usize>()
.saturating_add(headers_charge(&snapshot.request_headers))
.saturating_add(headers_charge(&snapshot.response_headers))
.saturating_add(network_observation_journal_charge(
&snapshot.network_observation_journal,
))
.saturating_add(snapshot.redirect_chain.iter().fold(0, |total, redirect| {
total
.saturating_add(512)
.saturating_add(string_charge(redirect.from_url.as_str()))
.saturating_add(string_charge(redirect.to_url.as_str()))
.saturating_add(headers_charge(&redirect.headers))
.saturating_add(
redirect
.request_extra_info
.as_ref()
.map(|extra| {
headers_charge(&extra.headers)
.saturating_add(cookie_query_report_charge(&extra.cookie_report))
})
.unwrap_or_default(),
)
.saturating_add(
redirect
.response_extra_info
.as_ref()
.map(|extra| {
headers_charge(&extra.headers).saturating_add(cookie_query_report_charge(
&extra.request_extra_info.cookie_report,
))
})
.unwrap_or_default(),
)
.saturating_add(
redirect
.request_cookie_report
.as_ref()
.map(cookie_query_report_charge)
.unwrap_or_default(),
)
.saturating_add(redirect.cookie_set_reports.len().saturating_mul(128))
}))
.saturating_add(
snapshot
.response_body
@@ -263,6 +302,66 @@ fn child_frame_document_network_charge(
)
}
fn network_observation_journal_charge(journal: &moli_fetch::NetworkObservationJournal) -> usize {
journal.exchanges().iter().fold(0, |total, exchange| {
let request = exchange.request();
total
.saturating_add(512)
.saturating_add(request.method().map(string_charge).unwrap_or_default())
.saturating_add(headers_charge(request.headers()))
.saturating_add(
request
.cookie_report()
.map(cookie_query_report_charge)
.unwrap_or_default(),
)
.saturating_add(
exchange
.response()
.map(|response| headers_charge(response.headers()))
.unwrap_or_default(),
)
})
}
fn cookie_query_report_charge(report: &moli_cookie_jar::StoredCookieQueryReport) -> usize {
report
.included_cookies
.iter()
.chain(&report.excluded_cookies)
.fold(0, |total, access| {
let cookie = &access.cookie;
total
.saturating_add(512)
.saturating_add(string_charge(&cookie.name))
.saturating_add(string_charge(&cookie.value))
.saturating_add(string_charge(&cookie.domain))
.saturating_add(string_charge(&cookie.path))
.saturating_add(
cookie
.partition_key
.as_ref()
.and_then(|key| key.top_level_site())
.map(string_charge)
.unwrap_or_default(),
)
.saturating_add(
access
.site_for_cookies_url
.as_ref()
.map(|url| string_charge(url.as_str()))
.unwrap_or_default(),
)
.saturating_add(
access
.top_frame_origin_url
.as_ref()
.map(|url| string_charge(url.as_str()))
.unwrap_or_default(),
)
})
}
fn runtime_inspector_message_transport_charge_bytes(
message: &crate::runtime::RendererRuntimeInspectorMessage,
) -> usize {
@@ -487,3 +586,37 @@ fn navigation_response_transport_charge_bytes(
.unwrap_or(0),
)
}
#[cfg(test)]
mod tests {
use super::network_observation_journal_charge;
#[test]
fn network_observation_journal_charge_includes_methods_and_headers() {
let journal = moli_fetch::NetworkObservationJournal::from_exchanges(vec![
moli_fetch::NetworkExchangeObservation::new(
moli_fetch::NetworkRequestObservation::new_with_method(
"POST",
vec![("X-Request".to_owned(), "request-value".to_owned())],
),
Some(moli_fetch::NetworkResponseObservation::new(
200,
vec![("X-Response".to_owned(), "response-value".to_owned())],
)),
),
]);
let empty_exchange = moli_fetch::NetworkObservationJournal::from_exchanges(vec![
moli_fetch::NetworkExchangeObservation::new(
moli_fetch::NetworkRequestObservation::new(Vec::new()),
None,
),
]);
assert!(
network_observation_journal_charge(&journal)
> network_observation_journal_charge(&empty_exchange),
"transport accounting must charge retained method and header strings"
);
}
}