use std::path::{Path, PathBuf};
use std::sync::mpsc;
use std::time::{Duration, Instant};
use crate::cdp::{CdpEvent, run_session};
use crate::http::captured::CapturedRow;
#[derive(Debug, Clone)]
pub struct Options {
pub workspace: PathBuf,
pub url: String,
pub max_seconds: Option<u64>,
pub idle_ms: u64,
pub verbose: bool,
}
impl Default for Options {
fn default() -> Self {
Self {
workspace: PathBuf::from("."),
url: String::new(),
max_seconds: None,
idle_ms: 2000,
verbose: true,
}
}
}
pub fn captured_log_path(workspace: &Path) -> PathBuf {
workspace.join(".rqst").join("captured").join("log.jsonl")
}
pub fn run(opts: Options) -> Result<usize, String> {
use std::fs;
use std::io::Write;
let log_path = captured_log_path(&opts.workspace);
if let Some(parent) = log_path.parent() {
fs::create_dir_all(parent).map_err(|e| format!("create log dir: {e}"))?;
}
let mut log = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&log_path)
.map_err(|e| format!("open {}: {e}", log_path.display()))?;
let profile_dir =
std::env::temp_dir().join(format!("mnml-proxy-profile-{}", std::process::id()));
fs::create_dir_all(&profile_dir).map_err(|e| format!("mkdir profile: {e}"))?;
let (event_tx, event_rx) = mpsc::channel::<CdpEvent>();
let (cmd_tx, cmd_rx) = mpsc::channel::<crate::cdp::CdpCommand>();
let url = opts.url.clone();
let pd = profile_dir.clone();
let session_handle = std::thread::spawn(move || {
run_session(&url, &pd, true, &event_tx, &cmd_rx);
});
let started = Instant::now();
let mut last_event = Instant::now();
let mut written: usize = 0;
loop {
match event_rx.recv_timeout(Duration::from_millis(100)) {
Ok(CdpEvent::Connected { ws_url }) => {
if opts.verbose {
eprintln!("mnml proxy: attached to {ws_url}");
}
}
Ok(CdpEvent::Message(v)) => {
last_event = Instant::now();
if let Some(method) = v.get("method").and_then(|m| m.as_str())
&& method == "Network.requestWillBeSent"
&& let Some(row) = decode_network_request(&v)
{
if opts.verbose {
eprintln!(" {} {}", row.method, row.url);
}
if let Ok(line) = serde_json::to_string(&row) {
let _ = writeln!(log, "{line}");
written += 1;
}
}
}
Ok(CdpEvent::Closed(reason)) => {
if opts.verbose {
eprintln!("mnml proxy: session closed — {reason}");
}
break;
}
Err(mpsc::RecvTimeoutError::Timeout) => { }
Err(mpsc::RecvTimeoutError::Disconnected) => break,
}
if let Some(cap) = opts.max_seconds
&& started.elapsed() >= Duration::from_secs(cap)
{
if opts.verbose {
eprintln!("mnml proxy: --seconds {cap} elapsed, stopping ({written} captured)");
}
break;
}
if written > 0 && last_event.elapsed() >= Duration::from_millis(opts.idle_ms) {
if opts.verbose {
eprintln!(
"mnml proxy: idle for {}ms, stopping ({written} captured)",
opts.idle_ms
);
}
break;
}
}
drop(cmd_tx); let _ = session_handle.join();
let _ = fs::remove_dir_all(&profile_dir);
Ok(written)
}
fn decode_network_request(v: &serde_json::Value) -> Option<CapturedRow> {
let p = v.get("params")?;
let request = p.get("request")?;
let request_id = p.get("requestId")?.as_str()?.to_string();
let method = request.get("method")?.as_str()?.to_string();
let url = request.get("url")?.as_str()?.to_string();
let headers: Vec<(String, String)> = request
.get("headers")
.and_then(|h| h.as_object())
.map(|obj| {
obj.iter()
.filter_map(|(k, v)| v.as_str().map(|s| (k.clone(), s.to_string())))
.collect()
})
.unwrap_or_default();
let body = request
.get("postData")
.and_then(|b| b.as_str())
.map(str::to_string);
let at = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
Some(CapturedRow {
at,
request_id,
method,
url,
headers,
body,
paused: false,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn decode_network_request_extracts_method_and_url() {
let v = serde_json::json!({
"method": "Network.requestWillBeSent",
"params": {
"requestId": "r1",
"request": {
"method": "POST",
"url": "https://api/x",
"headers": {"content-type": "application/json"},
"postData": "{\"k\":1}"
}
}
});
let row = decode_network_request(&v).unwrap();
assert_eq!(row.method, "POST");
assert_eq!(row.url, "https://api/x");
assert_eq!(row.body.as_deref(), Some("{\"k\":1}"));
assert!(row.headers.iter().any(|(k, _)| k == "content-type"));
}
#[test]
fn captured_log_path_under_workspace() {
let p = captured_log_path(Path::new("/tmp/x"));
assert!(p.ends_with(".rqst/captured/log.jsonl"));
}
#[test]
fn decode_returns_none_on_missing_fields() {
let v = serde_json::json!({"method": "Network.requestWillBeSent", "params": {}});
assert!(decode_network_request(&v).is_none());
}
}