use std::path::Path;
use crate::claude_peer::{read_registry, registry_dir, ClaudePeerSession, ClaudePeerStatus};
use crate::claude_relay::{send_through_relay, supercode_program, RelayReceipt, RELAY_NAME_PREFIX};
use crate::live_runtime::LiveRuntimeRecord;
use crate::mailbox::{
local_machine_name, mail_root, Envelope, IdleSubscription, MailAddress, Mailbox, ReplyVia,
};
use crate::runtime_mail::{deliver_to_runtime, RuntimeDelivery};
use crate::HarnessHomes;
pub const CODEX_HOOK_ARGUMENTS: &str = "message hook codex";
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Door {
Runtime(Box<LiveRuntimeRecord>),
Native(Box<ClaudePeerSession>),
Hook,
Stored,
Operator,
}
impl Door {
pub fn name(&self) -> &'static str {
match self {
Self::Runtime(_) => "runtime",
Self::Native(_) => "native",
Self::Hook => "hook",
Self::Stored => "stored",
Self::Operator => "operator",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum NoDoor {
OtherMachine(String),
NotRunning,
}
pub fn door_for(homes: &HarnessHomes, to: &MailAddress) -> Result<Door, NoDoor> {
if to.machine != local_machine_name() {
return Err(NoDoor::OtherMachine(to.machine.clone()));
}
if to.harness == "operator" {
return Ok(Door::Operator);
}
LiveSessions::read(homes)
.sessions
.into_iter()
.find(|session| &session.address == to)
.map(|session| session.door)
.ok_or(NoDoor::NotRunning)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LiveSession {
pub address: MailAddress,
pub name: String,
pub status: String,
pub door: Door,
pub pid: Option<u32>,
pub cwd: Option<std::path::PathBuf>,
pub tmux: Option<String>,
}
#[derive(Debug, Default)]
pub struct LiveSessions {
sessions: Vec<LiveSession>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Unresolved {
Stale(String),
Unknown(String),
}
impl LiveSessions {
pub fn read(homes: &HarnessHomes) -> Self {
let machine = local_machine_name();
let registry = read_registry(®istry_dir(homes));
let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
registry.iter().find(|session| {
session.session_id == record.source.session_id
|| session.session_id == record.runtime_session_id
})
};
let mut sessions: Vec<LiveSession> = crate::runtime_mail::controlled_runtimes()
.into_iter()
.filter_map(|record| {
let registered = registered(&record);
if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
return None;
}
let address =
MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
.ok()?;
let short: String = record.source.session_id.chars().take(8).collect();
let name = match registered {
Some(session) if !session.name.is_empty() => session.name.clone(),
_ => format!("{}-{short}", record.source.harness),
};
Some(LiveSession {
name: format!("{name}@{machine}"),
address,
status: "hosted".into(),
pid: Some(record.pid),
cwd: Some(record.source.workspace.clone()),
tmux: None,
door: Door::Runtime(Box::new(record)),
})
})
.collect();
let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
sessions.iter().any(|session| &session.address == address)
};
for session in registry {
if session.name.starts_with(RELAY_NAME_PREFIX) {
continue;
}
let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
continue;
};
if controlled(&address, &sessions) {
continue;
}
sessions.push(LiveSession {
address,
name: format!("{}@{machine}", session.name),
status: session
.status
.as_ref()
.map(|status| status.as_str().to_string())
.unwrap_or_else(|| "unknown".into()),
pid: Some(session.pid),
cwd: session.cwd.clone(),
tmux: session.tmux.clone(),
door: Door::Native(Box::new(session)),
});
}
let user_hook = codex_user_hook_installed();
for (path, status) in crate::codex_peer::live_rollouts(&homes.codex) {
let Some((session_id, None)) = crate::codex_peer::rollout_session(&path) else {
continue;
};
let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
continue;
};
if controlled(&address, &sessions) {
continue;
}
let cwd = crate::codex_peer::rollout_cwd(&path);
let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
sessions.push(LiveSession {
name: format!("{}@{machine}", codex_name(&session_id)),
address,
status: status.as_str().to_string(),
pid: None,
cwd,
tmux: None,
door: if hooked { Door::Hook } else { Door::Stored },
});
}
Self { sessions }
}
pub fn all(&self) -> &[LiveSession] {
&self.sessions
}
pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
self.sessions
.iter()
.find(|session| {
session.address.harness == harness && session.address.session_id == session_id
})
.map(|session| session.door.name())
}
pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
let machine = local_machine_name();
if let Ok(address) = MailAddress::parse(to) {
return self
.sessions
.iter()
.find(|session| session.address == address)
.ok_or_else(|| {
Unresolved::Stale(format!(
"{to} is no longer running. Nothing was sent. Run supercode message list \
for the live sessions."
))
});
}
let wanted = match to.split_once('@') {
Some((name, at)) if at == machine => name.to_string(),
Some((_, at)) => {
return Err(Unresolved::Unknown(format!(
"Not sent: {to} is on machine {at}, not this one. Nothing was sent."
)))
}
None => to.to_string(),
};
let short = |session: &LiveSession| {
session
.name
.split('@')
.next()
.unwrap_or_default()
.to_string()
};
let matching: Vec<&LiveSession> = self
.sessions
.iter()
.filter(|session| short(session) == wanted)
.collect();
if let [only] = matching.as_slice() {
return Ok(only);
}
let hint = if matching.len() > 1 {
format!(
" {} sessions are named {wanted}; use its address.",
matching.len()
)
} else {
let near: Vec<String> = self
.sessions
.iter()
.filter(|session| {
let name = short(session);
name.contains(&wanted)
|| wanted.contains(&name)
|| name
.chars()
.zip(wanted.chars())
.take_while(|(a, b)| a == b)
.count()
>= 4
})
.take(3)
.map(|session| {
format!(
"{} ({}, {})",
session.name, session.address.harness, session.status
)
})
.collect();
if near.is_empty() {
String::new()
} else {
format!(" Did you mean: {}?", near.join(", "))
}
};
Err(Unresolved::Unknown(format!(
"No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
was sent."
)))
}
}
pub fn has_message_tools(pid: u32) -> bool {
let Ok(output) = std::process::Command::new("ps")
.args(["-A", "-o", "ppid=,command="])
.output()
else {
return false;
};
String::from_utf8_lossy(&output.stdout).lines().any(|line| {
let line = line.trim_start();
let Some((ppid, command)) = line.split_once(' ') else {
return false;
};
ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
})
}
pub fn codex_name(session_id: &str) -> String {
format!("codex-{}", session_id.chars().take(8).collect::<String>())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Caller {
pub address: MailAddress,
pub name: String,
}
pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
shell tool.";
pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
let machine = local_machine_name();
let registry = read_registry(®istry_dir(homes));
let hosted = crate::runtime_mail::controlled_runtimes();
for &pid in pids {
if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
let short: String = record.source.session_id.chars().take(8).collect();
return Ok(Caller {
address: MailAddress::new(
&machine,
&record.source.harness,
&record.source.session_id,
)
.map_err(|error| error.to_string())?,
name: format!("{}-{short}@{machine}", record.source.harness),
});
}
if let Some(session) = registry.iter().find(|session| session.pid == pid) {
if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
if !claimed.is_empty() && claimed != session.session_id {
return Err(format!(
"{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
process {pid} above this command is session {})",
session.session_id
));
}
}
return Ok(Caller {
address: MailAddress::new(&machine, "claude-code", &session.session_id)
.map_err(|error| error.to_string())?,
name: format!("{}@{machine}", session.name),
});
}
if let Some(thread) = std::env::var("CODEX_THREAD_ID")
.ok()
.filter(|thread| !thread.is_empty())
.filter(|thread| crate::codex_peer::holds_session(pid, thread))
{
return Ok(Caller {
address: MailAddress::new(&machine, "codex", &thread)
.map_err(|error| error.to_string())?,
name: format!("{}@{machine}", codex_name(&thread)),
});
}
if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
return Ok(Caller {
address: MailAddress::new(&machine, "codex", &session_id)
.map_err(|error| error.to_string())?,
name: format!("{}@{machine}", codex_name(&session_id)),
});
}
}
#[cfg(windows)]
if let Some(session) = msys_cut_claim(®istry, pids) {
return Ok(Caller {
address: MailAddress::new(&machine, "claude-code", &session.session_id)
.map_err(|error| error.to_string())?,
name: format!("{}@{machine}", session.name),
});
}
Err(format!(
"{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
sessions in {} and {} hosted runtimes)",
registry.len(),
registry_dir(homes).display(),
hosted.len()
))
}
#[cfg(windows)]
fn msys_cut_claim<'a>(
registry: &'a [ClaudePeerSession],
pids: &[u32],
) -> Option<&'a ClaudePeerSession> {
let table = process_table();
let [.., shell, cut] = pids else {
return None;
};
if table.contains_key(cut) {
return None;
}
let name = table.get(shell)?.1.to_ascii_lowercase();
if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
return None;
}
let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
registry
.iter()
.find(|session| !claimed.is_empty() && session.session_id == claimed)
}
pub fn process_ancestry() -> Vec<u32> {
ancestry_of(std::process::id())
}
pub fn ancestry_of(pid: u32) -> Vec<u32> {
let parents = parent_pids();
let mut chain = vec![pid];
let mut current = pid;
while let Some(&parent) = parents.get(¤t) {
if parent <= 1 || chain.contains(&parent) {
break;
}
chain.push(parent);
current = parent;
}
chain
}
#[cfg(not(windows))]
fn parent_pids() -> std::collections::HashMap<u32, u32> {
let Ok(output) = std::process::Command::new("ps")
.args(["-axo", "pid=,ppid="])
.output()
else {
return Default::default();
};
String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|line| {
let mut fields = line.split_whitespace();
Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
})
.collect()
}
#[cfg(windows)]
fn parent_pids() -> std::collections::HashMap<u32, u32> {
process_table()
.into_iter()
.map(|(pid, (parent, _))| (pid, parent))
.collect()
}
#[cfg(windows)]
fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
use windows_sys::Win32::System::Diagnostics::ToolHelp::{
CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
TH32CS_SNAPPROCESS,
};
let mut parents = std::collections::HashMap::new();
let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
if snapshot == INVALID_HANDLE_VALUE {
return parents;
}
let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
while has_entry {
let length = entry
.szExeFile
.iter()
.position(|&unit| unit == 0)
.unwrap_or(entry.szExeFile.len());
let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
}
unsafe {
CloseHandle(snapshot);
}
parents
}
pub fn codex_hooks_path() -> std::path::PathBuf {
std::env::var_os("CODEX_HOME")
.map(std::path::PathBuf::from)
.or_else(|| {
supercode_interchange::user_home()
.map(std::path::PathBuf::into_os_string)
.map(|home| std::path::PathBuf::from(home).join(".codex"))
})
.unwrap_or_else(|| std::path::PathBuf::from(".codex"))
.join("hooks.json")
}
pub fn codex_user_hook_installed() -> bool {
std::fs::read_to_string(codex_hooks_path())
.is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
}
pub fn codex_project_hook_installed(cwd: &Path) -> bool {
cwd.ancestors().any(|directory| {
std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
.is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Delivered {
Steered,
Started,
Native {
busy: bool,
},
Hooked,
Queued,
Stored,
Operator,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Refused {
CannotQueueNative,
TooLong(usize),
}
pub const MAX_RELAYED_BYTES: usize = 100_000;
pub async fn deliver(
envelope: &Envelope,
to: &MailAddress,
door: &Door,
wake: bool,
notify_when_idle: bool,
) -> Result<Result<Delivered, Refused>, String> {
let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
let mut final_reply = false;
let delivered = match door {
Door::Runtime(record) => {
let mut sent = envelope.clone();
let answers = matches!(envelope.reply_via, ReplyVia::Command);
if answers {
sent.reply_via = ReplyVia::FinalMessage {
destination: envelope.from_name.clone(),
};
}
match deliver_to_runtime(record, sent.render(), wake).await? {
RuntimeDelivery::Steered => {
mailbox
.deliver_read(&sent)
.map_err(|error| error.to_string())?;
final_reply = answers;
Delivered::Steered
}
RuntimeDelivery::Started => {
mailbox
.deliver_read(&sent)
.map_err(|error| error.to_string())?;
final_reply = answers;
Delivered::Started
}
RuntimeDelivery::NotWoken => {
mailbox
.deliver(envelope)
.map_err(|error| error.to_string())?;
Delivered::Queued
}
}
}
Door::Native(session) => {
let busy = session.status != Some(ClaudePeerStatus::Idle);
if !wake && !busy {
return Ok(Err(Refused::CannotQueueNative));
}
let mut envelope = envelope.clone();
if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
envelope.reply_via = ReplyVia::Tool;
}
let envelope = &envelope;
let text = envelope.render();
if text.len() > MAX_RELAYED_BYTES {
return Ok(Err(Refused::TooLong(text.len())));
}
match send_through_relay(
&envelope.from,
&envelope.from_name,
&session.name,
text,
&envelope.id,
)
.await
{
RelayReceipt::Delivered { .. } => {
mailbox
.deliver_read(envelope)
.map_err(|error| error.to_string())?;
Delivered::Native { busy }
}
RelayReceipt::Failed { detail } => return Err(detail),
}
}
Door::Hook => {
mailbox
.deliver(envelope)
.map_err(|error| error.to_string())?;
Delivered::Hooked
}
Door::Stored => {
mailbox
.deliver(envelope)
.map_err(|error| error.to_string())?;
Delivered::Stored
}
Door::Operator => {
mailbox
.deliver(envelope)
.map_err(|error| error.to_string())?;
Delivered::Operator
}
};
let notice = notify_when_idle && !matches!(door, Door::Operator);
if notice || final_reply {
let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
subscription.notice = notice;
subscription.final_reply = final_reply;
mailbox
.subscribe_idle(&subscription)
.map_err(|error| error.to_string())?;
if let Ok(program) = supercode_program() {
crate::claude_relay::ensure_machine_daemon(&program)
.await
.ok();
}
}
Ok(Ok(delivered))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum UserTurn {
Steered,
Started,
Typed,
Waiting,
}
impl UserTurn {
pub const fn as_str(self) -> &'static str {
match self {
Self::Steered => "steered",
Self::Started => "started",
Self::Typed => "typed",
Self::Waiting => "waiting",
}
}
}
pub fn daemon_pane(session: &LiveSession) -> Option<String> {
let name = session.tmux.as_deref()?.split(':').next()?;
let rest = name.strip_prefix("p_")?;
(!rest.is_empty()
&& rest
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
.then(|| name.to_string())
}
pub async fn deliver_user_turn(
homes: &HarnessHomes,
envelope: &Envelope,
to: &MailAddress,
) -> Result<UserTurn, String> {
let session = LiveSessions::read(homes)
.sessions
.into_iter()
.find(|session| &session.address == to)
.ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
if let Door::Runtime(record) = &session.door {
let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
mailbox
.deliver_read(envelope)
.map_err(|error| error.to_string())?;
return Ok(match delivered {
RuntimeDelivery::Steered => UserTurn::Steered,
_ => UserTurn::Started,
});
}
let Some(pane) = daemon_pane(&session) else {
let name = session.name.split('@').next().unwrap_or(&session.name);
return Err(format!(
"{name} runs outside supercode, where nothing can speak as its user. Open it in a \
pane with `supercode open {name}`; nothing was sent."
));
};
mailbox
.deliver(envelope)
.map_err(|error| error.to_string())?;
let typed = type_user_turns(&mailbox, &pane).await;
Ok(if typed.contains(&envelope.id) {
UserTurn::Typed
} else {
UserTurn::Waiting
})
}
pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
let mut typed = Vec::new();
for waiting in mailbox.user_turns().unwrap_or_default() {
let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
break;
};
match submit_when_composer_empty(pane, &stored.envelope.body).await {
Ok(true) => {
mailbox.acknowledge(&stored).ok();
typed.push(stored.envelope.id.clone());
}
Ok(false) => {
mailbox.release(&stored).ok();
break;
}
Err(error) => {
mailbox.release(&stored).ok();
eprintln!(
"supercode: the user's turn {} for {} waits: {error}",
stored.envelope.id,
mailbox.address()
);
break;
}
}
}
typed
}
async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
let entry = crate::teams_entry().map_err(|error| error.to_string())?;
let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
.ok()
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| "node".into());
let output = tokio::process::Command::new(node)
.arg(entry)
.args(["input", pane, text, "--when-composer-empty"])
.stdin(std::process::Stdio::null())
.output()
.await
.map_err(|error| error.to_string())?;
if !output.status.success() {
return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
&output.stderr,
)));
}
let answer: serde_json::Value =
serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
Ok(answer["delivered"] == true)
}