use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use serde::{Deserialize, Serialize};
use crate::mailbox::{local_machine_name, mail_root, Envelope, MailAddress, MailKind, Mailbox};
pub const AGENT_HARNESS: &str = "agent";
pub const SESSION_THREAD_PREFIX: &str = "s-";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Agent {
pub name: String,
pub main_session: MailAddress,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub folder: Option<PathBuf>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub harness: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idle_minutes: Option<u64>,
}
impl Agent {
pub fn address(&self) -> std::io::Result<MailAddress> {
agent_address(&self.name)
}
}
pub fn agent_address(name: &str) -> std::io::Result<MailAddress> {
MailAddress::new(local_machine_name(), AGENT_HARNESS, name)
.map_err(|error| std::io::Error::other(error.to_string()))
}
fn agents_dir() -> PathBuf {
mail_root().join("agents")
}
fn threads_dir() -> PathBuf {
mail_root().join("threads")
}
fn valid_name(name: &str) -> bool {
!name.is_empty()
&& name
.chars()
.all(|character| character.is_ascii_alphanumeric() || "-_.".contains(character))
}
fn write_atomically(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let temporary = path.with_extension(format!("tmp.{}", std::process::id()));
std::fs::write(&temporary, bytes)?;
std::fs::rename(temporary, path)
}
pub fn declare(agent: &Agent) -> std::io::Result<()> {
if !valid_name(&agent.name) {
return Err(std::io::Error::other(format!(
"`{}` is not an agent name: use letters, digits, `-`, `_` and `.`",
agent.name
)));
}
let bytes = serde_json::to_vec_pretty(agent).map_err(std::io::Error::other)?;
write_atomically(&agents_dir().join(format!("{}.json", agent.name)), &bytes)
}
pub fn load(name: &str) -> Option<Agent> {
if !valid_name(name) {
return None;
}
let bytes = std::fs::read(agents_dir().join(format!("{name}.json"))).ok()?;
serde_json::from_slice(&bytes).ok()
}
pub fn all_agents() -> Vec<Agent> {
let Ok(entries) = std::fs::read_dir(agents_dir()) else {
return Vec::new();
};
let mut agents: Vec<Agent> = entries
.flatten()
.filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
.filter_map(|entry| std::fs::read(entry.path()).ok())
.filter_map(|bytes| serde_json::from_slice(&bytes).ok())
.collect();
agents.sort_by(|a, b| a.name.cmp(&b.name));
agents
}
fn owners_account_manager_file() -> PathBuf {
agents_dir().join("owners-account-manager")
}
pub fn set_owners_account_manager(name: &str) -> std::io::Result<()> {
write_atomically(
&owners_account_manager_file(),
format!("{name}\n").as_bytes(),
)
}
pub fn owners_account_manager() -> Option<Agent> {
let name = std::fs::read_to_string(owners_account_manager_file()).ok()?;
load(name.trim())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Role {
To,
Cc,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Participant {
pub address: MailAddress,
pub role: Role,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Thread {
pub id: String,
pub agent: String,
pub holder: MailAddress,
pub participants: Vec<Participant>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub markers: BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub surface: Option<String>,
pub created_at_ms: u64,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub holder_launched: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent: Option<String>,
}
impl Thread {
pub fn role_of(&self, address: &MailAddress) -> Option<Role> {
self.participants
.iter()
.find(|participant| &participant.address == address)
.map(|participant| participant.role)
}
pub fn join(&mut self, address: &MailAddress, role: Role) {
match self
.participants
.iter_mut()
.find(|participant| &participant.address == address)
{
Some(participant) => participant.role = role,
None => self.participants.push(Participant {
address: address.clone(),
role,
}),
}
}
pub fn delegated(&self) -> bool {
load(&self.agent).is_some_and(|agent| agent.main_session != self.holder)
}
}
fn thread_path(id: &str) -> Option<PathBuf> {
valid_name(id).then(|| threads_dir().join(format!("{id}.json")))
}
pub fn thread(id: &str) -> Option<Thread> {
let bytes = std::fs::read(thread_path(id)?).ok()?;
serde_json::from_slice(&bytes).ok()
}
pub fn all_threads() -> Vec<Thread> {
let Ok(entries) = std::fs::read_dir(threads_dir()) else {
return Vec::new();
};
let mut threads: Vec<Thread> = entries
.flatten()
.filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
.filter_map(|entry| std::fs::read(entry.path()).ok())
.filter_map(|bytes| serde_json::from_slice(&bytes).ok())
.collect();
threads.sort_by(|a, b| (a.created_at_ms, &a.id).cmp(&(b.created_at_ms, &b.id)));
threads
}
pub fn threads_of(name: &str) -> Vec<Thread> {
all_threads()
.into_iter()
.filter(|thread| thread.agent == name)
.collect()
}
pub fn thread_held_by(address: &MailAddress) -> Option<Thread> {
all_threads().into_iter().rev().find(|thread| {
&thread.holder == address
&& thread.parent.is_none()
&& load(&thread.agent).is_some_and(|agent| &agent.main_session != address)
})
}
pub fn thread_of_marker(marker: &str) -> Option<(Thread, String)> {
all_threads().into_iter().find_map(|thread| {
let id = thread.markers.get(marker)?.clone();
Some((thread, id))
})
}
pub fn agent_of_session(address: &MailAddress) -> Option<Agent> {
if let Some(agent) = all_agents()
.into_iter()
.find(|agent| &agent.main_session == address)
{
return Some(agent);
}
thread_held_by(address).and_then(|thread| load(&thread.agent))
}
pub fn update_thread(
id: &str,
create: impl FnOnce() -> Option<Thread>,
change: impl FnOnce(&mut Thread),
) -> std::io::Result<Option<Thread>> {
let Some(path) = thread_path(id) else {
return Ok(None);
};
std::fs::create_dir_all(threads_dir())?;
let lock = path.with_extension("lock");
let started = Instant::now();
loop {
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&lock)
{
Ok(_) => break,
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
let stale = std::fs::metadata(&lock)
.and_then(|meta| meta.modified())
.ok()
.and_then(|modified| modified.elapsed().ok())
.is_some_and(|age| age > Duration::from_secs(10));
if stale {
std::fs::remove_file(&lock).ok();
} else if started.elapsed() > Duration::from_secs(30) {
return Err(std::io::Error::new(
std::io::ErrorKind::TimedOut,
format!("thread {id} is locked by another writer"),
));
}
std::thread::sleep(Duration::from_millis(20));
}
Err(error) => return Err(error),
}
}
let result = (|| {
let current = std::fs::read(&path)
.ok()
.and_then(|bytes| serde_json::from_slice::<Thread>(&bytes).ok())
.or_else(create);
let Some(mut thread) = current else {
return Ok(None);
};
change(&mut thread);
let bytes = serde_json::to_vec_pretty(&thread).map_err(std::io::Error::other)?;
write_atomically(&path, &bytes)?;
Ok(Some(thread))
})();
std::fs::remove_file(&lock).ok();
result
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Recipient {
pub address: MailAddress,
pub wake: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Plan {
pub recipients: Vec<Recipient>,
pub agent: Option<MailAddress>,
pub thread: Option<Thread>,
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|elapsed| elapsed.as_millis() as u64)
.unwrap_or_default()
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Channel {
pub marker: Option<String>,
pub reply_marker: Option<String>,
pub surface: Option<String>,
}
fn record_marker(thread: &mut Thread, channel: &Channel, id: &str) {
if let Some(marker) = &channel.marker {
thread.markers.insert(marker.clone(), id.to_string());
}
}
pub fn plan(
envelope: &mut Envelope,
to: &MailAddress,
channel: &Channel,
) -> std::io::Result<Option<Plan>> {
let sender = envelope.from.clone();
let id = envelope.id.clone();
if envelope.in_reply_to.is_none() {
if let Some((thread, parent)) = channel.reply_marker.as_deref().and_then(thread_of_marker) {
envelope.in_reply_to = Some(parent);
envelope.thread = Some(thread.id);
}
}
let known = envelope
.in_reply_to
.as_ref()
.and(envelope.thread.as_deref())
.and_then(thread);
if let Some(known) = known {
let receiver = (to.harness != AGENT_HARNESS && to != &sender).then(|| to.clone());
let saved = update_thread(
&known.id,
|| None,
|thread| {
if thread.role_of(&sender).is_none() {
thread.join(&sender, Role::To);
}
if let Some(receiver) = &receiver {
if thread.role_of(receiver).is_none() {
thread.join(receiver, Role::To);
}
}
record_marker(thread, channel, &id);
},
)?
.unwrap_or(known);
envelope.thread = Some(saved.id.clone());
let recipients = saved
.participants
.iter()
.filter(|participant| participant.address != sender)
.map(|participant| Recipient {
address: participant.address.clone(),
wake: participant.role == Role::To,
})
.collect();
return Ok(Some(Plan {
recipients,
agent: None,
thread: Some(saved),
}));
}
if envelope.in_reply_to.is_none() {
if let Some(held) = thread_held_by(&sender) {
let agent = load(&held.agent);
let own_agent = to.harness == AGENT_HARNESS && to.session_id == held.agent;
let to_main = own_agent || agent.as_ref().is_some_and(|a| &a.main_session == to);
if to_main {
let main = agent.map(|a| a.main_session).unwrap_or_else(|| to.clone());
let saved = update_thread(
&held.id,
|| None,
|thread| record_marker(thread, channel, &id),
)?
.unwrap_or(held);
envelope.thread = Some(saved.id.clone());
return Ok(Some(Plan {
recipients: vec![Recipient {
address: main,
wake: true,
}],
agent: None,
thread: Some(saved),
}));
}
let root = id.clone();
let receiver = if to.harness == AGENT_HARNESS {
load(&to.session_id)
.map(|other| other.main_session)
.ok_or_else(|| {
std::io::Error::other(format!(
"no agent named {} is declared on this machine",
to.session_id
))
})?
} else {
to.clone()
};
let mut participants = vec![
Participant {
address: sender.clone(),
role: Role::To,
},
Participant {
address: receiver,
role: Role::To,
},
];
if let Some(agent) = &agent {
participants.push(Participant {
address: agent.main_session.clone(),
role: Role::Cc,
});
}
let saved = update_thread(
&root,
|| {
Some(Thread {
id: root.clone(),
agent: held.agent.clone(),
holder: sender.clone(),
participants,
markers: BTreeMap::new(),
surface: None,
created_at_ms: now_ms(),
holder_launched: false,
parent: Some(held.id.clone()),
})
},
|thread| record_marker(thread, channel, &id),
)?;
envelope.thread = Some(root);
let recipients = saved
.as_ref()
.map(|thread| {
thread
.participants
.iter()
.filter(|p| p.address != sender)
.map(|p| Recipient {
address: p.address.clone(),
wake: p.role == Role::To,
})
.collect()
})
.unwrap_or_default();
return Ok(Some(Plan {
recipients,
agent: None,
thread: saved,
}));
}
}
if envelope.in_reply_to.is_none() && envelope.voice_for.as_ref() == Some(to) {
if let Some(held) = thread_held_by(to) {
envelope.thread = Some(held.id.clone());
let mut recipients = vec![Recipient {
address: to.clone(),
wake: true,
}];
recipients.extend(
held.participants
.iter()
.filter(|p| p.role == Role::Cc && p.address != sender && &p.address != to)
.map(|p| Recipient {
address: p.address.clone(),
wake: false,
}),
);
return Ok(Some(Plan {
recipients,
agent: None,
thread: Some(held),
}));
}
}
if to.harness == AGENT_HARNESS {
let agent = load(&to.session_id).ok_or_else(|| {
std::io::Error::other(format!(
"no agent named {} is declared on this machine",
to.session_id
))
})?;
let main = agent.main_session.clone();
if main == sender {
return Err(std::io::Error::other(format!(
"that is your own agent's address ({}); you are its main session",
agent.name
)));
}
if agent_of_session(&sender).is_some_and(|own| own.name == agent.name) {
return Ok(Some(Plan {
recipients: vec![Recipient {
address: main,
wake: true,
}],
agent: None,
thread: None,
}));
}
envelope.thread = None;
let root = id.clone();
let surface = channel.surface.clone();
let saved = update_thread(
&root,
|| {
Some(Thread {
id: root.clone(),
agent: agent.name.clone(),
holder: main.clone(),
participants: vec![
Participant {
address: main.clone(),
role: Role::To,
},
Participant {
address: sender.clone(),
role: Role::To,
},
],
markers: BTreeMap::new(),
surface,
created_at_ms: now_ms(),
holder_launched: false,
parent: None,
})
},
|thread| record_marker(thread, channel, &id),
)?;
return Ok(Some(Plan {
recipients: vec![Recipient {
address: main,
wake: true,
}],
agent: Some(to.clone()),
thread: saved,
}));
}
if envelope.in_reply_to.is_none() {
if let Some(agent) = all_agents()
.into_iter()
.find(|agent| agent.main_session == sender)
{
if to != &sender && to.harness != AGENT_HARNESS {
let root = id.clone();
let saved = update_thread(
&root,
|| {
Some(Thread {
id: root.clone(),
agent: agent.name.clone(),
holder: sender.clone(),
participants: vec![
Participant {
address: sender.clone(),
role: Role::To,
},
Participant {
address: to.clone(),
role: Role::To,
},
],
markers: BTreeMap::new(),
surface: None,
created_at_ms: now_ms(),
holder_launched: false,
parent: None,
})
},
|thread| record_marker(thread, channel, &id),
)?;
return Ok(Some(Plan {
recipients: vec![Recipient {
address: to.clone(),
wake: true,
}],
agent: None,
thread: saved,
}));
}
}
}
Ok(None)
}
pub fn root_of_main(main: &MailAddress, envelope: &Envelope) -> std::io::Result<Option<Thread>> {
let Some(agent) = agent_of_session(main).filter(|agent| &agent.main_session == main) else {
return Ok(None);
};
if &envelope.from == main || envelope.thread_id() != envelope.id {
return Ok(None);
}
let sender = envelope.from.clone();
update_thread(
&envelope.id,
|| {
Some(Thread {
id: envelope.id.clone(),
agent: agent.name.clone(),
holder: main.clone(),
participants: vec![
Participant {
address: main.clone(),
role: Role::To,
},
Participant {
address: sender,
role: Role::To,
},
],
markers: BTreeMap::new(),
surface: None,
created_at_ms: now_ms(),
holder_launched: false,
parent: None,
})
},
|_| {},
)
}
pub fn delegate_to(
thread_id: &str,
holder: &MailAddress,
launched: bool,
) -> std::io::Result<Option<Thread>> {
let Some(current) = thread(thread_id) else {
return Ok(None);
};
let Some(agent) = load(¤t.agent) else {
return Ok(None);
};
update_thread(
thread_id,
|| None,
|thread| {
if holder == &agent.main_session && thread.holder_launched && &thread.holder != holder {
let previous = thread.holder.clone();
thread.participants.retain(|p| p.address != previous);
}
thread.holder = holder.clone();
thread.holder_launched = launched && holder != &agent.main_session;
thread.join(holder, Role::To);
if holder != &agent.main_session {
thread.join(&agent.main_session, Role::Cc);
}
},
)
}
pub const RESUME_WAIT: Duration = Duration::from_secs(45);
pub fn resumable(address: &MailAddress) -> bool {
all_threads()
.iter()
.any(|thread| &thread.holder == address && thread.holder_launched)
&& recorded_resume_arguments(address).is_some()
}
fn recorded_resume_arguments(address: &MailAddress) -> Option<Vec<String>> {
let query = crate::DiscoveryQuery {
harnesses: vec![crate::HarnessId::new(&address.harness)],
query: Some(address.session_id.clone()),
limit: Some(50),
..Default::default()
};
let locator = crate::sdk::discover_sessions(&query)
.ok()?
.into_iter()
.find(|row| row.locator.session_id == address.session_id)?
.locator;
let config = crate::harness_service::recorded_config(&locator).ok()?;
if config.get("session_id").and_then(|v| v.as_str()) != Some(address.session_id.as_str()) {
return None;
}
resume_arguments(&config)
}
fn resume_arguments(config: &serde_json::Value) -> Option<Vec<String>> {
let text = |key: &str| config.get(key).and_then(|v| v.as_str());
let mut args = Vec::new();
match text("harness")? {
"codex" => {
let approval = text("approval_policy")?;
if !["untrusted", "on-failure", "on-request", "never"].contains(&approval) {
return None;
}
let policy = config.get("sandbox_policy")?.as_object()?;
let kind = policy.get("type")?.as_str()?;
let allowed: &[&str] = match kind {
"read-only" | "danger-full-access" => &["type"],
"workspace-write" => &[
"type",
"writable_roots",
"network_access",
"exclude_tmpdir_env_var",
"exclude_slash_tmp",
],
_ => return None,
};
if policy.keys().any(|key| !allowed.contains(&key.as_str())) {
return None;
}
args.extend([
"--ask-for-approval".into(),
approval.into(),
"--sandbox".into(),
kind.into(),
]);
if kind == "workspace-write" {
let roots = policy
.get("writable_roots")
.cloned()
.unwrap_or(serde_json::json!([]));
if !roots
.as_array()?
.iter()
.all(|root| root.as_str().is_some_and(|r| Path::new(r).is_absolute()))
{
return None;
}
args.extend([
"-c".into(),
format!("sandbox_workspace_write.writable_roots={roots}"),
]);
for key in [
"network_access",
"exclude_tmpdir_env_var",
"exclude_slash_tmp",
] {
let value = match policy.get(key) {
Some(value) => value.as_bool()?,
None => key != "network_access",
};
args.extend([
"-c".into(),
format!("sandbox_workspace_write.{key}={value}"),
]);
}
}
}
"claude-code" => {
let mode = text("permission_mode")?;
if !["default", "acceptEdits", "plan", "bypassPermissions"].contains(&mode) {
return None;
}
args.extend(["--permission-mode".into(), mode.into()]);
}
_ => return None,
}
if let Some(model) = text("model").filter(|m| !m.is_empty()) {
args.extend(["--model".into(), model.into()]);
}
Some(args)
}
fn resumed_path(address: &MailAddress) -> PathBuf {
let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
agents_dir().join("resumed").join(&hash[..24])
}
pub fn resume(address: &MailAddress) -> std::io::Result<()> {
let mut recorded = recorded_resume_arguments(address).ok_or_else(|| {
std::io::Error::other(format!(
"{address} has no recorded permissions to resume with; its mail waits"
))
})?;
if address.harness == "claude-code" {
if let Some(thread) = all_threads()
.into_iter()
.find(|thread| &thread.holder == address)
{
let id: String = address.session_id.chars().take(6).collect();
recorded.extend(["--name".to_string(), format!("{}-{id}", thread.agent)]);
}
}
let program = crate::claude_relay::supercode_program().map_err(std::io::Error::other)?;
let marker = resumed_path(address);
if let Some(dir) = marker.parent() {
std::fs::create_dir_all(dir)?;
}
write_atomically(&marker, now_ms().to_string().as_bytes())?;
let output = std::process::Command::new(program)
.args(["open", &address.session_id, "--detach", "--"])
.args(&recorded)
.stdin(std::process::Stdio::null())
.output()?;
if output.status.success() {
Ok(())
} else {
Err(std::io::Error::other(crate::mailbox::error_line(
&String::from_utf8_lossy(&output.stderr),
)))
}
}
pub fn file_unread(address: &MailAddress, envelope: &Envelope) -> std::io::Result<()> {
Mailbox::open(&mail_root(), address)?
.deliver(envelope)
.map(|_| ())
}
pub fn register_folder_sessions(homes: &crate::HarnessHomes) {
let agents: Vec<Agent> = all_agents()
.into_iter()
.filter(|agent| agent.folder.is_some())
.collect();
if agents.is_empty() {
return;
}
let holders: Vec<MailAddress> = all_threads()
.into_iter()
.map(|thread| thread.holder)
.collect();
for session in crate::mail_route::LiveSessions::read(homes).all() {
let Some(cwd) = &session.cwd else { continue };
let Some(agent) = agents.iter().find(|agent| {
agent.folder.as_deref() == Some(cwd.as_path()) && agent.main_session != session.address
}) else {
continue;
};
if holders.contains(&session.address) {
continue;
}
let id = format!("{SESSION_THREAD_PREFIX}{}", session.address.session_id);
let holder = session.address.clone();
let main = agent.main_session.clone();
update_thread(
&id,
|| {
Some(Thread {
id: id.clone(),
agent: agent.name.clone(),
holder: holder.clone(),
participants: vec![
Participant {
address: holder.clone(),
role: Role::To,
},
Participant {
address: main.clone(),
role: Role::Cc,
},
],
markers: BTreeMap::new(),
surface: None,
created_at_ms: now_ms(),
holder_launched: false,
parent: None,
})
},
|_| {},
)
.ok();
}
}
fn typed_cursor_path(address: &MailAddress) -> PathBuf {
let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
agents_dir().join("typed").join(&hash[..24])
}
pub fn copy_typed_lines(homes: &crate::HarnessHomes) {
let agents = all_agents();
if agents.is_empty() {
return;
}
let threads = all_threads();
let account_manager = owners_account_manager();
let mut sessions: Vec<(MailAddress, String)> = agents
.iter()
.map(|agent| (agent.main_session.clone(), agent.name.clone()))
.collect();
for thread in &threads {
if !sessions
.iter()
.any(|(address, _)| address == &thread.holder)
{
sessions.push((thread.holder.clone(), thread.agent.clone()));
}
}
let now = now_ms();
for (session, agent_name) in sessions {
let cursor_path = typed_cursor_path(&session);
let Some(cursor) = std::fs::read_to_string(&cursor_path)
.ok()
.and_then(|text| text.trim().parse::<u64>().ok())
else {
write_atomically(&cursor_path, now.to_string().as_bytes()).ok();
continue;
};
if let Some(path) = crate::mail_transcript::transcript(homes, &session) {
let changed = std::fs::metadata(&path)
.and_then(|meta| meta.modified())
.ok()
.and_then(|modified| modified.duration_since(UNIX_EPOCH).ok())
.map(|since| since.as_millis() as u64);
if changed.is_some_and(|changed| changed <= cursor) {
continue;
}
}
let lines: Vec<(u64, String)> = if session.harness == "codex" {
Mailbox::open(&mail_root(), &session)
.and_then(|mailbox| mailbox.list())
.unwrap_or_default()
.into_iter()
.filter(|stored| {
stored.envelope.kind == MailKind::User
&& stored
.envelope
.id
.starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
&& stored.envelope.created_at_ms > cursor
})
.map(|stored| (stored.envelope.created_at_ms, stored.envelope.body))
.collect()
} else {
let Some(mail) = crate::mail_transcript::read(homes, &session) else {
continue;
};
let mut typed_by_supercode: std::collections::HashMap<String, usize> =
std::collections::HashMap::new();
for stored in Mailbox::open(&mail_root(), &session)
.and_then(|mailbox| mailbox.list())
.unwrap_or_default()
{
if stored.envelope.kind == MailKind::User
&& !stored
.envelope
.id
.starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
{
*typed_by_supercode
.entry(stored.envelope.body.trim().to_string())
.or_default() += 1;
}
}
mail.typed
.into_iter()
.filter(|line| !line.withdrawn && line.sent_at_ms > cursor)
.filter(|line| match typed_by_supercode.get_mut(line.text.trim()) {
Some(left) if *left > 0 => {
*left -= 1;
false
}
_ => true,
})
.map(|line| (line.sent_at_ms, line.text))
.collect()
};
let Some(newest) = lines.iter().map(|(at, _)| *at).max() else {
continue;
};
let held = threads
.iter()
.rev()
.find(|thread| thread.holder == session && thread.delegated());
let mut receivers: Vec<MailAddress> = held
.map(|thread| {
thread
.participants
.iter()
.filter(|p| p.role == Role::Cc && p.address != session)
.map(|p| p.address.clone())
.collect()
})
.unwrap_or_default();
if let Some(account_manager) = &account_manager {
if agent_name != account_manager.name
&& !receivers.contains(&account_manager.main_session)
{
receivers.push(account_manager.main_session.clone());
}
}
let Ok(user) = crate::mail_transcript::user_address(&session.machine) else {
continue;
};
for (sent_at_ms, text) in lines {
let copy = Envelope {
id: crate::mail_transcript::typed_line_id(&session, sent_at_ms, &text),
created_at_ms: sent_at_ms,
from: user.clone(),
from_name: format!("user@{} (typed into {session})", session.machine),
sender_identity: None,
kind: MailKind::Typed,
reply_via: crate::mailbox::ReplyVia::None,
in_reply_to: None,
in_reply_to_inferred: false,
thread: held.map(|thread| thread.id.clone()),
native_from: None,
voice_for: None,
subject: None,
body: text,
};
for receiver in &receivers {
if let Err(error) = file_unread(receiver, ©) {
eprintln!("supercode mail: the line typed into {session} was not copied to {receiver}: {error}");
}
}
}
write_atomically(&cursor_path, newest.to_string().as_bytes()).ok();
}
}
pub fn close_idle_threads(homes: &crate::HarnessHomes) {
let launched: Vec<(Thread, u64)> = all_threads()
.into_iter()
.filter(|thread| thread.holder_launched)
.filter_map(|thread| {
let minutes = load(&thread.agent)?.idle_minutes?;
Some((thread, minutes))
})
.collect();
if launched.is_empty() {
return;
}
let live = crate::mail_route::LiveSessions::read(homes);
let now = now_ms();
for (thread, minutes) in launched {
let Some(session) = live
.all()
.iter()
.find(|session| session.address == thread.holder)
else {
continue;
};
if session.status != "idle" {
continue;
}
let Some(last) = session.last_message_at_ms(homes) else {
continue;
};
let resumed = std::fs::read_to_string(resumed_path(&thread.holder))
.ok()
.and_then(|text| text.trim().parse::<u64>().ok())
.unwrap_or(0);
let last = last.max(resumed);
if now.saturating_sub(last) < minutes.saturating_mul(60_000) {
continue;
}
let Ok(program) = crate::claude_relay::supercode_program() else {
return;
};
std::process::Command::new(program)
.args(["close", &thread.holder.to_string()])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status()
.ok();
}
}