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