Skip to main content

rightkit_browser/
session.rs

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