use super::*;
pub(crate) fn scan_one_file_multi(
path: &Path,
targets: &[String],
ac: &aho_corasick::AhoCorasick,
always: &[usize],
) -> Result<Vec<(usize, ScanResult)>> {
let session_id = crate::subagent::session_id_from_path(path);
let is_subagent = crate::subagent::is_subagent_path(path);
let parent_session_id =
crate::subagent::parent_session_id_from_path(path).unwrap_or_else(|| session_id.clone());
let Some(mmap) = mmap_bytes(path)? else {
return Ok(Vec::new());
};
let bytes: &[u8] = &mmap;
if always.is_empty() && !ac.is_match(bytes) {
return Ok(Vec::new());
}
let mut present: Vec<usize> = ac
.find_overlapping_iter(bytes)
.map(|m| m.pattern().as_usize())
.collect();
present.extend_from_slice(always); present.sort_unstable();
present.dedup();
let (records, skipped) =
crate::parse::parse_candidates_parallel(bytes, line_is_recover_candidate);
let recs: Vec<&Record> = records.iter().map(|(_, r)| r).collect();
let turns = group_turn_indices_deduped(&recs, |r| *r);
let opaque = collect_opaque_commands(&session_id, &records, &turns);
let mut out = Vec::new();
for ti in present {
let events = extract_with_turns(&records, &turns, Some(&targets[ti]));
if !events.is_empty() {
out.push((
ti,
ScanResult {
session_id: session_id.clone(),
is_subagent,
parent_session_id: parent_session_id.clone(),
events,
opaque: opaque.clone(),
merged_line_origin: std::collections::BTreeMap::new(),
skipped_lines: skipped,
},
));
}
}
Ok(out)
}
pub(crate) struct BestReconstruction {
pub(crate) content: String,
pub(crate) known: usize,
pub(crate) total: usize,
pub(crate) boundaries: usize,
pub(crate) bash_file: usize,
pub(crate) bash_opaque: usize,
}
pub(crate) fn reconstruct_best(
scans: Vec<ScanResult>,
when: &str,
) -> Result<Option<BestReconstruction>> {
let merged = merge_groups_for_reconstruction(scans);
let mut best: Option<(BestReconstruction, Option<String>)> = None;
for s in &merged {
if s.events.is_empty() {
continue;
}
let cutoff = resolve_cutoff(when, &s.events)?;
let rep = replay(&s.events, cutoff);
let known = rep.final_buffer.known_lines();
if known.is_empty() {
continue;
}
let total = rep
.final_buffer
.seen_total_lines
.unwrap_or_else(|| known.last().map(|(n, _)| *n).unwrap_or(0));
let latest_ts = s
.events
.iter()
.filter_map(|e| e.timestamp_utc.clone())
.max();
let mut content = known
.iter()
.map(|(_, t)| t.as_str())
.collect::<Vec<_>>()
.join("\n");
if !content.is_empty() {
content.push('\n');
}
let fresher = match &best {
None => true,
Some((b, bts)) => (&latest_ts, known.len()) > (bts, b.known),
};
if fresher {
best = Some((
BestReconstruction {
content,
known: known.len(),
total,
boundaries: rep.boundaries.len(),
bash_file: rep.counts.bash,
bash_opaque: s.opaque.len(),
},
latest_ts,
));
}
}
Ok(best.map(|(b, _)| b))
}
pub(crate) fn write_batch_report(out_dir: &Path, outcomes: &[BatchOutcome]) -> Result<()> {
std::fs::create_dir_all(out_dir)
.with_context(|| format!("cannot create {}", out_dir.display()))?;
let mut body = String::from(
"status\tknown_lines\ttotal_lines\tboundaries\tbash_file\tbash_opaque\ttarget\twritten_to\n",
);
let (mut complete, mut partial, mut none, mut skipped) = (0usize, 0usize, 0usize, 0usize);
let mut flagged = 0usize;
for o in outcomes {
match o.status {
"complete" => complete += 1,
"partial" => partial += 1,
"no-history" => none += 1,
"skipped-exists" => skipped += 1,
_ => {}
}
if o.boundaries + o.bash_file + o.bash_opaque > 0 {
flagged += 1;
}
body.push_str(&format!(
"{}\t{}\t{}\t{}\t{}\t{}\t{}\t{}\n",
o.status,
o.known,
o.total,
o.boundaries,
o.bash_file,
o.bash_opaque,
o.target,
o.written
.as_ref()
.map(|p| p.display().to_string())
.unwrap_or_default()
));
}
let report = out_dir.join("recovery-report.tsv");
std::fs::write(&report, &body).with_context(|| format!("cannot write {}", report.display()))?;
println!(
"recovered {} file(s): {complete} complete · {partial} partial · {none} no-history · \
{skipped} skipped (already present)",
outcomes.len()
);
if flagged > 0 {
println!(
"{flagged} file(s) had integrity boundaries or unparsed mutating bash in their \
window (boundaries/bash_file/bash_opaque columns): complete means complete from \
the tool stream, not verified against disk"
);
}
println!("report: {}", report.display());
Ok(())
}
pub(crate) fn window_admits(
turn_index: usize,
ts: Option<&str>,
turn_range: Option<(usize, usize)>,
time_window: &TimeWindow,
) -> bool {
if let Some((lo, hi)) = turn_range {
if turn_index < lo || turn_index > hi {
return false;
}
}
time_window.contains(ts)
}
pub(crate) fn scan_one_file(path: &Path, target_file: Option<&str>) -> Result<ScanResult> {
let session_id = crate::subagent::session_id_from_path(path);
let is_subagent = crate::subagent::is_subagent_path(path);
let parent_session_id =
crate::subagent::parent_session_id_from_path(path).unwrap_or_else(|| session_id.clone());
let Some(mmap) = mmap_bytes(path)? else {
return Ok(ScanResult {
session_id,
is_subagent,
parent_session_id,
events: Vec::new(),
opaque: Vec::new(),
merged_line_origin: std::collections::BTreeMap::new(),
skipped_lines: 0,
});
};
let bytes: &[u8] = &mmap;
if let Some(t) = target_file {
let base = basename_of(t);
if raw_needle_safe(base) && memmem::find(bytes, base.as_bytes()).is_none() {
return Ok(ScanResult {
session_id,
is_subagent,
parent_session_id,
events: Vec::new(),
opaque: Vec::new(),
merged_line_origin: std::collections::BTreeMap::new(),
skipped_lines: 0,
});
}
}
let (records, skipped) =
crate::parse::parse_candidates_parallel(bytes, line_is_recover_candidate);
let recs: Vec<&Record> = records.iter().map(|(_, r)| r).collect();
let turns = group_turn_indices_deduped(&recs, |r| *r);
let events = extract_with_turns(&records, &turns, target_file);
let opaque = collect_opaque_commands(&session_id, &records, &turns);
Ok(ScanResult {
session_id,
is_subagent,
parent_session_id,
events,
opaque,
merged_line_origin: std::collections::BTreeMap::new(),
skipped_lines: skipped,
})
}
pub(crate) fn line_is_recover_candidate(line: &[u8]) -> bool {
static NEEDLES: std::sync::LazyLock<[memmem::Finder<'static>; 11]> =
std::sync::LazyLock::new(|| {
[
memmem::Finder::new(b"toolUseResult"),
memmem::Finder::new(b"Edit"),
memmem::Finder::new(b"Write"),
memmem::Finder::new(b"Read"),
memmem::Finder::new(b"Bash"),
memmem::Finder::new(b"PowerShell"),
memmem::Finder::new(b"filePath"),
memmem::Finder::new(b"file_path"),
memmem::Finder::new(b"file-history-snapshot"),
memmem::Finder::new(b"edited_text_file"),
memmem::Finder::new(b"tool_use_error"),
]
});
crate::parse::line_has_user_role_marker(line) || NEEDLES.iter().any(|f| f.find(line).is_some())
}
pub(crate) fn collect_opaque_commands(
session_id: &str,
records: &[(usize, Record)],
turns: &[Vec<usize>],
) -> Vec<OpaqueCommand> {
let mut out: Vec<OpaqueCommand> = Vec::new();
for (turn_index, idxs) in turns.iter().enumerate() {
for &i in idxs {
let (line_no, rec) = (records[i].0, &records[i].1);
let Some(blocks) = rec.blocks() else { continue };
for b in blocks {
let Block::ToolUse {
name: Some(name),
input: Some(input),
..
} = b
else {
continue;
};
let cmd = input.get("command").and_then(serde_json::Value::as_str);
if name == "Bash" {
let Some(cmd) = cmd else { continue };
for bm in crate::bash_mutations::parse_bash_mutations(cmd) {
if crate::bash_mutations::is_class_marker(&bm.path)
&& !matches!(bm.path.as_str(), "git:add" | "git:commit")
{
out.push(OpaqueCommand {
session_id: session_id.to_string(),
line_no,
turn_index,
timestamp_utc: rec.timestamp.clone(),
marker: bm.path,
});
}
}
} else if name == "PowerShell" && cmd.is_some() {
out.push(OpaqueCommand {
session_id: session_id.to_string(),
line_no,
turn_index,
timestamp_utc: rec.timestamp.clone(),
marker: "powershell".to_string(),
});
}
}
}
}
out
}
pub(crate) fn extract_with_turns(
records: &[(usize, Record)],
turns: &[Vec<usize>],
target_file: Option<&str>,
) -> Vec<FileEvent> {
let mut events: Vec<FileEvent> = Vec::new();
for (turn_index, idxs) in turns.iter().enumerate() {
let mut id_to_path: BTreeMap<String, String> = BTreeMap::new();
let mut ids_with_result: std::collections::HashSet<String> =
std::collections::HashSet::new();
let mut failed_ids: std::collections::HashSet<String> = std::collections::HashSet::new();
for &i in idxs {
let rec = &records[i].1;
collect_tool_use_paths(rec.blocks(), &mut id_to_path);
let has_structured_result = rec.tool_use_result.is_some();
if let Some(blocks) = rec.blocks() {
for b in blocks {
if let Block::ToolResult {
tool_use_id: Some(id),
is_error,
..
} = b
{
if has_structured_result {
ids_with_result.insert(id.clone());
}
if *is_error == Some(true) {
failed_ids.insert(id.clone());
}
}
}
}
}
for &i in idxs {
let (line_no, rec) = (&records[i].0, &records[i].1);
extract_from_record(
*line_no,
turn_index,
rec,
target_file,
&id_to_path,
&mut events,
);
extract_input_fallback(
*line_no,
turn_index,
rec,
target_file,
&ids_with_result,
&failed_ids,
&mut events,
);
}
}
events
}