Skip to main content

rightkit_browser/
session.rs

1use crate::config::{
2    reject_real_browser_profile, resolve_executable, 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
26/// Per-step wait before escalating a close (CDP -> SIGTERM -> SIGKILL).
27const GRACE: Duration = Duration::from_secs(5);
28
29enum ProfileGuard {
30    Temp(#[allow(dead_code)] tempfile::TempDir),
31    Named { _lock: std::fs::File },
32}
33
34/// One Chrome process with any number of pages.
35pub struct BrowserSession {
36    id: String,
37    browser: Mutex<Browser>,
38    pages: StdMutex<Vec<Arc<BrowserPage>>>,
39    current: StdMutex<Option<String>>,
40    page_seq: AtomicU64,
41    handler: StdMutex<Option<JoinHandle<()>>>,
42    watcher: StdMutex<Option<JoinHandle<()>>>,
43    /// Released (temporary profile deleted) once Chrome has exited, on shutdown or drop.
44    profile_guard: StdMutex<Option<ProfileGuard>>,
45    profile_dir: PathBuf,
46    download_dir: PathBuf,
47    upload_root: Option<PathBuf>,
48    headless: bool,
49    guard: Arc<Guard>,
50    pid: Option<u32>,
51    proc: Arc<StdMutex<Option<ChromeProc>>>,
52    proxy: Option<PinProxy>,
53    /// Set before an intentional stop so the CDP pump does not report a crash.
54    stopping: Arc<AtomicBool>,
55    stopped: AtomicBool,
56}
57
58impl BrowserSession {
59    pub async fn launch(
60        #[cfg_attr(not(windows), allow(unused_mut))] mut opts: LaunchOptions,
61    ) -> Result<Self> {
62        validate_extra_args(&opts.extra_args)?;
63        let exe = resolve_executable(&opts)?;
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_guard: StdMutex::new(Some(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: CDP `Browser.close` and up to 5 s for exit, then SIGTERM to
435    /// the process group and up to 5 s more, then SIGKILL (job termination on Windows).
436    /// Safe to call twice. A temporary profile is deleted once Chrome has exited.
437    pub async fn shutdown(&self) {
438        self.stop(false).await;
439    }
440
441    /// Cancel: kill the owned Chrome now and wait for it to exit. Helper
442    /// processes exit with their parent; none are detached.
443    pub async fn cancel(&self) {
444        self.stop(true).await;
445    }
446
447    async fn stop(&self, hard: bool) {
448        self.stopping.store(true, Ordering::SeqCst);
449        if !hard {
450            let graceful = async {
451                let _ = self.browser.lock().await.close().await;
452                let end = Instant::now() + GRACE;
453                loop {
454                    let done = match self.proc.lock().unwrap().as_mut() {
455                        Some(p) => matches!(p.child.try_wait(), Ok(Some(_)) | Err(_)),
456                        None => true,
457                    };
458                    if done || Instant::now() >= end {
459                        break;
460                    }
461                    tokio::time::sleep(Duration::from_millis(20)).await;
462                }
463            };
464            let _ = tokio::time::timeout(Duration::from_secs(6), graceful).await;
465        }
466        // Whatever is left: SIGTERM the process group with a bounded wait, SIGKILL last
467        // (cancel kills at once). Windows: the job object takes the whole tree.
468        let taken = self.proc.lock().unwrap().take();
469        if let Some(mut p) = taken {
470            let grace = if hard { Duration::ZERO } else { GRACE };
471            let _ = tokio::task::spawn_blocking(move || p.terminate_gracefully(grace)).await;
472        }
473        drop(self.profile_guard.lock().unwrap().take());
474        for p in self.pages.lock().unwrap().drain(..) {
475            p.abort_pumps();
476        }
477        if let Some(h) = self.handler.lock().unwrap().take() {
478            h.abort();
479        }
480        if let Some(w) = self.watcher.lock().unwrap().take() {
481            w.abort();
482        }
483        let reason = if hard { "cancelled" } else { "shutdown" };
484        self.emit_stopped(reason);
485    }
486
487    fn emit_stopped(&self, reason: &str) {
488        if !self.stopped.swap(true, Ordering::SeqCst) {
489            let pid = self.pid;
490            self.guard
491                .emit_kind(LifecycleKind::Stop, reason, |e| e.pid = pid);
492        }
493    }
494}
495
496/// Emit at most one `Crash` per session, never during an intentional stop.
497fn report_crash(
498    guard: &Guard,
499    stopping: &AtomicBool,
500    crashed: &AtomicBool,
501    status: Option<std::process::ExitStatus>,
502    pid: Option<u32>,
503) {
504    if stopping.load(Ordering::SeqCst) || crashed.swap(true, Ordering::SeqCst) {
505        return;
506    }
507    let reason = match status {
508        Some(s) => format!("chrome exited unexpectedly ({s})"),
509        None => "devtools connection lost".to_string(),
510    };
511    guard.emit_kind(LifecycleKind::Crash, reason, |e| e.pid = pid);
512}
513
514/// Exit status of the owned Chrome, if it has exited.
515fn chrome_exit(proc: &StdMutex<Option<ChromeProc>>) -> Option<std::process::ExitStatus> {
516    proc.lock()
517        .ok()?
518        .as_mut()
519        .and_then(|p| p.child.try_wait().ok().flatten())
520}
521
522impl Drop for BrowserSession {
523    fn drop(&mut self) {
524        self.stopping.store(true, Ordering::SeqCst);
525        for p in self.pages.lock().unwrap().iter() {
526            p.abort_pumps();
527        }
528        if let Some(h) = self.handler.lock().unwrap().take() {
529            h.abort();
530        }
531        if let Some(w) = self.watcher.lock().unwrap().take() {
532            w.abort();
533        }
534        // An abandoned session (cancelled caller, panic, early return) still closes
535        // gracefully: SIGTERM to Chrome's process group, a bounded wait, SIGKILL last.
536        if let Some(mut p) = self.proc.lock().unwrap().take() {
537            p.terminate_gracefully(GRACE);
538        }
539        // Chrome has exited: now the temporary profile can go.
540        drop(self.profile_guard.lock().unwrap().take());
541        self.emit_stopped("dropped");
542    }
543}
544
545fn completed_files(dir: &Path) -> HashSet<PathBuf> {
546    std::fs::read_dir(dir)
547        .map(|rd| {
548            rd.filter_map(|e| e.ok())
549                .map(|e| e.path())
550                .filter(|p| {
551                    p.is_file()
552                        && !p
553                            .extension()
554                            .is_some_and(|e| e == "crdownload" || e == "tmp")
555                })
556                .collect()
557        })
558        .unwrap_or_default()
559}