use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
use std::io::BufRead;
use std::path::{Path, PathBuf};
use serde::Serialize;
use crate::mailbox::{Envelope, MailAddress, MailKind, Mailbox, ReplyVia, StoredEnvelope};
use crate::HarnessHomes;
pub const TYPED_ID_PREFIX: &str = "u-";
pub const CHANNEL_ID_PREFIX: &str = "c-";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ChannelLine {
pub header: String,
pub from: String,
pub sent_at_ms: u64,
pub delivered_at_ms: Option<u64>,
pub text: String,
}
pub fn channel_line_id(address: &MailAddress, header: &str) -> String {
let hash = blake3::hash(format!("{address}\n{header}").as_bytes()).to_hex();
format!("{CHANNEL_ID_PREFIX}{}", &hash[..24])
}
fn channel_lines(text: &str, fallback_ms: u64, delivered_at_ms: Option<u64>) -> Vec<ChannelLine> {
let starts: Vec<usize> = text
.match_indices("[channel: ")
.map(|(at, _)| at)
.filter(|at| *at == 0 || text[..*at].ends_with('\n'))
.collect();
starts
.iter()
.enumerate()
.filter_map(|(n, start)| {
let end = starts.get(n + 1).copied().unwrap_or(text.len());
let line = text[*start..end].trim_end();
let header = &line[..line.find(']')? + 1];
let field = |name: &str| {
header
.split(" · ")
.find_map(|part| part.strip_prefix(name))
.map(|value| value.trim_end_matches(']').trim().to_string())
};
Some(ChannelLine {
header: header.to_string(),
from: field("from: ").unwrap_or_else(|| "a channel".to_string()),
sent_at_ms: field("at: ")
.and_then(|at| supercode_interchange::sidecar::rfc3339_to_ms(&at))
.and_then(|at| u64::try_from(at).ok())
.unwrap_or(fallback_ms),
delivered_at_ms,
text: line.to_string(),
})
})
.collect()
}
pub const ANSWER_ID_PREFIX: &str = "a-";
pub fn answer_id(address: &MailAddress, native_id: &str) -> String {
let hash = blake3::hash(format!("{address}\n{native_id}").as_bytes()).to_hex();
format!("{ANSWER_ID_PREFIX}{}", &hash[..24])
}
pub fn user_address(machine: &str) -> std::io::Result<MailAddress> {
MailAddress::new(machine.to_string(), "operator", "user")
.map_err(|error| std::io::Error::other(error.0))
}
pub fn typed_line_id(address: &MailAddress, sent_at_ms: u64, text: &str) -> String {
let hash = blake3::hash(format!("{address}\n{sent_at_ms}\n{text}").as_bytes()).to_hex();
format!("{TYPED_ID_PREFIX}{}", &hash[..24])
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TypedLine {
pub sent_at_ms: u64,
pub delivered_at_ms: Option<u64>,
pub withdrawn: bool,
pub text: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Answer {
pub native_id: String,
pub at_ms: u64,
pub text: String,
}
#[derive(Debug, Clone, Default)]
pub struct TranscriptMail {
pub typed: Vec<TypedLine>,
pub answers: Vec<Answer>,
pub channel: Vec<ChannelLine>,
pub delivered: HashMap<String, u64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum Delivery {
Delivered {
at_ms: u64,
},
Sent,
Withdrawn,
Unknown,
NotRead,
}
impl Delivery {
pub fn describe(&self) -> String {
match self {
Self::Delivered { at_ms } => {
supercode_interchange::sidecar::ms_to_rfc3339(*at_ms as i64)
}
Self::Sent => "not yet".to_string(),
Self::Withdrawn => "never (taken back)".to_string(),
Self::Unknown => "unknown".to_string(),
Self::NotRead => "not read here (supercode message show reads it)".to_string(),
}
}
}
pub fn transcript(homes: &HarnessHomes, address: &MailAddress) -> Option<PathBuf> {
if address.harness != "claude-code" {
return None;
}
let file = format!("{}.jsonl", address.session_id);
std::fs::read_dir(&homes.claude_code)
.ok()?
.flatten()
.map(|project| project.path().join(&file))
.find(|path| path.is_file())
}
fn at_ms(value: &serde_json::Value) -> u64 {
value
.as_str()
.and_then(supercode_interchange::sidecar::rfc3339_to_ms)
.and_then(|at| u64::try_from(at).ok())
.unwrap_or_default()
}
fn prompt_text(content: &serde_json::Value) -> Option<String> {
let text = match content {
serde_json::Value::String(text) => text.clone(),
serde_json::Value::Array(parts) => parts
.iter()
.filter(|part| part["type"] == "text")
.filter_map(|part| part["text"].as_str())
.collect::<Vec<_>>()
.join("\n"),
_ => return None,
};
(!text.trim().is_empty()).then_some(text)
}
fn rendered_mail_ids(line: &str) -> Vec<&str> {
let mut ids = Vec::new();
for mark in ["cross-session-message id=\\\"", "(message "] {
for (at, _) in line.match_indices(mark) {
let rest = &line[at + mark.len()..];
let end = rest
.find(|character: char| !(character.is_ascii_alphanumeric() || character == '-'))
.unwrap_or(rest.len());
let id = &rest[..end];
if id.starts_with("m-") {
ids.push(id);
}
}
}
ids
}
fn machine_made(text: &str) -> bool {
text.contains("<cross-session-message")
|| text.starts_with("[Cross-session")
|| text.starts_with("[channel: ")
|| text.starts_with("<task-notification>")
}
pub fn read_claude(path: &Path) -> std::io::Result<TranscriptMail> {
let reader = std::io::BufReader::new(std::fs::File::open(path)?);
let mut mail = TranscriptMail::default();
let mut queued: VecDeque<(String, Option<usize>)> = VecDeque::new();
let mut dequeued = 0_usize;
let mut absorbed: VecDeque<(String, Option<usize>)> = VecDeque::new();
let take = |waiting: &mut VecDeque<(String, Option<usize>)>, text: &str| {
let at = waiting.iter().position(|(words, _)| words == text)?;
waiting.remove(at).map(|(_, index)| index)
};
let mut not_typed: HashSet<usize> = HashSet::new();
let mut taken_in: HashMap<String, VecDeque<u64>> = HashMap::new();
let mut early: Vec<(String, Vec<String>, Option<usize>, u64)> = Vec::new();
for line in reader.lines() {
let line = line?;
let prompt = line.contains("\"promptSource\"") || line.contains("\"type\":\"user\"");
let queue = line.contains("\"queue-operation\"");
let attachment = line.contains("\"queued_command\"");
let answer = line.contains("\"stop_reason\":\"end_turn\"")
|| line.contains("\"stop_reason\":\"stop_sequence\"");
let rendered = line.contains("cross-session-message id=") || line.contains("(message m-");
if !(prompt || queue || attachment || rendered || answer) {
continue;
}
let Ok(record) = serde_json::from_str::<serde_json::Value>(&line) else {
continue;
};
if record["isSidechain"] == true {
continue;
}
let kind = record["type"].as_str().unwrap_or_default();
let mut at = at_ms(&record["timestamp"]);
if kind == "queue-operation" && record["operation"] == "remove" {
if let Some(text) = record["content"].as_str() {
if let Some(found) = early.iter().position(|(words, ..)| words == text) {
let (_, ids, index, _) = early.remove(found);
for id in ids {
mail.delivered.entry(id).or_insert(at);
}
if let Some(index) = index {
mail.typed[index].delivered_at_ms = Some(at);
}
continue;
}
taken_in.entry(text.to_string()).or_default().push_back(at);
}
}
let mut waits_for_remove = false;
if kind == "attachment" && record["attachment"]["type"] == "queued_command" {
match record["attachment"]["prompt"]
.as_str()
.and_then(|text| taken_in.get_mut(text))
.and_then(VecDeque::pop_front)
{
Some(taken) => at = taken,
None => waits_for_remove = true,
}
}
if waits_for_remove {
let attached = &record["attachment"];
let text = attached["prompt"].as_str().unwrap_or_default().to_string();
let human = attached["origin"]["kind"] == "human" || attached["humanTurn"] == true;
let ids = rendered_mail_ids(&line)
.into_iter()
.map(str::to_string)
.collect();
let index = match take(&mut queued, &text) {
Some(Some(index)) if !human => {
not_typed.insert(index);
None
}
Some(index) => index,
None if human && !text.trim().is_empty() => {
mail.typed.push(TypedLine {
sent_at_ms: at,
delivered_at_ms: None,
withdrawn: false,
text: text.clone(),
});
Some(mail.typed.len() - 1)
}
None => None,
};
early.push((text, ids, index, at));
continue;
}
if rendered && matches!(kind, "user" | "attachment") {
for id in rendered_mail_ids(&line) {
mail.delivered.entry(id.to_string()).or_insert(at);
}
}
match kind {
"assistant"
if answer
&& matches!(
record["message"]["stop_reason"].as_str(),
Some("end_turn" | "stop_sequence")
) =>
{
let (Some(native_id), Some(text)) = (
record["uuid"].as_str(),
prompt_text(&record["message"]["content"]),
) else {
continue;
};
mail.answers.push(Answer {
native_id: native_id.to_string(),
at_ms: at,
text,
});
}
"queue-operation" if record["operation"] == "dequeue" => dequeued += 1,
"queue-operation" => {
let Some(text) = record["content"].as_str() else {
continue;
};
match record["operation"].as_str() {
Some("enqueue") => {
let index = (!machine_made(text) && !text.trim().is_empty()).then(|| {
mail.typed.push(TypedLine {
sent_at_ms: at,
delivered_at_ms: None,
withdrawn: false,
text: text.to_string(),
});
mail.typed.len() - 1
});
queued.push_back((text.to_string(), index));
}
Some("remove") => {
let Some(index) = take(&mut queued, text) else {
continue;
};
match record["reason"].as_str() {
Some("absorbed_mid_turn" | "delivered_to_agent") => {
if let Some(index) = index {
mail.typed[index].delivered_at_ms = Some(at);
}
absorbed.push_back((text.to_string(), index));
}
_ => {
if let Some(index) = index {
mail.typed[index].withdrawn = true;
}
}
}
}
_ => {}
}
}
"attachment" => {
let attached = &record["attachment"];
let Some(text) = attached["prompt"].as_str() else {
continue;
};
if attached["type"] != "queued_command" || text.trim().is_empty() {
continue;
}
let human = attached["origin"]["kind"] == "human" || attached["humanTurn"] == true;
match take(&mut absorbed, text) {
Some(Some(index)) if !human => {
not_typed.insert(index);
}
Some(_) => {}
None if human => mail.typed.push(TypedLine {
sent_at_ms: at_ms(&attached["timestamp"]).min(at),
delivered_at_ms: Some(at),
withdrawn: false,
text: text.to_string(),
}),
None => {}
}
}
"user" => {
let Some(text) = prompt_text(&record["message"]["content"]) else {
continue;
};
let source = record["promptSource"].as_str();
let human = matches!(source, Some("typed" | "queued")) && record["isMeta"] != true;
if !human {
mail.channel.extend(channel_lines(&text, at, Some(at)));
}
let mut landing: Vec<Option<usize>> = Vec::new();
let let_go = std::mem::take(&mut dequeued);
if let_go == 0 {
landing.extend(take(&mut queued, &text));
}
let wanted = let_go;
while landing.len() < wanted {
let Some(at) = queued
.iter()
.position(|(words, _)| text.contains(words.as_str()))
else {
break;
};
landing.extend(queued.remove(at).map(|(_, index)| index));
}
if landing.is_empty() && wanted > 0 {
for _ in 0..wanted {
landing.extend(queued.pop_front().map(|(_, index)| index));
}
}
if !landing.is_empty() {
for index in landing.into_iter().flatten() {
if human {
mail.typed[index].delivered_at_ms = Some(at);
} else {
not_typed.insert(index);
}
}
continue;
}
if human {
mail.typed.push(TypedLine {
sent_at_ms: at,
delivered_at_ms: Some(at),
withdrawn: false,
text,
});
}
}
_ => {}
}
}
for (_, ids, index, at) in early {
for id in ids {
mail.delivered.entry(id).or_insert(at);
}
if let Some(index) = index {
mail.typed[index].delivered_at_ms.get_or_insert(at);
}
}
for (text, _) in &queued {
if text.starts_with("[channel: ") {
mail.channel.extend(channel_lines(text, 0, None));
}
}
let mut index = 0;
mail.typed.retain(|_| {
index += 1;
!not_typed.contains(&(index - 1))
});
Ok(mail)
}
pub fn read(homes: &HarnessHomes, address: &MailAddress) -> Option<TranscriptMail> {
read_claude(&transcript(homes, address)?).ok()
}
fn newest_before(items: &[(u64, String)], at: u64, inclusive: bool) -> Option<String> {
items
.iter()
.rev()
.find(|(when, _)| if inclusive { *when <= at } else { *when < at })
.map(|(_, id)| id.clone())
}
fn answer_times(address: &MailAddress, mail: &TranscriptMail) -> Vec<(u64, String)> {
let mut times: Vec<(u64, String)> = mail
.answers
.iter()
.map(|answer| (answer.at_ms, answer_id(address, &answer.native_id)))
.collect();
times.sort();
times
}
fn file_derived(
mailbox: &Mailbox,
filed: &HashMap<&str, &StoredEnvelope>,
envelope: Envelope,
) -> std::io::Result<bool> {
if let Some(stored) = filed.get(envelope.id.as_str()) {
if stored.envelope.in_reply_to == envelope.in_reply_to {
return Ok(false);
}
std::fs::remove_file(&stored.path).ok();
}
mailbox.file_read(&envelope)?;
Ok(true)
}
pub fn file_typed_lines(mailbox: &Mailbox, mail: &TranscriptMail) -> std::io::Result<usize> {
let address = mailbox.address();
let filed = mailbox.list()?;
let by_id: HashMap<&str, &StoredEnvelope> = filed
.iter()
.map(|stored| (stored.envelope.id.as_str(), stored))
.collect();
let answers = answer_times(address, mail);
let mut typed_by_supercode: HashMap<&str, usize> = HashMap::new();
for stored in &filed {
if stored.envelope.kind == MailKind::User
&& !stored.envelope.id.starts_with(TYPED_ID_PREFIX)
{
*typed_by_supercode
.entry(stored.envelope.body.trim())
.or_default() += 1;
}
}
let user = user_address(&address.machine)?;
let current: HashSet<String> = mail
.typed
.iter()
.map(|line| typed_line_id(address, line.sent_at_ms, &line.text))
.collect();
for stored in &filed {
if stored.envelope.id.starts_with(TYPED_ID_PREFIX) && !current.contains(&stored.envelope.id)
{
std::fs::remove_file(&stored.path).ok();
}
}
let mut count = 0;
for line in &mail.typed {
let id = typed_line_id(address, line.sent_at_ms, &line.text);
if !by_id.contains_key(id.as_str()) {
if let Some(left) = typed_by_supercode.get_mut(line.text.trim()) {
if *left > 0 {
*left -= 1;
continue;
}
}
}
let in_reply_to = newest_before(&answers, line.sent_at_ms, false);
count += usize::from(file_derived(
mailbox,
&by_id,
Envelope {
id,
created_at_ms: line.sent_at_ms,
from: user.clone(),
from_name: format!("user@{}", address.machine),
kind: MailKind::User,
reply_via: ReplyVia::None,
in_reply_to_inferred: in_reply_to.is_some(),
in_reply_to,
thread: None,
native_from: None,
body: line.text.clone(),
},
)?);
}
for line in &mail.channel {
let id = channel_line_id(address, &line.header);
let in_reply_to = newest_before(&answers, line.sent_at_ms, false);
count += usize::from(file_derived(
mailbox,
&by_id,
Envelope {
id,
created_at_ms: line.sent_at_ms,
from: MailAddress::new(address.machine.clone(), "operator", "channel")
.map_err(|error| std::io::Error::other(error.0))?,
from_name: line.from.clone(),
kind: MailKind::Channel,
reply_via: ReplyVia::None,
in_reply_to_inferred: in_reply_to.is_some(),
in_reply_to,
thread: None,
native_from: None,
body: line.text.clone(),
},
)?);
}
Ok(count)
}
pub fn file_answers(
root: &Path,
address: &MailAddress,
name: &str,
mail: &TranscriptMail,
) -> std::io::Result<usize> {
let mailbox = Mailbox::open(root, &user_address(&address.machine)?)?;
let listed = mailbox.list()?;
let mut filed: HashMap<&str, &StoredEnvelope> = HashMap::new();
for stored in &listed {
if stored.envelope.id.starts_with(ANSWER_ID_PREFIX)
&& stored.envelope.kind != MailKind::Answer
{
std::fs::remove_file(&stored.path).ok();
continue;
}
filed.insert(stored.envelope.id.as_str(), stored);
}
let mut people: Vec<(u64, String)> = mail
.typed
.iter()
.map(|line| {
(
line.sent_at_ms,
typed_line_id(address, line.sent_at_ms, &line.text),
)
})
.chain(
mail.channel
.iter()
.map(|line| (line.sent_at_ms, channel_line_id(address, &line.header))),
)
.chain(
Mailbox::open(root, address)?
.list()?
.into_iter()
.filter(|stored| {
stored.envelope.kind == MailKind::User
&& !stored.envelope.id.starts_with(TYPED_ID_PREFIX)
})
.map(|stored| (stored.envelope.created_at_ms, stored.envelope.id)),
)
.collect();
people.sort();
let mut count = 0;
for answer in &mail.answers {
let id = answer_id(address, &answer.native_id);
let in_reply_to = newest_before(&people, answer.at_ms, true);
count += usize::from(file_derived(
&mailbox,
&filed,
Envelope {
id,
created_at_ms: answer.at_ms,
from: address.clone(),
from_name: name.to_string(),
kind: MailKind::Answer,
reply_via: ReplyVia::None,
in_reply_to_inferred: in_reply_to.is_some(),
in_reply_to,
thread: None,
native_from: None,
body: answer.text.clone(),
},
)?);
}
Ok(count)
}
pub fn sent_by(root: &Path, address: &MailAddress) -> Vec<(MailAddress, StoredEnvelope)> {
let mut sent: Vec<(MailAddress, StoredEnvelope)> = crate::mailbox::all_mailboxes(root)
.into_iter()
.filter(|mailbox| mailbox.address() != address)
.flat_map(|mailbox| {
let to = mailbox.address().clone();
mailbox
.list()
.unwrap_or_default()
.into_iter()
.filter(|stored| &stored.envelope.from == address)
.map(move |stored| (to.clone(), stored))
})
.collect();
sent.sort_by_key(|(_, stored)| stored.envelope.created_at_ms);
sent
}
pub fn deliveries(
address: &MailAddress,
stored: &[StoredEnvelope],
mail: Option<&TranscriptMail>,
) -> BTreeMap<String, Delivery> {
let Some(mail) = mail else {
return stored
.iter()
.map(|stored| (stored.envelope.id.clone(), Delivery::Unknown))
.collect();
};
let typed: HashMap<String, &TypedLine> = mail
.typed
.iter()
.map(|line| (typed_line_id(address, line.sent_at_ms, &line.text), line))
.collect();
let channel: HashMap<String, &ChannelLine> = mail
.channel
.iter()
.map(|line| (channel_line_id(address, &line.header), line))
.collect();
let mut used: HashSet<usize> = HashSet::new();
let mut result = BTreeMap::new();
for stored in stored {
let envelope = &stored.envelope;
if typed.get(&envelope.id).is_some_and(|line| line.withdrawn) {
result.insert(envelope.id.clone(), Delivery::Withdrawn);
continue;
}
let at = if let Some(line) = typed.get(&envelope.id) {
line.delivered_at_ms
} else if let Some(line) = channel.get(&envelope.id) {
line.delivered_at_ms
} else if let Some(at) = mail
.delivered
.get(&envelope.id)
.or_else(|| mail.delivered.get(crate::mailbox::short_id(&envelope.id)))
{
Some(*at)
} else if envelope.kind == MailKind::User {
mail.typed
.iter()
.enumerate()
.find(|(index, line)| {
!used.contains(index)
&& line.sent_at_ms + 1_000 >= envelope.created_at_ms
&& line.text.trim() == envelope.body.trim()
})
.and_then(|(index, line)| {
used.insert(index);
line.delivered_at_ms
})
} else {
None
};
result.insert(
envelope.id.clone(),
match at {
Some(at_ms) => Delivery::Delivered { at_ms },
None => Delivery::Sent,
},
);
}
result
}
pub const MIN_PREFIX_DIGITS: usize = 4;
pub fn find_messages(root: &Path, id: &str) -> std::io::Result<Vec<(Mailbox, StoredEnvelope)>> {
let digits = id.split_once('-').map_or(0, |(_, digits)| digits.len());
if digits < MIN_PREFIX_DIGITS {
return Ok(Vec::new());
}
let mut found: Vec<(Mailbox, StoredEnvelope)> = Vec::new();
for mailbox in crate::mailbox::all_mailboxes(root) {
for stored in mailbox.find_prefix(id)? {
if !found
.iter()
.any(|(_, known)| known.envelope.id == stored.envelope.id)
{
found.push((mailbox.clone(), stored));
}
}
}
Ok(found)
}