The CDP event broadcast buffer (256 events) was too small for pages with many concurrent API requests, causing silent event drops. Modern SPAs routinely fire 100+ API calls during page load, generating 300+ CDP network events that would overflow the buffer between drain cycles. Changes: - Increase CDP broadcast buffer from 256 to 4096 (event channel) and 512 to 4096 (raw channel) - Reduce background drain interval from 500ms to 100ms - Handle Network.loadingFailed events in HAR recording - Enable Network.enable on cross-origin iframe sessions during HAR recording and request tracking - Allow Network events from iframe sessions through the session filter - Log a warning when buffer overflow occurs instead of silently dropping Fixes #1128 Co-authored-by: ctate <366502+ctate@users.noreply.github.com>
362 lines
13 KiB
Rust
362 lines
13 KiB
Rust
use std::collections::HashMap;
|
|
use std::io::Write;
|
|
use std::sync::atomic::{AtomicU64, Ordering};
|
|
use std::sync::Arc;
|
|
|
|
use futures_util::{SinkExt, StreamExt};
|
|
use serde_json::Value;
|
|
use tokio::sync::{broadcast, oneshot, Mutex};
|
|
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
|
|
use tokio_tungstenite::tungstenite::protocol::WebSocketConfig;
|
|
use tokio_tungstenite::tungstenite::Message;
|
|
|
|
use super::types::{CdpCommand, CdpEvent, CdpMessage};
|
|
|
|
type PendingMap = Arc<Mutex<HashMap<u64, oneshot::Sender<CdpMessage>>>>;
|
|
|
|
/// Interval between WebSocket ping frames sent to keep the connection alive
|
|
/// through intermediate proxies (reverse proxies, load balancers, service meshes).
|
|
const WS_KEEPALIVE_INTERVAL_SECS: u64 = 30;
|
|
|
|
/// Raw incoming CDP message (text) broadcast to all subscribers.
|
|
/// Used by the inspect proxy to forward responses and events to DevTools.
|
|
#[derive(Debug, Clone)]
|
|
pub struct RawCdpMessage {
|
|
pub text: String,
|
|
pub session_id: Option<String>,
|
|
}
|
|
|
|
pub struct CdpClient {
|
|
ws_tx: Arc<
|
|
Mutex<
|
|
futures_util::stream::SplitSink<
|
|
tokio_tungstenite::WebSocketStream<
|
|
tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>,
|
|
>,
|
|
Message,
|
|
>,
|
|
>,
|
|
>,
|
|
next_id: AtomicU64,
|
|
pending: PendingMap,
|
|
event_tx: broadcast::Sender<CdpEvent>,
|
|
raw_tx: broadcast::Sender<RawCdpMessage>,
|
|
_reader_handle: tokio::task::JoinHandle<()>,
|
|
_keepalive_handle: tokio::task::JoinHandle<()>,
|
|
}
|
|
|
|
impl CdpClient {
|
|
pub async fn connect(url: &str) -> Result<Self, String> {
|
|
Self::connect_with_headers(url, None).await
|
|
}
|
|
|
|
pub async fn connect_with_headers(
|
|
url: &str,
|
|
headers: Option<Vec<(String, String)>>,
|
|
) -> Result<Self, String> {
|
|
let mut request = url
|
|
.into_client_request()
|
|
.map_err(|e| format!("Invalid WebSocket URL: {}", e))?;
|
|
|
|
if let Some(hdrs) = headers {
|
|
let req_headers = request.headers_mut();
|
|
for (key, value) in hdrs {
|
|
if let (Ok(name), Ok(val)) = (
|
|
key.parse::<tokio_tungstenite::tungstenite::http::header::HeaderName>(),
|
|
value.parse::<tokio_tungstenite::tungstenite::http::header::HeaderValue>(),
|
|
) {
|
|
req_headers.insert(name, val);
|
|
}
|
|
}
|
|
}
|
|
|
|
let ws_config = WebSocketConfig {
|
|
max_message_size: None,
|
|
max_frame_size: None,
|
|
..Default::default()
|
|
};
|
|
|
|
let (ws_stream, _) =
|
|
tokio_tungstenite::connect_async_with_config(request, Some(ws_config), false)
|
|
.await
|
|
.map_err(|e| format!("CDP WebSocket connect failed: {}", e))?;
|
|
|
|
enable_tcp_keepalive(ws_stream.get_ref());
|
|
|
|
let (ws_tx, mut ws_rx) = ws_stream.split();
|
|
let ws_tx = Arc::new(Mutex::new(ws_tx));
|
|
|
|
let pending: PendingMap = Arc::new(Mutex::new(HashMap::new()));
|
|
let (event_tx, _) = broadcast::channel(4096);
|
|
let (raw_tx, _) = broadcast::channel(4096);
|
|
|
|
let pending_clone = pending.clone();
|
|
let event_tx_clone = event_tx.clone();
|
|
let raw_tx_clone = raw_tx.clone();
|
|
|
|
// Notify used to stop the keepalive task when the reader loop exits.
|
|
let (cancel_tx, mut cancel_rx) = tokio::sync::watch::channel(false);
|
|
|
|
let reader_handle = tokio::spawn(async move {
|
|
while let Some(msg) = ws_rx.next().await {
|
|
// Accept both Text and Binary frames — remote CDP proxies
|
|
// (e.g. Browserless) may send responses as Binary frames.
|
|
let msg = match msg {
|
|
Ok(Message::Text(text)) => text,
|
|
Ok(Message::Binary(data)) => match String::from_utf8(data) {
|
|
Ok(text) => text,
|
|
Err(_) => continue,
|
|
},
|
|
Ok(Message::Close(frame)) => {
|
|
if std::env::var("AGENT_BROWSER_DEBUG").is_ok() {
|
|
let reason = frame
|
|
.as_ref()
|
|
.map(|f| format!("code={}, reason={}", f.code, f.reason))
|
|
.unwrap_or_else(|| "no frame".to_string());
|
|
let _ =
|
|
writeln!(std::io::stderr(), "[cdp] WebSocket Close: {}", reason);
|
|
}
|
|
break;
|
|
}
|
|
Ok(Message::Pong(_)) => continue,
|
|
Ok(_) => continue,
|
|
Err(e) => {
|
|
if std::env::var("AGENT_BROWSER_DEBUG").is_ok() {
|
|
let _ = writeln!(std::io::stderr(), "[cdp] WebSocket Error: {}", e);
|
|
}
|
|
break;
|
|
}
|
|
};
|
|
|
|
// Broadcast raw message for inspect proxy subscribers before typed parse,
|
|
// so messages with negative IDs (used by the inspect proxy) are still delivered.
|
|
if raw_tx_clone.receiver_count() > 0 {
|
|
let session_id = serde_json::from_str::<serde_json::Value>(&msg)
|
|
.ok()
|
|
.and_then(|v| v.get("sessionId")?.as_str().map(String::from));
|
|
let _ = raw_tx_clone.send(RawCdpMessage {
|
|
text: msg.clone(),
|
|
session_id,
|
|
});
|
|
}
|
|
|
|
let parsed: CdpMessage = match serde_json::from_str(&msg) {
|
|
Ok(m) => m,
|
|
// Expected for inspect proxy messages with negative IDs
|
|
// (CdpMessage.id is u64); handled via raw broadcast above.
|
|
Err(_) => continue,
|
|
};
|
|
|
|
if let Some(id) = parsed.id {
|
|
// Response to a command
|
|
let mut pending = pending_clone.lock().await;
|
|
if let Some(tx) = pending.remove(&id) {
|
|
let _ = tx.send(parsed);
|
|
}
|
|
} else if let Some(ref method) = parsed.method {
|
|
// Event
|
|
let event = CdpEvent {
|
|
method: method.clone(),
|
|
params: parsed.params.clone().unwrap_or(Value::Null),
|
|
session_id: parsed.session_id.clone(),
|
|
};
|
|
let _ = event_tx_clone.send(event);
|
|
}
|
|
}
|
|
|
|
// Reader loop exited (connection closed or error). Drop all pending
|
|
// command senders so callers get an immediate channel-closed error
|
|
// instead of waiting for the 30-second timeout.
|
|
pending_clone.lock().await.clear();
|
|
|
|
// Stop the keepalive task — the connection is gone.
|
|
let _ = cancel_tx.send(true);
|
|
});
|
|
|
|
// Spawn a keepalive task that sends WebSocket Ping frames at a regular
|
|
// interval. This prevents intermediate proxies (Envoy, nginx, OpenResty,
|
|
// cloud load balancers) from closing idle WebSocket connections. If the
|
|
// send fails, the connection is dead and we stop pinging.
|
|
let keepalive_tx = ws_tx.clone();
|
|
let keepalive_handle = tokio::spawn(async move {
|
|
let interval = std::time::Duration::from_secs(WS_KEEPALIVE_INTERVAL_SECS);
|
|
loop {
|
|
tokio::select! {
|
|
_ = tokio::time::sleep(interval) => {}
|
|
_ = cancel_rx.changed() => break,
|
|
}
|
|
let mut tx = keepalive_tx.lock().await;
|
|
if tx.send(Message::Ping(Vec::new())).await.is_err() {
|
|
break;
|
|
}
|
|
}
|
|
});
|
|
|
|
Ok(Self {
|
|
ws_tx,
|
|
next_id: AtomicU64::new(1),
|
|
pending,
|
|
event_tx,
|
|
raw_tx,
|
|
_reader_handle: reader_handle,
|
|
_keepalive_handle: keepalive_handle,
|
|
})
|
|
}
|
|
|
|
pub async fn send_command(
|
|
&self,
|
|
method: &str,
|
|
params: Option<Value>,
|
|
session_id: Option<&str>,
|
|
) -> Result<Value, String> {
|
|
let id = self.next_id.fetch_add(1, Ordering::SeqCst);
|
|
|
|
let cmd = CdpCommand {
|
|
id,
|
|
method: method.to_string(),
|
|
params,
|
|
session_id: session_id.filter(|s| !s.is_empty()).map(|s| s.to_string()),
|
|
};
|
|
|
|
let json = serde_json::to_string(&cmd)
|
|
.map_err(|e| format!("Failed to serialize CDP command: {}", e))?;
|
|
|
|
let (tx, rx) = oneshot::channel();
|
|
|
|
{
|
|
let mut pending = self.pending.lock().await;
|
|
pending.insert(id, tx);
|
|
}
|
|
|
|
{
|
|
let mut ws_tx = self.ws_tx.lock().await;
|
|
ws_tx
|
|
.send(Message::Text(json))
|
|
.await
|
|
.map_err(|e| format!("Failed to send CDP command: {}", e))?;
|
|
}
|
|
|
|
let response = match tokio::time::timeout(std::time::Duration::from_secs(30), rx).await {
|
|
Ok(Ok(resp)) => resp,
|
|
Ok(Err(_)) => return Err("CDP response channel closed".to_string()),
|
|
Err(_) => {
|
|
self.pending.lock().await.remove(&id);
|
|
return Err(format!("CDP command timed out: {}", method));
|
|
}
|
|
};
|
|
|
|
if let Some(error) = response.error {
|
|
return Err(format!("CDP error ({}): {}", method, error));
|
|
}
|
|
|
|
Ok(response.result.unwrap_or(Value::Null))
|
|
}
|
|
|
|
pub fn subscribe(&self) -> broadcast::Receiver<CdpEvent> {
|
|
self.event_tx.subscribe()
|
|
}
|
|
|
|
/// Subscribe to all raw incoming CDP messages (responses + events).
|
|
/// Used by the inspect proxy to forward traffic to the DevTools frontend.
|
|
pub fn subscribe_raw(&self) -> broadcast::Receiver<RawCdpMessage> {
|
|
self.raw_tx.subscribe()
|
|
}
|
|
|
|
/// Create a lightweight handle for the inspect WebSocket proxy.
|
|
/// Contains only what's needed to forward messages bidirectionally.
|
|
pub fn inspect_handle(&self) -> InspectProxyHandle {
|
|
InspectProxyHandle {
|
|
ws_tx: self.ws_tx.clone(),
|
|
raw_tx: self.raw_tx.clone(),
|
|
}
|
|
}
|
|
|
|
pub async fn send_command_typed<P: serde::Serialize, R: serde::de::DeserializeOwned>(
|
|
&self,
|
|
method: &str,
|
|
params: &P,
|
|
session_id: Option<&str>,
|
|
) -> Result<R, String> {
|
|
let params_value = serde_json::to_value(params)
|
|
.map_err(|e| format!("Failed to serialize params: {}", e))?;
|
|
let result = self
|
|
.send_command(method, Some(params_value), session_id)
|
|
.await?;
|
|
serde_json::from_value(result)
|
|
.map_err(|e| format!("Failed to deserialize CDP response for {}: {}", method, e))
|
|
}
|
|
|
|
pub async fn send_command_no_params(
|
|
&self,
|
|
method: &str,
|
|
session_id: Option<&str>,
|
|
) -> Result<Value, String> {
|
|
self.send_command(method, None, session_id).await
|
|
}
|
|
|
|
/// Send raw JSON through the WebSocket without tracking a response.
|
|
/// Used by the inspect proxy to forward DevTools frontend messages.
|
|
pub async fn send_raw(&self, json: String) -> Result<(), String> {
|
|
let mut ws_tx = self.ws_tx.lock().await;
|
|
ws_tx
|
|
.send(Message::Text(json))
|
|
.await
|
|
.map_err(|e| format!("Failed to send raw CDP message: {}", e))
|
|
}
|
|
}
|
|
|
|
type WsTx = Arc<
|
|
Mutex<
|
|
futures_util::stream::SplitSink<
|
|
tokio_tungstenite::WebSocketStream<
|
|
tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>,
|
|
>,
|
|
Message,
|
|
>,
|
|
>,
|
|
>;
|
|
|
|
/// Lightweight handle for the inspect WebSocket proxy, holding only
|
|
/// the cloneable parts of CdpClient needed for bidirectional message forwarding.
|
|
pub struct InspectProxyHandle {
|
|
ws_tx: WsTx,
|
|
raw_tx: broadcast::Sender<RawCdpMessage>,
|
|
}
|
|
|
|
impl InspectProxyHandle {
|
|
pub async fn send_raw(&self, json: String) -> Result<(), String> {
|
|
let mut ws_tx = self.ws_tx.lock().await;
|
|
ws_tx
|
|
.send(Message::Text(json))
|
|
.await
|
|
.map_err(|e| format!("Failed to send raw CDP message: {}", e))
|
|
}
|
|
|
|
pub fn subscribe_raw(&self) -> broadcast::Receiver<RawCdpMessage> {
|
|
self.raw_tx.subscribe()
|
|
}
|
|
}
|
|
|
|
/// Enable TCP SO_KEEPALIVE on the underlying socket of a WebSocket connection.
|
|
/// This is best-effort: failures are silently ignored since the WebSocket-level
|
|
/// Ping keepalive provides the primary connection liveness mechanism.
|
|
fn enable_tcp_keepalive(stream: &tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>) {
|
|
let tcp_stream = match stream {
|
|
tokio_tungstenite::MaybeTlsStream::Plain(s) => s,
|
|
tokio_tungstenite::MaybeTlsStream::Rustls(s) => s.get_ref().0,
|
|
_ => return,
|
|
};
|
|
|
|
// SockRef borrows the fd without taking ownership.
|
|
let sock = socket2::SockRef::from(tcp_stream);
|
|
let keepalive = socket2::TcpKeepalive::new().with_time(std::time::Duration::from_secs(30));
|
|
|
|
// with_interval sets TCP_KEEPINTVL — the time between probes after the
|
|
// first keepalive probe goes unanswered. Available on most platforms
|
|
// (Linux, macOS, Windows, FreeBSD, etc.) but not OpenBSD or Haiku.
|
|
#[cfg(not(any(target_os = "openbsd", target_os = "haiku")))]
|
|
let keepalive = keepalive.with_interval(std::time::Duration::from_secs(10));
|
|
|
|
let _ = sock.set_tcp_keepalive(&keepalive);
|
|
}
|