use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use anyhow::{bail, Result};
use memchr::memmem;
use crate::model::Record;
use crate::parse::{
mmap_bytes, non_candidate_verdict, parse_line, scan_lines_parallel, LineVerdict,
};
use super::{
is_expired, now_utc, parse_header, read_inbox, read_ledger, read_message, states, InboxLine,
LedgerLine, Message, MessageState,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MsgVerdict {
Delivered,
IntentOnly,
Queued,
Held,
Expired,
Acked,
Refused,
}
impl MsgVerdict {
pub(crate) fn as_str(self) -> &'static str {
match self {
MsgVerdict::Delivered => "DELIVERED",
MsgVerdict::IntentOnly => "INTENT-ONLY",
MsgVerdict::Queued => "QUEUED",
MsgVerdict::Held => "HELD",
MsgVerdict::Expired => "EXPIRED",
MsgVerdict::Acked => "ACKED",
MsgVerdict::Refused => "REFUSED",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct Fact {
pub(crate) line: usize,
pub(crate) uuid: Option<String>,
}
#[derive(Debug, Clone)]
pub(crate) struct MsgReport {
pub(crate) id: String,
pub(crate) verdict: MsgVerdict,
pub(crate) lane: String,
pub(crate) session: String,
pub(crate) message: Option<Message>,
pub(crate) inbox: Option<InboxLine>,
pub(crate) state: MessageState,
pub(crate) emits: Vec<LedgerLine>,
pub(crate) fact: Option<Fact>,
}
impl MsgReport {
pub(crate) fn sort_key(&self) -> &str {
self.inbox
.as_ref()
.map_or_else(
|| self.state.first_ts_utc.as_deref(),
|i| Some(i.enqueued_utc.as_str()),
)
.unwrap_or("")
}
}
#[derive(Debug, Clone)]
pub(crate) struct LaneCtx {
pub(crate) lane: String,
pub(crate) session: String,
pub(crate) transcript: PathBuf,
pub(crate) root: PathBuf,
}
pub(crate) fn lane_from_target(target: &str) -> Result<LaneCtx> {
let token = if target.starts_with('@') {
target.to_string()
} else {
format!("@{target}")
};
let files = crate::path::resolve_session_files(
&[PathBuf::from(&token)],
crate::path::SubagentScope::TopLevelOnly,
crate::path::Caller::Other,
)?;
match files.len() {
0 => bail!("no transcript found for lane `{token}`"),
1 => lane_from_transcript(&files[0]),
_ => {
let ids: Vec<String> = files
.iter()
.map(|p| crate::subagent::session_id_from_path(p))
.collect();
bail!(
"`{token}` names {} lanes ({}) - a channel command operates on exactly one \
lane, so name it by its own transcript id",
files.len(),
ids.join(", ")
)
}
}
}
pub(crate) fn lane_from_transcript(transcript: &Path) -> Result<LaneCtx> {
let lane = crate::subagent::session_id_from_path(transcript);
if lane.is_empty() {
bail!("cannot read a lane id from {}", transcript.display());
}
let session =
crate::subagent::parent_session_id_from_path(transcript).unwrap_or_else(|| lane.clone());
let Some(sidecar) = session_sidecar_dir(transcript) else {
bail!(
"cannot locate the session sidecar directory for {}",
transcript.display()
);
};
Ok(LaneCtx {
lane,
session,
transcript: transcript.to_path_buf(),
root: super::channel_dir(&sidecar),
})
}
fn session_sidecar_dir(transcript: &Path) -> Option<PathBuf> {
for dir in transcript.ancestors() {
if dir.file_name().and_then(|s| s.to_str()) == Some("subagents") {
return dir.parent().map(Path::to_path_buf);
}
}
let stem = transcript.file_stem()?.to_str()?;
Some(transcript.parent()?.join(stem))
}
pub(crate) struct LaneView {
pub(crate) inboxes: BTreeMap<String, InboxLine>,
pub(crate) states: BTreeMap<String, MessageState>,
pub(crate) emits: BTreeMap<String, Vec<LedgerLine>>,
pub(crate) skipped_lines: usize,
}
impl LaneView {
pub(crate) fn ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self.inboxes.keys().cloned().collect();
for id in self.states.keys() {
if !self.inboxes.contains_key(id) {
ids.push(id.clone());
}
}
ids.sort();
ids.dedup();
ids
}
}
pub(crate) fn read_lane(ctx: &LaneCtx) -> Result<LaneView> {
let (inbox_lines, inbox_skipped) = read_inbox(&ctx.root, &ctx.lane)?;
let (ledger_lines, ledger_skipped) = read_ledger(&ctx.root, &ctx.lane)?;
let mut inboxes: BTreeMap<String, InboxLine> = BTreeMap::new();
for line in inbox_lines {
inboxes.entry(line.id.clone()).or_insert(line);
}
let mut emits: BTreeMap<String, Vec<LedgerLine>> = BTreeMap::new();
for line in &ledger_lines {
if matches!(line, LedgerLine::Emit { .. }) {
if let Some(id) = line.id() {
emits.entry(id.to_string()).or_default().push(line.clone());
}
}
}
Ok(LaneView {
inboxes,
states: states(&ledger_lines),
emits,
skipped_lines: inbox_skipped + ledger_skipped,
})
}
pub(crate) fn scan_facts(
transcript: &Path,
only: Option<&str>,
) -> Result<(BTreeMap<String, Fact>, usize)> {
let Some(map) = mmap_bytes(transcript)? else {
return Ok((BTreeMap::new(), 0));
};
let channel = memmem::Finder::new(Record::CSIFT_CHANNEL_NEEDLE.as_bytes());
let id_pattern = only.map(|id| format!("id={id}"));
let id_needle = id_pattern
.as_ref()
.map(|p| memmem::Finder::new(p.as_bytes()));
let (hits, skipped) = scan_lines_parallel(&map[..], |line, line_no| {
if channel.find(line).is_none()
|| id_needle.as_ref().is_some_and(|f| f.find(line).is_none())
{
return non_candidate_verdict(line);
}
match parse_line(line) {
Ok(Some(rec)) => match fact_id(&rec) {
Some(id) if only.is_none_or(|want| want == id) => LineVerdict::Keep((
id,
Fact {
line: line_no,
uuid: rec.uuid.clone(),
},
)),
_ => LineVerdict::Ignore,
},
Ok(None) => LineVerdict::Ignore,
Err(_) => LineVerdict::Skip,
}
});
let mut out: BTreeMap<String, Fact> = BTreeMap::new();
for (id, fact) in hits {
out.entry(id).or_insert(fact);
}
Ok((out, skipped))
}
fn fact_id(rec: &Record) -> Option<String> {
let text = rec.csift_channel_text()?;
parse_header(&text).map(|h| h.id)
}
pub(crate) fn reconcile(
ctx: &LaneCtx,
id: &str,
view: &LaneView,
fact: Option<Fact>,
) -> Result<Option<MsgReport>> {
let inbox = view.inboxes.get(id).cloned();
let state = view.states.get(id).cloned();
let message = read_message(&ctx.root, id)?;
if inbox.is_none() && state.is_none() {
return Ok(None);
}
let state = state.unwrap_or_default();
let emits = view.emits.get(id).cloned().unwrap_or_default();
let verdict = verdict_for(&state, inbox.as_ref(), fact.as_ref(), &now_utc());
Ok(Some(MsgReport {
id: id.to_string(),
verdict,
lane: ctx.lane.clone(),
session: ctx.session.clone(),
message,
inbox,
state,
emits,
fact,
}))
}
pub(crate) fn is_message_expired(
state: &MessageState,
inbox: Option<&InboxLine>,
now: &str,
) -> bool {
if state.expired {
return true;
}
inbox
.and_then(|i| i.expires_utc.as_deref())
.is_some_and(|deadline| is_expired(deadline, now))
}
fn verdict_for(
state: &MessageState,
inbox: Option<&InboxLine>,
fact: Option<&Fact>,
now: &str,
) -> MsgVerdict {
if state.acked {
return MsgVerdict::Acked;
}
let emitted = !state.emitted_parts.is_empty();
if !emitted && !state.refused_reasons.is_empty() {
return MsgVerdict::Refused;
}
if !emitted && is_message_expired(state, inbox, now) {
return MsgVerdict::Expired;
}
if !emitted && !state.held_reasons.is_empty() {
return MsgVerdict::Held;
}
if emitted {
return if fact.is_some() {
MsgVerdict::Delivered
} else {
MsgVerdict::IntentOnly
};
}
MsgVerdict::Queued
}