Skip to main content

rightkit_browser/
session.rs

1use crate::config::{
2    find_chrome, reject_real_browser_profile, validate_extra_args, validate_profile_name,
3    LaunchOptions, ProfileSpec,
4};
5use crate::error::{cdp, BrowserError, Result};
6use crate::page::BrowserPage;
7use crate::policy::{Guard, LifecycleKind};
8use crate::proc::ChromeProc;
9use crate::proxy::{PinProxy, PinnedConnection};
10use chromiumoxide::browser::Browser;
11use chromiumoxide::cdp::browser_protocol::browser::{
12    PermissionDescriptor, PermissionSetting, ResetPermissionsParams, SetDownloadBehaviorBehavior,
13    SetDownloadBehaviorParams, SetPermissionParams,
14};
15use futures::StreamExt;
16use std::collections::HashSet;
17use std::path::{Path, PathBuf};
18use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
19use std::sync::{Arc, Mutex as StdMutex};
20use std::time::{Duration, Instant};
21use tokio::sync::Mutex;
22use tokio::task::JoinHandle;
23
24static SESSION_SEQ: AtomicU64 = AtomicU64::new(0);
25
26enum ProfileGuard {
27    Temp(#[allow(dead_code)] tempfile::TempDir),
28    Named { _lock: std::fs::File },
29}
30
31/// One Chrome process with any number of pages.
32pub struct BrowserSession {
33    id: String,
34    browser: Mutex<Browser>,
35    pages: StdMutex<Vec<Arc<BrowserPage>>>,
36    current: StdMutex<Option<String>>,
37    page_seq: AtomicU64,
38    handler: StdMutex<Option<JoinHandle<()>>>,
39    watcher: StdMutex<Option<JoinHandle<()>>>,
40    _profile: ProfileGuard,
41    profile_dir: PathBuf,
42    download_dir: PathBuf,
43    upload_root: Option<PathBuf>,
44    headless: bool,
45    guard: Arc<Guard>,
46    pid: Option<u32>,
47    proc: Arc<StdMutex<Option<ChromeProc>>>,
48    proxy: Option<PinProxy>,
49    /// Set before an intentional stop so the CDP pump does not report a crash.
50    stopping: Arc<AtomicBool>,
51    stopped: AtomicBool,
52}
53
54impl BrowserSession {
55    pub async fn launch(
56        #[cfg_attr(not(windows), allow(unused_mut))] mut opts: LaunchOptions,
57    ) -> Result<Self> {
58        validate_extra_args(&opts.extra_args)?;
59        let exe = opts
60            .chrome_path
61            .clone()
62            .or_else(find_chrome)
63            .ok_or(BrowserError::ChromeNotFound)?;
64        let (profile_guard, profile_dir, default_dl) = match &opts.profile {
65            ProfileSpec::Temporary => {
66                let t = tempfile::tempdir()?;
67                let p = t.path().to_path_buf();
68                let dl = p.join("downloads");
69                (ProfileGuard::Temp(t), p, dl)
70            }
71            ProfileSpec::Named { root, name } => {
72                validate_profile_name(name)?;
73                reject_real_browser_profile(root)?;
74                let dir = root.join(name);
75                std::fs::create_dir_all(&dir)?;
76                let lock = std::fs::OpenOptions::new()
77                    .create(true)
78                    .write(true)
79                    .truncate(false)
80                    .open(dir.join(".rightkit-browser.lock"))?;
81                match lock.try_lock() {
82                    Ok(()) => {}
83                    Err(std::fs::TryLockError::WouldBlock) => {
84                        return Err(BrowserError::ProfileInUse(name.clone()))
85                    }
86                    Err(std::fs::TryLockError::Error(e)) => return Err(e.into()),
87                }
88                let dl = dir.join("rightkit-downloads");
89                (ProfileGuard::Named { _lock: lock }, dir, dl)
90            }
91        };
92        let download_dir = opts.download_dir.clone().unwrap_or(default_dl);
93        std::fs::create_dir_all(&download_dir)?;
94
95        let id = format!("bs-{}", SESSION_SEQ.fetch_add(1, Ordering::Relaxed));
96        let guard = Arc::new(Guard {
97            session_id: id.clone(),
98            admission: opts.admission.clone(),
99            network: opts.network.clone(),
100            events: opts.on_event.clone(),
101        });
102        let proxy = if opts.network.unrestricted {
103            None
104        } else {
105            Some(PinProxy::start(guard.clone()).await?)
106        };
107
108        let mut args: Vec<String> = [
109            "--remote-debugging-port=0",
110            "--no-first-run",
111            "--no-default-browser-check",
112            "--disable-background-networking",
113            "--disable-background-timer-throttling",
114            "--disable-backgrounding-occluded-windows",
115            "--disable-breakpad",
116            "--disable-client-side-phishing-detection",
117            "--disable-component-extensions-with-background-pages",
118            "--disable-default-apps",
119            "--disable-dev-shm-usage",
120            "--disable-extensions",
121            "--disable-hang-monitor",
122            "--disable-ipc-flooding-protection",
123            "--disable-prompt-on-repost",
124            "--disable-renderer-backgrounding",
125            "--disable-sync",
126            "--force-color-profile=srgb",
127            "--metrics-recording-only",
128            "--password-store=basic",
129            "--use-mock-keychain",
130            // No unmanaged windows: a pop-up would be a page the policy never saw.
131            "--block-new-web-contents",
132        ]
133        .iter()
134        .map(|s| s.to_string())
135        .collect();
136        args.push(format!("--user-data-dir={}", profile_dir.display()));
137        args.push(format!(
138            "--window-size={},{}",
139            opts.viewport.0, opts.viewport.1
140        ));
141        args.push(if opts.headless {
142            "--headless=new".into()
143        } else {
144            "--hide-crash-restore-bubble".into()
145        });
146        if opts.mute_audio {
147            args.push("--mute-audio".into());
148        }
149        if let Some(p) = &proxy {
150            // Every connection, loopback included, goes through the pinning proxy;
151            // Chrome itself resolves nothing but the proxy address.
152            args.push(format!("--proxy-server=http://{}", p.addr));
153            args.push("--proxy-bypass-list=<-loopback>".into());
154            args.push("--host-resolver-rules=MAP * ~NOTFOUND , EXCLUDE 127.0.0.1".into());
155            // Keep cross-site frames in the page process so interception covers them.
156            args.push("--disable-site-isolation-trials".into());
157            // No UDP path around the proxy: WebRTC only over proxied routes, no QUIC.
158            args.push("--force-webrtc-ip-handling-policy=disable_non_proxied_udp".into());
159            args.push("--disable-quic".into());
160        }
161        for a in &opts.extra_args {
162            args.push(format!("--{}", a.trim_start_matches('-')));
163        }
164        args.push("about:blank".into());
165        let launch_timeout = opts.launch_timeout;
166        #[cfg(windows)]
167        let job = opts.windows_job.take();
168        #[cfg(not(windows))]
169        let job = None;
170        let proc = tokio::task::spawn_blocking(move || {
171            crate::proc::spawn_chrome(&exe, &args, launch_timeout, job)
172        })
173        .await
174        .map_err(|e| BrowserError::Launch(e.to_string()))??;
175        let pid = Some(proc.child.id());
176        let (browser, mut handler) = Browser::connect(proc.ws_url.clone())
177            .await
178            .map_err(|e| BrowserError::Launch(e.to_string()))?;
179        let proc = Arc::new(StdMutex::new(Some(proc)));
180        let stopping = Arc::new(AtomicBool::new(false));
181        let crashed = Arc::new(AtomicBool::new(false));
182        let (pump_proc, pump_stopping, pump_guard, pump_crashed) = (
183            proc.clone(),
184            stopping.clone(),
185            guard.clone(),
186            crashed.clone(),
187        );
188        let pump = tokio::spawn(async move {
189            // Transient handler errors are tolerated while Chrome is alive; the
190            // stream ending (or Chrome exiting) ends the session.
191            let mut errors = 0u32;
192            while let Some(ev) = handler.next().await {
193                if ev.is_err() {
194                    errors += 1;
195                    if errors > 50 || chrome_exit(&pump_proc).is_some() {
196                        break;
197                    }
198                } else {
199                    errors = 0;
200                }
201            }
202            // Give the OS a moment to reap so the reason carries the exit status.
203            let mut status = None;
204            for _ in 0..25 {
205                status = chrome_exit(&pump_proc);
206                if status.is_some() || pump_stopping.load(Ordering::SeqCst) {
207                    break;
208                }
209                tokio::time::sleep(Duration::from_millis(20)).await;
210            }
211            report_crash(&pump_guard, &pump_stopping, &pump_crashed, status, pid);
212        });
213        // The CDP stream does not always end when Chrome dies, so also watch the process.
214        let (w_proc, w_stopping, w_crashed, w_guard) = (
215            proc.clone(),
216            stopping.clone(),
217            crashed.clone(),
218            guard.clone(),
219        );
220        let watcher = tokio::spawn(async move {
221            loop {
222                tokio::time::sleep(Duration::from_millis(200)).await;
223                if w_stopping.load(Ordering::SeqCst) || w_crashed.load(Ordering::SeqCst) {
224                    return;
225                }
226                if let Some(status) = chrome_exit(&w_proc) {
227                    report_crash(&w_guard, &w_stopping, &w_crashed, Some(status), pid);
228                    return;
229                }
230            }
231        });
232        browser
233            .execute(
234                SetDownloadBehaviorParams::builder()
235                    .behavior(SetDownloadBehaviorBehavior::Allow)
236                    .download_path(download_dir.to_string_lossy().to_string())
237                    .build()
238                    .map_err(BrowserError::Launch)?,
239            )
240            .await
241            .map_err(cdp)?;
242
243        let mode = if opts.headless { "headless" } else { "headed" };
244        guard.emit_kind(LifecycleKind::Start, format!("launched ({mode})"), |e| {
245            e.pid = pid
246        });
247        Ok(Self {
248            id,
249            browser: Mutex::new(browser),
250            pages: StdMutex::new(Vec::new()),
251            current: StdMutex::new(None),
252            page_seq: AtomicU64::new(0),
253            handler: StdMutex::new(Some(pump)),
254            watcher: StdMutex::new(Some(watcher)),
255            _profile: profile_guard,
256            profile_dir,
257            download_dir,
258            upload_root: opts.upload_root,
259            headless: opts.headless,
260            guard,
261            pid,
262            proc,
263            proxy,
264            stopping,
265            stopped: AtomicBool::new(false),
266        })
267    }
268
269    pub fn id(&self) -> &str {
270        &self.id
271    }
272    /// OS process id of the Chrome child this session owns.
273    pub fn process_id(&self) -> Option<u32> {
274        self.pid
275    }
276    /// Connections the pinning proxy made, with the address it pinned for each.
277    pub fn pinned_connections(&self) -> Vec<PinnedConnection> {
278        self.proxy
279            .as_ref()
280            .map(|p| p.log.lock().unwrap().clone())
281            .unwrap_or_default()
282    }
283    pub fn is_headless(&self) -> bool {
284        self.headless
285    }
286    pub fn profile_dir(&self) -> &Path {
287        &self.profile_dir
288    }
289    pub fn download_dir(&self) -> &Path {
290        &self.download_dir
291    }
292
293    /// Open a new page (tab), make it current, and navigate to `url`.
294    pub async fn new_page(&self, url: &str) -> Result<Arc<BrowserPage>> {
295        let inner = self
296            .browser
297            .lock()
298            .await
299            .new_page("about:blank")
300            .await
301            .map_err(cdp)?;
302        let id = format!("p{}", self.page_seq.fetch_add(1, Ordering::Relaxed) + 1);
303        let page = Arc::new(
304            BrowserPage::attach(
305                id.clone(),
306                inner,
307                self.upload_root.clone(),
308                self.guard.clone(),
309            )
310            .await?,
311        );
312        self.pages.lock().unwrap().push(page.clone());
313        *self.current.lock().unwrap() = Some(id);
314        if url != "about:blank" {
315            page.goto(url).await?;
316        }
317        Ok(page)
318    }
319
320    pub fn pages(&self) -> Vec<Arc<BrowserPage>> {
321        self.pages.lock().unwrap().clone()
322    }
323
324    pub fn page(&self, id: &str) -> Result<Arc<BrowserPage>> {
325        self.pages
326            .lock()
327            .unwrap()
328            .iter()
329            .find(|p| p.id() == id)
330            .cloned()
331            .ok_or_else(|| BrowserError::UnknownPage(id.into()))
332    }
333
334    /// The page actions default to.
335    pub fn current(&self) -> Result<Arc<BrowserPage>> {
336        let id = self
337            .current
338            .lock()
339            .unwrap()
340            .clone()
341            .ok_or_else(|| BrowserError::UnknownPage("<none>".into()))?;
342        self.page(&id)
343    }
344
345    /// Make `id` current and bring it to the front.
346    pub async fn select_page(&self, id: &str) -> Result<Arc<BrowserPage>> {
347        let p = self.page(id)?;
348        p.bring_to_front().await?;
349        *self.current.lock().unwrap() = Some(id.into());
350        Ok(p)
351    }
352
353    pub async fn close_page(&self, id: &str) -> Result<()> {
354        let p = self.page(id)?;
355        self.pages.lock().unwrap().retain(|x| x.id() != id);
356        {
357            let mut cur = self.current.lock().unwrap();
358            if cur.as_deref() == Some(id) {
359                *cur = self
360                    .pages
361                    .lock()
362                    .unwrap()
363                    .last()
364                    .map(|p| p.id().to_string());
365            }
366        }
367        p.close_inner().await
368    }
369
370    /// Completed files currently in the download directory.
371    pub fn downloads(&self) -> HashSet<PathBuf> {
372        completed_files(&self.download_dir)
373    }
374
375    /// Wait for a download not in `before` to finish (no `.crdownload`, stable size).
376    pub async fn wait_download(
377        &self,
378        before: &HashSet<PathBuf>,
379        timeout: Duration,
380    ) -> Result<PathBuf> {
381        let deadline = Instant::now() + timeout;
382        loop {
383            let now = completed_files(&self.download_dir);
384            if let Some(p) = now.difference(before).next() {
385                let size = std::fs::metadata(p).map(|m| m.len()).unwrap_or(0);
386                tokio::time::sleep(Duration::from_millis(100)).await;
387                if std::fs::metadata(p).map(|m| m.len()).unwrap_or(1) == size {
388                    return Ok(p.clone());
389                }
390            }
391            if Instant::now() >= deadline {
392                return Err(BrowserError::Timeout("download did not complete".into()));
393            }
394            tokio::time::sleep(Duration::from_millis(100)).await;
395        }
396    }
397
398    /// Set a W3C permission (`geolocation`, `notifications`, `clipboard-read`,
399    /// `midi`, ...) to granted or denied for `origin`, or every origin when `None`.
400    pub async fn set_permission(
401        &self,
402        origin: Option<&str>,
403        name: &str,
404        granted: bool,
405    ) -> Result<()> {
406        let mut params = SetPermissionParams::new(
407            PermissionDescriptor::new(name),
408            if granted {
409                PermissionSetting::Granted
410            } else {
411                PermissionSetting::Denied
412            },
413        );
414        params.origin = origin.map(str::to_string);
415        self.browser
416            .lock()
417            .await
418            .execute(params)
419            .await
420            .map_err(cdp)?;
421        Ok(())
422    }
423
424    pub async fn reset_permissions(&self) -> Result<()> {
425        self.browser
426            .lock()
427            .await
428            .execute(ResetPermissionsParams::default())
429            .await
430            .map_err(cdp)?;
431        Ok(())
432    }
433
434    /// Close Chrome gracefully, killing it if it does not exit within 5 seconds.
435    /// Safe to call twice. Temporary profiles are deleted on drop.
436    pub async fn shutdown(&self) {
437        self.stop(false).await;
438    }
439
440    /// Cancel: kill the owned Chrome now and wait for it to exit. Helper
441    /// processes exit with their parent; none are detached.
442    pub async fn cancel(&self) {
443        self.stop(true).await;
444    }
445
446    async fn stop(&self, hard: bool) {
447        self.stopping.store(true, Ordering::SeqCst);
448        if !hard {
449            let graceful = async {
450                let _ = self.browser.lock().await.close().await;
451                let end = Instant::now() + Duration::from_secs(5);
452                loop {
453                    let done = match self.proc.lock().unwrap().as_mut() {
454                        Some(p) => matches!(p.child.try_wait(), Ok(Some(_)) | Err(_)),
455                        None => true,
456                    };
457                    if done || Instant::now() >= end {
458                        break;
459                    }
460                    tokio::time::sleep(Duration::from_millis(20)).await;
461                }
462            };
463            let _ = tokio::time::timeout(Duration::from_secs(6), graceful).await;
464        }
465        // Whatever is left (or everything, on cancel) dies now: process group on
466        // Unix, job object on Windows. Dropping `ChromeProc` also stops the watchdog.
467        if let Some(mut p) = self.proc.lock().unwrap().take() {
468            let _ = p.child.terminate_tree();
469            if let Some(w) = p._watchdog.as_mut() {
470                let _ = w.terminate_tree();
471            }
472        }
473        for p in self.pages.lock().unwrap().drain(..) {
474            p.abort_pumps();
475        }
476        if let Some(h) = self.handler.lock().unwrap().take() {
477            h.abort();
478        }
479        if let Some(w) = self.watcher.lock().unwrap().take() {
480            w.abort();
481        }
482        let reason = if hard { "cancelled" } else { "shutdown" };
483        self.emit_stopped(reason);
484    }
485
486    fn emit_stopped(&self, reason: &str) {
487        if !self.stopped.swap(true, Ordering::SeqCst) {
488            let pid = self.pid;
489            self.guard
490                .emit_kind(LifecycleKind::Stop, reason, |e| e.pid = pid);
491        }
492    }
493}
494
495/// Emit at most one `Crash` per session, never during an intentional stop.
496fn report_crash(
497    guard: &Guard,
498    stopping: &AtomicBool,
499    crashed: &AtomicBool,
500    status: Option<std::process::ExitStatus>,
501    pid: Option<u32>,
502) {
503    if stopping.load(Ordering::SeqCst) || crashed.swap(true, Ordering::SeqCst) {
504        return;
505    }
506    let reason = match status {
507        Some(s) => format!("chrome exited unexpectedly ({s})"),
508        None => "devtools connection lost".to_string(),
509    };
510    guard.emit_kind(LifecycleKind::Crash, reason, |e| e.pid = pid);
511}
512
513/// Exit status of the owned Chrome, if it has exited.
514fn chrome_exit(proc: &StdMutex<Option<ChromeProc>>) -> Option<std::process::ExitStatus> {
515    proc.lock()
516        .ok()?
517        .as_mut()
518        .and_then(|p| p.child.try_wait().ok().flatten())
519}
520
521impl Drop for BrowserSession {
522    fn drop(&mut self) {
523        self.stopping.store(true, Ordering::SeqCst);
524        for p in self.pages.lock().unwrap().iter() {
525            p.abort_pumps();
526        }
527        if let Some(h) = self.handler.lock().unwrap().take() {
528            h.abort();
529        }
530        if let Some(w) = self.watcher.lock().unwrap().take() {
531            w.abort();
532        }
533        // Dropping `ChromeProc` terminates Chrome's process tree, so an abandoned
534        // session (cancelled caller, panic, early return) never leaves Chrome running.
535        if let Some(mut p) = self.proc.lock().unwrap().take() {
536            let _ = p.child.terminate_tree();
537        }
538        self.emit_stopped("dropped");
539    }
540}
541
542fn completed_files(dir: &Path) -> HashSet<PathBuf> {
543    std::fs::read_dir(dir)
544        .map(|rd| {
545            rd.filter_map(|e| e.ok())
546                .map(|e| e.path())
547                .filter(|p| {
548                    p.is_file()
549                        && !p
550                            .extension()
551                            .is_some_and(|e| e == "crdownload" || e == "tmp")
552                })
553                .collect()
554        })
555        .unwrap_or_default()
556}