mod events;
use std::collections::{HashMap, HashSet, VecDeque};
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Context, Poll};
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD as BASE64;
use futures::{Stream, StreamExt};
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::warn;
use zendriver_transport::{AccountedRawEvent, SessionHandle};
use crate::url_matcher::UrlMatcher;
use events::{
DataReceived, EventSourceMessage, LoadingFailed, RequestIdOnly, RequestWillBeSent,
ResponseReceived, WebSocketCreated, WebSocketFrameEvent,
};
const CHANNEL_CAP: usize = 1024;
const MAX_TRACKED: usize = 10_000;
#[derive(Debug, Clone)]
pub enum NetworkEvent {
Http(NetworkExchange),
HttpData {
request_id: String,
chunk: Vec<u8>,
},
WebSocketOpen {
request_id: String,
url: String,
},
WebSocketFrame {
request_id: String,
direction: FrameDirection,
opcode: u8,
payload: String,
},
WebSocketClose {
request_id: String,
},
EventSourceMessage {
request_id: String,
event_name: String,
event_id: String,
data: String,
},
DeliveryBoundary(NetworkDeliveryBoundary),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum NetworkDeliveryBoundary {
Lagged {
missed: u64,
generation: u64,
},
Reconnected {
previous: u64,
generation: u64,
},
Disconnected {
generation: u64,
},
CorrelationEvicted {
url: String,
},
DecodeFailed,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FrameDirection {
Sent,
Received,
}
#[derive(Debug, Clone)]
pub struct MonitoredRequest {
pub url: String,
pub method: String,
pub headers: HashMap<String, String>,
pub post_data: Option<String>,
}
#[derive(Debug, Clone)]
pub struct MonitoredResponse {
pub status: u16,
pub status_text: String,
pub headers: HashMap<String, String>,
pub mime_type: String,
}
#[derive(Clone)]
pub struct NetworkExchange {
pub request: MonitoredRequest,
pub response: Option<MonitoredResponse>,
pub error: Option<String>,
pub(crate) request_id: String,
pub(crate) session: SessionHandle,
}
impl std::fmt::Debug for NetworkExchange {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("NetworkExchange")
.field("request", &self.request)
.field("response", &self.response)
.field("error", &self.error)
.finish()
}
}
impl NetworkExchange {
#[must_use]
pub fn request_id(&self) -> &str {
&self.request_id
}
#[must_use]
pub fn status(&self) -> Option<u16> {
self.response.as_ref().map(|r| r.status)
}
#[must_use]
pub fn is_success(&self) -> bool {
matches!(self.status(), Some(s) if (200..300).contains(&s))
}
pub async fn body(&self) -> crate::Result<Vec<u8>> {
let res = self
.session
.call(
"Network.getResponseBody",
serde_json::json!({ "requestId": self.request_id }),
)
.await
.map_err(|e| crate::ZendriverError::NetworkMonitor(format!("getResponseBody: {e}")))?;
let body = res
.get("body")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
if res
.get("base64Encoded")
.and_then(serde_json::Value::as_bool)
.unwrap_or(false)
{
BASE64
.decode(body)
.map_err(|e| crate::ZendriverError::NetworkMonitor(format!("base64: {e}")))
} else {
Ok(body.as_bytes().to_vec())
}
}
pub async fn text(&self) -> crate::Result<String> {
Ok(String::from_utf8_lossy(&self.body().await?).into_owned())
}
}
pub struct MonitorBuilder {
session: SessionHandle,
url_pattern: Option<UrlMatcher>,
stream_bodies: bool,
}
impl MonitorBuilder {
pub(crate) fn new(session: SessionHandle) -> Self {
Self {
session,
url_pattern: None,
stream_bodies: false,
}
}
#[must_use]
pub fn url_pattern(mut self, pattern: impl Into<UrlMatcher>) -> Self {
self.url_pattern = Some(pattern.into());
self
}
#[must_use]
pub fn stream_bodies(mut self, enabled: bool) -> Self {
self.stream_bodies = enabled;
self
}
pub async fn start(self) -> crate::Result<NetworkMonitor> {
let (tx, rx) = mpsc::channel(CHANNEL_CAP);
let cancel = CancellationToken::new();
let task = tokio::spawn(run_monitor(
self.session,
self.url_pattern,
self.stream_bodies,
tx,
cancel.clone(),
));
Ok(NetworkMonitor {
rx,
cancel,
_task: task,
})
}
}
pub struct NetworkMonitor {
rx: mpsc::Receiver<NetworkEvent>,
cancel: CancellationToken,
_task: JoinHandle<()>,
}
impl NetworkMonitor {
pub fn stop(self) {
self.cancel.cancel();
}
}
impl Drop for NetworkMonitor {
fn drop(&mut self) {
self.cancel.cancel();
}
}
impl Stream for NetworkMonitor {
type Item = NetworkEvent;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<NetworkEvent>> {
self.rx.poll_recv(cx)
}
}
type PartialExchange = (MonitoredRequest, Option<MonitoredResponse>);
async fn run_monitor(
session: SessionHandle,
filter: Option<UrlMatcher>,
stream_bodies: bool,
tx: mpsc::Sender<NetworkEvent>,
cancel: CancellationToken,
) {
let session_id = session.session_id().to_string();
let warned_stream_unsupported = Arc::new(AtomicBool::new(false));
let mut events = session.connection().subscribe_raw_accounted();
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
let mut streaming: HashSet<String> = HashSet::new();
let enable_session = session.clone();
tokio::spawn(async move {
if let Err(e) = enable_session
.call("Network.enable", serde_json::json!({}))
.await
{
warn!(error = %e, "network monitor: Network.enable failed; events may be inactive");
}
});
loop {
tokio::select! {
() = cancel.cancelled() => return,
next = events.next() => {
let Some(acc) = next else { return };
match acc {
AccountedRawEvent::Event { event: ev, .. } => {
if ev.session_id.as_deref() != Some(session_id.as_str()) {
continue;
}
match ev.method.as_str() {
"Network.requestWillBeSent" => {
let Ok(p) = serde_json::from_value::<RequestWillBeSent>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
urls.insert(p.request_id.clone(), p.request.url.clone());
track_order(&mut order, &urls, p.request_id.clone());
if urls.len() > MAX_TRACKED {
let evicted = evict_oldest(&mut urls, &mut partial, &mut order);
if let Some(url) = evicted {
if emit_boundary(&tx, NetworkDeliveryBoundary::CorrelationEvicted { url })
.await
{
return;
}
}
}
if stream_bodies
&& !warned_stream_unsupported.load(Ordering::Relaxed)
&& !streaming.contains(&p.request_id)
&& filter_allows(filter.as_ref(), Some(p.request.url.as_str()))
{
streaming.insert(p.request_id.clone());
spawn_stream_resource_content(
&session,
&tx,
&warned_stream_unsupported,
p.request_id.clone(),
);
}
let req = MonitoredRequest {
url: p.request.url,
method: p.request.method,
headers: p.request.headers,
post_data: p.request.post_data,
};
partial.insert(p.request_id, (req, None));
}
"Network.responseReceived" => {
let Ok(p) = serde_json::from_value::<ResponseReceived>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if let Some(entry) = partial.get_mut(&p.request_id) {
entry.1 = Some(MonitoredResponse {
status: p.response.status,
status_text: p.response.status_text,
headers: p.response.headers,
mime_type: p.response.mime_type,
});
}
}
"Network.dataReceived" if stream_bodies => {
let Ok(p) = serde_json::from_value::<DataReceived>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
let Some(data) = p.data else { continue };
match BASE64.decode(&data) {
Ok(chunk) => {
if tx
.send(NetworkEvent::HttpData { request_id: p.request_id, chunk })
.await
.is_err()
{
return;
}
}
Err(_) => {
if emit_decode_failed(&tx).await {
return;
}
}
}
}
"Network.loadingFinished" => {
let Ok(p) = serde_json::from_value::<RequestIdOnly>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if let Some((req, resp)) = partial.remove(&p.request_id) {
if filter_allows(filter.as_ref(), Some(&req.url)) {
let exchange = NetworkExchange {
request: req,
response: resp,
error: None,
request_id: p.request_id.clone(),
session: session.clone(),
};
if tx.send(NetworkEvent::Http(exchange)).await.is_err() {
return;
}
}
}
urls.remove(&p.request_id);
streaming.remove(&p.request_id);
}
"Network.loadingFailed" => {
let Ok(p) = serde_json::from_value::<LoadingFailed>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if let Some((req, resp)) = partial.remove(&p.request_id) {
if filter_allows(filter.as_ref(), Some(&req.url)) {
let exchange = NetworkExchange {
request: req,
response: resp,
error: Some(p.error_text),
request_id: p.request_id.clone(),
session: session.clone(),
};
if tx.send(NetworkEvent::Http(exchange)).await.is_err() {
return;
}
}
}
urls.remove(&p.request_id);
streaming.remove(&p.request_id);
}
"Network.webSocketCreated" => {
let Ok(p) = serde_json::from_value::<WebSocketCreated>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
urls.insert(p.request_id.clone(), p.url.clone());
track_order(&mut order, &urls, p.request_id.clone());
if urls.len() > MAX_TRACKED {
let evicted = evict_oldest(&mut urls, &mut partial, &mut order);
if let Some(url) = evicted {
if emit_boundary(&tx, NetworkDeliveryBoundary::CorrelationEvicted { url })
.await
{
return;
}
}
}
if filter_allows(filter.as_ref(), Some(&p.url))
&& tx
.send(NetworkEvent::WebSocketOpen {
request_id: p.request_id,
url: p.url,
})
.await
.is_err()
{
return;
}
}
"Network.webSocketFrameSent" | "Network.webSocketFrameReceived" => {
let direction = if ev.method.ends_with("Sent") {
FrameDirection::Sent
} else {
FrameDirection::Received
};
let Ok(p) = serde_json::from_value::<WebSocketFrameEvent>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if filter_allows(filter.as_ref(), urls.get(&p.request_id).map(String::as_str))
&& tx
.send(NetworkEvent::WebSocketFrame {
request_id: p.request_id,
direction,
opcode: p.response.opcode,
payload: p.response.payload_data,
})
.await
.is_err()
{
return;
}
}
"Network.webSocketClosed" => {
let Ok(p) = serde_json::from_value::<RequestIdOnly>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if filter_allows(filter.as_ref(), urls.get(&p.request_id).map(String::as_str))
&& tx
.send(NetworkEvent::WebSocketClose {
request_id: p.request_id.clone(),
})
.await
.is_err()
{
return;
}
urls.remove(&p.request_id);
}
"Network.eventSourceMessageReceived" => {
let Ok(p) = serde_json::from_value::<EventSourceMessage>(ev.params) else {
if emit_decode_failed(&tx).await {
return;
}
continue;
};
if filter_allows(filter.as_ref(), urls.get(&p.request_id).map(String::as_str))
&& tx
.send(NetworkEvent::EventSourceMessage {
request_id: p.request_id,
event_name: p.event_name,
event_id: p.event_id,
data: p.data,
})
.await
.is_err()
{
return;
}
}
_ => {}
}
}
AccountedRawEvent::Lagged { generation, missed } => {
partial.clear();
urls.clear();
streaming.clear();
if emit_boundary(&tx, NetworkDeliveryBoundary::Lagged { missed, generation }).await {
return;
}
}
AccountedRawEvent::Reconnected { previous, generation } => {
partial.clear();
urls.clear();
streaming.clear();
if emit_boundary(&tx, NetworkDeliveryBoundary::Reconnected { previous, generation }).await {
return;
}
}
AccountedRawEvent::Disconnected { generation } => {
partial.clear();
urls.clear();
streaming.clear();
let _ = emit_boundary(&tx, NetworkDeliveryBoundary::Disconnected { generation }).await;
return;
}
#[allow(unreachable_patterns)]
_ => {
partial.clear();
urls.clear();
streaming.clear();
if emit_boundary(&tx, NetworkDeliveryBoundary::Unknown).await {
return;
}
}
}
}
}
}
}
async fn emit_boundary(tx: &mpsc::Sender<NetworkEvent>, boundary: NetworkDeliveryBoundary) -> bool {
tx.send(NetworkEvent::DeliveryBoundary(boundary))
.await
.is_err()
}
async fn emit_decode_failed(tx: &mpsc::Sender<NetworkEvent>) -> bool {
emit_boundary(tx, NetworkDeliveryBoundary::DecodeFailed).await
}
fn spawn_stream_resource_content(
session: &SessionHandle,
tx: &mpsc::Sender<NetworkEvent>,
warned: &Arc<AtomicBool>,
request_id: String,
) {
let session = session.clone();
let tx = tx.clone();
let warned = Arc::clone(warned);
tokio::spawn(async move {
match session
.call(
"Network.streamResourceContent",
serde_json::json!({ "requestId": request_id }),
)
.await
{
Ok(res) => {
let buffered = res
.get("bufferedData")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
if buffered.is_empty() {
return;
}
if let Ok(chunk) = BASE64.decode(buffered) {
if !chunk.is_empty() {
let _ = tx.send(NetworkEvent::HttpData { request_id, chunk }).await;
}
}
}
Err(e) => {
if !warned.swap(true, Ordering::Relaxed) {
warn!(
error = %e,
"network monitor: Network.streamResourceContent failed (needs Chrome ~124+); \
stream_bodies falling back to whole-body capture for this and future requests"
);
}
}
}
});
}
fn filter_allows(filter: Option<&UrlMatcher>, url: Option<&str>) -> bool {
match filter {
None => true,
Some(m) => url.is_some_and(|u| m.matches(u)),
}
}
fn track_order(order: &mut VecDeque<String>, urls: &HashMap<String, String>, id: String) {
order.push_back(id);
if order.len() > MAX_TRACKED * 2 {
order.retain(|k| urls.contains_key(k));
}
}
fn evict_oldest(
urls: &mut HashMap<String, String>,
partial: &mut HashMap<String, PartialExchange>,
order: &mut VecDeque<String>,
) -> Option<String> {
while let Some(id) = order.pop_front() {
if let Some(url) = urls.remove(&id) {
partial.remove(&id);
warn!("network monitor correlation map exceeded {MAX_TRACKED}; evicting oldest entry");
return Some(url);
}
}
None
}
#[cfg(test)]
#[allow(clippy::panic, clippy::unwrap_used)]
mod tests {
use std::time::Duration;
use serde_json::json;
use zendriver_transport::testing::MockConnection;
use super::*;
const SID: &str = "S1";
async fn spawn_monitor(
filter: Option<UrlMatcher>,
) -> (
NetworkMonitor,
MockConnection,
zendriver_transport::Connection,
) {
spawn_monitor_with(MockConnection::pair(), filter).await
}
async fn spawn_monitor_with(
pair: (MockConnection, zendriver_transport::Connection),
filter: Option<UrlMatcher>,
) -> (
NetworkMonitor,
MockConnection,
zendriver_transport::Connection,
) {
let (mut mock, conn) = pair;
let session = SessionHandle::new(conn.clone(), SID);
let mut builder = MonitorBuilder::new(session);
if let Some(f) = filter {
builder = builder.url_pattern(f);
}
let monitor = builder.start().await.unwrap();
let id = mock.expect_cmd("Network.enable").await;
mock.reply(id, json!({})).await;
(monitor, mock, conn)
}
async fn next_event(monitor: &mut NetworkMonitor) -> NetworkEvent {
tokio::time::timeout(Duration::from_secs(2), monitor.next())
.await
.expect("timed out waiting for a NetworkEvent")
.expect("monitor stream ended unexpectedly")
}
async fn assert_no_event(monitor: &mut NetworkMonitor) {
let res = tokio::time::timeout(Duration::from_millis(300), monitor.next()).await;
assert!(res.is_err(), "expected no event, got {res:?}");
}
#[tokio::test]
async fn http_request_correlates_to_one_exchange() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "1",
"request": {
"url": "https://example.com/api/users",
"method": "GET",
"headers": { "Accept": "application/json" }
}
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({
"requestId": "1",
"response": {
"status": 200,
"statusText": "OK",
"mimeType": "application/json"
}
}),
SID,
)
.await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "1" }), SID)
.await;
let event = next_event(&mut monitor).await;
let NetworkEvent::Http(exchange) = event else {
panic!("expected NetworkEvent::Http, got {event:?}");
};
assert_eq!(exchange.request.url, "https://example.com/api/users");
assert_eq!(exchange.request.method, "GET");
assert_eq!(exchange.status(), Some(200));
assert!(exchange.is_success());
assert!(exchange.error.is_none());
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn loading_failed_emits_error_exchange() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "7",
"request": { "url": "https://example.com/boom", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.loadingFailed",
json!({ "requestId": "7", "errorText": "net::ERR_ABORTED" }),
SID,
)
.await;
let event = next_event(&mut monitor).await;
let NetworkEvent::Http(exchange) = event else {
panic!("expected NetworkEvent::Http, got {event:?}");
};
assert_eq!(exchange.request.url, "https://example.com/boom");
assert!(exchange.response.is_none());
assert_eq!(exchange.status(), None);
assert_eq!(exchange.error.as_deref(), Some("net::ERR_ABORTED"));
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn ws_frames_emit_tagged_events() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.webSocketCreated",
json!({ "requestId": "ws1", "url": "wss://echo.example.com/socket" }),
SID,
)
.await;
mock.emit_event_for_session(
"Network.webSocketFrameSent",
json!({ "requestId": "ws1", "response": { "opcode": 1, "payloadData": "ping" } }),
SID,
)
.await;
mock.emit_event_for_session(
"Network.webSocketFrameReceived",
json!({ "requestId": "ws1", "response": { "opcode": 1, "payloadData": "pong" } }),
SID,
)
.await;
mock.emit_event_for_session(
"Network.webSocketClosed",
json!({ "requestId": "ws1" }),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::WebSocketOpen { request_id, url } => {
assert_eq!(request_id, "ws1");
assert_eq!(url, "wss://echo.example.com/socket");
}
other => panic!("expected WebSocketOpen, got {other:?}"),
}
match next_event(&mut monitor).await {
NetworkEvent::WebSocketFrame {
request_id,
direction,
opcode,
payload,
} => {
assert_eq!(request_id, "ws1");
assert_eq!(direction, FrameDirection::Sent);
assert_eq!(opcode, 1);
assert_eq!(payload, "ping");
}
other => panic!("expected WebSocketFrame(Sent), got {other:?}"),
}
match next_event(&mut monitor).await {
NetworkEvent::WebSocketFrame {
direction, payload, ..
} => {
assert_eq!(direction, FrameDirection::Received);
assert_eq!(payload, "pong");
}
other => panic!("expected WebSocketFrame(Received), got {other:?}"),
}
match next_event(&mut monitor).await {
NetworkEvent::WebSocketClose { request_id } => assert_eq!(request_id, "ws1"),
other => panic!("expected WebSocketClose, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn event_source_message_emits_event() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.eventSourceMessageReceived",
json!({
"requestId": "sse1",
"eventName": "update",
"eventId": "42",
"data": "tick"
}),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::EventSourceMessage {
request_id,
event_name,
event_id,
data,
} => {
assert_eq!(request_id, "sse1");
assert_eq!(event_name, "update");
assert_eq!(event_id, "42");
assert_eq!(data, "tick");
}
other => panic!("expected EventSourceMessage, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
async fn spawn_monitor_streaming(
filter: Option<UrlMatcher>,
) -> (
NetworkMonitor,
MockConnection,
zendriver_transport::Connection,
) {
let (mut mock, conn) = MockConnection::pair();
let session = SessionHandle::new(conn.clone(), SID);
let mut builder = MonitorBuilder::new(session).stream_bodies(true);
if let Some(f) = filter {
builder = builder.url_pattern(f);
}
let monitor = builder.start().await.unwrap();
let id = mock.expect_cmd("Network.enable").await;
mock.reply(id, json!({})).await;
(monitor, mock, conn)
}
#[tokio::test]
async fn stream_bodies_prepends_buffered_data_then_streams_data_received() {
let (mut monitor, mut mock, conn) = spawn_monitor_streaming(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s1",
"request": { "url": "https://example.com/big", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "s1", "response": { "status": 200 } }),
SID,
)
.await;
let id = mock.expect_cmd("Network.streamResourceContent").await;
assert_eq!(mock.last_sent()["params"]["requestId"], "s1");
mock.reply(id, json!({ "bufferedData": BASE64.encode("hello-") }))
.await;
match next_event(&mut monitor).await {
NetworkEvent::HttpData { request_id, chunk } => {
assert_eq!(request_id, "s1");
assert_eq!(chunk, b"hello-");
}
other => panic!("expected HttpData (bufferedData), got {other:?}"),
}
mock.emit_event_for_session(
"Network.dataReceived",
json!({
"requestId": "s1",
"timestamp": 1.0,
"dataLength": 5,
"encodedDataLength": 5,
"data": BASE64.encode("world")
}),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::HttpData { request_id, chunk } => {
assert_eq!(request_id, "s1");
assert_eq!(chunk, b"world");
}
other => panic!("expected HttpData (dataReceived), got {other:?}"),
}
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "s1" }), SID)
.await;
match next_event(&mut monitor).await {
NetworkEvent::Http(exchange) => {
assert_eq!(exchange.request_id(), "s1");
assert_eq!(exchange.request.url, "https://example.com/big");
}
other => panic!("expected NetworkEvent::Http, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn stream_bodies_empty_buffered_data_emits_no_leading_chunk() {
let (mut monitor, mut mock, conn) = spawn_monitor_streaming(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s2",
"request": { "url": "https://example.com/empty-buffer", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "s2", "response": { "status": 200 } }),
SID,
)
.await;
let id = mock.expect_cmd("Network.streamResourceContent").await;
mock.reply(id, json!({ "bufferedData": "" })).await;
mock.emit_event_for_session(
"Network.dataReceived",
json!({
"requestId": "s2",
"timestamp": 1.0,
"dataLength": 6,
"encodedDataLength": 6,
"data": BASE64.encode("chunk1")
}),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::HttpData { request_id, chunk } => {
assert_eq!(request_id, "s2");
assert_eq!(chunk, b"chunk1", "the only chunk — no empty leading one");
}
other => panic!("expected HttpData, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn stream_resource_content_error_degrades_without_breaking_monitor() {
let (mut monitor, mut mock, conn) = spawn_monitor_streaming(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s3",
"request": { "url": "https://example.com/old-chrome", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "s3", "response": { "status": 200 } }),
SID,
)
.await;
let id = mock.expect_cmd("Network.streamResourceContent").await;
mock.reply_err(id, -32601, "'Network.streamResourceContent' wasn't found")
.await;
assert_no_event(&mut monitor).await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "s3" }), SID)
.await;
match next_event(&mut monitor).await {
NetworkEvent::Http(exchange) => {
assert_eq!(exchange.request.url, "https://example.com/old-chrome");
}
other => panic!("expected NetworkEvent::Http, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn stream_bodies_false_never_enables_or_emits_http_data() {
let (mut monitor, mut mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s4",
"request": { "url": "https://example.com/opt-out", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "s4", "response": { "status": 200 } }),
SID,
)
.await;
assert!(
mock.try_recv_cmd().is_none(),
"stream_bodies: false must never issue Network.streamResourceContent"
);
mock.emit_event_for_session(
"Network.dataReceived",
json!({
"requestId": "s4",
"timestamp": 1.0,
"dataLength": 4,
"encodedDataLength": 4,
"data": BASE64.encode("data")
}),
SID,
)
.await;
assert_no_event(&mut monitor).await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "s4" }), SID)
.await;
match next_event(&mut monitor).await {
NetworkEvent::Http(exchange) => {
assert_eq!(exchange.request.url, "https://example.com/opt-out");
}
other => panic!("expected NetworkEvent::Http, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn stream_bodies_dedups_repeated_request_will_be_sent_for_same_request() {
let (mut monitor, mut mock, conn) = spawn_monitor_streaming(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s5",
"request": { "url": "https://example.com/dup", "method": "GET" }
}),
SID,
)
.await;
let id = mock.expect_cmd("Network.streamResourceContent").await;
mock.reply(id, json!({ "bufferedData": "" })).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "s5",
"request": { "url": "https://example.com/dup-redirected", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "s5", "response": { "status": 200 } }),
SID,
)
.await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "s5" }), SID)
.await;
let _ = next_event(&mut monitor).await; assert!(
mock.try_recv_cmd().is_none(),
"a duplicate requestWillBeSent must not re-issue Network.streamResourceContent"
);
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn dropping_monitor_cancels_correlator_task() {
let (monitor, mock, conn) = spawn_monitor(None).await;
let cancel = monitor.cancel.clone();
assert!(!cancel.is_cancelled());
drop(monitor);
assert!(cancel.is_cancelled(), "Drop must cancel the correlator");
drop(mock);
conn.shutdown();
}
#[tokio::test]
async fn url_filter_drops_unmatched() {
let (mut monitor, mock, conn) = spawn_monitor(Some("/api/".into())).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "2",
"request": { "url": "https://example.com/static/app.js", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "2" }), SID)
.await;
assert_no_event(&mut monitor).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "3",
"request": { "url": "https://example.com/api/orders", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.responseReceived",
json!({ "requestId": "3", "response": { "status": 201 } }),
SID,
)
.await;
mock.emit_event_for_session("Network.loadingFinished", json!({ "requestId": "3" }), SID)
.await;
let event = next_event(&mut monitor).await;
let NetworkEvent::Http(exchange) = event else {
panic!("expected NetworkEvent::Http, got {event:?}");
};
assert_eq!(exchange.request.url, "https://example.com/api/orders");
assert_eq!(exchange.status(), Some(201));
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn events_for_other_sessions_are_ignored() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "x",
"request": { "url": "https://other.example.com/api/x", "method": "GET" }
}),
"OTHER",
)
.await;
mock.emit_event_for_session(
"Network.loadingFinished",
json!({ "requestId": "x" }),
"OTHER",
)
.await;
assert_no_event(&mut monitor).await;
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn lagged_mid_exchange_clears_partial_and_emits_boundary() {
let (mut monitor, mock, conn) =
spawn_monitor_with(MockConnection::pair_with_accounted_capacity(2), None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "mid1",
"request": { "url": "https://example.com/mid", "method": "GET" }
}),
SID,
)
.await;
tokio::time::sleep(Duration::from_millis(50)).await;
for i in 0..5u32 {
mock.emit_event("Test.dummy", json!({ "i": i })).await;
}
match next_event(&mut monitor).await {
NetworkEvent::DeliveryBoundary(NetworkDeliveryBoundary::Lagged {
generation,
missed,
}) => {
assert_eq!(generation, 1);
assert!(missed > 0, "expected a nonzero missed count, got {missed}");
}
other => panic!("expected DeliveryBoundary::Lagged, got {other:?}"),
}
mock.emit_event_for_session(
"Network.loadingFinished",
json!({ "requestId": "mid1" }),
SID,
)
.await;
assert_no_event(&mut monitor).await;
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn reconnected_mid_exchange_clears_partial_and_emits_boundary() {
let (mut monitor, mut mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "recon1",
"request": { "url": "https://example.com/recon", "method": "GET" }
}),
SID,
)
.await;
tokio::time::sleep(Duration::from_millis(50)).await;
mock.reconnect(&conn);
match next_event(&mut monitor).await {
NetworkEvent::DeliveryBoundary(NetworkDeliveryBoundary::Reconnected {
previous,
generation,
}) => {
assert_eq!(previous, 1);
assert_eq!(generation, 2);
}
other => panic!("expected DeliveryBoundary::Reconnected, got {other:?}"),
}
mock.emit_event_for_session(
"Network.loadingFinished",
json!({ "requestId": "recon1" }),
SID,
)
.await;
assert_no_event(&mut monitor).await;
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn correlation_cap_exceeded_emits_correlation_evicted() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
for i in 0..=MAX_TRACKED {
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": format!("r{i}"),
"request": { "url": format!("https://example.com/{i}"), "method": "GET" }
}),
SID,
)
.await;
}
match next_event(&mut monitor).await {
NetworkEvent::DeliveryBoundary(NetworkDeliveryBoundary::CorrelationEvicted { url }) => {
assert!(
url.starts_with("https://example.com/"),
"evicted url should be one of the inserted entries, got {url}"
);
}
other => panic!("expected DeliveryBoundary::CorrelationEvicted, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn malformed_payload_emits_decode_failed_without_raw_payload() {
let (mut monitor, mock, conn) = spawn_monitor(None).await;
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({ "requestId": "bad1" }),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::DeliveryBoundary(NetworkDeliveryBoundary::DecodeFailed) => {}
other => panic!("expected DeliveryBoundary::DecodeFailed, got {other:?}"),
}
mock.emit_event_for_session(
"Network.requestWillBeSent",
json!({
"requestId": "good1",
"request": { "url": "https://example.com/ok", "method": "GET" }
}),
SID,
)
.await;
mock.emit_event_for_session(
"Network.loadingFinished",
json!({ "requestId": "good1" }),
SID,
)
.await;
match next_event(&mut monitor).await {
NetworkEvent::Http(exchange) => {
assert_eq!(exchange.request.url, "https://example.com/ok");
}
other => panic!("expected NetworkEvent::Http, got {other:?}"),
}
monitor.stop();
conn.shutdown();
}
#[tokio::test]
async fn disconnected_emits_boundary_and_ends_monitor_task() {
let (mut monitor, mock, _conn) = spawn_monitor(None).await;
mock.disconnect();
match next_event(&mut monitor).await {
NetworkEvent::DeliveryBoundary(NetworkDeliveryBoundary::Disconnected {
generation,
}) => {
assert_eq!(generation, 1);
}
other => panic!("expected DeliveryBoundary::Disconnected, got {other:?}"),
}
let next = tokio::time::timeout(Duration::from_secs(2), monitor.next())
.await
.expect("monitor stream did not end within 2s after Disconnected");
assert!(
next.is_none(),
"expected the monitor stream to end after Disconnected, got {next:?}"
);
}
async fn make_exchange(status: Option<u16>, error: Option<&str>) -> NetworkExchange {
let (_mock, conn) = MockConnection::pair();
let session = SessionHandle::new(conn, "test-session");
let req = MonitoredRequest {
url: "https://example.com/api".into(),
method: "GET".into(),
headers: HashMap::new(),
post_data: None,
};
let resp = status.map(|s| MonitoredResponse {
status: s,
status_text: "OK".into(),
headers: HashMap::new(),
mime_type: "application/json".into(),
});
NetworkExchange {
request: req,
response: resp,
error: error.map(ToOwned::to_owned),
request_id: "r1".into(),
session,
}
}
#[tokio::test]
async fn status_returns_none_when_no_response() {
let ex = make_exchange(None, None).await;
assert!(ex.status().is_none());
assert!(!ex.is_success());
}
#[tokio::test]
async fn status_returns_some_for_200() {
let ex = make_exchange(Some(200), None).await;
assert_eq!(ex.status(), Some(200));
assert!(ex.is_success());
}
#[tokio::test]
async fn status_304_is_not_success() {
let ex = make_exchange(Some(304), None).await;
assert!(!ex.is_success());
}
#[tokio::test]
async fn status_404_is_not_success() {
let ex = make_exchange(Some(404), None).await;
assert!(!ex.is_success());
}
#[tokio::test]
async fn debug_does_not_include_session_field() {
let ex = make_exchange(Some(200), None).await;
let s = format!("{ex:?}");
assert!(s.contains("NetworkExchange"));
assert!(s.contains("request"));
assert!(s.contains("response"));
assert!(!s.contains("session"));
}
#[tokio::test]
async fn error_field_is_set_on_failed_exchange() {
let ex = make_exchange(None, Some("net::ERR_ABORTED")).await;
assert_eq!(ex.error.as_deref(), Some("net::ERR_ABORTED"));
}
#[test]
fn frame_direction_copy_and_eq() {
let d = FrameDirection::Sent;
let d2 = d;
assert_eq!(d, d2);
assert_ne!(FrameDirection::Sent, FrameDirection::Received);
}
#[test]
fn network_event_debug_roundtrip() {
let ev = NetworkEvent::WebSocketOpen {
request_id: "r1".into(),
url: "wss://echo.example.com".into(),
};
let s = format!("{ev:?}");
assert!(s.contains("WebSocketOpen"));
assert!(s.contains("wss://echo.example.com"));
}
#[test]
fn network_delivery_boundary_variants_construct_and_debug() {
let variants = [
NetworkDeliveryBoundary::Lagged {
missed: 3,
generation: 1,
},
NetworkDeliveryBoundary::Reconnected {
previous: 1,
generation: 2,
},
NetworkDeliveryBoundary::Disconnected { generation: 1 },
NetworkDeliveryBoundary::CorrelationEvicted {
url: "https://example.com/evicted".into(),
},
NetworkDeliveryBoundary::DecodeFailed,
NetworkDeliveryBoundary::Unknown,
];
for v in &variants {
let cloned = v.clone();
assert_eq!(v, &cloned);
let ev = NetworkEvent::DeliveryBoundary(cloned);
let s = format!("{ev:?}");
assert!(s.contains("DeliveryBoundary"), "got {s}");
}
}
fn partial_entry(url: &str) -> PartialExchange {
(
MonitoredRequest {
url: url.into(),
method: "GET".into(),
headers: HashMap::new(),
post_data: None,
},
None,
)
}
#[test]
fn evict_oldest_drops_partial_and_mirrored_url() {
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
partial.insert("req1".into(), partial_entry("https://example.com/a"));
urls.insert("req1".into(), "https://example.com/a".into());
track_order(&mut order, &urls, "req1".into());
let evicted = evict_oldest(&mut urls, &mut partial, &mut order);
assert_eq!(
evicted.as_deref(),
Some("https://example.com/a"),
"returns the evicted entry's URL so the caller can report CorrelationEvicted"
);
assert!(partial.is_empty(), "partial entry must be evicted");
assert!(
urls.is_empty(),
"the partial entry's mirrored url must be evicted too"
);
}
#[test]
fn evict_oldest_falls_back_to_urls_only_entry() {
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
urls.insert("ws1".into(), "wss://echo.example.com".into());
track_order(&mut order, &urls, "ws1".into());
let evicted = evict_oldest(&mut urls, &mut partial, &mut order);
assert_eq!(evicted.as_deref(), Some("wss://echo.example.com"));
assert!(urls.is_empty(), "urls-only entry must be evicted");
}
#[test]
fn evict_oldest_evicts_in_insertion_order() {
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
for (id, url) in [
("req1", "https://example.com/1"),
("req2", "https://example.com/2"),
("req3", "https://example.com/3"),
] {
urls.insert(id.into(), url.into());
track_order(&mut order, &urls, id.into());
}
assert_eq!(
evict_oldest(&mut urls, &mut partial, &mut order).as_deref(),
Some("https://example.com/1"),
"oldest (req1) evicted first"
);
assert_eq!(
evict_oldest(&mut urls, &mut partial, &mut order).as_deref(),
Some("https://example.com/2"),
"then req2"
);
assert_eq!(
evict_oldest(&mut urls, &mut partial, &mut order).as_deref(),
Some("https://example.com/3"),
"then req3"
);
}
#[test]
fn evict_oldest_skips_tombstones_for_completed_requests() {
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
urls.insert("req1".into(), "https://example.com/1".into());
track_order(&mut order, &urls, "req1".into());
urls.insert("req2".into(), "https://example.com/2".into());
track_order(&mut order, &urls, "req2".into());
urls.remove("req1");
assert_eq!(
evict_oldest(&mut urls, &mut partial, &mut order).as_deref(),
Some("https://example.com/2"),
"req1 tombstone skipped; oldest live (req2) evicted"
);
}
#[test]
fn evict_oldest_on_empty_is_a_noop() {
let mut partial: HashMap<String, PartialExchange> = HashMap::new();
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
let evicted = evict_oldest(&mut urls, &mut partial, &mut order);
assert!(evicted.is_none());
assert!(partial.is_empty() && urls.is_empty());
}
#[test]
fn track_order_stays_bounded_despite_tombstones() {
let mut urls: HashMap<String, String> = HashMap::new();
let mut order: VecDeque<String> = VecDeque::new();
for i in 0..(MAX_TRACKED * 2 + 5) {
let id = format!("req{i}");
urls.insert(id.clone(), "https://example.com/x".into());
track_order(&mut order, &urls, id.clone());
urls.remove(&id); }
assert!(
order.len() <= MAX_TRACKED * 2,
"order must stay bounded via compaction; grew to {}",
order.len()
);
}
}