use super::*;
#[allow(clippy::too_many_arguments)]
pub(crate) fn reconstruct_and_match(
path: &Path,
records: &[Kept],
args: &SearchArgs,
matcher: &Matcher,
turn_range: Option<crate::text::RangeSpec>,
time_window: &TimeWindow,
address: Option<&AddressSet>,
want_siblings: bool,
spawn_map: &HashMap<PathBuf, Option<Arc<DiscoveredSpawns>>>,
inner_parallel: bool,
head_is_fork: bool,
) -> (Vec<Exchange>, usize, ChainCounts) {
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 chain = crate::model::Chain::build_by(records, |k| &k.rec, None);
let survivor_lines: HashMap<usize, usize> = (0..records.len())
.filter_map(|i| chain.replay_of(i).map(|s| (s, records[s].line_no)))
.collect();
let index_turns = crate::model::group_turn_indices_chained(records, |k| &k.rec, &chain);
let plan_index = PlanIndex::from_records(records.iter().map(|k| &k.rec));
let filter = args.label_filter();
let summarize_index = if wants_compaction_pairing(&filter) {
crate::model::SummarizeIndex::from_records(records.iter().map(|k| &k.rec))
} else {
crate::model::SummarizeIndex::default()
};
let tool_names = build_tool_name_index(records);
let (use_ids, result_ids) = tool_pair_ids(records);
let spawn_lookup = spawn_map
.get(&discovery_root_for(path))
.and_then(|o| o.as_deref());
let first_opener_line = if head_is_fork {
None
} else {
records
.iter()
.find(|k| k.rec.opens_turn())
.map(|k| k.line_no)
};
let resume_prompts = resume_prompt_uuids(records);
let env = ClassifyEnv {
owner_id: &session_id,
is_subagent,
parent_id: &parent_session_id,
first_opener_line,
spawn: spawn_lookup.map(|s| s as &(dyn SpawnLookup + Sync)),
resume_prompts: &resume_prompts,
summarize: &summarize_index,
};
let excerpt_max = if args.no_truncate || address.is_some() {
usize::MAX
} else {
EXCERPT_MAX
};
let turn_bounds = turn_range.map(|spec| spec.resolve(index_turns.len(), false));
let want_diff = !(args.count_only || args.sessions_with_matches || args.count_by.is_some());
let build_exchange = |turn_index: usize, idxs: &[usize], abandoned: bool| -> Option<Exchange> {
if !abandoned {
if let Some((lo, hi)) = turn_bounds {
if turn_index < lo || turn_index > hi {
return None;
}
}
}
let turn = Turn {
index: turn_index,
records: idxs.iter().map(|&i| &records[i]).collect(),
indices: idxs.to_vec(),
};
let (mut hits, hit_idxs) = collect_turn_hits(
&turn,
&chain,
&survivor_lines,
filter,
matcher,
time_window,
args.resolve_persisted,
excerpt_max,
&plan_index,
&tool_names,
address,
&env,
);
if hits.is_empty() {
return None;
}
let (mut siblings, siblings_hidden) = if want_siblings {
collect_turn_siblings(
&turn,
&hit_idxs,
args.resolve_persisted,
excerpt_max,
&plan_index,
&tool_names,
&env,
)
} else {
(Vec::new(), 0)
};
for h in hits.iter_mut().chain(siblings.iter_mut()) {
set_pairing(h, &use_ids, &result_ids);
}
let record_uuids = turn
.records
.iter()
.filter_map(|k| k.rec.uuid.clone())
.collect();
let started_utc = turn
.records
.first()
.and_then(|k| k.rec.timestamp.clone())
.or_else(|| hits.iter().find_map(|h| h.timestamp_utc.clone()));
let turn_line_nos: Vec<usize> = turn
.records
.iter()
.map(|k| k.line_no)
.filter(|&n| n > 0)
.collect();
let turn_lines = match (turn_line_nos.iter().min(), turn_line_nos.iter().max()) {
(Some(&a), Some(&b)) => (a, b),
_ => (0, 0),
};
let head = idxs.first().copied();
let abandoned_kind = head.filter(|_| abandoned).and_then(|i| chain.kind(i));
let abandoned_root_line = head
.filter(|_| abandoned)
.and_then(|i| chain.abandoned_root(i))
.map(|r| records[r].line_no);
let draft_diff = if abandoned && want_diff {
head.and_then(|i| draft_diff_for(records, i, &chain, &plan_index))
} else {
None
};
Some(Exchange {
session_id: session_id.clone(),
is_subagent,
parent_session_id: parent_session_id.clone(),
turn_index: (!abandoned).then_some(turn.index),
started_utc,
hits,
siblings,
siblings_hidden,
turn_lines,
record_uuids,
superseded_draft: abandoned,
abandoned_kind,
abandoned_root_line,
draft_diff,
})
};
const PAR_TURNS_MIN_RECORDS: usize = 1024;
let mut out: Vec<Exchange> = if inner_parallel && records.len() >= PAR_TURNS_MIN_RECORDS {
index_turns
.par_iter()
.enumerate()
.filter_map(|(i, idxs)| build_exchange(i, idxs, false))
.collect()
} else {
index_turns
.iter()
.enumerate()
.filter_map(|(i, idxs)| build_exchange(i, idxs, false))
.collect()
};
if address.is_some() || turn_bounds.is_none() {
let mut branches: BTreeMap<usize, Vec<usize>> = BTreeMap::new();
for i in 0..records.len() {
if let Some(root) = chain.abandoned_root(i) {
branches.entry(root).or_default().push(i);
}
}
for (_, idxs) in branches {
out.extend(build_exchange(0, &idxs, true));
}
}
let counts = ChainCounts {
abandoned_records: chain.abandoned_records,
drafts: chain.drafts,
rewound_turns: chain.rewound_turns,
replay_copies: chain.replay_copies,
boundary_cut_line: chain.boundary_cut.map(|i| records[i].line_no),
leaf_source: Some(chain.leaf_source.as_str()),
};
(out, index_turns.len(), counts)
}
fn draft_diff_for(
records: &[Kept],
draft_idx: usize,
chain: &crate::model::Chain,
plan_index: &PlanIndex,
) -> Option<DraftDiff> {
let sent = records.get(chain.superseding(draft_idx)?)?;
let draft_text = records
.get(draft_idx)?
.rec
.reconstructed_user_text(Some(plan_index))?;
let sent_text = sent.rec.reconstructed_user_text(Some(plan_index))?;
let d = crate::chardiff::char_diff(&draft_text, &sent_text);
let sent_chars = sent_text.chars().count();
let pct = (sent_chars > 0).then(|| d.chars as f64 * 100.0 / sent_chars as f64);
Some(DraftDiff {
superseding_line: sent.line_no,
superseding_uuid: sent.rec.uuid.clone(),
chars: d.chars,
pct,
exact: d.exact,
})
}
pub(crate) fn sibling_cap(class: Class) -> Option<usize> {
if class == Class::AgentThinkingNarration {
return Some(1);
}
let path = class.path();
if path.starts_with("agent.thinking") {
Some(2)
} else if path.starts_with("agent.tool.") {
Some(3)
} else if path.starts_with("harness") {
Some(2)
} else {
None }
}
pub(crate) struct Turn<'a> {
pub(crate) index: usize,
pub(crate) records: Vec<&'a Kept>,
pub(crate) indices: Vec<usize>,
}
#[derive(Debug, Default)]
pub(crate) struct DiscoveredSpawns {
pub(crate) by_tool_use_id: HashMap<String, String>,
pub(crate) by_name: HashMap<String, String>,
}
impl SpawnLookup for DiscoveredSpawns {
fn child_for_spawn_tool_use_id(&self, tool_use_id: &str) -> Option<String> {
self.by_tool_use_id.get(tool_use_id).cloned()
}
fn child_for_spawn_name(&self, name: &str) -> Option<String> {
self.by_name.get(name).cloned()
}
}
pub(crate) fn parent_session_jsonl(path: &Path) -> Option<PathBuf> {
for anc in path.ancestors() {
if anc.file_name().and_then(|n| n.to_str()) == Some("subagents") {
return anc.parent().map(|d| d.with_extension("jsonl"));
}
}
None
}
pub(crate) fn discovery_root_for(path: &Path) -> PathBuf {
if is_subagent_path(path) {
parent_session_jsonl(path).unwrap_or_else(|| path.to_path_buf())
} else {
path.to_path_buf()
}
}
pub(crate) fn build_spawn_lookup(discovery_root: &Path) -> Option<DiscoveredSpawns> {
let subs = discover_subagents(discovery_root).ok()?;
if subs.is_empty() {
return None;
}
let mut out = DiscoveredSpawns::default();
for s in subs {
if let Some(tuid) = s.spawn_tool_use_id {
out.by_tool_use_id.entry(tuid).or_insert(s.agent_id.clone());
}
if let Some(name) = s.name {
out.by_name.entry(name).or_insert(s.agent_id.clone());
}
}
if out.by_tool_use_id.is_empty() && out.by_name.is_empty() {
return None;
}
Some(out)
}
pub(crate) struct ClassifyEnv<'a> {
pub(crate) owner_id: &'a str,
pub(crate) is_subagent: bool,
pub(crate) parent_id: &'a str,
pub(crate) first_opener_line: Option<usize>,
pub(crate) spawn: Option<&'a (dyn SpawnLookup + Sync)>,
pub(crate) resume_prompts: &'a HashSet<String>,
pub(crate) summarize: &'a crate::model::SummarizeIndex,
}
impl ClassifyEnv<'_> {
pub(crate) fn ctx_for(&self, kept: &Kept) -> ClassifyCtx<'_> {
ClassifyCtx {
owner_id: Some(self.owner_id),
owner_name: None,
is_subagent: self.is_subagent,
parent_id: Some(self.parent_id),
is_transcript_opener: self.is_subagent
&& kept.line_no != 0
&& Some(kept.line_no) == self.first_opener_line,
spawn: self.spawn.map(|s| s as &dyn SpawnLookup),
resume_prompt_uuids: Some(self.resume_prompts),
summarize: Some(self.summarize),
}
}
}
pub(crate) fn wants_compaction_pairing(filter: &LabelFilter<'_>) -> bool {
filter.selected(Class::CompactionBoundary.path())
|| filter.selected(Class::CompactionSummary.path())
}
pub(crate) fn resume_prompt_uuids(records: &[Kept]) -> HashSet<String> {
let mut out = HashSet::new();
for k in records {
if k.rec.is_resume_prompt() {
if let Some(uuid) = k.rec.uuid.as_deref() {
out.insert(uuid.to_string());
}
}
}
out
}
pub(crate) fn build_tool_name_index(records: &[Kept]) -> HashMap<String, String> {
let mut map: HashMap<String, String> = HashMap::new();
for k in records {
if let Some(blocks) = k.rec.blocks() {
for b in blocks {
if let Block::ToolUse {
id: Some(id),
name: Some(name),
..
} = b
{
map.entry(id.clone()).or_insert_with(|| name.clone());
}
}
}
}
map
}
pub(crate) fn tool_pair_ids(records: &[Kept]) -> (HashSet<String>, HashSet<String>) {
let mut uses = HashSet::new();
let mut results = HashSet::new();
for k in records {
if let Some(blocks) = k.rec.blocks() {
for b in blocks {
match b {
Block::ToolUse { id: Some(id), .. } => {
uses.insert(id.clone());
}
Block::ToolResult {
tool_use_id: Some(id),
..
} => {
results.insert(id.clone());
}
_ => {}
}
}
}
}
(uses, results)
}
pub(crate) fn set_pairing(h: &mut Hit, use_ids: &HashSet<String>, result_ids: &HashSet<String>) {
let Some(id) = h.tool_use_id.as_deref() else {
return;
};
h.pair = match h.class {
Some(Class::AgentToolUse | Class::CommSent | Class::CommSignal) => {
Some(if result_ids.contains(id) {
Pairing::Paired
} else {
Pairing::PendingNoResult
})
}
Some(Class::AgentToolResult | Class::CommInbox) => Some(if use_ids.contains(id) {
Pairing::Paired
} else {
Pairing::OrphanResult
}),
_ => None,
};
}