Skip to main content

oxicode_agent/tools/browse/
oxibrowser_backend.rs

1//! oxibrowser-core backend for the browser engine.
2//!
3//! Implements `BrowserEngine` and `BrowserTab` using the pure-Rust
4//! `oxibrowser-core` headless browser. Only compiled with
5//! `#[cfg(feature = "native-browser")]`.
6
7use super::config::BrowseConfig;
8use super::engine::{
9    BrowseProgress, BrowserEngine, BrowserError, BrowserTab as BrowserTabTrait, PageContent,
10    TabCallbackRegistry,
11};
12use async_trait::async_trait;
13use serde_json::Value;
14use std::sync::Arc;
15use tokio::sync::Mutex;
16use tokio::sync::broadcast::error::RecvError;
17use tokio::task::JoinHandle;
18
19/// Extract the `tab_id` from any `BrowserEvent` variant.
20fn extract_event_tab_id(event: &oxibrowser_core::BrowserEvent) -> uuid::Uuid {
21    match event {
22        oxibrowser_core::BrowserEvent::NavigationStarted { tab_id, .. }
23        | oxibrowser_core::BrowserEvent::WaitingForSelector { tab_id, .. }
24        | oxibrowser_core::BrowserEvent::DocumentReady { tab_id, .. }
25        | oxibrowser_core::BrowserEvent::ScreenshotCaptured { tab_id, .. }
26        | oxibrowser_core::BrowserEvent::PdfExported { tab_id, .. } => *tab_id,
27        // `BrowserEvent` is `#[non_exhaustive]`; fall through for forward-compat.
28        _ => uuid::Uuid::nil(),
29    }
30}
31
32/// Convert an `oxibrowser_core::BrowserEvent` into a `BrowseProgress`.
33///
34/// Returns `None` for unknown variants (forward-compatible with
35/// future `BrowserEvent` additions).
36fn browse_progress_from_event(event: &oxibrowser_core::BrowserEvent) -> Option<BrowseProgress> {
37    use oxibrowser_core::BrowserEvent::*;
38    match event {
39        NavigationStarted { url, .. } => {
40            Some(BrowseProgress::NavigationStarted { url: url.clone() })
41        }
42        WaitingForSelector {
43            selector,
44            timeout_ms,
45            ..
46        } => Some(BrowseProgress::WaitingForSelector {
47            selector: selector.clone(),
48            timeout_ms: *timeout_ms,
49        }),
50        DocumentReady {
51            final_url,
52            title,
53            status,
54            total_bytes,
55            total_duration,
56            ..
57        } => Some(BrowseProgress::DocumentReady {
58            url: final_url.clone(),
59            title: title.clone(),
60            status: *status,
61            bytes: *total_bytes,
62            duration_ms: total_duration.as_millis() as u64,
63        }),
64        ScreenshotCaptured {
65            bytes,
66            viewport_width,
67            duration,
68            ..
69        } => Some(BrowseProgress::ScreenshotCaptured {
70            bytes: *bytes,
71            width: *viewport_width,
72            duration_ms: duration.as_millis() as u64,
73        }),
74        PdfExported {
75            bytes,
76            viewport_width,
77            duration,
78            ..
79        } => Some(BrowseProgress::PdfExported {
80            bytes: *bytes,
81            width: *viewport_width,
82            duration_ms: duration.as_millis() as u64,
83        }),
84        _ => None,
85    }
86}
87
88// ── OxicodeBrowserEngine ──────────────────────────────────────────────────────────
89
90/// Browser engine powered by `oxibrowser-core`.
91///
92/// Spins a background task in its constructor that drains the browser's
93/// event stream and invokes whatever callback is currently installed in
94/// `progress_forwarder()`. The task exits gracefully when the browser
95/// is dropped (the broadcast sender is dropped → `RecvError::Closed`).
96///
97/// Single-tenant — see `BrowseTool::execution_mode`.
98pub struct OxicodeBrowserEngine {
99    browser: oxibrowser_core::Browser,
100    config: BrowseConfig,
101    /// Shared per-tab callback registry.
102    progress: Arc<TabCallbackRegistry>,
103    /// Background task that drains browser events into the forwarder.
104    /// Held so we can `await` it on `close()` for clean shutdown.
105    event_task: Mutex<Option<JoinHandle<()>>>,
106}
107
108impl OxicodeBrowserEngine {
109    /// Create a new engine with default config.
110    pub async fn new() -> Result<Self, BrowserError> {
111        Self::with_config(BrowseConfig::default()).await
112    }
113
114    /// Create a new engine with custom config.
115    ///
116    /// Propagates `BrowseConfig` fields (user_agent, obey_robots, js_timeout_ms)
117    /// to the underlying `oxibrowser-core` `BrowserConfig`.
118    pub async fn with_config(config: BrowseConfig) -> Result<Self, BrowserError> {
119        let mut browser_config = oxibrowser_core::BrowserConfig::headless();
120
121        // Propagate SDK-level settings to the browser engine
122        if let Some(ref ua) = config.user_agent {
123            browser_config.user_agent = ua.clone();
124        }
125        browser_config.obey_robots = config.obey_robots;
126        browser_config.js_timeout_ms = config.js_timeout_ms;
127
128        let browser = oxibrowser_core::Browser::new(browser_config)
129            .await
130            .map_err(|e| BrowserError::Backend(format!("Failed to create browser: {}", e)))?;
131
132        // Spawn the event-drain task. It lives for the lifetime of the engine:
133        // when the browser (and thus its event_tx) is dropped, the task's
134        // receiver returns `RecvError::Closed` and the task exits cleanly.
135        let progress = Arc::new(TabCallbackRegistry::new());
136        let mut events_rx = browser.subscribe_events();
137        let progress_clone = Arc::clone(&progress);
138        let event_task = tokio::spawn(async move {
139            loop {
140                match events_rx.recv().await {
141                    Ok(event) => {
142                        let tab_id = extract_event_tab_id(&event);
143                        // Enrich context FIRST so the String callback
144                        // below reads the enriched context_cell.
145                        if let Some(bp) = browse_progress_from_event(&event) {
146                            progress_clone.invoke_browse(&tab_id, bp);
147                        }
148                        progress_clone.invoke(&tab_id, event.short_label());
149                    }
150                    Err(RecvError::Lagged(skipped)) => {
151                        tracing::debug!(
152                            skipped = skipped,
153                            "oxibrowser event subscriber lagged; some events were dropped"
154                        );
155                    }
156                    Err(RecvError::Closed) => {
157                        break;
158                    }
159                }
160            }
161        });
162
163        Ok(Self {
164            browser,
165            config,
166            progress,
167            event_task: Mutex::new(Some(event_task)),
168        })
169    }
170}
171
172impl Default for OxicodeBrowserEngine {
173    fn default() -> Self {
174        // Default cannot be async, so use blocking runtime.
175        // Prefer `OxicodeBrowserEngine::new().await` in async contexts.
176        // SAFETY: `Runtime::new()` cannot fail with default config; and
177        // `block_on(Self::new())` panics rather than returning a half-built
178        // engine because `Default` has no Result channel. A failing browser
179        // init is an environment error (no Chrome/backend) that the caller
180        // should handle via `OxicodeBrowserEngine::new().await` instead.
181        #[allow(clippy::expect_used)]
182        let rt = tokio::runtime::Runtime::new().expect("failed to create tokio runtime");
183        #[allow(clippy::expect_used)]
184        rt.block_on(Self::new())
185            .expect("Failed to create default OxicodeBrowserEngine")
186    }
187}
188
189#[async_trait]
190impl BrowserEngine for OxicodeBrowserEngine {
191    async fn new_tab(&self) -> Result<Box<dyn BrowserTabTrait>, BrowserError> {
192        let tab = self
193            .browser
194            .new_tab()
195            .await
196            .map_err(|e| BrowserError::Backend(format!("Failed to create tab: {}", e)))?;
197        let tab_id = tab.tab_id();
198        Ok(Box::new(OxicodeTab {
199            inner: tab,
200            config: self.config.clone(),
201            tab_id,
202            registry: Arc::clone(&self.progress),
203        }))
204    }
205
206    async fn close(&self) -> Result<(), BrowserError> {
207        // Close the browser first. After this returns, the browser's internal
208        // event_tx is dropped — but the broadcast channel itself stays alive
209        // because the spawned event task holds its own sender clone. We need
210        // to cancel the task explicitly to make `close()` mean "fully shut
211        // down". The task will then exit with no further events forwarded.
212        self.browser
213            .close()
214            .await
215            .map_err(|e| BrowserError::Backend(format!("Browser close failed: {}", e)))?;
216
217        if let Some(handle) = self.event_task.lock().await.take() {
218            handle.abort();
219            let _ = handle.await; // ignore JoinError from abort
220        }
221        Ok(())
222    }
223
224    async fn is_alive(&self) -> bool {
225        self.browser.is_open()
226    }
227
228    fn callback_registry(&self) -> Arc<TabCallbackRegistry> {
229        Arc::clone(&self.progress)
230    }
231}
232
233// ── OxicodeTab ────────────────────────────────────────────────────────────────────
234
235/// A single browser tab backed by `oxibrowser-core`.
236#[allow(dead_code)] // config kept for future per-tab settings
237pub struct OxicodeTab {
238    inner: oxibrowser_core::Tab,
239    config: BrowseConfig,
240    /// Stable tab identity from `oxibrowser_core::Tab::tab_id()`.
241    tab_id: uuid::Uuid,
242    /// Shared per-tab callback registry.
243    registry: Arc<TabCallbackRegistry>,
244}
245
246impl OxicodeTab {
247    /// Register a progress callback for this tab.
248    pub fn set_progress_callback(&self, cb: crate::tools::ProgressCallback) {
249        self.registry.set(self.tab_id, cb);
250    }
251
252    /// Remove the progress callback for this tab.
253    pub fn clear_progress_callback(&self) {
254        self.registry.clear(&self.tab_id);
255    }
256
257    /// Register a structured browse progress callback for this tab.
258    pub fn set_browse_progress_callback_impl(&self, cb: super::engine::BrowseProgressCallback) {
259        self.registry.set_browse(self.tab_id, cb);
260    }
261
262    /// Return this tab's stable ID.
263    pub fn tab_id(&self) -> uuid::Uuid {
264        self.tab_id
265    }
266}
267
268#[async_trait]
269impl BrowserTabTrait for OxicodeTab {
270    async fn goto(&self, url: &str) -> Result<PageContent, BrowserError> {
271        let page = self
272            .inner
273            .goto(url)
274            .await
275            .map_err(|e| BrowserError::Navigation(e.to_string()))?;
276        Ok(browse_result_to_page_content(page))
277    }
278
279    async fn click(&self, selector: &str) -> Result<(), BrowserError> {
280        self.inner
281            .click(selector)
282            .await
283            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
284    }
285
286    async fn type_(&self, selector: &str, text: &str) -> Result<(), BrowserError> {
287        self.inner
288            .r#type(selector, text)
289            .await
290            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
291    }
292
293    async fn fill(&self, selector: &str, value: &str) -> Result<(), BrowserError> {
294        self.inner
295            .fill(selector, value)
296            .await
297            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
298    }
299
300    async fn press(&self, combo: &str) -> Result<(), BrowserError> {
301        self.inner
302            .press(combo)
303            .await
304            .map_err(|e| BrowserError::Evaluation(e.to_string()))
305    }
306
307    async fn wait_for(&self, selector: &str, timeout_ms: u64) -> Result<(), BrowserError> {
308        self.inner
309            .wait_for(selector, timeout_ms)
310            .await
311            .map_err(|e| BrowserError::Timeout(e.to_string()))
312    }
313    /// Native structured-wait override — maps our portable
314    /// [`super::engine::BrowseWaitCondition`] to `oxibrowser_core::tab::WaitCondition` so
315    /// `NetworkIdle` / `DomContentLoaded` / `Load` honour real in-flight
316    /// traffic semantics (Playwright/Puppeteer "networkidle" parity).
317    async fn wait_for_condition(
318        &self,
319        cond: &super::engine::BrowseWaitCondition,
320        timeout_ms: u64,
321    ) -> Result<(), BrowserError> {
322        use super::engine::BrowseWaitCondition as Bwc;
323        let mapped = match cond {
324            Bwc::Visible(s) => oxibrowser_core::tab::WaitCondition::Visible(s.clone()),
325            Bwc::NetworkIdle => oxibrowser_core::tab::WaitCondition::NetworkIdle,
326            Bwc::DomContentLoaded => oxibrowser_core::tab::WaitCondition::DomContentLoaded,
327            Bwc::Load => oxibrowser_core::tab::WaitCondition::Load,
328        };
329        self.inner
330            .wait_for_condition(mapped, timeout_ms)
331            .await
332            .map_err(|e| BrowserError::Timeout(e.to_string()))
333    }
334
335    async fn content(&self) -> Result<PageContent, BrowserError> {
336        let page = self
337            .inner
338            .content()
339            .await
340            .map_err(|e| BrowserError::Backend(e.to_string()))?;
341        Ok(browse_result_to_page_content(page))
342    }
343    /// omp `observe()` parity — runs the JS accessibility-surface synthesis
344    /// via `evaluate()` and parses the result into a [`super::engine::Observation`].
345    /// Returns the page's visible, interactive elements with stable
346    /// `data-oxicode-ref` selectors (no coordinates — boa only approximates
347    /// layout geometry).
348    async fn observe(&self) -> Result<super::engine::Observation, BrowserError> {
349        let page = self
350            .inner
351            .content()
352            .await
353            .map_err(|e| BrowserError::Backend(e.to_string()))?;
354        let value = self
355            .inner
356            .evaluate(super::helpers::JS_OBSERVE)
357            .await
358            .map_err(|e| BrowserError::Evaluation(e.to_string()))?;
359        Ok(super::engine::Observation {
360            url: page.url,
361            title: page.title,
362            elements: super::helpers::parse_observed_elements(value),
363        })
364    }
365
366    async fn query_all(&self, selector: &str) -> Result<Vec<String>, BrowserError> {
367        self.inner
368            .query_all(selector)
369            .await
370            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
371    }
372
373    async fn evaluate(&self, js: &str) -> Result<Value, BrowserError> {
374        self.inner
375            .evaluate(js)
376            .await
377            .map_err(|e| BrowserError::Evaluation(e.to_string()))
378    }
379
380    async fn screenshot(&self, width: u32) -> Result<Vec<u8>, BrowserError> {
381        self.inner
382            .screenshot(width)
383            .await
384            .map_err(|e| BrowserError::Screenshot(e.to_string()))
385    }
386
387    async fn print_to_pdf(&self, width: u32) -> Result<Vec<u8>, BrowserError> {
388        self.inner
389            .print_to_pdf(width)
390            .await
391            .map_err(|e| BrowserError::Pdf(e.to_string()))
392    }
393
394    async fn close(&self) -> Result<(), BrowserError> {
395        self.inner
396            .close()
397            .await
398            .map_err(|e| BrowserError::TabClosed(e.to_string()))
399    }
400
401    // ── Navigation — oxibrowser native history management ──────────────
402
403    async fn back(&self) -> Result<PageContent, BrowserError> {
404        let page = self
405            .inner
406            .back()
407            .await
408            .map_err(|e| BrowserError::Navigation(e.to_string()))?;
409        Ok(browse_result_to_page_content(page))
410    }
411
412    async fn forward(&self) -> Result<PageContent, BrowserError> {
413        let page = self
414            .inner
415            .forward()
416            .await
417            .map_err(|e| BrowserError::Navigation(e.to_string()))?;
418        Ok(browse_result_to_page_content(page))
419    }
420
421    async fn reload(&self) -> Result<PageContent, BrowserError> {
422        let page = self
423            .inner
424            .reload()
425            .await
426            .map_err(|e| BrowserError::Navigation(e.to_string()))?;
427        Ok(browse_result_to_page_content(page))
428    }
429
430    // ── Form interaction — oxibrowser native implementations ──────────
431
432    async fn select_option(&self, selector: &str, value: &str) -> Result<(), BrowserError> {
433        self.inner
434            .select_option(selector, value)
435            .await
436            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
437    }
438
439    async fn check(&self, selector: &str) -> Result<(), BrowserError> {
440        self.inner
441            .check(selector)
442            .await
443            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
444    }
445
446    async fn uncheck(&self, selector: &str) -> Result<(), BrowserError> {
447        self.inner
448            .uncheck(selector)
449            .await
450            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
451    }
452
453    // ── Advanced interaction — oxibrowser native ──────────────────────
454
455    async fn clear(&self, selector: &str) -> Result<(), BrowserError> {
456        self.inner
457            .clear_input(selector)
458            .await
459            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
460    }
461
462    async fn hover(&self, selector: &str) -> Result<(), BrowserError> {
463        self.inner
464            .hover(selector)
465            .await
466            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
467    }
468
469    async fn double_click(&self, selector: &str) -> Result<(), BrowserError> {
470        self.inner
471            .double_click(selector)
472            .await
473            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
474    }
475
476    async fn right_click(&self, selector: &str) -> Result<(), BrowserError> {
477        self.inner
478            .right_click(selector)
479            .await
480            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
481    }
482
483    async fn scroll(&self, delta_x: f64, delta_y: f64) -> Result<(), BrowserError> {
484        self.inner
485            .scroll(delta_x, delta_y)
486            .await
487            .map_err(|e| BrowserError::Evaluation(e.to_string()))
488    }
489
490    async fn scroll_into_view(&self, selector: &str) -> Result<(), BrowserError> {
491        self.inner
492            .scroll_into_view(selector, true)
493            .await
494            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
495    }
496
497    async fn drag(&self, from_selector: &str, to_selector: &str) -> Result<(), BrowserError> {
498        self.inner
499            .drag(from_selector, to_selector)
500            .await
501            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
502    }
503
504    async fn upload_file(&self, selector: &str, path: &str) -> Result<(), BrowserError> {
505        self.inner
506            .upload_file(selector, path)
507            .await
508            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
509    }
510
511    async fn get_value(&self, selector: &str) -> Result<String, BrowserError> {
512        self.inner
513            .get_value(selector)
514            .await
515            .map_err(|e| BrowserError::ElementNotFound(e.to_string()))
516    }
517
518    async fn evaluate_await(&self, js: &str) -> Result<Value, BrowserError> {
519        self.inner
520            .evaluate_await(js)
521            .await
522            .map_err(|e| BrowserError::Evaluation(e.to_string()))
523    }
524
525    fn is_closed(&self) -> bool {
526        self.inner.is_closed()
527    }
528
529    fn tab_id(&self) -> uuid::Uuid {
530        self.tab_id
531    }
532
533    fn as_any(&self) -> &dyn std::any::Any {
534        self
535    }
536
537    fn clear_progress_callback(&self) {
538        self.registry.clear(&self.tab_id);
539    }
540
541    fn set_browse_progress_callback(&self, cb: super::engine::BrowseProgressCallback) {
542        self.set_browse_progress_callback_impl(cb);
543    }
544}
545
546// ── Helpers ───────────────────────────────────────────────────────────────────
547
548/// Convert an `oxibrowser_core::BrowseResult` into our portable `PageContent`.
549fn browse_result_to_page_content(page: oxibrowser_core::BrowseResult) -> PageContent {
550    PageContent {
551        url: page.url.clone(),
552        title: page.title.clone(),
553        status: page.status,
554        markdown: page.markdown.clone(),
555        html: page.html.clone(),
556    }
557}
558
559#[cfg(test)]
560mod tests {
561    use super::*;
562    use std::sync::Mutex as StdMutex;
563    use std::sync::atomic::{AtomicUsize, Ordering};
564    use std::time::Duration;
565
566    /// End-to-end: the engine's background task should drain browser events
567    /// and invoke the callback installed in `progress_forwarder()`.
568    ///
569    /// We use a `data:` URL so the test does not require network access.
570    #[tokio::test]
571    async fn engine_forwards_browser_events_to_progress_callback() {
572        let engine = OxicodeBrowserEngine::new().await.unwrap();
573        let registry = engine.callback_registry();
574        let received: Arc<StdMutex<Vec<String>>> = Arc::new(StdMutex::new(Vec::new()));
575        let received_clone = Arc::clone(&received);
576
577        // Open a tab first to get its tab_id
578        let tab = engine.new_tab().await.unwrap();
579        let tab_id = tab
580            .as_any()
581            .downcast_ref::<OxicodeTab>()
582            .map(|t| t.tab_id())
583            .unwrap_or_default();
584
585        registry.set(
586            tab_id,
587            oxicode_ai::progress_callback(move |msg: String| {
588                received_clone.lock().unwrap().push(msg);
589            }),
590        );
591
592        // Navigate to a data: URL.
593        let _ = tab
594            .goto("data:text/html,<title>Hi</title><p>Hello</p>")
595            .await
596            .unwrap();
597
598        // Give the background task a moment to drain the broadcast channel.
599        tokio::time::sleep(Duration::from_millis(50)).await;
600
601        let got = received.lock().unwrap().clone();
602        assert!(
603            got.iter().any(|s| s.starts_with("Opening")),
604            "expected 'Opening …' event, got {got:?}"
605        );
606        assert!(
607            got.iter().any(|s| s.contains("Loaded")),
608            "expected 'Loaded …' event, got {got:?}"
609        );
610
611        let _ = tab.close().await;
612        let _ = engine.close().await;
613    }
614
615    /// Replacing the callback should drop the old one. Two callbacks should
616    /// not both fire for the same event.
617    #[tokio::test]
618    async fn engine_replaces_progress_callback_cleanly() {
619        let engine = OxicodeBrowserEngine::new().await.unwrap();
620        let registry = engine.callback_registry();
621        let count_a = Arc::new(AtomicUsize::new(0));
622        let count_b = Arc::new(AtomicUsize::new(0));
623
624        // Open tab to get its tab_id
625        let tab = engine.new_tab().await.unwrap();
626        let tab_id = tab
627            .as_any()
628            .downcast_ref::<OxicodeTab>()
629            .map(|t| t.tab_id())
630            .unwrap_or_default();
631
632        let ca = Arc::clone(&count_a);
633        registry.set(
634            tab_id,
635            oxicode_ai::progress_callback(move |_| {
636                ca.fetch_add(1, Ordering::SeqCst);
637            }),
638        );
639
640        let _ = tab.goto("data:text/html,<title>A</title>").await.unwrap();
641        tokio::time::sleep(Duration::from_millis(50)).await;
642        let a_after_first = count_a.load(Ordering::SeqCst);
643        assert!(a_after_first > 0, "callback A should have fired");
644
645        // Replace with B.
646        let cb_clone = Arc::clone(&count_b);
647        registry.set(
648            tab_id,
649            oxicode_ai::progress_callback(move |_| {
650                cb_clone.fetch_add(1, Ordering::SeqCst);
651            }),
652        );
653
654        let _ = tab.goto("data:text/html,<title>B</title>").await.unwrap();
655        tokio::time::sleep(Duration::from_millis(50)).await;
656
657        let a_final = count_a.load(Ordering::SeqCst);
658        let b_final = count_b.load(Ordering::SeqCst);
659        assert_eq!(
660            a_final, a_after_first,
661            "callback A should not fire after being replaced"
662        );
663        assert!(b_final > 0, "callback B should have fired");
664
665        let _ = tab.close().await;
666        let _ = engine.close().await;
667    }
668
669    /// End-to-end: `invoke_browse` should fire the structured
670    /// `BrowseProgressCallback` with `DocumentReady` carrying the page title
671    /// and HTTP status. This is the key T2 integration test for
672    /// `BrowseProgress` propagation.
673    #[tokio::test]
674    async fn engine_forwards_browse_progress_to_callback() {
675        use crate::tools::browse::BrowseProgress;
676
677        let engine = OxicodeBrowserEngine::new().await.unwrap();
678        let registry = engine.callback_registry();
679        let received: Arc<StdMutex<Vec<BrowseProgress>>> = Arc::new(StdMutex::new(Vec::new()));
680        let received_clone = Arc::clone(&received);
681
682        let tab = engine.new_tab().await.unwrap();
683        let tab_id = tab.tab_id();
684
685        registry.set_browse(
686            tab_id,
687            Arc::new(move |bp: BrowseProgress| {
688                received_clone.lock().unwrap().push(bp);
689            }),
690        );
691
692        // Navigate to a data: URL — must produce DocumentReady.
693        let _ = tab
694            .goto("data:text/html,<title>Hi</title><p>Hello</p>")
695            .await
696            .unwrap();
697
698        // Allow drain task to process events.
699        tokio::time::sleep(Duration::from_millis(100)).await;
700
701        let events = received.lock().unwrap().clone();
702        assert!(
703            events
704                .iter()
705                .any(|bp| matches!(bp, BrowseProgress::DocumentReady { status: 200, .. })),
706            "expected DocumentReady with status 200, got {events:?}"
707        );
708        let doc_ready = events.iter().find_map(|bp| match bp {
709            BrowseProgress::DocumentReady {
710                title,
711                bytes,
712                duration_ms,
713                ..
714            } => Some((title.clone(), *bytes, *duration_ms)),
715            _ => None,
716        });
717        let (title, bytes, duration_ms) = doc_ready.expect("DocumentReady present");
718        assert_eq!(title, "Hi");
719        assert!(
720            bytes > 0,
721            "bytes should be > 0 for non-empty page, got {bytes}"
722        );
723        assert!(
724            duration_ms < 30_000,
725            "duration_ms should be reasonable, got {duration_ms}"
726        );
727
728        let _ = tab.close().await;
729        let _ = engine.close().await;
730    }
731
732    /// Open two tabs, register per-tab browse callbacks, and verify each
733    /// callback receives only its own tab's `BrowseProgress` events.
734    #[tokio::test]
735    async fn engine_routes_browse_progress_by_tab_id() {
736        use crate::tools::browse::BrowseProgress;
737
738        let engine = OxicodeBrowserEngine::new().await.unwrap();
739        let registry = engine.callback_registry();
740
741        let received_a: Arc<StdMutex<Vec<BrowseProgress>>> = Arc::new(StdMutex::new(Vec::new()));
742        let received_b: Arc<StdMutex<Vec<BrowseProgress>>> = Arc::new(StdMutex::new(Vec::new()));
743        let ra = Arc::clone(&received_a);
744        let rb = Arc::clone(&received_b);
745
746        let tab_a = engine.new_tab().await.unwrap();
747        let tab_b = engine.new_tab().await.unwrap();
748        let tid_a = tab_a.tab_id();
749        let tid_b = tab_b.tab_id();
750
751        registry.set_browse(
752            tid_a,
753            Arc::new(move |bp: BrowseProgress| {
754                ra.lock().unwrap().push(bp);
755            }),
756        );
757        registry.set_browse(
758            tid_b,
759            Arc::new(move |bp: BrowseProgress| {
760                rb.lock().unwrap().push(bp);
761            }),
762        );
763
764        let _ = tab_a
765            .goto("data:text/html,<title>OnlyA</title>")
766            .await
767            .unwrap();
768        let _ = tab_b
769            .goto("data:text/html,<title>OnlyB</title>")
770            .await
771            .unwrap();
772
773        tokio::time::sleep(Duration::from_millis(100)).await;
774
775        let got_a = received_a.lock().unwrap().clone();
776        let got_b = received_b.lock().unwrap().clone();
777
778        let a_titles: Vec<&str> = got_a
779            .iter()
780            .filter_map(|bp| match bp {
781                BrowseProgress::DocumentReady { title, .. } => Some(title.as_str()),
782                _ => None,
783            })
784            .collect();
785        let b_titles: Vec<&str> = got_b
786            .iter()
787            .filter_map(|bp| match bp {
788                BrowseProgress::DocumentReady { title, .. } => Some(title.as_str()),
789                _ => None,
790            })
791            .collect();
792
793        assert!(
794            a_titles.contains(&"OnlyA"),
795            "A should have OnlyA, got {a_titles:?}"
796        );
797        assert!(!a_titles.contains(&"OnlyB"), "A should NOT have OnlyB");
798        assert!(
799            b_titles.contains(&"OnlyB"),
800            "B should have OnlyB, got {b_titles:?}"
801        );
802        assert!(!b_titles.contains(&"OnlyA"), "B should NOT have OnlyA");
803
804        let _ = tab_a.close().await;
805        let _ = tab_b.close().await;
806        let _ = engine.close().await;
807    }
808
809    /// Open two tabs in one engine, register two callbacks, navigate each.
810    /// Assert each callback fires only for its own tab's events.
811    #[tokio::test]
812    async fn engine_routes_events_by_tab_id_concurrent() {
813        let engine = OxicodeBrowserEngine::new().await.unwrap();
814        let registry = engine.callback_registry();
815
816        let received_a: Arc<StdMutex<Vec<String>>> = Arc::new(StdMutex::new(Vec::new()));
817        let received_b: Arc<StdMutex<Vec<String>>> = Arc::new(StdMutex::new(Vec::new()));
818        let received_a_clone = Arc::clone(&received_a);
819        let received_b_clone = Arc::clone(&received_b);
820
821        // Open two tabs
822        let tab_a = engine.new_tab().await.unwrap();
823        let tab_b = engine.new_tab().await.unwrap();
824        let tab_id_a = tab_a.tab_id();
825        let tab_id_b = tab_b.tab_id();
826        assert_ne!(tab_id_a, tab_id_b, "two tabs must have distinct IDs");
827
828        // Register per-tab callbacks
829        registry.set(
830            tab_id_a,
831            oxicode_ai::progress_callback(move |msg: String| {
832                received_a_clone.lock().unwrap().push(msg);
833            }),
834        );
835        registry.set(
836            tab_id_b,
837            oxicode_ai::progress_callback(move |msg: String| {
838                received_b_clone.lock().unwrap().push(msg);
839            }),
840        );
841
842        // Navigate tab A
843        let _ = tab_a
844            .goto("data:text/html,<title>TabA</title>")
845            .await
846            .unwrap();
847        tokio::time::sleep(Duration::from_millis(50)).await;
848
849        // Navigate tab B
850        let _ = tab_b
851            .goto("data:text/html,<title>TabB</title>")
852            .await
853            .unwrap();
854        tokio::time::sleep(Duration::from_millis(50)).await;
855
856        let got_a = received_a.lock().unwrap().clone();
857        let got_b = received_b.lock().unwrap().clone();
858
859        // Each tab should have received its own events
860        assert!(
861            got_a.iter().any(|s| s.contains("TabA")),
862            "tab A callback should have received TabA events, got {got_a:?}"
863        );
864        assert!(
865            got_b.iter().any(|s| s.contains("TabB")),
866            "tab B callback should have received TabB events, got {got_b:?}"
867        );
868        // Cross-contamination check: A's callback should NOT have B's events
869        assert!(
870            !got_a.iter().any(|s| s.contains("TabB")),
871            "tab A callback should NOT have received TabB events, got {got_a:?}"
872        );
873        assert!(
874            !got_b.iter().any(|s| s.contains("TabA")),
875            "tab B callback should NOT have received TabA events, got {got_b:?}"
876        );
877
878        let _ = tab_a.close().await;
879        let _ = tab_b.close().await;
880        let _ = engine.close().await;
881    }
882}