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::{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
26enum ProfileGuard {
27 Temp(#[allow(dead_code)] tempfile::TempDir),
28 Named { _lock: std::fs::File },
29}
30
31pub 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 watcher: StdMutex<Option<JoinHandle<()>>>,
40 _profile: ProfileGuard,
41 profile_dir: PathBuf,
42 download_dir: PathBuf,
43 upload_root: Option<PathBuf>,
44 headless: bool,
45 guard: Arc<Guard>,
46 pid: Option<u32>,
47 proc: Arc<StdMutex<Option<ChromeProc>>>,
48 proxy: Option<PinProxy>,
49 stopping: Arc<AtomicBool>,
51 stopped: AtomicBool,
52}
53
54impl BrowserSession {
55 pub async fn launch(
56 #[cfg_attr(not(windows), allow(unused_mut))] mut opts: LaunchOptions,
57 ) -> Result<Self> {
58 validate_extra_args(&opts.extra_args)?;
59 let exe = opts
60 .chrome_path
61 .clone()
62 .or_else(find_chrome)
63 .ok_or(BrowserError::ChromeNotFound)?;
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 "--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 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 args.push("--disable-site-isolation-trials".into());
157 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 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 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 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: 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 pub fn process_id(&self) -> Option<u32> {
274 self.pid
275 }
276 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 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 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 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 pub fn downloads(&self) -> HashSet<PathBuf> {
372 completed_files(&self.download_dir)
373 }
374
375 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 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 pub async fn shutdown(&self) {
437 self.stop(false).await;
438 }
439
440 pub async fn cancel(&self) {
443 self.stop(true).await;
444 }
445
446 async fn stop(&self, hard: bool) {
447 self.stopping.store(true, Ordering::SeqCst);
448 if !hard {
449 let graceful = async {
450 let _ = self.browser.lock().await.close().await;
451 let end = Instant::now() + Duration::from_secs(5);
452 loop {
453 let done = match self.proc.lock().unwrap().as_mut() {
454 Some(p) => matches!(p.child.try_wait(), Ok(Some(_)) | Err(_)),
455 None => true,
456 };
457 if done || Instant::now() >= end {
458 break;
459 }
460 tokio::time::sleep(Duration::from_millis(20)).await;
461 }
462 };
463 let _ = tokio::time::timeout(Duration::from_secs(6), graceful).await;
464 }
465 if let Some(mut p) = self.proc.lock().unwrap().take() {
468 let _ = p.child.terminate_tree();
469 if let Some(w) = p._watchdog.as_mut() {
470 let _ = w.terminate_tree();
471 }
472 }
473 for p in self.pages.lock().unwrap().drain(..) {
474 p.abort_pumps();
475 }
476 if let Some(h) = self.handler.lock().unwrap().take() {
477 h.abort();
478 }
479 if let Some(w) = self.watcher.lock().unwrap().take() {
480 w.abort();
481 }
482 let reason = if hard { "cancelled" } else { "shutdown" };
483 self.emit_stopped(reason);
484 }
485
486 fn emit_stopped(&self, reason: &str) {
487 if !self.stopped.swap(true, Ordering::SeqCst) {
488 let pid = self.pid;
489 self.guard
490 .emit_kind(LifecycleKind::Stop, reason, |e| e.pid = pid);
491 }
492 }
493}
494
495fn report_crash(
497 guard: &Guard,
498 stopping: &AtomicBool,
499 crashed: &AtomicBool,
500 status: Option<std::process::ExitStatus>,
501 pid: Option<u32>,
502) {
503 if stopping.load(Ordering::SeqCst) || crashed.swap(true, Ordering::SeqCst) {
504 return;
505 }
506 let reason = match status {
507 Some(s) => format!("chrome exited unexpectedly ({s})"),
508 None => "devtools connection lost".to_string(),
509 };
510 guard.emit_kind(LifecycleKind::Crash, reason, |e| e.pid = pid);
511}
512
513fn chrome_exit(proc: &StdMutex<Option<ChromeProc>>) -> Option<std::process::ExitStatus> {
515 proc.lock()
516 .ok()?
517 .as_mut()
518 .and_then(|p| p.child.try_wait().ok().flatten())
519}
520
521impl Drop for BrowserSession {
522 fn drop(&mut self) {
523 self.stopping.store(true, Ordering::SeqCst);
524 for p in self.pages.lock().unwrap().iter() {
525 p.abort_pumps();
526 }
527 if let Some(h) = self.handler.lock().unwrap().take() {
528 h.abort();
529 }
530 if let Some(w) = self.watcher.lock().unwrap().take() {
531 w.abort();
532 }
533 if let Some(mut p) = self.proc.lock().unwrap().take() {
536 let _ = p.child.terminate_tree();
537 }
538 self.emit_stopped("dropped");
539 }
540}
541
542fn completed_files(dir: &Path) -> HashSet<PathBuf> {
543 std::fs::read_dir(dir)
544 .map(|rd| {
545 rd.filter_map(|e| e.ok())
546 .map(|e| e.path())
547 .filter(|p| {
548 p.is_file()
549 && !p
550 .extension()
551 .is_some_and(|e| e == "crdownload" || e == "tmp")
552 })
553 .collect()
554 })
555 .unwrap_or_default()
556}