diff --git a/README.md b/README.md index 004128a..aa8b8de 100644 --- a/README.md +++ b/README.md @@ -153,15 +153,15 @@ enforces several configurable limits: Increase `--max-input-size` if you work with captures larger than 2 GiB (e.g. `--max-input-size 8GiB`). The other limits rarely need tuning. -Mitmproxy flow files are now processed incrementally — memory usage stays bounded -regardless of input size. HAR files are still loaded in full (streaming HAR is planned). +Both mitmproxy flow files and HAR files are processed incrementally — memory usage +stays bounded regardless of input size. ## Supported Formats | Format | Versions | Extension | |--------|----------|-----------| | mitmproxy flow dumps | v19, v20, v21 | `.flow` | -| HAR (HTTP Archive) | 1.2 | `.har` | +| HAR (HTTP Archive) | 1.2 (incrementally parsed) | `.har` | Format is auto-detected from file content. Use `--format` to override. diff --git a/src/har_reader.rs b/src/har_reader.rs index d3ab6a6..6d516fd 100644 --- a/src/har_reader.rs +++ b/src/har_reader.rs @@ -1,33 +1,61 @@ +use std::fs::File; +use std::io::{BufRead, BufReader, Read}; use std::path::Path; +use base64::Engine; +use serde::Deserialize; use tracing::{debug, warn}; use crate::error::{Error, Result}; use crate::types::CapturedRequest; -use base64::Engine; -fn strip_bom(input: &[u8]) -> &[u8] { - if input.starts_with(&[0xEF, 0xBB, 0xBF]) { - &input[3..] - } else { - input - } +#[derive(Deserialize)] +struct StreamingHarEntry { + request: StreamingHarRequest, + response: StreamingHarResponse, } -fn convert_headers(headers: &[har::v1_2::Headers]) -> Vec<(String, String)> { - headers - .iter() - .map(|h| (h.name.clone(), h.value.clone())) - .collect() +#[derive(Deserialize)] +struct StreamingHarRequest { + method: String, + url: String, + #[serde(default)] + headers: Vec, + #[serde(rename = "postData", default)] + post_data: Option, } -fn decode_body(content: &har::v1_2::Content) -> Option> { - let text = content.text.as_deref()?; - if content.encoding.as_deref() == Some("base64") { - base64::engine::general_purpose::STANDARD.decode(text).ok() - } else { - Some(text.as_bytes().to_vec()) - } +#[derive(Deserialize)] +struct StreamingHarResponse { + status: i64, + #[serde(rename = "statusText", default)] + status_text: String, + #[serde(default)] + headers: Vec, + #[serde(default)] + content: StreamingHarContent, +} + +#[derive(Deserialize)] +struct StreamingHarHeader { + name: String, + value: String, +} + +#[derive(Deserialize, Default)] +struct StreamingHarPostData { + #[serde(default)] + text: Option, +} + +#[derive(Deserialize, Default)] +struct StreamingHarContent { + #[serde(default)] + text: Option, + #[serde(rename = "mimeType", default)] + mime_type: Option, + #[serde(default)] + encoding: Option, } pub struct HarFlowWrapper { @@ -43,32 +71,49 @@ pub struct HarFlowWrapper { } impl HarFlowWrapper { - fn from_entry(entry: &har::v1_2::Entries) -> Self { - let req = &entry.request; - let resp = &entry.response; - - let request_body = req + fn from_streaming_entry(entry: StreamingHarEntry) -> Self { + let request_body = entry + .request .post_data - .as_ref() - .and_then(|pd| pd.text.as_deref()) - .map(|t| t.as_bytes().to_vec()); + .and_then(|pd| pd.text) + .map(String::into_bytes); - let response_content_type = resp.content.mime_type.clone(); + let response_content_type = entry.response.content.mime_type.clone(); + let response_body = decode_streaming_body(&entry.response.content); Self { - url: req.url.clone(), - method: req.method.clone(), - request_headers: convert_headers(&req.headers), + url: entry.request.url, + method: entry.request.method, + request_headers: entry + .request + .headers + .into_iter() + .map(|h| (h.name, h.value)) + .collect(), request_body, - response_status: resp.status as u16, - response_reason: resp.status_text.clone(), - response_headers: convert_headers(&resp.headers), - response_body: decode_body(&resp.content), + response_status: entry.response.status as u16, + response_reason: entry.response.status_text, + response_headers: entry + .response + .headers + .into_iter() + .map(|h| (h.name, h.value)) + .collect(), + response_body, response_content_type, } } } +fn decode_streaming_body(content: &StreamingHarContent) -> Option> { + let text = content.text.as_deref()?; + if content.encoding.as_deref() == Some("base64") { + base64::engine::general_purpose::STANDARD.decode(text).ok() + } else { + Some(text.as_bytes().to_vec()) + } +} + impl CapturedRequest for HarFlowWrapper { fn get_url(&self) -> &str { &self.url @@ -107,61 +152,266 @@ impl CapturedRequest for HarFlowWrapper { } } -fn parse_har_bytes(bytes: &[u8]) -> Result>> { - let clean = strip_bom(bytes); - let har_doc = har::from_slice(clean).map_err(|e| Error::HarParse(e.to_string()))?; +fn read_byte(reader: &mut impl Read) -> Result> { + let mut buf = [0u8; 1]; + match reader.read(&mut buf)? { + 0 => Ok(None), + _ => Ok(Some(buf[0])), + } +} - let entries = match &har_doc.log { - har::Spec::V1_2(log) => &log.entries, - har::Spec::V1_3(_) => { - return Err(Error::HarParse("HAR v1.3 not yet supported".into())); +fn skip_ws_byte(reader: &mut impl Read) -> Result> { + loop { + match read_byte(reader)? { + None => return Ok(None), + Some(b) if b.is_ascii_whitespace() => continue, + Some(b) => return Ok(Some(b)), } - }; + } +} - let requests: Vec> = entries - .iter() - .map(|e| Box::new(HarFlowWrapper::from_entry(e)) as Box) - .collect(); +fn strip_bom_from_reader(reader: &mut BufReader) -> Result<()> { + let buf = reader.fill_buf()?; + if buf.starts_with(&[0xEF, 0xBB, 0xBF]) { + reader.consume(3); + } + Ok(()) +} - Ok(requests) +/// Scan forward through JSON until positioned just past `"entries": [`. +/// Tracks string boundaries so `"entries"` inside a value is not mistaken for the key. +fn find_entries_array_start(reader: &mut impl Read) -> Result<()> { + let target = b"\"entries\""; + let mut in_string = false; + let mut escape_next = false; + let mut match_pos: usize = 0; + + loop { + let byte = read_byte(reader)? + .ok_or_else(|| Error::HarParse("unexpected EOF: entries array not found".into()))?; + + if escape_next { + escape_next = false; + match_pos = 0; + continue; + } + + if byte == b'\\' && in_string { + escape_next = true; + match_pos = 0; + continue; + } + + if byte == b'"' { + in_string = !in_string; + } + + if byte == target[match_pos] { + match_pos += 1; + if match_pos == target.len() { + let colon = skip_ws_byte(reader)? + .ok_or_else(|| Error::HarParse("unexpected EOF after entries key".into()))?; + if colon == b':' { + let bracket = skip_ws_byte(reader)?.ok_or_else(|| { + Error::HarParse("unexpected EOF expecting entries array".into()) + })?; + if bracket == b'[' { + return Ok(()); + } + } + match_pos = 0; + } + } else if byte == target[0] { + match_pos = 1; + } else { + match_pos = 0; + } + } } -pub fn read_har_file(path: &Path) -> Result>> { - if path.is_dir() { - let mut all = Vec::new(); - let mut entries: Vec<_> = std::fs::read_dir(path)? - .filter_map(|e| e.ok()) - .filter(|e| { - e.path() - .extension() - .is_some_and(|ext| ext.eq_ignore_ascii_case("har")) - }) - .collect(); - entries.sort_by_key(|e| e.path()); - for entry in entries { - match std::fs::read(entry.path()) { - Ok(bytes) => match parse_har_bytes(&bytes) { - Ok(requests) => { - debug!(path = %entry.path().display(), count = requests.len(), "Parsed HAR file"); - all.extend(requests); +/// Read a balanced JSON object after the opening `{` has been consumed. +/// Tracks nesting depth across braces/brackets and handles string escapes. +fn read_json_object(reader: &mut impl Read) -> Result> { + let mut buf = Vec::with_capacity(4096); + buf.push(b'{'); + let mut depth: i32 = 1; + let mut in_string = false; + let mut escape_next = false; + + loop { + let byte = read_byte(reader)? + .ok_or_else(|| Error::HarParse("unexpected EOF inside entry object".into()))?; + buf.push(byte); + + if escape_next { + escape_next = false; + continue; + } + + if in_string { + match byte { + b'\\' => escape_next = true, + b'"' => in_string = false, + _ => {} + } + continue; + } + + match byte { + b'"' => in_string = true, + b'{' | b'[' => depth += 1, + b'}' | b']' => { + depth -= 1; + if depth == 0 { + return Ok(buf); + } + } + _ => {} + } + } +} + +pub struct HarStreamIter { + reader: BufReader, + done: bool, + entry_index: usize, +} + +impl HarStreamIter { + fn new(path: &Path) -> Result { + let file = File::open(path)?; + let mut reader = BufReader::with_capacity(64 * 1024, file); + + strip_bom_from_reader(&mut reader)?; + find_entries_array_start(&mut reader)?; + + Ok(Self { + reader, + done: false, + entry_index: 0, + }) + } +} + +impl Iterator for HarStreamIter { + type Item = Result>; + + fn next(&mut self) -> Option { + if self.done { + return None; + } + + let byte = match skip_ws_byte(&mut self.reader) { + Ok(Some(b)) => b, + Ok(None) => { + self.done = true; + return None; + } + Err(e) => { + self.done = true; + return Some(Err(e)); + } + }; + + if byte == b']' { + self.done = true; + return None; + } + + let byte = if byte == b',' { + match skip_ws_byte(&mut self.reader) { + Ok(Some(b)) => b, + Ok(None) => { + self.done = true; + return None; + } + Err(e) => { + self.done = true; + return Some(Err(e)); + } + } + } else { + byte + }; + + if byte == b']' { + self.done = true; + return None; + } + + if byte != b'{' { + self.done = true; + return Some(Err(Error::HarParse(format!( + "expected '{{' at start of entry {}, got '{}'", + self.entry_index, byte as char + )))); + } + + match read_json_object(&mut self.reader) { + Ok(buf) => { + let idx = self.entry_index; + self.entry_index += 1; + match serde_json::from_slice::(&buf) { + Ok(entry) => { + let wrapper = HarFlowWrapper::from_streaming_entry(entry); + Some(Ok(Box::new(wrapper) as Box)) } Err(e) => { - warn!(path = %entry.path().display(), error = %e, "Skipping unparseable HAR file"); + warn!(entry = idx, error = %e, "Failed to parse HAR entry"); + Some(Err(Error::HarParse(format!("entry {idx}: {e}")))) } - }, - Err(e) => { - warn!(path = %entry.path().display(), error = %e, "Skipping unreadable HAR file"); } } + Err(e) => { + self.done = true; + Some(Err(e)) + } } - Ok(all) - } else { - let bytes = std::fs::read(path)?; - debug!(path = %path.display(), "Parsing HAR file"); - parse_har_bytes(&bytes) } } +type RequestIter = Box>>>; + +pub fn stream_har_file(path: &Path) -> Result { + if path.is_dir() { + return stream_har_dir(path); + } + debug!(path = %path.display(), "Streaming HAR file"); + let iter = HarStreamIter::new(path)?; + Ok(Box::new(iter)) +} + +fn stream_har_dir(path: &Path) -> Result { + let mut dir_entries: Vec<_> = std::fs::read_dir(path)? + .filter_map(|e| e.ok()) + .filter(|e| { + e.path() + .extension() + .is_some_and(|ext| ext.eq_ignore_ascii_case("har")) + }) + .collect(); + dir_entries.sort_by_key(|e| e.path()); + + let iter = dir_entries + .into_iter() + .flat_map(|entry| match HarStreamIter::new(&entry.path()) { + Ok(it) => { + debug!(path = %entry.path().display(), "Streaming HAR file from directory"); + Box::new(it) as Box>>> + } + Err(e) => { + warn!(path = %entry.path().display(), error = %e, "Skipping unparseable HAR file"); + Box::new(std::iter::empty()) + } + }); + + Ok(Box::new(iter)) +} + +pub fn read_har_file(path: &Path) -> Result>> { + stream_har_file(path)?.collect() +} + pub fn har_heuristic(path: &Path) -> bool { if path.is_dir() { return false; @@ -169,14 +419,18 @@ pub fn har_heuristic(path: &Path) -> bool { let Ok(file) = std::fs::File::open(path) else { return false; }; - use std::io::Read; + use std::io::Read as _; let mut buf = [0u8; 4096]; let mut reader = std::io::BufReader::new(file); let n = match reader.read(&mut buf) { Ok(n) => n, Err(_) => return false, }; - let clean = strip_bom(&buf[..n]); + let clean = if buf[..n].starts_with(&[0xEF, 0xBB, 0xBF]) { + &buf[3..n] + } else { + &buf[..n] + }; clean .iter() .find(|b| !b.is_ascii_whitespace()) @@ -304,11 +558,134 @@ mod tests { let mut with_bom = vec![0xEF, 0xBB, 0xBF]; with_bom.extend_from_slice(&original); - let requests = parse_har_bytes(&with_bom).unwrap(); + let tmp = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(tmp.path(), &with_bom).unwrap(); + + let requests = read_har_file(tmp.path()).unwrap(); assert_eq!(requests.len(), 1); assert_eq!( requests[0].get_url(), "https://api.example.com/api/v1/users" ); } + + #[test] + fn stream_har_does_not_materialize_all() { + use std::io::Write; + + let mut har = String::from( + r#"{"log":{"version":"1.2","creator":{"name":"test","version":"1.0"},"entries":["#, + ); + let entry_count = 20; + for i in 0..entry_count { + if i > 0 { + har.push(','); + } + har.push_str(&format!( + r#"{{"request":{{"method":"GET","url":"https://example.com/api/item/{i}","headers":[]}},"response":{{"status":200,"statusText":"OK","headers":[],"content":{{"text":"{{\"id\":{i}}}"}}}}}} +"#, + )); + } + har.push_str("]}}\n"); + + let tmp = tempfile::NamedTempFile::new().unwrap(); + tmp.as_file().write_all(har.as_bytes()).unwrap(); + tmp.as_file().sync_all().unwrap(); + + let mut iter = stream_har_file(tmp.path()).unwrap(); + + let first = iter.next().unwrap().unwrap(); + assert_eq!(first.get_url(), "https://example.com/api/item/0"); + assert_eq!(first.get_method(), "GET"); + + let rest: Vec<_> = iter.collect::>>().unwrap(); + assert_eq!(rest.len(), entry_count - 1); + } + + #[test] + fn stream_matches_read_for_fixtures() { + for name in &["simple.har", "multi.har", "base64_body.har"] { + let path = fixture(name); + let collected: Vec<_> = stream_har_file(&path) + .unwrap() + .collect::>>() + .unwrap(); + + let direct = read_har_file(&path).unwrap(); + + assert_eq!( + collected.len(), + direct.len(), + "entry count mismatch for {name}" + ); + for (i, (a, b)) in collected.iter().zip(direct.iter()).enumerate() { + assert_eq!(a.get_url(), b.get_url(), "{name} entry {i} url"); + assert_eq!(a.get_method(), b.get_method(), "{name} entry {i} method"); + assert_eq!( + a.get_response_status_code(), + b.get_response_status_code(), + "{name} entry {i} status" + ); + assert_eq!( + a.get_response_body(), + b.get_response_body(), + "{name} entry {i} body" + ); + assert_eq!( + a.get_request_body(), + b.get_request_body(), + "{name} entry {i} req body" + ); + assert_eq!( + a.get_response_content_type(), + b.get_response_content_type(), + "{name} entry {i} content-type" + ); + } + } + } + + #[test] + fn stream_malformed_entry_returns_error() { + use std::io::Write; + + let har = r#"{"log":{"version":"1.2","entries":[ + {"request":{"method":"GET","url":"https://ok.example.com","headers":[]},"response":{"status":200,"statusText":"OK","headers":[],"content":{}}}, + {"INVALID JSON STRUCTURE": true}, + {"request":{"method":"POST","url":"https://also-ok.example.com","headers":[]},"response":{"status":201,"statusText":"Created","headers":[],"content":{}}} + ]}}"#; + + let tmp = tempfile::NamedTempFile::new().unwrap(); + tmp.as_file().write_all(har.as_bytes()).unwrap(); + tmp.as_file().sync_all().unwrap(); + + let results: Vec<_> = stream_har_file(tmp.path()).unwrap().collect(); + + assert!(results[0].is_ok()); + assert_eq!( + results[0].as_ref().unwrap().get_url(), + "https://ok.example.com" + ); + + assert!(results[1].is_err()); + + assert!(results[2].is_ok()); + assert_eq!( + results[2].as_ref().unwrap().get_url(), + "https://also-ok.example.com" + ); + } + + #[test] + fn stream_empty_entries_array() { + use std::io::Write; + + let har = r#"{"log":{"version":"1.2","entries":[]}}"#; + let tmp = tempfile::NamedTempFile::new().unwrap(); + tmp.as_file().write_all(har.as_bytes()).unwrap(); + tmp.as_file().sync_all().unwrap(); + + let results: Vec<_> = stream_har_file(tmp.path()).unwrap().collect(); + assert!(results.is_empty()); + } } diff --git a/src/main.rs b/src/main.rs index 0b94374..4a3c740 100644 --- a/src/main.rs +++ b/src/main.rs @@ -187,25 +187,23 @@ fn stream_input( } } InputFormat::Har => { - debug!(path = %path.display(), "Reading as HAR format"); - // TODO: HAR streaming is planned for a future release. - let requests = har_reader::read_har_file(path).context("failed to read HAR file")?; - Ok(Box::new(requests.into_iter().map(Ok))) + debug!(path = %path.display(), "Streaming as HAR format"); + let iter = har_reader::stream_har_file(path).context("failed to stream HAR file")?; + Ok(Box::new(iter)) } InputFormat::Auto => { if path.is_dir() { debug!(path = %path.display(), "Auto-detecting format for directory"); let mitmproxy_result = mitmproxy_reader::stream_mitmproxy_dir(path); - let har_result = har_reader::read_har_file(path); + let har_result = har_reader::stream_har_file(path); match (mitmproxy_result, har_result) { - (Ok(m_iter), Ok(h_vec)) => { - let combined: RequestIter = - Box::new(m_iter.chain(h_vec.into_iter().map(Ok))); + (Ok(m_iter), Ok(h_iter)) => { + let combined: RequestIter = Box::new(m_iter.chain(h_iter)); Ok(combined) } (Ok(m_iter), Err(_)) => Ok(m_iter), - (Err(_), Ok(h_vec)) => Ok(Box::new(h_vec.into_iter().map(Ok))), + (Err(_), Ok(h_iter)) => Ok(h_iter), (Err(e1), Err(_e2)) => { Err(e1).context("failed to read directory as mitmproxy or HAR") } @@ -226,17 +224,17 @@ fn stream_input( Ok(Box::new(iter)) } else if hs > ms { info!(path = %path.display(), "Auto-detected as HAR format"); - let requests = har_reader::read_har_file(path) - .context("detected as HAR format but failed to parse")?; - Ok(Box::new(requests.into_iter().map(Ok))) + let iter = har_reader::stream_har_file(path) + .context("detected as HAR format but failed to stream")?; + Ok(Box::new(iter)) } else if ms > 0 { warn!(path = %path.display(), "Ambiguous format detection, trying mitmproxy first"); match mitmproxy_reader::stream_mitmproxy_file(path) { Ok(iter) => Ok(Box::new(iter)), Err(_) => { - let requests = har_reader::read_har_file(path) + let iter = har_reader::stream_har_file(path) .context("failed to parse as either mitmproxy or HAR")?; - Ok(Box::new(requests.into_iter().map(Ok))) + Ok(Box::new(iter)) } } } else {