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