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
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 _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 "--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 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 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 pub fn process_id(&self) -> Option<u32> {
214 self.pid
215 }
216 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 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 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 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 pub fn downloads(&self) -> HashSet<PathBuf> {
312 completed_files(&self.download_dir)
313 }
314
315 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 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 pub async fn shutdown(&self) {
377 self.stop(false).await;
378 }
379
380 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 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 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}