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::{BrowserEvent, Guard};
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::{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    _profile: ProfileGuard,
40    profile_dir: PathBuf,
41    download_dir: PathBuf,
42    upload_root: Option<PathBuf>,
43    headless: bool,
44    guard: Arc<Guard>,
45    pid: Option<u32>,
46    proc: StdMutex<Option<ChromeProc>>,
47    proxy: Option<PinProxy>,
48    stopped: std::sync::atomic::AtomicBool,
49}
50
51impl BrowserSession {
52    pub async fn launch(opts: LaunchOptions) -> Result<Self> {
53        validate_extra_args(&opts.extra_args)?;
54        let exe = opts
55            .chrome_path
56            .clone()
57            .or_else(find_chrome)
58            .ok_or(BrowserError::ChromeNotFound)?;
59        let (profile_guard, profile_dir, default_dl) = match &opts.profile {
60            ProfileSpec::Temporary => {
61                let t = tempfile::tempdir()?;
62                let p = t.path().to_path_buf();
63                let dl = p.join("downloads");
64                (ProfileGuard::Temp(t), p, dl)
65            }
66            ProfileSpec::Named { root, name } => {
67                validate_profile_name(name)?;
68                reject_real_browser_profile(root)?;
69                let dir = root.join(name);
70                std::fs::create_dir_all(&dir)?;
71                let lock = std::fs::OpenOptions::new()
72                    .create(true)
73                    .write(true)
74                    .truncate(false)
75                    .open(dir.join(".rightkit-browser.lock"))?;
76                match lock.try_lock() {
77                    Ok(()) => {}
78                    Err(std::fs::TryLockError::WouldBlock) => {
79                        return Err(BrowserError::ProfileInUse(name.clone()))
80                    }
81                    Err(std::fs::TryLockError::Error(e)) => return Err(e.into()),
82                }
83                let dl = dir.join("rightkit-downloads");
84                (ProfileGuard::Named { _lock: lock }, dir, dl)
85            }
86        };
87        let download_dir = opts.download_dir.clone().unwrap_or(default_dl);
88        std::fs::create_dir_all(&download_dir)?;
89
90        let id = format!("bs-{}", SESSION_SEQ.fetch_add(1, Ordering::Relaxed));
91        let guard = Arc::new(Guard {
92            session_id: id.clone(),
93            admission: opts.admission.clone(),
94            network: opts.network.clone(),
95            events: opts.on_event.clone(),
96        });
97        let proxy = if opts.network.unrestricted {
98            None
99        } else {
100            Some(PinProxy::start(guard.clone()).await?)
101        };
102
103        let mut args: Vec<String> = [
104            "--remote-debugging-port=0",
105            "--no-first-run",
106            "--no-default-browser-check",
107            "--disable-background-networking",
108            "--disable-background-timer-throttling",
109            "--disable-backgrounding-occluded-windows",
110            "--disable-breakpad",
111            "--disable-client-side-phishing-detection",
112            "--disable-component-extensions-with-background-pages",
113            "--disable-default-apps",
114            "--disable-dev-shm-usage",
115            "--disable-extensions",
116            "--disable-hang-monitor",
117            "--disable-ipc-flooding-protection",
118            "--disable-prompt-on-repost",
119            "--disable-renderer-backgrounding",
120            "--disable-sync",
121            "--force-color-profile=srgb",
122            "--metrics-recording-only",
123            "--password-store=basic",
124            "--use-mock-keychain",
125            // No unmanaged windows: a pop-up would be a page the policy never saw.
126            "--block-new-web-contents",
127        ]
128        .iter()
129        .map(|s| s.to_string())
130        .collect();
131        args.push(format!("--user-data-dir={}", profile_dir.display()));
132        args.push(format!(
133            "--window-size={},{}",
134            opts.viewport.0, opts.viewport.1
135        ));
136        args.push(if opts.headless {
137            "--headless=new".into()
138        } else {
139            "--hide-crash-restore-bubble".into()
140        });
141        if opts.mute_audio {
142            args.push("--mute-audio".into());
143        }
144        if let Some(p) = &proxy {
145            // Every connection, loopback included, goes through the pinning proxy;
146            // Chrome itself resolves nothing but the proxy address.
147            args.push(format!("--proxy-server=http://{}", p.addr));
148            args.push("--proxy-bypass-list=<-loopback>".into());
149            args.push("--host-resolver-rules=MAP * ~NOTFOUND , EXCLUDE 127.0.0.1".into());
150            // Keep cross-site frames in the page process so interception covers them.
151            args.push("--disable-site-isolation-trials".into());
152        }
153        for a in &opts.extra_args {
154            args.push(format!("--{}", a.trim_start_matches('-')));
155        }
156        args.push("about:blank".into());
157        let launch_timeout = opts.launch_timeout;
158        let proc = tokio::task::spawn_blocking(move || {
159            crate::proc::spawn_chrome(&exe, &args, launch_timeout)
160        })
161        .await
162        .map_err(|e| BrowserError::Launch(e.to_string()))??;
163        let pid = Some(proc.child.id());
164        let (browser, mut handler) = Browser::connect(proc.ws_url.clone())
165            .await
166            .map_err(|e| BrowserError::Launch(e.to_string()))?;
167        let pump = tokio::spawn(async move {
168            while let Some(ev) = handler.next().await {
169                if ev.is_err() {
170                    break;
171                }
172            }
173        });
174        browser
175            .execute(
176                SetDownloadBehaviorParams::builder()
177                    .behavior(SetDownloadBehaviorBehavior::Allow)
178                    .download_path(download_dir.to_string_lossy().to_string())
179                    .build()
180                    .map_err(BrowserError::Launch)?,
181            )
182            .await
183            .map_err(cdp)?;
184
185        guard.emit(BrowserEvent::Started {
186            session_id: id.clone(),
187            pid,
188        });
189        Ok(Self {
190            id,
191            browser: Mutex::new(browser),
192            pages: StdMutex::new(Vec::new()),
193            current: StdMutex::new(None),
194            page_seq: AtomicU64::new(0),
195            handler: StdMutex::new(Some(pump)),
196            _profile: profile_guard,
197            profile_dir,
198            download_dir,
199            upload_root: opts.upload_root,
200            headless: opts.headless,
201            guard,
202            pid,
203            proc: StdMutex::new(Some(proc)),
204            proxy,
205            stopped: std::sync::atomic::AtomicBool::new(false),
206        })
207    }
208
209    pub fn id(&self) -> &str {
210        &self.id
211    }
212    /// OS process id of the Chrome child this session owns.
213    pub fn process_id(&self) -> Option<u32> {
214        self.pid
215    }
216    /// Connections the pinning proxy made, with the address it pinned for each.
217    pub fn pinned_connections(&self) -> Vec<PinnedConnection> {
218        self.proxy
219            .as_ref()
220            .map(|p| p.log.lock().unwrap().clone())
221            .unwrap_or_default()
222    }
223    pub fn is_headless(&self) -> bool {
224        self.headless
225    }
226    pub fn profile_dir(&self) -> &Path {
227        &self.profile_dir
228    }
229    pub fn download_dir(&self) -> &Path {
230        &self.download_dir
231    }
232
233    /// Open a new page (tab), make it current, and navigate to `url`.
234    pub async fn new_page(&self, url: &str) -> Result<Arc<BrowserPage>> {
235        let inner = self
236            .browser
237            .lock()
238            .await
239            .new_page("about:blank")
240            .await
241            .map_err(cdp)?;
242        let id = format!("p{}", self.page_seq.fetch_add(1, Ordering::Relaxed) + 1);
243        let page = Arc::new(
244            BrowserPage::attach(
245                id.clone(),
246                inner,
247                self.upload_root.clone(),
248                self.guard.clone(),
249            )
250            .await?,
251        );
252        self.pages.lock().unwrap().push(page.clone());
253        *self.current.lock().unwrap() = Some(id);
254        if url != "about:blank" {
255            page.goto(url).await?;
256        }
257        Ok(page)
258    }
259
260    pub fn pages(&self) -> Vec<Arc<BrowserPage>> {
261        self.pages.lock().unwrap().clone()
262    }
263
264    pub fn page(&self, id: &str) -> Result<Arc<BrowserPage>> {
265        self.pages
266            .lock()
267            .unwrap()
268            .iter()
269            .find(|p| p.id() == id)
270            .cloned()
271            .ok_or_else(|| BrowserError::UnknownPage(id.into()))
272    }
273
274    /// The page actions default to.
275    pub fn current(&self) -> Result<Arc<BrowserPage>> {
276        let id = self
277            .current
278            .lock()
279            .unwrap()
280            .clone()
281            .ok_or_else(|| BrowserError::UnknownPage("<none>".into()))?;
282        self.page(&id)
283    }
284
285    /// Make `id` current and bring it to the front.
286    pub async fn select_page(&self, id: &str) -> Result<Arc<BrowserPage>> {
287        let p = self.page(id)?;
288        p.bring_to_front().await?;
289        *self.current.lock().unwrap() = Some(id.into());
290        Ok(p)
291    }
292
293    pub async fn close_page(&self, id: &str) -> Result<()> {
294        let p = self.page(id)?;
295        self.pages.lock().unwrap().retain(|x| x.id() != id);
296        {
297            let mut cur = self.current.lock().unwrap();
298            if cur.as_deref() == Some(id) {
299                *cur = self
300                    .pages
301                    .lock()
302                    .unwrap()
303                    .last()
304                    .map(|p| p.id().to_string());
305            }
306        }
307        p.close_inner().await
308    }
309
310    /// Completed files currently in the download directory.
311    pub fn downloads(&self) -> HashSet<PathBuf> {
312        completed_files(&self.download_dir)
313    }
314
315    /// Wait for a download not in `before` to finish (no `.crdownload`, stable size).
316    pub async fn wait_download(
317        &self,
318        before: &HashSet<PathBuf>,
319        timeout: Duration,
320    ) -> Result<PathBuf> {
321        let deadline = Instant::now() + timeout;
322        loop {
323            let now = completed_files(&self.download_dir);
324            if let Some(p) = now.difference(before).next() {
325                let size = std::fs::metadata(p).map(|m| m.len()).unwrap_or(0);
326                tokio::time::sleep(Duration::from_millis(100)).await;
327                if std::fs::metadata(p).map(|m| m.len()).unwrap_or(1) == size {
328                    return Ok(p.clone());
329                }
330            }
331            if Instant::now() >= deadline {
332                return Err(BrowserError::Timeout("download did not complete".into()));
333            }
334            tokio::time::sleep(Duration::from_millis(100)).await;
335        }
336    }
337
338    /// Set a W3C permission (`geolocation`, `notifications`, `clipboard-read`,
339    /// `midi`, ...) to granted or denied for `origin`, or every origin when `None`.
340    pub async fn set_permission(
341        &self,
342        origin: Option<&str>,
343        name: &str,
344        granted: bool,
345    ) -> Result<()> {
346        let mut params = SetPermissionParams::new(
347            PermissionDescriptor::new(name),
348            if granted {
349                PermissionSetting::Granted
350            } else {
351                PermissionSetting::Denied
352            },
353        );
354        params.origin = origin.map(str::to_string);
355        self.browser
356            .lock()
357            .await
358            .execute(params)
359            .await
360            .map_err(cdp)?;
361        Ok(())
362    }
363
364    pub async fn reset_permissions(&self) -> Result<()> {
365        self.browser
366            .lock()
367            .await
368            .execute(ResetPermissionsParams::default())
369            .await
370            .map_err(cdp)?;
371        Ok(())
372    }
373
374    /// Close Chrome gracefully, killing it if it does not exit within 5 seconds.
375    /// Safe to call twice. Temporary profiles are deleted on drop.
376    pub async fn shutdown(&self) {
377        self.stop(false).await;
378    }
379
380    /// Cancel: kill the owned Chrome now and wait for it to exit. Helper
381    /// processes exit with their parent; none are detached.
382    pub async fn cancel(&self) {
383        self.stop(true).await;
384    }
385
386    async fn stop(&self, hard: bool) {
387        if !hard {
388            let graceful = async {
389                let _ = self.browser.lock().await.close().await;
390                let end = Instant::now() + Duration::from_secs(5);
391                loop {
392                    let done = match self.proc.lock().unwrap().as_mut() {
393                        Some(p) => matches!(p.child.try_wait(), Ok(Some(_)) | Err(_)),
394                        None => true,
395                    };
396                    if done || Instant::now() >= end {
397                        break;
398                    }
399                    tokio::time::sleep(Duration::from_millis(20)).await;
400                }
401            };
402            let _ = tokio::time::timeout(Duration::from_secs(6), graceful).await;
403        }
404        // Whatever is left (or everything, on cancel) dies now: process group on
405        // Unix, job object on Windows. Dropping `ChromeProc` also stops the watchdog.
406        if let Some(mut p) = self.proc.lock().unwrap().take() {
407            let _ = p.child.terminate_tree();
408            if let Some(w) = p._watchdog.as_mut() {
409                let _ = w.terminate_tree();
410            }
411        }
412        for p in self.pages.lock().unwrap().drain(..) {
413            p.abort_pumps();
414        }
415        if let Some(h) = self.handler.lock().unwrap().take() {
416            h.abort();
417        }
418        if !self.stopped.swap(true, Ordering::SeqCst) {
419            self.guard.emit(BrowserEvent::Stopped {
420                session_id: self.id.clone(),
421            });
422        }
423    }
424}
425
426impl Drop for BrowserSession {
427    fn drop(&mut self) {
428        for p in self.pages.lock().unwrap().iter() {
429            p.abort_pumps();
430        }
431        if let Some(h) = self.handler.lock().unwrap().take() {
432            h.abort();
433        }
434        // Dropping `ChromeProc` terminates Chrome's process tree, so an abandoned
435        // session (cancelled caller, panic, early return) never leaves Chrome running.
436        if let Some(mut p) = self.proc.lock().unwrap().take() {
437            let _ = p.child.terminate_tree();
438        }
439        if !self.stopped.swap(true, Ordering::SeqCst) {
440            self.guard.emit(BrowserEvent::Stopped {
441                session_id: self.id.clone(),
442            });
443        }
444    }
445}
446
447fn completed_files(dir: &Path) -> HashSet<PathBuf> {
448    std::fs::read_dir(dir)
449        .map(|rd| {
450            rd.filter_map(|e| e.ok())
451                .map(|e| e.path())
452                .filter(|p| {
453                    p.is_file()
454                        && !p
455                            .extension()
456                            .is_some_and(|e| e == "crdownload" || e == "tmp")
457                })
458                .collect()
459        })
460        .unwrap_or_default()
461}