Skip to main content

browser_automation_cli/native/cdp/
client.rs

1// SPDX-License-Identifier: MIT OR Apache-2.0
2//! CDP client over chromiumoxide (single connection — no dual WebSocket).
3#![allow(missing_docs)]
4//!
5//! Chrome one-shot: `Browser::launch` only.
6//! Lightpanda / attach path: `Browser::connect` only.
7//! FORBIDDEN: second `tokio-tungstenite` attach to the same browser.
8//!
9//! # Workload
10//!
11//! **I/O-bound** CDP WebSocket. Multi-page listener attach fans out with
12//! [`crate::concurrency::join_bounded`] after releasing `browser.lock`
13//! (rules: never hold a lock across unbounded sequential awaits when pages
14//! are independent).
15
16use std::borrow::Cow;
17use std::sync::Arc;
18
19use chromiumoxide::browser::Browser;
20use chromiumoxide::cdp::browser_protocol::fetch::EventRequestPaused;
21use chromiumoxide::cdp::browser_protocol::network::{
22    EventLoadingFailed, EventLoadingFinished, EventRequestWillBeSent,
23};
24use chromiumoxide::cdp::browser_protocol::page::{
25    EventDomContentEventFired, EventJavascriptDialogOpening, EventLoadEventFired,
26    EventScreencastFrame,
27};
28use chromiumoxide::cdp::browser_protocol::tracing::{EventDataCollected, EventTracingComplete};
29use chromiumoxide::cdp::js_protocol::heap_profiler::{
30    EventAddHeapSnapshotChunk, EventReportHeapSnapshotProgress,
31};
32use chromiumoxide::cdp::js_protocol::runtime::EventConsoleApiCalled;
33use chromiumoxide::error::CdpError;
34use chromiumoxide::page::Page;
35use chromiumoxide::types::{Command, Method, MethodId};
36use chromiumoxide::Handler;
37use futures::StreamExt;
38use serde::Serialize;
39use serde_json::Value;
40use tokio::sync::{broadcast, Mutex};
41use tokio::task::JoinHandle;
42
43use super::types::CdpEvent;
44
45/// Dynamic CDP command for `Browser::execute` / `Page::execute`.
46#[derive(Debug, Clone)]
47struct RawCdpCommand {
48    method: String,
49    params: Value,
50}
51
52impl Serialize for RawCdpCommand {
53    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
54    where
55        S: serde::Serializer,
56    {
57        match &self.params {
58            Value::Null => {
59                use serde::ser::SerializeMap;
60                let map = serializer.serialize_map(Some(0))?;
61                map.end()
62            }
63            other => other.serialize(serializer),
64        }
65    }
66}
67
68impl Method for RawCdpCommand {
69    fn identifier(&self) -> MethodId {
70        Cow::Owned(self.method.clone())
71    }
72}
73
74impl Command for RawCdpCommand {
75    type Response = Value;
76}
77
78/// CDP client wrapping a shared chromiumoxide [`Browser`].
79///
80/// # Interior mutability
81///
82/// `browser` uses **`tokio::sync::Mutex`** because guards are held across
83/// `.await` points (`Browser::execute`, `pages()`, `event_listener`). A
84/// `std::sync::Mutex` here would block the async runtime (rules: never hold
85/// std mutex across `.await`). The mutex is not exposed in the public agent
86/// JSON API — only as an internal handle for FINALIZE.
87///
88/// # Ownership
89///
90/// Holds the event-handler task and shared browser mutex — do not discard
91/// without FINALIZE (`#[must_use]`).
92#[must_use = "CdpClient owns the CDP connection and handler tasks"]
93pub struct CdpClient {
94    browser: Arc<Mutex<Browser>>,
95    event_tx: broadcast::Sender<CdpEvent>,
96    _handler: JoinHandle<()>,
97    _event_forwarders: Vec<JoinHandle<()>>,
98}
99
100impl CdpClient {
101    /// Build client from an already-launched or connected chromiumoxide Browser + handler.
102    pub async fn from_browser(browser: Browser, mut handler: Handler) -> Result<Self, String> {
103        let handler_task = tokio::spawn(async move {
104            while let Some(h) = handler.next().await {
105                if h.is_err() {
106                    break;
107                }
108            }
109        });
110
111        let (event_tx, _) = broadcast::channel(4096);
112
113        let browser = Arc::new(Mutex::new(browser));
114        let event_forwarders = spawn_event_forwarders(browser.clone(), event_tx.clone()).await?;
115
116        Ok(Self {
117            browser,
118            event_tx,
119            _handler: handler_task,
120            _event_forwarders: event_forwarders,
121        })
122    }
123
124    /// Attach via chromiumoxide `Browser::connect` (lightpanda only).
125    pub async fn connect(url: &str) -> Result<Self, String> {
126        Self::connect_with_headers(url, None).await
127    }
128
129    /// Headers are ignored on the oxide path (chromiumoxide connect has no custom WS headers API).
130    pub async fn connect_with_headers(
131        url: &str,
132        _headers: Option<Vec<(String, String)>>,
133    ) -> Result<Self, String> {
134        let (browser, handler) = Browser::connect(url)
135            .await
136            .map_err(|e| format!("CDP Browser::connect failed: {e}"))?;
137        Self::from_browser(browser, handler).await
138    }
139
140    /// Shared browser handle (for FINALIZE close/wait/kill).
141    pub fn browser(&self) -> Arc<Mutex<Browser>> {
142        self.browser.clone()
143    }
144
145    pub async fn send_command(
146        &self,
147        method: &str,
148        params: Option<Value>,
149        session_id: Option<&str>,
150    ) -> Result<Value, String> {
151        let cmd = RawCdpCommand {
152            method: method.to_string(),
153            params: params.unwrap_or(Value::Null),
154        };
155
156        let result = if let Some(sid) = session_id.filter(|s| !s.is_empty()) {
157            let page = self.page_for_session(sid).await?;
158            page.execute(cmd)
159                .await
160                .map_err(|e| format_cdp_err(method, &e))?
161        } else {
162            let browser = self.browser.lock().await;
163            browser
164                .execute(cmd)
165                .await
166                .map_err(|e| format_cdp_err(method, &e))?
167        };
168
169        Ok(result.result)
170    }
171
172    pub fn subscribe(&self) -> broadcast::Receiver<CdpEvent> {
173        self.event_tx.subscribe()
174    }
175
176    pub async fn send_command_typed<P: serde::Serialize, R: serde::de::DeserializeOwned>(
177        &self,
178        method: &str,
179        params: &P,
180        session_id: Option<&str>,
181    ) -> Result<R, String> {
182        let params_value = serde_json::to_value(params)
183            .map_err(|e| format!("Failed to serialize params: {}", e))?;
184        let result = self
185            .send_command(method, Some(params_value), session_id)
186            .await?;
187        serde_json::from_value(result)
188            .map_err(|e| format!("Failed to deserialize CDP response for {}: {}", method, e))
189    }
190
191    pub async fn send_command_no_params(
192        &self,
193        method: &str,
194        session_id: Option<&str>,
195    ) -> Result<Value, String> {
196        self.send_command(method, None, session_id).await
197    }
198
199    /// Best-effort command (still awaits oxide execute).
200    pub async fn send_command_no_wait(
201        &self,
202        method: &str,
203        params: Option<Value>,
204        session_id: Option<&str>,
205    ) -> Result<(), String> {
206        let _ = self.send_command(method, params, session_id).await;
207        Ok(())
208    }
209
210    async fn page_for_session(&self, session_id: &str) -> Result<Page, String> {
211        let browser = self.browser.lock().await;
212        let pages = browser
213            .pages()
214            .await
215            .map_err(|e| format!("Browser::pages failed: {e}"))?;
216        for page in pages {
217            if page.session_id().as_ref() == session_id {
218                return Ok(page);
219            }
220        }
221        // Fallback: first page if only one exists (session id mismatch after attach).
222        let pages = browser
223            .pages()
224            .await
225            .map_err(|e| format!("Browser::pages failed: {e}"))?;
226        let n = pages.len();
227        if n == 1 {
228            if let Some(page) = pages.into_iter().next() {
229                return Ok(page);
230            }
231        }
232        Err(format!(
233            "No chromiumoxide Page for session_id={session_id} (pages={n})"
234        ))
235    }
236}
237
238/// Format a chromiumoxide error for agent-facing messages (Display only → borrow).
239fn format_cdp_err(method: &str, e: &CdpError) -> String {
240    format!("CDP error ({method}): {e}")
241}
242
243/// Forward a typed CDP event stream onto the shared broadcast channel.
244///
245/// # Macro policy (`rules_rust_macros`)
246///
247/// Prefer **generics + monomorphization** over `macro_rules!` when the only
248/// variation is a type parameter and a method string. A previous local
249/// `macro_rules! fwd` expanded identical bodies for each CDP event type; that
250/// is exactly what a generic function does without hygiene / follow-set /
251/// double-evaluation concerns.
252///
253/// Browser-level listeners do not expose a session id; lifecycle accepts `None`.
254fn spawn_cdp_event_forwarder<T, St>(
255    mut stream: St,
256    method: &'static str,
257    event_tx: broadcast::Sender<CdpEvent>,
258) -> JoinHandle<()>
259where
260    T: serde::Serialize + Send + Sync + 'static,
261    St: futures::Stream<Item = Arc<T>> + Send + Unpin + 'static,
262{
263    tokio::spawn(async move {
264        while let Some(ev) = stream.next().await {
265            let params = serde_json::to_value(ev.as_ref()).unwrap_or(Value::Null);
266            let _ = event_tx.send(CdpEvent {
267                method: method.to_string(),
268                params,
269                session_id: None,
270            });
271        }
272    })
273}
274
275/// Subscribe to one browser-level CDP event type and spawn a forwarder task.
276async fn attach_browser_event_forwarder<T>(
277    browser: &Browser,
278    method: &'static str,
279    event_tx: broadcast::Sender<CdpEvent>,
280) -> Result<JoinHandle<()>, String>
281where
282    T: chromiumoxide::cdp::IntoEventKind + serde::Serialize + Unpin + 'static,
283{
284    let stream = browser
285        .event_listener::<T>()
286        .await
287        .map_err(|e| format!("event_listener {method}: {e}"))?;
288    Ok(spawn_cdp_event_forwarder(stream, method, event_tx))
289}
290
291/// Subscribe to one page-level CDP event type and spawn a forwarder task.
292async fn attach_page_event_forwarder<T>(
293    page: &Page,
294    method: &'static str,
295    event_tx: broadcast::Sender<CdpEvent>,
296) -> Result<(), String>
297where
298    T: chromiumoxide::cdp::IntoEventKind + serde::Serialize + Unpin + 'static,
299{
300    let stream = page
301        .event_listener::<T>()
302        .await
303        .map_err(|e| format!("page {method} listener: {e}"))?;
304    // Page-scoped tasks are fire-and-forget for the session lifetime (same as
305    // pre-refactor); browser-level handles are retained on `CdpClient`.
306    let _handle = spawn_cdp_event_forwarder(stream, method, event_tx);
307    Ok(())
308}
309
310async fn spawn_event_forwarders(
311    browser: Arc<Mutex<Browser>>,
312    event_tx: broadcast::Sender<CdpEvent>,
313) -> Result<Vec<JoinHandle<()>>, String> {
314    let mut handles = Vec::with_capacity(13);
315    let b = browser.lock().await;
316
317    handles.push(
318        attach_browser_event_forwarder::<EventLoadEventFired>(
319            &b,
320            "Page.loadEventFired",
321            event_tx.clone(),
322        )
323        .await?,
324    );
325    handles.push(
326        attach_browser_event_forwarder::<EventDomContentEventFired>(
327            &b,
328            "Page.domContentEventFired",
329            event_tx.clone(),
330        )
331        .await?,
332    );
333    handles.push(
334        attach_browser_event_forwarder::<EventRequestWillBeSent>(
335            &b,
336            "Network.requestWillBeSent",
337            event_tx.clone(),
338        )
339        .await?,
340    );
341    handles.push(
342        attach_browser_event_forwarder::<EventLoadingFinished>(
343            &b,
344            "Network.loadingFinished",
345            event_tx.clone(),
346        )
347        .await?,
348    );
349    handles.push(
350        attach_browser_event_forwarder::<EventLoadingFailed>(
351            &b,
352            "Network.loadingFailed",
353            event_tx.clone(),
354        )
355        .await?,
356    );
357    handles.push(
358        attach_browser_event_forwarder::<EventRequestPaused>(
359            &b,
360            "Fetch.requestPaused",
361            event_tx.clone(),
362        )
363        .await?,
364    );
365    handles.push(
366        attach_browser_event_forwarder::<EventJavascriptDialogOpening>(
367            &b,
368            "Page.javascriptDialogOpening",
369            event_tx.clone(),
370        )
371        .await?,
372    );
373    // Console API (context7/docs-rs): required for --capture-console.
374    handles.push(
375        attach_browser_event_forwarder::<EventConsoleApiCalled>(
376            &b,
377            "Runtime.consoleAPICalled",
378            event_tx.clone(),
379        )
380        .await?,
381    );
382    // Heap / tracing / screencast: required for heap take, perf stop, screencast frames.
383    handles.push(
384        attach_browser_event_forwarder::<EventAddHeapSnapshotChunk>(
385            &b,
386            "HeapProfiler.addHeapSnapshotChunk",
387            event_tx.clone(),
388        )
389        .await?,
390    );
391    handles.push(
392        attach_browser_event_forwarder::<EventReportHeapSnapshotProgress>(
393            &b,
394            "HeapProfiler.reportHeapSnapshotProgress",
395            event_tx.clone(),
396        )
397        .await?,
398    );
399    handles.push(
400        attach_browser_event_forwarder::<EventDataCollected>(
401            &b,
402            "Tracing.dataCollected",
403            event_tx.clone(),
404        )
405        .await?,
406    );
407    handles.push(
408        attach_browser_event_forwarder::<EventTracingComplete>(
409            &b,
410            "Tracing.tracingComplete",
411            event_tx.clone(),
412        )
413        .await?,
414    );
415    handles.push(
416        attach_browser_event_forwarder::<EventScreencastFrame>(
417            &b,
418            "Page.screencastFrame",
419            event_tx.clone(),
420        )
421        .await?,
422    );
423
424    drop(b);
425    Ok(handles)
426}
427
428impl CdpClient {
429    /// Attach page-level console listeners (context7 pattern: page.event_listener).
430    /// Complements browser-level forwarders when Runtime events are page-scoped.
431    pub async fn attach_page_console_forwarders(&self) -> Result<(), String> {
432        self.attach_page_event_forwarders_console().await
433    }
434
435    /// Page-level Network.requestWillBeSent (page-scoped CDP events).
436    pub async fn attach_page_network_forwarders(&self) -> Result<(), String> {
437        let pages = {
438            let browser = self.browser.lock().await;
439            browser
440                .pages()
441                .await
442                .map_err(|e| format!("Browser::pages for network listeners: {e}"))?
443        };
444        let event_tx = self.event_tx.clone();
445        let limit = crate::concurrency::effective_limit_capped(8);
446        let futs: Vec<_> = pages
447            .into_iter()
448            .map(|page| {
449                let event_tx = event_tx.clone();
450                async move {
451                    attach_page_event_forwarder::<EventRequestWillBeSent>(
452                        &page,
453                        "Network.requestWillBeSent",
454                        event_tx,
455                    )
456                    .await
457                }
458            })
459            .collect();
460        let results = crate::concurrency::join_bounded(futs, limit).await;
461        for r in results {
462            r?;
463        }
464        Ok(())
465    }
466
467    async fn attach_page_event_forwarders_console(&self) -> Result<(), String> {
468        let pages = {
469            let browser = self.browser.lock().await;
470            browser
471                .pages()
472                .await
473                .map_err(|e| format!("Browser::pages for console listeners: {e}"))?
474        };
475        let event_tx = self.event_tx.clone();
476        let limit = crate::concurrency::effective_limit_capped(8);
477        let futs: Vec<_> = pages
478            .into_iter()
479            .map(|page| {
480                let event_tx = event_tx.clone();
481                async move {
482                    attach_page_event_forwarder::<EventConsoleApiCalled>(
483                        &page,
484                        "Runtime.consoleAPICalled",
485                        event_tx,
486                    )
487                    .await
488                }
489            })
490            .collect();
491        let results = crate::concurrency::join_bounded(futs, limit).await;
492        for r in results {
493            r?;
494        }
495        Ok(())
496    }
497
498    /// Page-scoped CDP events (heap chunks, screencast frames, JS dialogs).
499    /// Browser-level listeners miss target-session events; attach after pages exist.
500    ///
501    /// Multi-page attach is I/O-bound → [`join_bounded`] after releasing the
502    /// browser lock (PAR-53).
503    pub async fn attach_page_session_forwarders(&self) -> Result<(), String> {
504        let pages = {
505            let browser = self.browser.lock().await;
506            browser
507                .pages()
508                .await
509                .map_err(|e| format!("Browser::pages for session listeners: {e}"))?
510        };
511        let event_tx = self.event_tx.clone();
512        let limit = crate::concurrency::effective_limit_capped(8);
513        let futs: Vec<_> = pages
514            .into_iter()
515            .map(|page| {
516                let event_tx = event_tx.clone();
517                async move {
518                    attach_page_event_forwarder::<EventAddHeapSnapshotChunk>(
519                        &page,
520                        "HeapProfiler.addHeapSnapshotChunk",
521                        event_tx.clone(),
522                    )
523                    .await?;
524                    attach_page_event_forwarder::<EventReportHeapSnapshotProgress>(
525                        &page,
526                        "HeapProfiler.reportHeapSnapshotProgress",
527                        event_tx.clone(),
528                    )
529                    .await?;
530                    attach_page_event_forwarder::<EventScreencastFrame>(
531                        &page,
532                        "Page.screencastFrame",
533                        event_tx.clone(),
534                    )
535                    .await?;
536                    // Page-scoped dialog open (required for eval auto-accept).
537                    attach_page_event_forwarder::<EventJavascriptDialogOpening>(
538                        &page,
539                        "Page.javascriptDialogOpening",
540                        event_tx,
541                    )
542                    .await?;
543                    Ok::<(), String>(())
544                }
545            })
546            .collect();
547        let results = crate::concurrency::join_bounded(futs, limit).await;
548        for r in results {
549            r?;
550        }
551        Ok(())
552    }
553}
554
555#[cfg(test)]
556mod tests {
557    use super::*;
558    use futures::stream;
559    use serde::Serialize;
560
561    #[derive(Debug, Serialize)]
562    struct DummyEvent {
563        n: u32,
564    }
565
566    #[tokio::test]
567    async fn cdp_event_forwarder_serializes_and_publishes() {
568        let (tx, mut rx) = broadcast::channel(4);
569        let stream = stream::iter(vec![Arc::new(DummyEvent { n: 7 })]);
570        let handle = spawn_cdp_event_forwarder(stream, "Test.event", tx);
571        let ev = rx.recv().await.expect("event delivered");
572        assert_eq!(ev.method, "Test.event");
573        assert_eq!(ev.params["n"], 7);
574        assert!(ev.session_id.is_none());
575        handle.await.expect("forwarder task");
576    }
577}