From 8ce7ba891abe70b3d17393335277a367f7d8e565 Mon Sep 17 00:00:00 2001 From: KeyCode17 Date: Sat, 30 May 2026 09:34:33 +0700 Subject: [PATCH 1/3] refactor(cdp): split adapter setup (connect/browser_arc/Debug) into sibling file Co-Authored-By: Claude Sonnet 4.6 --- .../infrastructure/chromiumoxide_adapter.rs | 43 ++--------------- .../chromiumoxide_adapter_setup.rs | 48 +++++++++++++++++++ ras-cdp/src/infrastructure/mod.rs | 1 + 3 files changed, 52 insertions(+), 40 deletions(-) create mode 100644 ras-cdp/src/infrastructure/chromiumoxide_adapter_setup.rs diff --git a/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs b/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs index c8eae14..09154fe 100644 --- a/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs +++ b/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs @@ -3,12 +3,9 @@ use std::time::Duration; use async_trait::async_trait; use chromiumoxide::Browser; -use chromiumoxide::handler::HandlerConfig; -use futures::StreamExt; use ras_errors::AppError; use ras_types::{BackendNodeId, ContextId, TargetId}; use tokio::sync::Mutex; -use tracing::warn; use url::Url; use crate::domain::repository::{BrowserPort, ScreenshotFormat}; @@ -28,43 +25,9 @@ use crate::infrastructure::mouse_input::{ use crate::infrastructure::timeout::within; pub struct ChromiumoxideAdapter { - browser: Arc>, - cdp_url: Url, - request_timeout: Duration, -} - -impl std::fmt::Debug for ChromiumoxideAdapter { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("ChromiumoxideAdapter") - .field("cdp_url", &self.cdp_url) - .finish() - } -} - -impl ChromiumoxideAdapter { - #[must_use] - pub fn browser_arc(&self) -> Arc> { - Arc::clone(&self.browser) - } - - pub async fn connect(cdp_url: Url, request_timeout: Duration) -> Result { - let cfg = HandlerConfig::default(); - let (browser, mut handler) = Browser::connect_with_config(cdp_url.as_str(), cfg) - .await - .map_err(|e| AppError::BrowserDisconnected(format!("cdp connect: {e}")))?; - tokio::spawn(async move { - while let Some(ev) = handler.next().await { - if let Err(e) = ev { - warn!(error = %e, "cdp handler event error"); - } - } - }); - Ok(Self { - browser: Arc::new(Mutex::new(browser)), - cdp_url, - request_timeout, - }) - } + pub(crate) browser: Arc>, + pub(crate) cdp_url: Url, + pub(crate) request_timeout: Duration, } macro_rules! op { diff --git a/ras-cdp/src/infrastructure/chromiumoxide_adapter_setup.rs b/ras-cdp/src/infrastructure/chromiumoxide_adapter_setup.rs new file mode 100644 index 0000000..75618cd --- /dev/null +++ b/ras-cdp/src/infrastructure/chromiumoxide_adapter_setup.rs @@ -0,0 +1,48 @@ +//! Construction and [`std::fmt::Debug`] implementation for [`ChromiumoxideAdapter`]. + +use std::sync::Arc; +use std::time::Duration; + +use chromiumoxide::Browser; +use chromiumoxide::handler::HandlerConfig; +use futures::StreamExt; +use ras_errors::AppError; +use tokio::sync::Mutex; +use tracing::warn; +use url::Url; + +use super::chromiumoxide_adapter::ChromiumoxideAdapter; + +impl std::fmt::Debug for ChromiumoxideAdapter { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ChromiumoxideAdapter") + .field("cdp_url", &self.cdp_url) + .finish() + } +} + +impl ChromiumoxideAdapter { + #[must_use] + pub fn browser_arc(&self) -> Arc> { + Arc::clone(&self.browser) + } + + pub async fn connect(cdp_url: Url, request_timeout: Duration) -> Result { + let cfg = HandlerConfig::default(); + let (browser, mut handler) = Browser::connect_with_config(cdp_url.as_str(), cfg) + .await + .map_err(|e| AppError::BrowserDisconnected(format!("cdp connect: {e}")))?; + tokio::spawn(async move { + while let Some(ev) = handler.next().await { + if let Err(e) = ev { + warn!(error = %e, "cdp handler event error"); + } + } + }); + Ok(Self { + browser: Arc::new(Mutex::new(browser)), + cdp_url, + request_timeout, + }) + } +} diff --git a/ras-cdp/src/infrastructure/mod.rs b/ras-cdp/src/infrastructure/mod.rs index b302e3a..16abf07 100644 --- a/ras-cdp/src/infrastructure/mod.rs +++ b/ras-cdp/src/infrastructure/mod.rs @@ -1,5 +1,6 @@ pub mod cdp_ext; pub mod chromiumoxide_adapter; +pub mod chromiumoxide_adapter_setup; pub mod chromiumoxide_helpers; pub mod chromiumoxide_input; pub mod context_ops; From 3bc9e223db49fa6ba57ac2f529ddc11162f9f4a5 Mon Sep 17 00:00:00 2001 From: KeyCode17 Date: Sat, 30 May 2026 09:39:26 +0700 Subject: [PATCH 2/3] feat(cdp): CDP->BrowserEvent producer per tab (attach_events) with live e2e Co-Authored-By: Claude Sonnet 4.6 --- ras-cdp/Cargo.toml | 3 +- ras-cdp/src/domain/repository.rs | 14 ++++ .../infrastructure/chromiumoxide_adapter.rs | 11 +++ ras-cdp/src/infrastructure/event_pump.rs | 77 +++++++++++++++++++ ras-cdp/src/infrastructure/mod.rs | 1 + ras-cdp/tests/e2e_event_pump.rs | 57 ++++++++++++++ 6 files changed, 162 insertions(+), 1 deletion(-) create mode 100644 ras-cdp/src/infrastructure/event_pump.rs create mode 100644 ras-cdp/tests/e2e_event_pump.rs diff --git a/ras-cdp/Cargo.toml b/ras-cdp/Cargo.toml index 1770aee..02e3441 100644 --- a/ras-cdp/Cargo.toml +++ b/ras-cdp/Cargo.toml @@ -27,7 +27,8 @@ tracing = { workspace = true } url = { workspace = true } [dev-dependencies] -tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "time"] } url = { workspace = true } serde_json = { workspace = true } ras-types = { workspace = true } +ras-events = { workspace = true } diff --git a/ras-cdp/src/domain/repository.rs b/ras-cdp/src/domain/repository.rs index 00b4e58..37b042a 100644 --- a/ras-cdp/src/domain/repository.rs +++ b/ras-cdp/src/domain/repository.rs @@ -1,5 +1,8 @@ +use std::sync::Arc; + use async_trait::async_trait; use ras_errors::AppError; +use ras_events::EventBus; use ras_types::{BackendNodeId, ContextId, TargetId}; use serde::{Deserialize, Serialize}; use url::Url; @@ -92,4 +95,15 @@ pub trait BrowserPort: Send + Sync + 'static { "list_targets_in not supported by this BrowserPort".into(), )) } + + /// Forward this target's browser events to `bus` (per-session isolation). + /// + /// Default: no-op. Override in adapters that support live CDP streams. + async fn attach_events( + &self, + _target: &TargetId, + _bus: Arc, + ) -> Result<(), AppError> { + Ok(()) + } } diff --git a/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs b/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs index 09154fe..495f9a8 100644 --- a/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs +++ b/ras-cdp/src/infrastructure/chromiumoxide_adapter.rs @@ -8,6 +8,8 @@ use ras_types::{BackendNodeId, ContextId, TargetId}; use tokio::sync::Mutex; use url::Url; +use ras_events::EventBus; + use crate::domain::repository::{BrowserPort, ScreenshotFormat}; use crate::domain::viewport::Viewport; use crate::infrastructure::cdp_ext::{ @@ -160,4 +162,13 @@ impl BrowserPort for ChromiumoxideAdapter { async fn list_targets_in(&self, ctx: &ContextId) -> Result, AppError> { list_targets_in(&self.browser, ctx).await } + + async fn attach_events( + &self, + target: &TargetId, + bus: std::sync::Arc, + ) -> Result<(), AppError> { + let page = page_for(&self.browser, target).await?; + crate::infrastructure::event_pump::attach(&page, target.clone(), bus).await + } } diff --git a/ras-cdp/src/infrastructure/event_pump.rs b/ras-cdp/src/infrastructure/event_pump.rs new file mode 100644 index 0000000..38c0fc6 --- /dev/null +++ b/ras-cdp/src/infrastructure/event_pump.rs @@ -0,0 +1,77 @@ +//! CDP event forwarding: arms listeners on a [`chromiumoxide::Page`] and +//! publishes mapped [`BrowserEvent`] values to the supplied bus. + +use std::sync::Arc; + +use chromiumoxide::Page; +use chromiumoxide::cdp::browser_protocol::page::{ + DialogType, EventFrameNavigated, EventJavascriptDialogOpening, +}; +use futures::StreamExt; +use ras_errors::AppError; +use ras_events::{BrowserEvent, DialogKind, EventBus}; +use ras_types::TargetId; +use url::Url; + +/// Arm CDP listeners on `page` and forward mapped events to `bus`. +/// +/// Listeners are armed before returning; the spawned tasks run until the +/// page closes (streams end). +pub(crate) async fn attach( + page: &Page, + target: TargetId, + bus: Arc, +) -> Result<(), AppError> { + let nav_stream = page + .event_listener::() + .await + .map_err(|e| AppError::ActionFailed(format!("listen frameNavigated: {e}")))?; + + let nav_target = target.clone(); + let nav_bus = bus.clone(); + tokio::spawn(async move { + let mut s = nav_stream; + while let Some(ev) = s.next().await { + if ev.frame.parent_id.is_some() { + continue; + } + if let Ok(url) = Url::parse(&ev.frame.url) { + let _ = nav_bus + .publish(BrowserEvent::NavigationCompleted { + target: nav_target.clone(), + url, + }) + .await; + } + } + }); + + let dialog_stream = page + .event_listener::() + .await + .map_err(|e| AppError::ActionFailed(format!("listen dialog: {e}")))?; + + tokio::spawn(async move { + let mut s = dialog_stream; + while let Some(ev) = s.next().await { + let kind = map_dialog_kind(&ev.r#type); + let _ = bus + .publish(BrowserEvent::DialogOpened { + kind, + message: ev.message.clone(), + }) + .await; + } + }); + + Ok(()) +} + +fn map_dialog_kind(t: &DialogType) -> DialogKind { + match t { + DialogType::Alert => DialogKind::Alert, + DialogType::Confirm => DialogKind::Confirm, + DialogType::Prompt => DialogKind::Prompt, + DialogType::Beforeunload => DialogKind::BeforeUnload, + } +} diff --git a/ras-cdp/src/infrastructure/mod.rs b/ras-cdp/src/infrastructure/mod.rs index 16abf07..4c46136 100644 --- a/ras-cdp/src/infrastructure/mod.rs +++ b/ras-cdp/src/infrastructure/mod.rs @@ -4,5 +4,6 @@ pub mod chromiumoxide_adapter_setup; pub mod chromiumoxide_helpers; pub mod chromiumoxide_input; pub mod context_ops; +pub mod event_pump; pub mod mouse_input; pub mod timeout; diff --git a/ras-cdp/tests/e2e_event_pump.rs b/ras-cdp/tests/e2e_event_pump.rs new file mode 100644 index 0000000..92d0dad --- /dev/null +++ b/ras-cdp/tests/e2e_event_pump.rs @@ -0,0 +1,57 @@ +use std::sync::Arc; +use std::time::Duration; + +use ras_cdp::BrowserPort; +use ras_cdp::infrastructure::chromiumoxide_adapter::ChromiumoxideAdapter; +use ras_events::{BroadcastBus, BrowserEvent, EventBus}; +use url::Url; + +fn cdp_url() -> Option { + std::env::var("CDP_URL").ok()?.parse().ok() +} + +#[tokio::test] +#[ignore] +async fn navigation_event_reaches_the_bus() { + let url = cdp_url().expect("set CDP_URL"); + let a = ChromiumoxideAdapter::connect(url, Duration::from_secs(30)) + .await + .expect("connect"); + + let ctx = a.create_context().await.expect("ctx"); + let tab = a + .new_target_in(&ctx, &"about:blank".parse().expect("about")) + .await + .expect("tab"); + + let bus: Arc = Arc::new(BroadcastBus::default()); + let mut rx = bus.subscribe(); + a.attach_events(&tab, bus.clone()).await.expect("attach"); + + let dest: Url = std::env::var("TEST_URL") + .unwrap_or_else(|_| "http://127.0.0.1:8732/".to_string()) + .parse() + .expect("dest"); + a.navigate(&tab, &dest).await.expect("navigate"); + + let got = tokio::time::timeout(Duration::from_secs(8), async { + loop { + match rx.recv().await { + Ok(BrowserEvent::NavigationCompleted { target, .. }) if target == tab => { + break true; + } + Ok(_) => continue, + Err(_) => break false, + } + } + }) + .await + .unwrap_or(false); + + assert!( + got, + "expected a NavigationCompleted event for the tab on the bus" + ); + + a.close_context(&ctx).await.expect("close"); +} From 6a1854a485878e85ee4d746f1a916e6f8cd4cd95 Mon Sep 17 00:00:00 2001 From: KeyCode17 Date: Sat, 30 May 2026 09:40:06 +0700 Subject: [PATCH 3/3] =?UTF-8?q?chore:=20bump=20to=203.5.0=20(Phase=204=20?= =?UTF-8?q?=E2=80=94=20event=20producer)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Cargo.lock | 76 +++++++++++++++++++++++++++--------------------------- Cargo.toml | 2 +- 2 files changed, 39 insertions(+), 39 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d26a8f1..5cd44d9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1465,7 +1465,7 @@ dependencies = [ [[package]] name = "ras-agent" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "chrono", @@ -1497,7 +1497,7 @@ dependencies = [ [[package]] name = "ras-browser" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-cdp", @@ -1515,7 +1515,7 @@ dependencies = [ [[package]] name = "ras-cdp" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "chromiumoxide", @@ -1534,7 +1534,7 @@ dependencies = [ [[package]] name = "ras-cli" -version = "3.4.0" +version = "3.5.0" dependencies = [ "anyhow", "clap", @@ -1564,7 +1564,7 @@ dependencies = [ [[package]] name = "ras-cloud" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1581,7 +1581,7 @@ dependencies = [ [[package]] name = "ras-config" -version = "3.4.0" +version = "3.5.0" dependencies = [ "dotenvy", "once_cell", @@ -1595,7 +1595,7 @@ dependencies = [ [[package]] name = "ras-cosmium" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-cdp", @@ -1612,7 +1612,7 @@ dependencies = [ [[package]] name = "ras-daemon" -version = "3.4.0" +version = "3.5.0" dependencies = [ "anyhow", "dotenvy", @@ -1632,7 +1632,7 @@ dependencies = [ [[package]] name = "ras-dom" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "chromiumoxide", @@ -1652,7 +1652,7 @@ dependencies = [ [[package]] name = "ras-errors" -version = "3.4.0" +version = "3.5.0" dependencies = [ "serde", "thiserror", @@ -1660,7 +1660,7 @@ dependencies = [ [[package]] name = "ras-events" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-broadcast", "async-trait", @@ -1678,7 +1678,7 @@ dependencies = [ [[package]] name = "ras-filesystem" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1693,7 +1693,7 @@ dependencies = [ [[package]] name = "ras-judge" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "image", @@ -1708,7 +1708,7 @@ dependencies = [ [[package]] name = "ras-llm" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1722,7 +1722,7 @@ dependencies = [ [[package]] name = "ras-llm-anthropic" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "chrono", @@ -1745,7 +1745,7 @@ dependencies = [ [[package]] name = "ras-llm-bedrock" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1762,7 +1762,7 @@ dependencies = [ [[package]] name = "ras-llm-cerebras" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1780,7 +1780,7 @@ dependencies = [ [[package]] name = "ras-llm-cloud" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1797,7 +1797,7 @@ dependencies = [ [[package]] name = "ras-llm-deepseek" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1815,7 +1815,7 @@ dependencies = [ [[package]] name = "ras-llm-google" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1832,7 +1832,7 @@ dependencies = [ [[package]] name = "ras-llm-groq" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1850,7 +1850,7 @@ dependencies = [ [[package]] name = "ras-llm-langchain" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1867,7 +1867,7 @@ dependencies = [ [[package]] name = "ras-llm-mistral" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1885,7 +1885,7 @@ dependencies = [ [[package]] name = "ras-llm-oci" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1902,7 +1902,7 @@ dependencies = [ [[package]] name = "ras-llm-ollama" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1919,7 +1919,7 @@ dependencies = [ [[package]] name = "ras-llm-openai" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1936,7 +1936,7 @@ dependencies = [ [[package]] name = "ras-llm-openrouter" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1954,7 +1954,7 @@ dependencies = [ [[package]] name = "ras-llm-vercel" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1972,7 +1972,7 @@ dependencies = [ [[package]] name = "ras-mcp" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -1989,7 +1989,7 @@ dependencies = [ [[package]] name = "ras-recording" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "image", @@ -2004,7 +2004,7 @@ dependencies = [ [[package]] name = "ras-sandbox" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -2017,7 +2017,7 @@ dependencies = [ [[package]] name = "ras-skills" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -2034,7 +2034,7 @@ dependencies = [ [[package]] name = "ras-telemetry" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -2048,7 +2048,7 @@ dependencies = [ [[package]] name = "ras-tokens" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "chrono", @@ -2066,7 +2066,7 @@ dependencies = [ [[package]] name = "ras-tools" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "indexmap", @@ -2088,7 +2088,7 @@ dependencies = [ [[package]] name = "ras-types" -version = "3.4.0" +version = "3.5.0" dependencies = [ "chrono", "indexmap", @@ -2104,7 +2104,7 @@ dependencies = [ [[package]] name = "ras-validation" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-errors", @@ -2119,7 +2119,7 @@ dependencies = [ [[package]] name = "ras-watchdogs" -version = "3.4.0" +version = "3.5.0" dependencies = [ "async-trait", "ras-browser", diff --git a/Cargo.toml b/Cargo.toml index 30e9b30..2b31b31 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -51,7 +51,7 @@ exclude = ["examples", "tests", "docs", "scripts"] default-members = ["ras-cli", "ras-daemon"] [workspace.package] -version = "3.4.0" +version = "3.5.0" edition = "2024" rust-version = "1.95.0" license = "MIT"