Skip to main content

glass/browser/
cdp.rs

1//! Low-level Chrome DevTools Protocol (CDP) client.
2//!
3//! Provides a multiplexed WebSocket connection to a Chrome debugging
4//! endpoint, including command dispatch, event subscription, and
5//! session-scoped target management.
6
7use futures_util::{SinkExt, StreamExt};
8use serde::de::DeserializeOwned;
9use serde::{Deserialize, Serialize};
10use serde_json::Value;
11use std::collections::HashMap;
12use std::error::Error;
13use std::fmt::{Display, Formatter};
14use std::path::Path;
15use std::sync::{
16    Arc,
17    atomic::{AtomicU64, Ordering},
18};
19use std::time::Duration;
20use tokio::sync::{broadcast, mpsc, oneshot};
21use tokio_tungstenite::tungstenite::Message;
22use tracing::{debug, info, warn};
23
24#[derive(Clone, Default)]
25struct CdpRoute {
26    target_id: Option<String>,
27    session_id: Option<String>,
28    context_id: Option<i64>,
29    frame_id: Option<String>,
30}
31
32tokio::task_local! {
33    static OPERATION_ROUTE: CdpRoute;
34    static CDP_WAIT_SCOPE: Arc<AtomicU64>;
35}
36
37/// A CDP method call request.
38#[derive(Debug, Serialize)]
39pub struct CdpRequest {
40    pub id: u64,
41    pub method: String,
42    #[serde(skip_serializing_if = "Option::is_none")]
43    pub params: Option<Value>,
44    #[serde(rename = "sessionId", skip_serializing_if = "Option::is_none")]
45    pub session_id: Option<String>,
46}
47
48/// A protocol or transport error returned by a CDP connection.
49#[derive(Debug, Clone, Deserialize, Serialize)]
50pub struct CdpError {
51    pub code: i64,
52    pub message: String,
53    #[serde(default, skip_serializing_if = "Option::is_none")]
54    pub data: Option<Value>,
55    #[serde(skip)]
56    kind: CdpErrorKind,
57}
58
59#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
60enum CdpErrorKind {
61    #[default]
62    Protocol,
63    Transport,
64    ResponseTimeout,
65}
66
67impl CdpError {
68    fn transport(message: impl Into<String>) -> Self {
69        Self {
70            code: -32_000,
71            message: message.into(),
72            data: None,
73            kind: CdpErrorKind::Transport,
74        }
75    }
76
77    fn response_timeout(timeout: Duration) -> Self {
78        Self {
79            code: -32_000,
80            message: format!(
81                "CDP response timeout after {} seconds",
82                timeout.as_secs_f64()
83            ),
84            data: None,
85            kind: CdpErrorKind::ResponseTimeout,
86        }
87    }
88
89    fn decode(error: impl std::fmt::Display) -> Self {
90        Self::transport(format!("failed to decode typed CDP response: {error}"))
91    }
92
93    /// Whether this error is the locally typed expiry of an unanswered CDP request.
94    pub fn is_response_timeout(&self) -> bool {
95        self.kind == CdpErrorKind::ResponseTimeout
96    }
97}
98
99impl Display for CdpError {
100    fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
101        write!(formatter, "CDP error {}: {}", self.code, self.message)
102    }
103}
104
105impl Error for CdpError {}
106
107/// Lightweight notification of a CDP event.
108///
109/// This default event stream intentionally omits payloads. Subscribe to
110/// [`CdpClient::subscribe_events_with_params`] only when a caller needs them.
111#[derive(Debug, Clone)]
112pub struct CdpEvent {
113    pub method: String,
114}
115
116/// CDP event with its full JSON payload for explicit diagnostic or media use.
117#[derive(Debug, Clone)]
118pub struct CdpEventWithParams {
119    pub method: String,
120    pub params: Value,
121    pub session_id: Option<String>,
122}
123
124/// Typed result envelope returned by `Runtime.evaluate` and
125/// `Runtime.callFunctionOn`.
126#[derive(Debug, Clone, Serialize, Deserialize)]
127pub struct RuntimeEvaluateResponse {
128    pub result: RuntimeRemoteObject,
129    #[serde(default, rename = "exceptionDetails")]
130    pub exception_details: Option<Value>,
131}
132
133/// The bounded subset of a Runtime remote object used by session operations.
134#[derive(Debug, Clone, Serialize, Deserialize)]
135pub struct RuntimeRemoteObject {
136    #[serde(rename = "type")]
137    #[serde(default)]
138    pub object_type: String,
139    #[serde(default)]
140    pub value: Option<Value>,
141    #[serde(default, rename = "objectId")]
142    pub object_id: Option<String>,
143    #[serde(default)]
144    pub description: Option<String>,
145}
146
147/// Typed response from `Page.navigate`.
148#[derive(Debug, Clone, Default, Serialize, Deserialize)]
149pub struct PageNavigateResponse {
150    #[serde(default, rename = "frameId")]
151    pub frame_id: Option<String>,
152    #[serde(default, rename = "loaderId")]
153    pub loader_id: Option<String>,
154    #[serde(default, rename = "errorText")]
155    pub error_text: Option<String>,
156}
157
158/// Typed response envelope from `Accessibility.getFullAXTree`.
159///
160/// Accessibility node properties are intentionally retained as JSON because
161/// Chrome adds protocol fields across versions; the response envelope and
162/// node collection are still validated at the CDP boundary.
163#[derive(Debug, Clone, Default, Serialize, Deserialize)]
164pub struct AccessibilityTreeResponse {
165    #[serde(default)]
166    pub nodes: Vec<Value>,
167}
168
169/// Typed response envelope from `DOM.getDocument` and
170/// `DOM.getFlattenedDocument`.
171#[derive(Debug, Clone, Default, Serialize, Deserialize)]
172pub struct DomDocumentResponse {
173    pub root: Value,
174    #[serde(default)]
175    pub nodes: Vec<Value>,
176}
177
178#[derive(Debug, Clone, Default, Deserialize)]
179struct DomFrameOwnerResponse {
180    #[serde(default, rename = "backendNodeId")]
181    backend_node_id: Option<i64>,
182}
183
184#[derive(Debug, Clone, Default, Deserialize)]
185struct DomBoxModelResponse {
186    model: DomBoxModel,
187}
188
189#[derive(Debug, Clone, Default, Deserialize)]
190struct DomBoxModel {
191    #[serde(default)]
192    content: Option<Vec<Option<f64>>>,
193}
194
195#[derive(Debug)]
196pub struct CdpScreencastFrame {
197    pub data: String,
198    pub metadata: Value,
199    pub session_id: Option<String>,
200}
201
202struct ScreencastSink {
203    session_id: Option<String>,
204    sender: mpsc::Sender<CdpScreencastFrame>,
205}
206
207#[derive(Debug, Deserialize)]
208struct IncomingMessage {
209    #[serde(default)]
210    id: Option<u64>,
211    #[serde(default)]
212    result: Option<Value>,
213    #[serde(default)]
214    error: Option<CdpError>,
215    #[serde(default)]
216    method: Option<String>,
217    #[serde(default, rename = "sessionId")]
218    session_id: Option<String>,
219}
220
221#[derive(Debug, Deserialize)]
222struct IncomingEventParams {
223    #[serde(default)]
224    params: Value,
225}
226
227enum Command {
228    Request {
229        id: u64,
230        json: String,
231        response: oneshot::Sender<Result<Value, CdpError>>,
232    },
233    Cancel {
234        id: u64,
235    },
236    FireAndForget {
237        json: String,
238    },
239    Close,
240}
241
242struct PendingRequestGuard {
243    tx: mpsc::UnboundedSender<Command>,
244    id: u64,
245    armed: bool,
246}
247
248impl PendingRequestGuard {
249    fn disarm(&mut self) {
250        self.armed = false;
251    }
252}
253
254impl Drop for PendingRequestGuard {
255    fn drop(&mut self) {
256        if self.armed {
257            let _ = self.tx.send(Command::Cancel { id: self.id });
258        }
259    }
260}
261
262/// A multiplexed connection to Chrome DevTools Protocol.
263///
264/// The connection task owns the WebSocket and the pending request map. This
265/// keeps response routing and event delivery in one place while allowing
266/// multiple commands to be in flight at once.
267#[derive(Clone)]
268pub struct CdpClient {
269    tx: mpsc::UnboundedSender<Command>,
270    next_id: Arc<AtomicU64>,
271    events: broadcast::Sender<CdpEvent>,
272    payload_events: broadcast::Sender<CdpEventWithParams>,
273    screencast_sink: Arc<std::sync::Mutex<Option<ScreencastSink>>>,
274    screencast_received: Arc<AtomicU64>,
275    screencast_dropped: Arc<AtomicU64>,
276    cdp_wait_nanos: Arc<AtomicU64>,
277    timeout: Duration,
278    active_route: Arc<std::sync::Mutex<CdpRoute>>,
279}
280
281struct RemoteObjectBatchGuard {
282    cdp: CdpClient,
283    object_ids: Vec<String>,
284}
285
286impl RemoteObjectBatchGuard {
287    fn new(cdp: CdpClient, array_id: String) -> Self {
288        Self {
289            cdp,
290            object_ids: vec![array_id],
291        }
292    }
293
294    async fn cleanup(&mut self) -> Result<(), CdpError> {
295        let object_ids = std::mem::take(&mut self.object_ids);
296        let mut first_error = None;
297        for object_id in object_ids {
298            if let Err(error) = self.cdp.release_object(&object_id).await
299                && first_error.is_none()
300            {
301                first_error = Some(error);
302            }
303        }
304        first_error.map_or(Ok(()), Err)
305    }
306}
307
308impl Drop for RemoteObjectBatchGuard {
309    fn drop(&mut self) {
310        let cdp = self.cdp.clone();
311        let object_ids = std::mem::take(&mut self.object_ids);
312        tokio::spawn(async move {
313            for object_id in object_ids {
314                let _ = cdp.release_object(&object_id).await;
315            }
316        });
317    }
318}
319
320impl CdpClient {
321    /// Number of CDP requests allocated by this connection so far.
322    ///
323    /// This is a monotonic diagnostic counter and does not retain request
324    /// payloads or alter command routing.
325    pub fn request_count(&self) -> u64 {
326        self.next_id.load(Ordering::Relaxed).saturating_sub(1)
327    }
328
329    /// Total time spent awaiting CDP responses on this connection.
330    ///
331    /// This is a monotonic diagnostic counter used by the benchmark to
332    /// separate protocol wait time from Glass-side orchestration work.
333    pub fn cdp_wait_nanos(&self) -> u64 {
334        self.cdp_wait_nanos.load(Ordering::Relaxed)
335    }
336
337    /// Measure CDP response wait time initiated by one async operation.
338    ///
339    /// The scope is task-local, so detached cleanup work does not get charged
340    /// to the operation that caused it.
341    pub async fn measure_cdp_wait<F>(&self, future: F) -> (F::Output, u64)
342    where
343        F: std::future::Future,
344    {
345        let wait_nanos = Arc::new(AtomicU64::new(0));
346        let output = CDP_WAIT_SCOPE.scope(Arc::clone(&wait_nanos), future).await;
347        (output, wait_nanos.load(Ordering::Relaxed))
348    }
349
350    /// Connect to a Chrome CDP page WebSocket using the default timeout.
351    pub async fn connect(ws_url: &str) -> Result<Self, Box<dyn Error>> {
352        Self::connect_with_timeout(ws_url, Duration::from_secs(30)).await
353    }
354
355    /// Connect to a Chrome CDP page WebSocket with a custom command timeout.
356    pub async fn connect_with_timeout(
357        ws_url: &str,
358        timeout: Duration,
359    ) -> Result<Self, Box<dyn Error>> {
360        info!(%ws_url, "connecting to CDP");
361        let (ws_stream, _) = tokio_tungstenite::connect_async(ws_url).await?;
362        let (mut write, mut read) = ws_stream.split();
363        let (tx, mut rx) = mpsc::unbounded_channel::<Command>();
364        let (event_tx, _) = broadcast::channel::<CdpEvent>(128);
365        let (payload_event_tx, _) = broadcast::channel::<CdpEventWithParams>(128);
366        let actor_events = event_tx.clone();
367        let actor_payload_events = payload_event_tx.clone();
368        let screencast_sink = Arc::new(std::sync::Mutex::new(None));
369        let actor_screencast_sink = Arc::clone(&screencast_sink);
370        let screencast_received = Arc::new(AtomicU64::new(0));
371        let actor_screencast_received = Arc::clone(&screencast_received);
372        let screencast_dropped = Arc::new(AtomicU64::new(0));
373        let actor_screencast_dropped = Arc::clone(&screencast_dropped);
374        let cdp_wait_nanos = Arc::new(AtomicU64::new(0));
375        let actor_tx = tx.clone();
376        let actor_next_id = Arc::new(AtomicU64::new(1));
377        let next_id = Arc::clone(&actor_next_id);
378
379        tokio::spawn(async move {
380            let mut pending: HashMap<u64, oneshot::Sender<Result<Value, CdpError>>> =
381                HashMap::new();
382            let mut close_reason = "CDP connection closed".to_string();
383
384            loop {
385                tokio::select! {
386                    command = rx.recv() => {
387                        match command {
388                            Some(Command::Request { id, json, response }) => {
389                                pending.insert(id, response);
390                                if let Err(error) = write.send(Message::Text(json.into())).await {
391                                    close_reason = format!("CDP write failed: {error}");
392                                    break;
393                                }
394                            }
395                            Some(Command::Cancel { id }) => {
396                                pending.remove(&id);
397                            }
398                            Some(Command::FireAndForget { json }) => {
399                                if let Err(error) = write.send(Message::Text(json.into())).await {
400                                    close_reason = format!("CDP write failed: {error}");
401                                    break;
402                                }
403                            }
404                            Some(Command::Close) | None => {
405                                let _ = write.send(Message::Close(None)).await;
406                                close_reason = "CDP connection closed by client".to_string();
407                                break;
408                            }
409                        }
410                    }
411                    message = read.next() => {
412                        match message {
413                            Some(Ok(Message::Text(text))) => {
414                                handle_incoming_message(
415                                    &mut pending,
416                                    &actor_events,
417                                    &actor_payload_events,
418                                    ScreencastDispatch {
419                                        sink: &actor_screencast_sink,
420                                        received: &actor_screencast_received,
421                                        dropped: &actor_screencast_dropped,
422                                        command_tx: &actor_tx,
423                                        next_id: &actor_next_id,
424                                    },
425                                    text.as_ref(),
426                                );
427                            }
428                            Some(Ok(Message::Binary(bytes))) => {
429                                match std::str::from_utf8(bytes.as_ref()) {
430                                    Ok(text) => handle_incoming_message(
431                                        &mut pending,
432                                        &actor_events,
433                                        &actor_payload_events,
434                                        ScreencastDispatch {
435                                            sink: &actor_screencast_sink,
436                                            received: &actor_screencast_received,
437                                            dropped: &actor_screencast_dropped,
438                                            command_tx: &actor_tx,
439                                            next_id: &actor_next_id,
440                                        },
441                                        text,
442                                    ),
443                                    Err(error) => warn!(%error, "ignoring non-UTF-8 CDP frame"),
444                                }
445                            }
446                            Some(Ok(Message::Ping(payload))) => {
447                                if let Err(error) = write.send(Message::Pong(payload)).await {
448                                    close_reason = format!("CDP pong failed: {error}");
449                                    break;
450                                }
451                            }
452                            Some(Ok(Message::Close(_))) => {
453                                close_reason = "CDP server closed the connection".to_string();
454                                break;
455                            }
456                            Some(Ok(_)) => {}
457                            Some(Err(error)) => {
458                                close_reason = format!("CDP read failed: {error}");
459                                break;
460                            }
461                            None => break,
462                        }
463                    }
464                }
465            }
466
467            let error = CdpError::transport(close_reason);
468            for (_, response) in pending.drain() {
469                let _ = response.send(Err(error.clone()));
470            }
471        });
472
473        Ok(Self {
474            tx,
475            next_id,
476            events: event_tx,
477            payload_events: payload_event_tx,
478            screencast_sink,
479            screencast_received,
480            screencast_dropped,
481            cdp_wait_nanos,
482            timeout,
483            active_route: Arc::new(std::sync::Mutex::new(CdpRoute::default())),
484        })
485    }
486
487    /// Subscribe to method-only CDP events such as page lifecycle events.
488    pub fn subscribe_events(&self) -> broadcast::Receiver<CdpEvent> {
489        self.events.subscribe()
490    }
491
492    /// Subscribe to CDP events with payloads for explicit media or diagnostic work.
493    ///
494    /// Payload JSON is parsed only while this stream has at least one receiver.
495    pub fn subscribe_events_with_params(&self) -> broadcast::Receiver<CdpEventWithParams> {
496        self.payload_events.subscribe()
497    }
498
499    pub fn open_screencast_channel(
500        &self,
501        session_id: Option<String>,
502    ) -> Result<mpsc::Receiver<CdpScreencastFrame>, CdpError> {
503        let mut sink = self
504            .screencast_sink
505            .lock()
506            .map_err(|_| CdpError::transport("screencast sink lock poisoned"))?;
507        if sink.is_some() {
508            return Err(CdpError::transport("a screencast scope is already active"));
509        }
510        let (sender, receiver) = mpsc::channel(2);
511        *sink = Some(ScreencastSink { session_id, sender });
512        self.screencast_received.store(0, Ordering::Relaxed);
513        self.screencast_dropped.store(0, Ordering::Relaxed);
514        Ok(receiver)
515    }
516
517    pub fn close_screencast_channel(&self) -> (u64, u64) {
518        if let Ok(mut sink) = self.screencast_sink.lock() {
519            *sink = None;
520        }
521        (
522            self.screencast_received.load(Ordering::Relaxed),
523            self.screencast_dropped.load(Ordering::Relaxed),
524        )
525    }
526
527    pub fn screencast_stats(&self) -> (u64, u64) {
528        (
529            self.screencast_received.load(Ordering::Relaxed),
530            self.screencast_dropped.load(Ordering::Relaxed),
531        )
532    }
533
534    pub fn current_session_id(&self) -> Option<String> {
535        self.current_route().session_id
536    }
537
538    pub async fn set_domain_enabled_for(
539        &self,
540        session_id: Option<String>,
541        domain: &str,
542        enabled: bool,
543    ) -> Result<(), CdpError> {
544        let method = format!("{domain}.{}", if enabled { "enable" } else { "disable" });
545        self.send_routed(&method, None, session_id, self.timeout)
546            .await?;
547        Ok(())
548    }
549
550    /// Send a CDP command and wait for the response.
551    pub async fn send(&self, method: &str, params: Option<Value>) -> Result<Value, CdpError> {
552        let session_id = self.current_route().session_id;
553        self.send_routed(method, params, session_id, self.timeout)
554            .await
555    }
556
557    /// Send a command and deserialize its response into a caller-owned type.
558    ///
559    /// The untyped [`send`](Self::send) method remains available at the CDP
560    /// boundary, while session code can use this helper for stable response
561    /// shapes without repeatedly indexing raw `serde_json::Value` objects.
562    pub async fn send_typed<R: DeserializeOwned>(
563        &self,
564        method: &str,
565        params: Option<Value>,
566    ) -> Result<R, CdpError> {
567        let value = self.send(method, params).await?;
568        serde_json::from_value(value).map_err(CdpError::decode)
569    }
570
571    /// Typed variant of [`send_browser`](Self::send_browser).
572    pub async fn send_browser_typed<R: DeserializeOwned>(
573        &self,
574        method: &str,
575        params: Option<Value>,
576    ) -> Result<R, CdpError> {
577        let value = self.send_browser(method, params).await?;
578        serde_json::from_value(value).map_err(CdpError::decode)
579    }
580
581    /// Typed variant of [`send_to_session`](Self::send_to_session).
582    pub async fn send_to_session_typed<R: DeserializeOwned>(
583        &self,
584        session_id: &str,
585        method: &str,
586        params: Option<Value>,
587    ) -> Result<R, CdpError> {
588        let value = self.send_to_session(session_id, method, params).await?;
589        serde_json::from_value(value).map_err(CdpError::decode)
590    }
591
592    pub async fn send_browser(
593        &self,
594        method: &str,
595        params: Option<Value>,
596    ) -> Result<Value, CdpError> {
597        self.send_routed(method, params, None, self.timeout).await
598    }
599
600    pub async fn send_to_session(
601        &self,
602        session_id: &str,
603        method: &str,
604        params: Option<Value>,
605    ) -> Result<Value, CdpError> {
606        self.send_routed(method, params, Some(session_id.to_string()), self.timeout)
607            .await
608    }
609
610    /// Send one operation-scoped command without changing the connection's
611    /// default response deadline or active route.
612    pub async fn send_with_timeout(
613        &self,
614        method: &str,
615        params: Option<Value>,
616        timeout: Duration,
617    ) -> Result<Value, CdpError> {
618        let session_id = self.current_route().session_id;
619        self.send_routed(method, params, session_id, timeout).await
620    }
621
622    async fn send_routed(
623        &self,
624        method: &str,
625        params: Option<Value>,
626        session_id: Option<String>,
627        timeout: Duration,
628    ) -> Result<Value, CdpError> {
629        let id = self.next_id.fetch_add(1, Ordering::Relaxed);
630        let request = CdpRequest {
631            id,
632            method: method.to_string(),
633            params,
634            session_id,
635        };
636        let json = serde_json::to_string(&request)
637            .map_err(|error| CdpError::transport(format!("failed to encode request: {error}")))?;
638        let (response_tx, response_rx) = oneshot::channel();
639
640        self.tx
641            .send(Command::Request {
642                id,
643                json,
644                response: response_tx,
645            })
646            .map_err(|_| CdpError::transport("CDP connection task is unavailable"))?;
647        let mut pending_guard = PendingRequestGuard {
648            tx: self.tx.clone(),
649            id,
650            armed: true,
651        };
652        let started = std::time::Instant::now();
653
654        let result = match tokio::time::timeout(timeout, response_rx).await {
655            Ok(Ok(result)) => {
656                pending_guard.disarm();
657                result
658            }
659            Ok(Err(_)) => {
660                pending_guard.disarm();
661                Err(CdpError::transport("CDP response channel closed"))
662            }
663            Err(_) => Err(CdpError::response_timeout(timeout)),
664        };
665        let elapsed_nanos = started.elapsed().as_nanos().min(u64::MAX as u128) as u64;
666        self.cdp_wait_nanos
667            .fetch_add(elapsed_nanos, Ordering::Relaxed);
668        if let Ok(scope) = CDP_WAIT_SCOPE.try_with(Arc::clone) {
669            scope.fetch_add(elapsed_nanos, Ordering::Relaxed);
670        }
671        result
672    }
673
674    pub fn set_active_session(&self, session_id: Option<String>) {
675        let mut route = self
676            .active_route
677            .lock()
678            .unwrap_or_else(|poison| poison.into_inner());
679        route.session_id = session_id;
680        route.context_id = None;
681        route.frame_id = None;
682    }
683
684    pub fn set_active_context(&self, context_id: Option<i64>) {
685        self.active_route
686            .lock()
687            .unwrap_or_else(|poison| poison.into_inner())
688            .context_id = context_id;
689    }
690
691    pub fn set_active_frame_context(&self, frame_id: Option<String>, context_id: Option<i64>) {
692        let mut route = self
693            .active_route
694            .lock()
695            .unwrap_or_else(|poison| poison.into_inner());
696        route.frame_id = frame_id;
697        route.context_id = context_id;
698    }
699
700    pub fn set_active_route(
701        &self,
702        session_id: Option<String>,
703        frame_id: Option<String>,
704        context_id: Option<i64>,
705    ) {
706        let mut route = self
707            .active_route
708            .lock()
709            .unwrap_or_else(|poison| poison.into_inner());
710        route.session_id = session_id;
711        route.frame_id = frame_id;
712        route.context_id = context_id;
713    }
714
715    pub fn set_active_target_route(
716        &self,
717        target_id: Option<String>,
718        session_id: Option<String>,
719        frame_id: Option<String>,
720        context_id: Option<i64>,
721    ) {
722        *self
723            .active_route
724            .lock()
725            .unwrap_or_else(|poison| poison.into_inner()) = CdpRoute {
726            target_id,
727            session_id,
728            frame_id,
729            context_id,
730        };
731    }
732
733    pub fn operation_identity(&self) -> Option<(String, String)> {
734        let route = self.current_route();
735        Some((route.target_id?, route.frame_id?))
736    }
737
738    pub fn set_active_frame(&self, frame_id: Option<String>) {
739        self.active_route
740            .lock()
741            .unwrap_or_else(|poison| poison.into_inner())
742            .frame_id = frame_id;
743    }
744
745    pub fn active_frame(&self) -> Option<String> {
746        self.current_route().frame_id
747    }
748
749    fn current_route(&self) -> CdpRoute {
750        OPERATION_ROUTE.try_with(Clone::clone).unwrap_or_else(|_| {
751            self.active_route
752                .lock()
753                .unwrap_or_else(|poison| poison.into_inner())
754                .clone()
755        })
756    }
757
758    pub async fn with_current_route<F: std::future::Future>(&self, future: F) -> F::Output {
759        if OPERATION_ROUTE.try_with(|_| ()).is_ok() {
760            future.await
761        } else {
762            OPERATION_ROUTE.scope(self.current_route(), future).await
763        }
764    }
765
766    pub async fn with_current_target_route<F: std::future::Future>(&self, future: F) -> F::Output {
767        let mut route = self.current_route();
768        route.context_id = None;
769        OPERATION_ROUTE.scope(route, future).await
770    }
771
772    /// Return the selected child frame viewport origin in target coordinates.
773    pub async fn frame_viewport_offset(&self, frame_id: &str) -> Result<(f64, f64), CdpError> {
774        let owner: DomFrameOwnerResponse = self
775            .send_typed(
776                "DOM.getFrameOwner",
777                Some(serde_json::json!({"frameId": frame_id})),
778            )
779            .await?;
780        let backend_node_id = owner
781            .backend_node_id
782            .ok_or_else(|| CdpError::transport("frame owner contained no backend node ID"))?;
783        let model: DomBoxModelResponse = self
784            .send_typed(
785                "DOM.getBoxModel",
786                Some(serde_json::json!({"backendNodeId": backend_node_id})),
787            )
788            .await?;
789        let content = model
790            .model
791            .content
792            .filter(|quad| quad.len() >= 2)
793            .ok_or_else(|| CdpError::transport("frame owner contained no content quad"))?;
794        let x = content[0].ok_or_else(|| CdpError::transport("frame owner x was not numeric"))?;
795        let y = content[1].ok_or_else(|| CdpError::transport("frame owner y was not numeric"))?;
796        Ok((x, y))
797    }
798
799    /// Navigate to a URL.
800    pub async fn navigate(&self, url: &str) -> Result<PageNavigateResponse, CdpError> {
801        self.send_typed("Page.navigate", Some(serde_json::json!({ "url": url })))
802            .await
803    }
804
805    /// Take a screenshot and return its base64-encoded image data.
806    pub async fn screenshot(&self, format: &str) -> Result<String, CdpError> {
807        self.screenshot_with_params(serde_json::json!({
808            "format": format,
809            "optimizeForSpeed": true
810        }))
811        .await
812    }
813
814    pub async fn screenshot_with_params(&self, params: Value) -> Result<String, CdpError> {
815        let mut result = self.send("Page.captureScreenshot", Some(params)).await?;
816        match result.get_mut("data").map(Value::take) {
817            Some(Value::String(data)) => Ok(data),
818            _ => Err(CdpError::transport(
819                "CDP screenshot response contained no data",
820            )),
821        }
822    }
823
824    pub async fn get_layout_metrics(&self) -> Result<Value, CdpError> {
825        self.send("Page.getLayoutMetrics", None).await
826    }
827
828    /// Get the accessibility tree.
829    pub async fn get_accessibility_tree(&self) -> Result<AccessibilityTreeResponse, CdpError> {
830        let frame_id = self.current_route().frame_id;
831        self.send_typed(
832            "Accessibility.getFullAXTree",
833            frame_id.map(|frame_id| serde_json::json!({"frameId": frame_id})),
834        )
835        .await
836    }
837
838    /// Get a flattened document tree including shadow DOM content.
839    /// Uses `pierce: true` to include open shadow roots up to the given depth.
840    pub async fn get_flattened_document(
841        &self,
842        depth: i64,
843    ) -> Result<DomDocumentResponse, CdpError> {
844        self.send_typed(
845            "DOM.getFlattenedDocument",
846            Some(serde_json::json!({ "depth": depth, "pierce": true })),
847        )
848        .await
849    }
850
851    /// Get the full document tree for an explicit deep-DOM inspection.
852    pub async fn get_deep_document(&self) -> Result<DomDocumentResponse, CdpError> {
853        self.send_typed("DOM.getDocument", Some(serde_json::json!({ "depth": -1 })))
854            .await
855    }
856
857    /// Get only the document root for operations that do not need descendants.
858    pub async fn get_document_root(&self) -> Result<DomDocumentResponse, CdpError> {
859        self.send_typed("DOM.getDocument", Some(serde_json::json!({ "depth": 0 })))
860            .await
861    }
862
863    /// Compatibility alias for the explicit deep-DOM request.
864    pub async fn get_document(&self) -> Result<DomDocumentResponse, CdpError> {
865        self.get_deep_document().await
866    }
867
868    /// Query a CSS selector and return the matching node.
869    pub async fn query_selector(&self, selector: &str) -> Result<Value, CdpError> {
870        let document = self.get_document_root().await?;
871        let root_id = document.root["nodeId"]
872            .as_i64()
873            .ok_or_else(|| CdpError::transport("DOM document response contained no root nodeId"))?;
874        self.send(
875            "DOM.querySelector",
876            Some(serde_json::json!({ "nodeId": root_id, "selector": selector })),
877        )
878        .await
879    }
880
881    /// Resolve a DOM/backend node to a reusable remote object.
882    pub async fn resolve_node_object(
883        &self,
884        node_id: Option<i64>,
885        backend_node_id: Option<i64>,
886    ) -> Result<String, CdpError> {
887        let mut params = serde_json::Map::new();
888        if let Some(node_id) = node_id {
889            params.insert("nodeId".to_string(), Value::from(node_id));
890        }
891        if let Some(backend_node_id) = backend_node_id {
892            params.insert("backendNodeId".to_string(), Value::from(backend_node_id));
893        }
894        let resolved = self
895            .send("DOM.resolveNode", Some(Value::Object(params)))
896            .await?;
897        resolved["object"]["objectId"]
898            .as_str()
899            .map(str::to_string)
900            .ok_or_else(|| CdpError::transport("DOM.resolveNode returned no objectId"))
901    }
902
903    /// Translate one already-resolved frontend node into its immutable backend
904    /// identity without repeating the caller's locator query.
905    pub async fn backend_node_id_for_node(&self, node_id: i64) -> Result<i64, CdpError> {
906        let described = self
907            .send(
908                "DOM.describeNode",
909                Some(serde_json::json!({"nodeId": node_id, "depth": 0})),
910            )
911            .await?;
912        described["node"]["backendNodeId"]
913            .as_i64()
914            .filter(|id| *id > 0)
915            .ok_or_else(|| CdpError::transport("DOM.describeNode returned no backendNodeId"))
916    }
917
918    /// Invoke a function on a previously resolved remote object.
919    pub async fn call_on_object(
920        &self,
921        object_id: &str,
922        function_declaration: &str,
923    ) -> Result<RuntimeEvaluateResponse, CdpError> {
924        self.send_typed(
925            "Runtime.callFunctionOn",
926            Some(serde_json::json!({
927                "objectId": object_id,
928                "functionDeclaration": function_declaration,
929                "returnByValue": true,
930                "awaitPromise": true
931            })),
932        )
933        .await
934    }
935
936    pub async fn release_object(&self, object_id: &str) -> Result<Value, CdpError> {
937        self.send(
938            "Runtime.releaseObject",
939            Some(serde_json::json!({ "objectId": object_id })),
940        )
941        .await
942    }
943
944    pub async fn release_object_for_session(
945        &self,
946        session_id: &str,
947        object_id: &str,
948    ) -> Result<Value, CdpError> {
949        self.send_to_session(
950            session_id,
951            "Runtime.releaseObject",
952            Some(serde_json::json!({"objectId": object_id})),
953        )
954        .await
955    }
956
957    /// Resolve a page-produced, bounded remote element array into DOM node IDs.
958    ///
959    /// The expression must return an Array with at most `limit` elements and a
960    /// numeric `glassCount` property containing the total logical match count.
961    pub async fn bounded_element_query(
962        &self,
963        expression: &str,
964        limit: usize,
965    ) -> Result<(usize, Vec<i64>), CdpError> {
966        // DOM.requestNode only returns frontend node IDs after the document has
967        // been requested in this CDP session.
968        self.get_document_root().await?;
969        let context_id = self.current_route().context_id;
970        let mut params = serde_json::json!({
971            "expression": expression,
972            "returnByValue": false,
973            "awaitPromise": true
974        });
975        if let Some(context_id) = context_id {
976            params["contextId"] = Value::from(context_id);
977        }
978        let evaluated = self.send("Runtime.evaluate", Some(params)).await?;
979        if evaluated.get("exceptionDetails").is_some() {
980            return Err(CdpError::transport("element query evaluation failed"));
981        }
982        let array_id = evaluated["result"]["objectId"]
983            .as_str()
984            .ok_or_else(|| CdpError::transport("element query returned no remote array"))?;
985        let mut remote_objects = RemoteObjectBatchGuard::new(self.clone(), array_id.to_string());
986        let properties = self
987            .send(
988                "Runtime.getProperties",
989                Some(serde_json::json!({
990                    "objectId": array_id,
991                    "ownProperties": true
992                })),
993            )
994            .await?;
995        let mut count = 0;
996        let mut objects = Vec::with_capacity(limit);
997        for property in properties["result"].as_array().into_iter().flatten() {
998            if property["name"].as_str() == Some("glassCount") {
999                count = property["value"]["value"].as_u64().unwrap_or(0) as usize;
1000                continue;
1001            }
1002            if property["name"]
1003                .as_str()
1004                .and_then(|name| name.parse::<usize>().ok())
1005                .is_some_and(|index| index < limit)
1006                && let Some(object_id) = property["value"]["objectId"].as_str()
1007            {
1008                remote_objects.object_ids.push(object_id.to_string());
1009                objects.push(object_id.to_string());
1010            }
1011        }
1012        let mut node_ids = Vec::with_capacity(objects.len());
1013        for object_id in objects {
1014            let requested = self
1015                .send(
1016                    "DOM.requestNode",
1017                    Some(serde_json::json!({ "objectId": object_id })),
1018                )
1019                .await;
1020            let requested = requested?;
1021            if let Some(node_id) = requested["nodeId"].as_i64().filter(|id| *id != 0) {
1022                node_ids.push(node_id);
1023            }
1024        }
1025        remote_objects.cleanup().await?;
1026        Ok((count, node_ids))
1027    }
1028
1029    /// Get the bounding box of a DOM node.
1030    pub async fn get_box_model(&self, node_id: i64) -> Result<Value, CdpError> {
1031        self.get_box_model_inner(Some(node_id), None).await
1032    }
1033
1034    /// Get the bounding box of a backend DOM node from an accessibility tree.
1035    pub async fn get_box_model_for_backend(&self, backend_node_id: i64) -> Result<Value, CdpError> {
1036        self.get_box_model_inner(None, Some(backend_node_id)).await
1037    }
1038
1039    /// Ask Chrome to scroll a DOM node into view only when it is necessary.
1040    ///
1041    /// The browser owns the visibility decision, avoiding a separate layout
1042    /// probe and avoiding a scroll when the target is already actionable.
1043    pub async fn scroll_into_view_if_needed(
1044        &self,
1045        node_id: Option<i64>,
1046        backend_node_id: Option<i64>,
1047    ) -> Result<Value, CdpError> {
1048        let mut params = serde_json::Map::new();
1049        if let Some(node_id) = node_id {
1050            params.insert("nodeId".to_string(), Value::from(node_id));
1051        }
1052        if let Some(backend_node_id) = backend_node_id {
1053            params.insert("backendNodeId".to_string(), Value::from(backend_node_id));
1054        }
1055        if params.is_empty() {
1056            return Err(CdpError::transport(
1057                "scrollIntoViewIfNeeded requires a nodeId or backendNodeId",
1058            ));
1059        }
1060        self.send("DOM.scrollIntoViewIfNeeded", Some(Value::Object(params)))
1061            .await
1062    }
1063
1064    async fn get_box_model_inner(
1065        &self,
1066        node_id: Option<i64>,
1067        backend_node_id: Option<i64>,
1068    ) -> Result<Value, CdpError> {
1069        let mut params = serde_json::Map::new();
1070        if let Some(node_id) = node_id {
1071            params.insert("nodeId".to_string(), Value::from(node_id));
1072        }
1073        if let Some(backend_node_id) = backend_node_id {
1074            params.insert("backendNodeId".to_string(), Value::from(backend_node_id));
1075        }
1076        self.send("DOM.getBoxModel", Some(Value::Object(params)))
1077            .await
1078    }
1079
1080    /// Evaluate JavaScript in the page.
1081    pub async fn evaluate(&self, expression: &str) -> Result<RuntimeEvaluateResponse, CdpError> {
1082        let context_id = self.current_route().context_id;
1083        self.evaluate_in_context(expression, context_id).await
1084    }
1085
1086    pub async fn evaluate_in_context(
1087        &self,
1088        expression: &str,
1089        context_id: Option<i64>,
1090    ) -> Result<RuntimeEvaluateResponse, CdpError> {
1091        let mut params = serde_json::json!({
1092            "expression": expression,
1093            "returnByValue": true,
1094            "awaitPromise": true
1095        });
1096        if let Some(context_id) = context_id {
1097            params["contextId"] = Value::from(context_id);
1098        }
1099        self.send_typed("Runtime.evaluate", Some(params)).await
1100    }
1101
1102    /// Insert text into the currently focused element.
1103    pub async fn insert_text(&self, text: &str) -> Result<Value, CdpError> {
1104        self.send(
1105            "Input.insertText",
1106            Some(serde_json::json!({ "text": text })),
1107        )
1108        .await
1109    }
1110
1111    /// Dispatch a mouse event via CDP Input.
1112    pub async fn dispatch_mouse_event(
1113        &self,
1114        event_type: &str,
1115        x: f64,
1116        y: f64,
1117        button: Option<&str>,
1118        click_count: Option<u32>,
1119    ) -> Result<Value, CdpError> {
1120        let mut params = serde_json::json!({
1121            "type": event_type,
1122            "x": x,
1123            "y": y,
1124        });
1125        if let Some(button) = button {
1126            params["button"] = Value::from(button);
1127        }
1128        if let Some(click_count) = click_count {
1129            params["clickCount"] = Value::from(click_count);
1130        }
1131        self.send("Input.dispatchMouseEvent", Some(params)).await
1132    }
1133
1134    /// Dispatch one mouse event with a caller-owned response window.
1135    ///
1136    /// This does not alter the connection default used by ordinary input.
1137    pub async fn dispatch_mouse_event_with_timeout(
1138        &self,
1139        event_type: &str,
1140        x: f64,
1141        y: f64,
1142        button: Option<&str>,
1143        click_count: Option<u32>,
1144        timeout: Duration,
1145    ) -> Result<Value, CdpError> {
1146        let mut params = serde_json::json!({"type": event_type, "x": x, "y": y});
1147        if let Some(button) = button {
1148            params["button"] = Value::from(button);
1149        }
1150        if let Some(click_count) = click_count {
1151            params["clickCount"] = Value::from(click_count);
1152        }
1153        self.send_with_timeout("Input.dispatchMouseEvent", Some(params), timeout)
1154            .await
1155    }
1156
1157    /// Dispatch a keyboard event via CDP Input.
1158    pub async fn dispatch_key_event(
1159        &self,
1160        event_type: &str,
1161        key: &str,
1162        code: &str,
1163    ) -> Result<Value, CdpError> {
1164        self.send(
1165            "Input.dispatchKeyEvent",
1166            Some(serde_json::json!({
1167                "type": event_type,
1168                "key": key,
1169                "code": code,
1170                "text": if event_type == "keyDown" { key } else { "" }
1171            })),
1172        )
1173        .await
1174    }
1175
1176    pub async fn dispatch_key_event_with_modifiers(
1177        &self,
1178        event_type: &str,
1179        key: &str,
1180        code: &str,
1181        text: &str,
1182        modifiers: i64,
1183    ) -> Result<Value, CdpError> {
1184        let virtual_key_code = match key {
1185            "Backspace" => 8,
1186            "Tab" => 9,
1187            "Enter" => 13,
1188            "Escape" => 27,
1189            "Delete" => 46,
1190            _ if key.len() == 1 => key.as_bytes()[0].to_ascii_uppercase() as i64,
1191            _ => 0,
1192        };
1193        self.send(
1194            "Input.dispatchKeyEvent",
1195            Some(serde_json::json!({
1196                "type": event_type,
1197                "key": key,
1198                "code": code,
1199                "text": text,
1200                "modifiers": modifiers,
1201                "windowsVirtualKeyCode": virtual_key_code,
1202                "nativeVirtualKeyCode": virtual_key_code
1203            })),
1204        )
1205        .await
1206    }
1207
1208    /// Invoke Blink's platform-independent editing command on the focused node.
1209    pub async fn dispatch_select_all(&self) -> Result<Value, CdpError> {
1210        self.send(
1211            "Input.dispatchKeyEvent",
1212            Some(serde_json::json!({
1213                "type": "rawKeyDown",
1214                "key": "a",
1215                "code": "KeyA",
1216                "commands": ["selectAll"]
1217            })),
1218        )
1219        .await
1220    }
1221
1222    pub async fn set_file_input_files(
1223        &self,
1224        node_id: Option<i64>,
1225        backend_node_id: Option<i64>,
1226        files: &[String],
1227    ) -> Result<Value, CdpError> {
1228        let mut params = serde_json::json!({"files": files});
1229        if let Some(node_id) = node_id {
1230            params["nodeId"] = Value::from(node_id);
1231        }
1232        if let Some(backend_node_id) = backend_node_id {
1233            params["backendNodeId"] = Value::from(backend_node_id);
1234        }
1235        self.send("DOM.setFileInputFiles", Some(params)).await
1236    }
1237
1238    /// Scroll the current page by a delta in CSS pixels.
1239    pub async fn scroll_by(&self, dx: f64, dy: f64) -> Result<Value, CdpError> {
1240        let expression = format!(
1241            "window.scrollBy({:.4}, {:.4}); window.scrollX + ',' + window.scrollY",
1242            dx, dy
1243        );
1244        serde_json::to_value(self.evaluate(&expression).await?).map_err(CdpError::decode)
1245    }
1246
1247    pub async fn get_cookies(&self) -> Result<Value, CdpError> {
1248        self.send("Network.getCookies", None).await
1249    }
1250
1251    pub async fn set_cookies(&self, cookies: Value) -> Result<Value, CdpError> {
1252        self.send(
1253            "Network.setCookies",
1254            Some(serde_json::json!({ "cookies": cookies })),
1255        )
1256        .await
1257    }
1258
1259    pub async fn clear_browser_cookies(&self) -> Result<(), CdpError> {
1260        self.send("Network.clearBrowserCookies", None).await?;
1261        Ok(())
1262    }
1263    pub async fn enable_page(&self) -> Result<(), CdpError> {
1264        self.send("Page.enable", None).await?;
1265        Ok(())
1266    }
1267
1268    /// Enable the only event domains required to wait for navigation and
1269    /// invalidate compact observations.
1270    pub async fn enable_observation_events(&self) -> Result<(), CdpError> {
1271        self.enable_page().await?;
1272        self.enable_dom().await?;
1273        Ok(())
1274    }
1275
1276    pub async fn enable_observation_events_for(&self, session_id: &str) -> Result<(), CdpError> {
1277        self.send_to_session(session_id, "Page.enable", None)
1278            .await?;
1279        self.send_to_session(session_id, "DOM.enable", None).await?;
1280        Ok(())
1281    }
1282
1283    pub async fn enable_runtime(&self) -> Result<(), CdpError> {
1284        self.send("Runtime.enable", None).await?;
1285        Ok(())
1286    }
1287
1288    pub async fn disable_runtime(&self) -> Result<(), CdpError> {
1289        self.send("Runtime.disable", None).await?;
1290        Ok(())
1291    }
1292
1293    pub async fn enable_log(&self) -> Result<(), CdpError> {
1294        self.send("Log.enable", None).await?;
1295        Ok(())
1296    }
1297
1298    pub async fn disable_log(&self) -> Result<(), CdpError> {
1299        self.send("Log.disable", None).await?;
1300        Ok(())
1301    }
1302
1303    pub async fn enable_network(&self) -> Result<(), CdpError> {
1304        self.send("Network.enable", None).await?;
1305        Ok(())
1306    }
1307
1308    pub async fn disable_network(&self) -> Result<(), CdpError> {
1309        self.send("Network.disable", None).await?;
1310        Ok(())
1311    }
1312
1313    pub async fn handle_javascript_dialog(&self, accept: bool) -> Result<Value, CdpError> {
1314        self.send(
1315            "Page.handleJavaScriptDialog",
1316            Some(serde_json::json!({"accept": accept})),
1317        )
1318        .await
1319    }
1320
1321    pub async fn set_download_behavior(
1322        &self,
1323        behavior: &str,
1324        download_path: Option<&Path>,
1325        events_enabled: bool,
1326    ) -> Result<Value, CdpError> {
1327        let mut params = serde_json::json!({
1328            "behavior": behavior,
1329            "eventsEnabled": events_enabled
1330        });
1331        if let Some(path) = download_path {
1332            params["downloadPath"] = Value::from(path.to_string_lossy().into_owned());
1333        }
1334        self.send_browser("Browser.setDownloadBehavior", Some(params))
1335            .await
1336    }
1337
1338    pub async fn enable_dom(&self) -> Result<(), CdpError> {
1339        self.send("DOM.enable", None).await?;
1340        Ok(())
1341    }
1342
1343    pub async fn enable_accessibility(&self) -> Result<(), CdpError> {
1344        self.send("Accessibility.enable", None).await?;
1345        Ok(())
1346    }
1347
1348    /// Ask an owned Chrome browser to close itself before process-level
1349    /// shutdown. This gives profile-backed state a chance to flush cleanly.
1350    pub async fn close_browser(&self) -> Result<(), CdpError> {
1351        self.send_browser("Browser.close", None).await?;
1352        Ok(())
1353    }
1354
1355    /// Override device metrics for viewport emulation.
1356    pub async fn set_device_metrics_override(
1357        &self,
1358        width: i64,
1359        height: i64,
1360        device_scale_factor: f64,
1361        mobile: bool,
1362    ) -> Result<Value, CdpError> {
1363        self.send(
1364            "Emulation.setDeviceMetricsOverride",
1365            Some(serde_json::json!({
1366                "width": width,
1367                "height": height,
1368                "deviceScaleFactor": device_scale_factor,
1369                "mobile": mobile,
1370            })),
1371        )
1372        .await
1373    }
1374
1375    /// Clear device metrics override.
1376    pub async fn clear_device_metrics_override(&self) -> Result<Value, CdpError> {
1377        self.send("Emulation.clearDeviceMetricsOverride", None)
1378            .await
1379    }
1380
1381    /// Ask the connection task to close its WebSocket.
1382    pub async fn close(&self) {
1383        let _ = self.tx.send(Command::Close);
1384    }
1385}
1386
1387struct ScreencastDispatch<'a> {
1388    sink: &'a std::sync::Mutex<Option<ScreencastSink>>,
1389    received: &'a AtomicU64,
1390    dropped: &'a AtomicU64,
1391    command_tx: &'a mpsc::UnboundedSender<Command>,
1392    next_id: &'a AtomicU64,
1393}
1394
1395fn handle_incoming_message(
1396    pending: &mut HashMap<u64, oneshot::Sender<Result<Value, CdpError>>>,
1397    events: &broadcast::Sender<CdpEvent>,
1398    payload_events: &broadcast::Sender<CdpEventWithParams>,
1399    screencast: ScreencastDispatch<'_>,
1400    text: &str,
1401) {
1402    let message: IncomingMessage = match serde_json::from_str(text) {
1403        Ok(message) => message,
1404        Err(error) => {
1405            warn!(%error, "ignoring malformed CDP message");
1406            return;
1407        }
1408    };
1409
1410    if let Some(id) = message.id {
1411        if let Some(response) = pending.remove(&id) {
1412            let result = match message.error {
1413                Some(error) => Err(error),
1414                None => Ok(message.result.unwrap_or(Value::Null)),
1415            };
1416            let _ = response.send(result);
1417        } else {
1418            debug!(id, "received CDP response with no pending request");
1419        }
1420        return;
1421    }
1422
1423    if let Some(method) = message.method {
1424        if method == "Page.screencastFrame" {
1425            match serde_json::from_str::<IncomingEventParams>(text) {
1426                Ok(mut payload) => {
1427                    let frame_session_id = payload.params["sessionId"].as_u64();
1428                    if let Some(frame_session_id) = frame_session_id {
1429                        let id = screencast.next_id.fetch_add(1, Ordering::Relaxed);
1430                        let mut ack = serde_json::json!({
1431                            "id": id,
1432                            "method": "Page.screencastFrameAck",
1433                            "params": {"sessionId": frame_session_id}
1434                        });
1435                        if let Some(session_id) = message.session_id.as_deref() {
1436                            ack["sessionId"] = Value::from(session_id);
1437                        }
1438                        let _ = screencast.command_tx.send(Command::FireAndForget {
1439                            json: ack.to_string(),
1440                        });
1441                    }
1442                    let data = payload.params["data"].take();
1443                    let metadata = payload.params["metadata"].take();
1444                    let frame = match data {
1445                        Value::String(data) => Some(CdpScreencastFrame {
1446                            data,
1447                            metadata,
1448                            session_id: message.session_id,
1449                        }),
1450                        _ => None,
1451                    };
1452                    if let Some(frame) = frame {
1453                        let sink = screencast.sink.lock().expect("screencast sink poisoned");
1454                        if let Some(sink) = sink.as_ref()
1455                            && sink.session_id == frame.session_id
1456                        {
1457                            screencast.received.fetch_add(1, Ordering::Relaxed);
1458                            if frame.data.len() > 32 * 1024 * 1024
1459                                || sink.sender.try_send(frame).is_err()
1460                            {
1461                                screencast.dropped.fetch_add(1, Ordering::Relaxed);
1462                            }
1463                        }
1464                    }
1465                }
1466                Err(error) => warn!(%error, "ignoring malformed screencast payload"),
1467            }
1468            let _ = events.send(CdpEvent { method });
1469            return;
1470        }
1471        let _ = events.send(CdpEvent {
1472            method: method.clone(),
1473        });
1474        if payload_events.receiver_count() > 0 {
1475            match serde_json::from_str::<IncomingEventParams>(text) {
1476                Ok(payload) => {
1477                    let _ = payload_events.send(CdpEventWithParams {
1478                        method,
1479                        params: payload.params,
1480                        session_id: message.session_id,
1481                    });
1482                }
1483                Err(error) => warn!(%error, "ignoring malformed CDP event payload"),
1484            }
1485        }
1486    }
1487}
1488
1489/// Exercise the production CDP envelope and event payload decoders.
1490#[cfg(feature = "fuzzing")]
1491#[doc(hidden)]
1492pub fn fuzz_incoming_message(text: &str) {
1493    if let Ok(message) = serde_json::from_str::<IncomingMessage>(text)
1494        && message.method.is_some()
1495    {
1496        let _ = serde_json::from_str::<IncomingEventParams>(text);
1497    }
1498}
1499
1500#[cfg(test)]
1501mod tests {
1502    use super::*;
1503    use futures_util::{SinkExt, StreamExt};
1504    use tokio::net::TcpListener;
1505    use tokio_tungstenite::accept_async;
1506
1507    #[tokio::test]
1508    async fn routes_concurrent_responses_by_id_and_delivers_events() {
1509        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1510        let address = listener.local_addr().unwrap();
1511
1512        let server = tokio::spawn(async move {
1513            let (stream, _) = listener.accept().await.unwrap();
1514            let websocket = accept_async(stream).await.unwrap();
1515            let (mut write, mut read) = websocket.split();
1516
1517            let first = read.next().await.unwrap().unwrap();
1518            let second = read.next().await.unwrap().unwrap();
1519            let first: Value = match first {
1520                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1521                _ => panic!("expected text frame"),
1522            };
1523            let second: Value = match second {
1524                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1525                _ => panic!("expected text frame"),
1526            };
1527
1528            write
1529                .send(Message::Text(
1530                    serde_json::json!({
1531                        "method": "Page.loadEventFired",
1532                        "params": {"frameId": "main"}
1533                    })
1534                    .to_string()
1535                    .into(),
1536                ))
1537                .await
1538                .unwrap();
1539
1540            for request in [second, first] {
1541                write
1542                    .send(Message::Text(
1543                        serde_json::json!({
1544                            "id": request["id"],
1545                            "result": {"method": request["method"]}
1546                        })
1547                        .to_string()
1548                        .into(),
1549                    ))
1550                    .await
1551                    .unwrap();
1552            }
1553        });
1554
1555        let client =
1556            CdpClient::connect_with_timeout(&format!("ws://{address}"), Duration::from_secs(2))
1557                .await
1558                .unwrap();
1559        let mut events = client.subscribe_events();
1560
1561        let (first, second) =
1562            tokio::join!(client.send("first", None), client.send("second", None),);
1563        assert_eq!(first.unwrap()["method"], "first");
1564        assert_eq!(second.unwrap()["method"], "second");
1565        assert_eq!(events.recv().await.unwrap().method, "Page.loadEventFired");
1566
1567        client.close().await;
1568        server.await.unwrap();
1569    }
1570
1571    #[tokio::test]
1572    async fn operation_route_is_immutable_across_selection_changes() {
1573        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1574        let address = listener.local_addr().unwrap();
1575        let server = tokio::spawn(async move {
1576            let (stream, _) = listener.accept().await.unwrap();
1577            let mut websocket = accept_async(stream).await.unwrap();
1578            for expected_session in ["old", "old", "new"] {
1579                let request = websocket.next().await.unwrap().unwrap();
1580                let request: Value = match request {
1581                    Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1582                    _ => panic!("expected text frame"),
1583                };
1584                assert_eq!(request["sessionId"], expected_session);
1585                websocket
1586                    .send(Message::Text(
1587                        serde_json::json!({"id": request["id"], "result": {}})
1588                            .to_string()
1589                            .into(),
1590                    ))
1591                    .await
1592                    .unwrap();
1593            }
1594        });
1595        let client = CdpClient::connect(&format!("ws://{address}"))
1596            .await
1597            .unwrap();
1598        client.set_active_target_route(
1599            Some("old-target".to_string()),
1600            Some("old".to_string()),
1601            Some("old-frame".to_string()),
1602            None,
1603        );
1604        client
1605            .with_current_route(async {
1606                client.send("first", None).await.unwrap();
1607                client.set_active_target_route(
1608                    Some("new-target".to_string()),
1609                    Some("new".to_string()),
1610                    Some("new-frame".to_string()),
1611                    None,
1612                );
1613                assert_eq!(
1614                    client.operation_identity(),
1615                    Some(("old-target".to_string(), "old-frame".to_string()))
1616                );
1617                client.send("second", None).await.unwrap();
1618            })
1619            .await;
1620        client.send("third", None).await.unwrap();
1621        client.close().await;
1622        server.await.unwrap();
1623    }
1624
1625    #[tokio::test]
1626    async fn delivers_event_payloads_only_to_opt_in_subscribers() {
1627        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1628        let address = listener.local_addr().unwrap();
1629        let server = tokio::spawn(async move {
1630            let (stream, _) = listener.accept().await.unwrap();
1631            let mut websocket = accept_async(stream).await.unwrap();
1632            let request = websocket.next().await.unwrap().unwrap();
1633            let request: Value = match request {
1634                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1635                _ => panic!("expected text frame"),
1636            };
1637
1638            for index in 0..5 {
1639                let route = if index == 0 { "foreign" } else { "wanted" };
1640                websocket
1641                    .send(Message::Text(
1642                        serde_json::json!({
1643                            "method": "Page.screencastFrame",
1644                            "sessionId": route,
1645                            "params": {"sessionId": 9, "data": format!("frame-{index}")}
1646                        })
1647                        .to_string()
1648                        .into(),
1649                    ))
1650                    .await
1651                    .unwrap();
1652                let ack = websocket.next().await.unwrap().unwrap();
1653                let ack: Value = match ack {
1654                    Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1655                    _ => panic!("expected text frame"),
1656                };
1657                assert_eq!(ack["method"], "Page.screencastFrameAck");
1658                assert_eq!(ack["params"]["sessionId"], 9);
1659                assert_eq!(ack["sessionId"], route);
1660            }
1661            websocket
1662                .send(Message::Text(
1663                    serde_json::json!({"id": request["id"], "result": {}})
1664                        .to_string()
1665                        .into(),
1666                ))
1667                .await
1668                .unwrap();
1669        });
1670
1671        let client = CdpClient::connect(&format!("ws://{address}"))
1672            .await
1673            .unwrap();
1674        let mut methods = client.subscribe_events();
1675        let mut payloads = client.subscribe_events_with_params();
1676        let mut frames = client
1677            .open_screencast_channel(Some("wanted".to_string()))
1678            .unwrap();
1679        client.send("test.ready", None).await.unwrap();
1680
1681        assert_eq!(methods.recv().await.unwrap().method, "Page.screencastFrame");
1682        assert_eq!(frames.recv().await.unwrap().data, "frame-1");
1683        assert_eq!(frames.recv().await.unwrap().data, "frame-2");
1684        assert!(payloads.try_recv().is_err());
1685        assert_eq!(client.screencast_stats(), (4, 2));
1686
1687        client.close().await;
1688        server.await.unwrap();
1689    }
1690
1691    #[tokio::test]
1692    async fn requests_fast_screenshot_encoding_and_moves_the_payload() {
1693        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1694        let address = listener.local_addr().unwrap();
1695        let server = tokio::spawn(async move {
1696            let (stream, _) = listener.accept().await.unwrap();
1697            let mut websocket = accept_async(stream).await.unwrap();
1698            let request = websocket.next().await.unwrap().unwrap();
1699            let request: Value = match request {
1700                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1701                _ => panic!("expected text frame"),
1702            };
1703            assert_eq!(request["method"], "Page.captureScreenshot");
1704            assert_eq!(request["params"]["format"], "png");
1705            assert_eq!(request["params"]["optimizeForSpeed"], true);
1706            websocket
1707                .send(Message::Text(
1708                    serde_json::json!({
1709                        "id": request["id"],
1710                        "result": {"data": "cG5n"}
1711                    })
1712                    .to_string()
1713                    .into(),
1714                ))
1715                .await
1716                .unwrap();
1717        });
1718
1719        let client = CdpClient::connect(&format!("ws://{address}"))
1720            .await
1721            .unwrap();
1722        assert_eq!(client.screenshot("png").await.unwrap(), "cG5n");
1723        client.close().await;
1724        server.await.unwrap();
1725    }
1726
1727    #[tokio::test]
1728    async fn selector_lookup_fetches_only_the_document_root() {
1729        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1730        let address = listener.local_addr().unwrap();
1731        let server = tokio::spawn(async move {
1732            let (stream, _) = listener.accept().await.unwrap();
1733            let mut websocket = accept_async(stream).await.unwrap();
1734
1735            let root_request = websocket.next().await.unwrap().unwrap();
1736            let root_request: Value = match root_request {
1737                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1738                _ => panic!("expected text frame"),
1739            };
1740            assert_eq!(root_request["method"], "DOM.getDocument");
1741            assert_eq!(root_request["params"], serde_json::json!({ "depth": 0 }));
1742            websocket
1743                .send(Message::Text(
1744                    serde_json::json!({
1745                        "id": root_request["id"],
1746                        "result": {"root": {"nodeId": 42}}
1747                    })
1748                    .to_string()
1749                    .into(),
1750                ))
1751                .await
1752                .unwrap();
1753
1754            let selector_request = websocket.next().await.unwrap().unwrap();
1755            let selector_request: Value = match selector_request {
1756                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1757                _ => panic!("expected text frame"),
1758            };
1759            assert_eq!(selector_request["method"], "DOM.querySelector");
1760            assert_eq!(
1761                selector_request["params"],
1762                serde_json::json!({ "nodeId": 42, "selector": "#save" })
1763            );
1764            websocket
1765                .send(Message::Text(
1766                    serde_json::json!({
1767                        "id": selector_request["id"],
1768                        "result": {"nodeId": 7}
1769                    })
1770                    .to_string()
1771                    .into(),
1772                ))
1773                .await
1774                .unwrap();
1775        });
1776
1777        let client = CdpClient::connect(&format!("ws://{address}"))
1778            .await
1779            .unwrap();
1780        assert_eq!(client.query_selector("#save").await.unwrap()["nodeId"], 7);
1781        client.close().await;
1782        server.await.unwrap();
1783    }
1784
1785    #[tokio::test]
1786    async fn backend_identity_describes_the_existing_frontend_node_without_a_second_query() {
1787        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1788        let address = listener.local_addr().unwrap();
1789        let server = tokio::spawn(async move {
1790            let (stream, _) = listener.accept().await.unwrap();
1791            let mut websocket = accept_async(stream).await.unwrap();
1792            let request = websocket.next().await.unwrap().unwrap();
1793            let request: Value = match request {
1794                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1795                _ => panic!("expected text frame"),
1796            };
1797            assert_eq!(request["method"], "DOM.describeNode");
1798            assert_eq!(
1799                request["params"],
1800                serde_json::json!({"nodeId": 17, "depth": 0})
1801            );
1802            websocket
1803                .send(Message::Text(
1804                    serde_json::json!({
1805                        "id": request["id"],
1806                        "result": {"node": {"backendNodeId": 91}}
1807                    })
1808                    .to_string()
1809                    .into(),
1810                ))
1811                .await
1812                .unwrap();
1813            assert!(
1814                tokio::time::timeout(Duration::from_millis(25), websocket.next())
1815                    .await
1816                    .is_err(),
1817                "backend translation must not repeat the selector query"
1818            );
1819        });
1820        let client = CdpClient::connect(&format!("ws://{address}"))
1821            .await
1822            .unwrap();
1823        assert_eq!(client.backend_node_id_for_node(17).await.unwrap(), 91);
1824        server.await.unwrap();
1825        client.close().await;
1826    }
1827
1828    #[tokio::test]
1829    async fn backend_identity_rejects_a_describe_response_without_backend_id() {
1830        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1831        let address = listener.local_addr().unwrap();
1832        let server = tokio::spawn(async move {
1833            let (stream, _) = listener.accept().await.unwrap();
1834            let mut websocket = accept_async(stream).await.unwrap();
1835            let request = websocket.next().await.unwrap().unwrap();
1836            let request: Value = match request {
1837                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1838                _ => panic!("expected text frame"),
1839            };
1840            websocket
1841                .send(Message::Text(
1842                    serde_json::json!({"id": request["id"], "result": {"node": {}}})
1843                        .to_string()
1844                        .into(),
1845                ))
1846                .await
1847                .unwrap();
1848        });
1849        let client = CdpClient::connect(&format!("ws://{address}"))
1850            .await
1851            .unwrap();
1852        let error = client.backend_node_id_for_node(17).await.unwrap_err();
1853        assert!(error.message.contains("no backendNodeId"));
1854        client.close().await;
1855        server.await.unwrap();
1856    }
1857
1858    #[tokio::test]
1859    async fn scroll_into_view_uses_the_backend_node_without_a_layout_probe() {
1860        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1861        let address = listener.local_addr().unwrap();
1862        let server = tokio::spawn(async move {
1863            let (stream, _) = listener.accept().await.unwrap();
1864            let mut websocket = accept_async(stream).await.unwrap();
1865            let request = websocket.next().await.unwrap().unwrap();
1866            let request: Value = match request {
1867                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1868                _ => panic!("expected text frame"),
1869            };
1870            assert_eq!(request["method"], "DOM.scrollIntoViewIfNeeded");
1871            assert_eq!(
1872                request["params"],
1873                serde_json::json!({ "backendNodeId": 42 })
1874            );
1875            websocket
1876                .send(Message::Text(
1877                    serde_json::json!({"id": request["id"], "result": {}})
1878                        .to_string()
1879                        .into(),
1880                ))
1881                .await
1882                .unwrap();
1883        });
1884
1885        let client = CdpClient::connect(&format!("ws://{address}"))
1886            .await
1887            .unwrap();
1888        client
1889            .scroll_into_view_if_needed(None, Some(42))
1890            .await
1891            .unwrap();
1892        client.close().await;
1893        server.await.unwrap();
1894    }
1895
1896    #[tokio::test]
1897    async fn observation_event_setup_enables_only_page_and_dom() {
1898        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1899        let address = listener.local_addr().unwrap();
1900        let server = tokio::spawn(async move {
1901            let (stream, _) = listener.accept().await.unwrap();
1902            let mut websocket = accept_async(stream).await.unwrap();
1903            let mut methods = Vec::new();
1904
1905            for _ in 0..2 {
1906                let request = websocket.next().await.unwrap().unwrap();
1907                let request: Value = match request {
1908                    Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1909                    _ => panic!("expected text frame"),
1910                };
1911                methods.push(request["method"].as_str().unwrap().to_string());
1912                websocket
1913                    .send(Message::Text(
1914                        serde_json::json!({"id": request["id"], "result": {}})
1915                            .to_string()
1916                            .into(),
1917                    ))
1918                    .await
1919                    .unwrap();
1920            }
1921
1922            assert_eq!(methods, ["Page.enable", "DOM.enable"]);
1923        });
1924
1925        let client = CdpClient::connect(&format!("ws://{address}"))
1926            .await
1927            .unwrap();
1928        client.enable_observation_events().await.unwrap();
1929        client.close().await;
1930        server.await.unwrap();
1931    }
1932
1933    #[tokio::test]
1934    async fn sends_browser_close_for_owned_session_shutdown() {
1935        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1936        let address = listener.local_addr().unwrap();
1937        let server = tokio::spawn(async move {
1938            let (stream, _) = listener.accept().await.unwrap();
1939            let mut websocket = accept_async(stream).await.unwrap();
1940            let request = websocket.next().await.unwrap().unwrap();
1941            let request: Value = match request {
1942                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1943                _ => panic!("expected text frame"),
1944            };
1945            assert_eq!(request["method"], "Browser.close");
1946            assert!(request.get("params").is_none());
1947            websocket
1948                .send(Message::Text(
1949                    serde_json::json!({"id": request["id"], "result": {}})
1950                        .to_string()
1951                        .into(),
1952                ))
1953                .await
1954                .unwrap();
1955        });
1956
1957        let client = CdpClient::connect(&format!("ws://{address}"))
1958            .await
1959            .unwrap();
1960        client.close_browser().await.unwrap();
1961        client.close().await;
1962        server.await.unwrap();
1963    }
1964
1965    #[tokio::test]
1966    async fn returns_a_timeout_when_the_server_does_not_respond() {
1967        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1968        let address = listener.local_addr().unwrap();
1969        let server = tokio::spawn(async move {
1970            let (stream, _) = listener.accept().await.unwrap();
1971            let (_write, mut read) = accept_async(stream).await.unwrap().split();
1972            let _ = read.next().await;
1973            tokio::time::sleep(Duration::from_millis(200)).await;
1974        });
1975
1976        let client =
1977            CdpClient::connect_with_timeout(&format!("ws://{address}"), Duration::from_millis(50))
1978                .await
1979                .unwrap();
1980        let error = client.send("never", None).await.unwrap_err();
1981        assert!(error.message.contains("timeout"));
1982        assert!(error.is_response_timeout());
1983        client.close().await;
1984        server.await.unwrap();
1985    }
1986
1987    #[tokio::test]
1988    async fn operation_timeout_is_short_and_does_not_change_the_connection_default() {
1989        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
1990        let address = listener.local_addr().unwrap();
1991        let server = tokio::spawn(async move {
1992            let (stream, _) = listener.accept().await.unwrap();
1993            let mut websocket = accept_async(stream).await.unwrap();
1994            let first = websocket.next().await.unwrap().unwrap();
1995            let first: Value = match first {
1996                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
1997                _ => panic!("expected text frame"),
1998            };
1999            assert_eq!(first["method"], "short");
2000            let second = websocket.next().await.unwrap().unwrap();
2001            let second: Value = match second {
2002                Message::Text(text) => serde_json::from_str(text.as_ref()).unwrap(),
2003                _ => panic!("expected text frame"),
2004            };
2005            assert_eq!(second["method"], "ordinary");
2006            websocket
2007                .send(Message::Text(
2008                    serde_json::json!({"id": second["id"], "result": {"ok": true}})
2009                        .to_string()
2010                        .into(),
2011                ))
2012                .await
2013                .unwrap();
2014        });
2015
2016        let client =
2017            CdpClient::connect_with_timeout(&format!("ws://{address}"), Duration::from_millis(500))
2018                .await
2019                .unwrap();
2020        let started = tokio::time::Instant::now();
2021        let error = client
2022            .send_with_timeout("short", None, Duration::from_millis(20))
2023            .await
2024            .unwrap_err();
2025        assert!(error.is_response_timeout());
2026        assert!(started.elapsed() < Duration::from_millis(200));
2027        assert_eq!(client.send("ordinary", None).await.unwrap()["ok"], true);
2028        client.close().await;
2029        server.await.unwrap();
2030    }
2031
2032    #[test]
2033    fn protocol_errors_are_not_typed_as_response_timeouts() {
2034        let error: CdpError = serde_json::from_value(serde_json::json!({
2035            "code": -32000,
2036            "message": "some protocol failure"
2037        }))
2038        .unwrap();
2039        assert!(!error.is_response_timeout());
2040    }
2041}