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
27const 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 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
78pub 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 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 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 "--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 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 args.push("--disable-site-isolation-trials".into());
171 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 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 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 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 pub fn process_id(&self) -> Option<u32> {
288 self.pid
289 }
290 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 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 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 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 pub fn downloads(&self) -> HashSet<PathBuf> {
386 completed_files(&self.download_dir)
387 }
388
389 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 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 pub async fn shutdown(&self) {
452 self.stop(false).await;
453 }
454
455 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 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
510fn 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
528fn 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 if let Some(mut p) = self.proc.lock().unwrap().take() {
551 p.terminate_gracefully(GRACE);
552 }
553 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 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 #[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}