use std::io::Write as _;
use serde_json::{Value, json};
use tokio::io::{AsyncBufReadExt, BufReader};
use crate::browser::{self, BrowserOptions};
use crate::cdp::client::CdpClient;
use crate::cli::Cli;
use crate::commands;
use crate::pipe_dispatch::EmulationRecovery;
use crate::session::{self, SessionStore};
pub async fn run_pipe(cli: &Cli) -> Result<(), crate::BoxError> {
let mut session = match open_session(cli).await {
Ok(session) => session,
Err(error) => return terminal_startup_error("pipe", &error),
};
let mut emulation_recovery =
EmulationRecovery::new(&session.client, &session.store, &cli.browser, &cli.page).await;
let stdin = BufReader::new(tokio::io::stdin());
let mut lines = stdin.lines();
let mut history: Vec<crate::macros_record::Observed> = Vec::new();
let processing: Result<(), crate::BoxError> = async {
loop {
let line = match lines.next_line().await {
Ok(Some(line)) => line,
Ok(None) => break,
Err(error) => return Err(format!("Failed to read pipe input: {error}").into()),
};
let line = line.trim().to_string();
if line.is_empty() {
continue;
}
let mut cmd: Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(e) => {
emit(&json!({"ok": false, "error": format!("Invalid JSON: {e}")}))?;
continue;
}
};
let record_path = match take_record_path(&mut cmd) {
Ok(path) => path,
Err(e) => {
emit(&json!({"ok": false, "error": e}))?;
continue;
}
};
if let Some(ref path) = record_path
&& let Err(e) = commands::record::start_recording(path)
{
emit(
&json!({"ok": false, "error": format!("{e}"), "hint": "Check the --record path's directory exists and is writable."}),
)?;
continue;
}
if cmd.get("cmd").and_then(Value::as_str) == Some("macro") {
let answer = crate::macros_cmd::dispatch_pipe(&cmd, &history)
.unwrap_or_else(|e| json!({"ok": false, "error": e.to_string()}));
emit(&answer)?;
continue;
}
let mut response = dispatch_on(&mut session, cli, &cmd, &mut emulation_recovery).await;
if let Some(ref path) = record_path
&& let Err(e) = commands::record::log_entry(path, &cmd, &response)
{
response["recording_error"] = json!(format!("{e}"));
}
let snapshot = session
.store
.browsers
.get(&cli.browser)
.and_then(|b| b.pages.get(&cli.page))
.and_then(|p| p.last_snapshot.clone());
history.push(crate::macros_record::Observed::read_with_snapshot(
&cmd,
&response,
snapshot.as_deref(),
));
emit(&response)?;
}
Ok(())
}
.await;
finalize_session(&mut session.store, processing, "pipe")
}
pub struct Session {
pub store: SessionStore,
pub browser_client: CdpClient,
pub client: CdpClient,
pub target_id: String,
pub policy: crate::run_helpers::ReportPolicy,
}
pub async fn open_session(cli: &Cli) -> Result<Session, crate::BoxError> {
let mut store = session::load_session()?;
let want_headless = !cli.headed;
let requested_proxy =
browser::normalized_proxy_option(cli.connect.as_deref(), cli.proxy_server.as_deref())?;
let requested_chrome_args =
browser::normalized_chrome_args_option(cli.connect.as_deref(), &cli.chrome_args)?;
let effective_proxy = requested_proxy.or_else(|| {
store
.browsers
.get(&cli.browser)
.and_then(|b| b.proxy_server.clone())
});
let effective_chrome_args =
crate::chrome_args::effective_chrome_args(&store, &cli.browser, &requested_chrome_args);
let (conn, browser_client) = connect_browser(
&mut store,
cli,
want_headless,
effective_proxy.clone(),
effective_chrome_args.clone(),
)
.await?;
let http_endpoint = conn
.http_endpoint
.as_deref()
.ok_or("No HTTP endpoint available. Cannot resolve page WebSocket URL.")?;
let target_id = {
let browser_session = session::ensure_browser(
&mut store,
&cli.browser,
&conn.ws_endpoint,
conn.pid,
want_headless,
effective_proxy,
effective_chrome_args,
);
crate::run_helpers::resolve_page_target(&browser_client, browser_session, &cli.page).await?
};
session::save_session(&mut store)?;
let page_ws = browser::get_page_ws_url(http_endpoint, &target_id).await?;
let client = CdpClient::connect(&page_ws).await?;
client.set_call_timeout(std::time::Duration::from_secs(cli.timeout));
client.enable("Page").await?;
commands::console::inject(&client).await;
if cli.stealth {
crate::setup::apply_stealth(&client).await;
} else {
client.enable("Runtime").await?;
}
let dialog_policy = crate::setup::DialogPolicy::parse(&cli.dialog)?;
client.spawn_dialog_handler(dialog_policy, cli.dialog_text.clone());
let policy = report_policy(cli)?;
Ok(Session {
store,
browser_client,
client,
target_id,
policy,
})
}
pub async fn dispatch_on(
session: &mut Session,
cli: &Cli,
cmd: &Value,
emulation_recovery: &mut EmulationRecovery,
) -> Value {
if let Some(response) = emulation_recovery.refusal_for(cmd) {
return response;
}
let mut ctx = crate::page_ctx::PageCtx {
client: &session.client,
browser_client: &session.browser_client,
store: &mut session.store,
browser: &cli.browser,
page: &cli.page,
target_id: &session.target_id,
timeout: cli.timeout,
max_depth: cli.max_depth,
report: session.policy,
};
let response = crate::pipe_dispatch::dispatch_single(&mut ctx, cmd, emulation_recovery).await;
emulation_recovery.update_after(cmd, &response);
response
}
pub async fn run_replay(
cli: &Cli,
file: &str,
vars: Option<&[String]>,
) -> Result<(), crate::BoxError> {
let content = std::fs::read_to_string(file)
.map_err(|e| format!("Cannot read replay file '{file}': {e}"))?;
let replacements: Vec<(&str, &str)> = vars
.unwrap_or(&[])
.iter()
.filter_map(|pair| pair.split_once('='))
.collect();
let mut session = match open_session(cli).await {
Ok(session) => session,
Err(error) => return terminal_startup_error("replay", &error),
};
let mut emulation_recovery =
EmulationRecovery::new(&session.client, &session.store, &cli.browser, &cli.page).await;
let processing: Result<(), crate::BoxError> = async {
for line in content.lines() {
let line = line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
let mut resolved = line.to_string();
for (key, val) in &replacements {
resolved = resolved.replace(&format!("{{{{{key}}}}}"), val);
}
let parsed: Value = serde_json::from_str(&resolved)
.map_err(|e| format!("Invalid JSON in replay: {e}"))?;
let mut cmd = if parsed.get("cmd").is_some_and(Value::is_object)
&& parsed.get("response").is_some()
{
parsed.get("cmd").cloned().unwrap_or_default()
} else {
parsed
};
let _ = take_record_path(&mut cmd);
let response = dispatch_on(&mut session, cli, &cmd, &mut emulation_recovery).await;
emit(&response)?;
}
Ok(())
}
.await;
finalize_session(&mut session.store, processing, "replay")
}
fn take_record_path(cmd: &mut Value) -> Result<Option<String>, String> {
let Some(value) = cmd.as_object_mut().and_then(|map| map.remove("_record")) else {
return Ok(None);
};
match value {
Value::String(path) => Ok(Some(path)),
Value::Null => Ok(None),
other => Err(format!("\"_record\" must be a file path, got {other}")),
}
}
fn report_policy(cli: &Cli) -> Result<crate::run_helpers::ReportPolicy, crate::BoxError> {
Ok(crate::run_helpers::ReportPolicy {
changes: cli.verdict == "auto",
budget: cli.budget,
on_intercept: crate::hit_test::OnIntercept::parse(&cli.on_intercept)?,
})
}
fn emit(value: &Value) -> Result<(), crate::BoxError> {
let line = serde_json::to_string(value)?;
let stdout = std::io::stdout();
let mut handle = stdout.lock();
writeln!(handle, "{line}")?;
handle.flush()?;
Ok(())
}
fn terminal_error(phase: &str, error: &str) -> Value {
json!({"ok": false, "terminal": true, "phase": phase, "error": error})
}
fn terminal_startup_error(phase: &str, error: &crate::BoxError) -> Result<(), crate::BoxError> {
let message = format!("Failed to start {phase} session: {error}");
let delivery = emit(&terminal_error("startup", &message));
match delivery {
Ok(()) => Err(message.into()),
Err(delivery_error) => Err(format!(
"{message}; also failed to deliver terminal error: {delivery_error}"
)
.into()),
}
}
fn finalize_session(
store: &mut SessionStore,
processing: Result<(), crate::BoxError>,
phase: &str,
) -> Result<(), crate::BoxError> {
let saved = session::save_session(store);
let message = match processing {
Ok(()) => match saved {
Ok(()) => return Ok(()),
Err(error) => format!("Failed to persist {phase} session at end of input: {error}"),
},
Err(error) => match saved {
Ok(()) => format!("{phase} session ended early: {error}"),
Err(save_error) => format!(
"{phase} session ended early: {error}; also failed to persist it: {save_error}"
),
},
};
let delivery = emit(&terminal_error("finalize", &message));
match delivery {
Ok(()) => Err(message.into()),
Err(delivery_error) => Err(format!(
"{message}; also failed to deliver terminal error: {delivery_error}"
)
.into()),
}
}
async fn connect_browser(
store: &mut SessionStore,
cli: &Cli,
want_headless: bool,
effective_proxy: Option<String>,
effective_chrome_args: Vec<String>,
) -> Result<(browser::BrowserConnection, CdpClient), crate::BoxError> {
if let Some(existing) = store.browsers.get(&cli.browser) {
let mode_matches = existing.headless == want_headless;
let ws = &existing.ws_endpoint;
let http = browser::extract_http_from_ws(ws);
if mode_matches {
if let Ok(client) = CdpClient::connect(ws).await {
session::ensure_proxy_compatible(existing, effective_proxy.as_deref())?;
session::ensure_chrome_args_compatible(existing, &effective_chrome_args)?;
let conn = browser::BrowserConnection {
ws_endpoint: ws.clone(),
http_endpoint: Some(http),
pid: existing.pid,
};
client.set_call_timeout(std::time::Duration::from_secs(cli.timeout));
return Ok((conn, client));
}
} else if let Some(pid) = existing.pid {
browser::kill_and_await_exit(&cli.browser, pid)?;
}
store.browsers.remove(&cli.browser);
}
let opts = BrowserOptions {
name: cli.browser.clone(),
headless: want_headless,
ignore_https_errors: cli.ignore_https_errors,
stealth: cli.stealth,
connect: cli.connect.clone(),
proxy_server: effective_proxy,
copy_cookies: cli.copy_cookies,
chrome_args: effective_chrome_args,
};
let conn = browser::resolve_browser(&opts).await?;
let client = CdpClient::connect(&conn.ws_endpoint).await?;
client.set_call_timeout(std::time::Duration::from_secs(cli.timeout));
Ok((conn, client))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn terminal_failures_are_machine_readable_failures() {
let response = terminal_error("finalize", "session store is read-only");
assert_eq!(response["ok"], false);
assert_eq!(response["terminal"], true);
assert_eq!(response["phase"], "finalize");
assert_eq!(response["error"], "session store is read-only");
}
}