browser_automation_cli/native/cdp/
client.rs1#![allow(missing_docs)]
4use 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#[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#[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 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 pub async fn connect(url: &str) -> Result<Self, String> {
126 Self::connect_with_headers(url, None).await
127 }
128
129 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 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 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 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
238fn format_cdp_err(method: &str, e: &CdpError) -> String {
240 format!("CDP error ({method}): {e}")
241}
242
243fn 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
275async 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
291async 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 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 handles.push(
375 attach_browser_event_forwarder::<EventConsoleApiCalled>(
376 &b,
377 "Runtime.consoleAPICalled",
378 event_tx.clone(),
379 )
380 .await?,
381 );
382 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 pub async fn attach_page_console_forwarders(&self) -> Result<(), String> {
432 self.attach_page_event_forwarders_console().await
433 }
434
435 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 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 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}