use std::collections::{BTreeMap, HashMap};
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use crate::claude_peer::{read_registry, registry_dir, ClaudePeerSession};
use crate::mailbox::{mail_root, Envelope, MailAddress, MailKind, ReplyVia};
use crate::HarnessHomes;
pub const RELAY_MODEL: &str = "haiku";
pub const RELAY_NAME_PREFIX: &str = "sc-";
pub const RELAY_STATUS_LINE: &str = "relay, not the session's status";
pub const SEND_TIMEOUT: Duration = Duration::from_secs(90);
const RELAY_TOOLS: &str = "ListAgents,SendMessage";
const NATIVE_PEER_PREAMBLE: &str = "Another Claude session sent a message:";
const NATIVE_PEER_OPENING: &str = "<cross-session-message ";
const NATIVE_PEER_CLOSING: &str = "</cross-session-message>";
const NATIVE_IDLE_NOTICE: &str = "[Cross-session idle notice]";
const NATIVE_DELIVERY_NOTICE: &str = "[Cross-session delivery notice]";
pub fn relay_name(name: &str, machine: &str) -> String {
let base = name.strip_prefix(RELAY_NAME_PREFIX).unwrap_or(name);
format!("{RELAY_NAME_PREFIX}{base}-on-{machine}")
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueuedSend {
pub to: String,
pub message: String,
}
pub fn gate_decision(hook_input: &Value, queued: Option<&QueuedSend>) -> Option<String> {
let tool = hook_input
.get("tool_name")
.and_then(Value::as_str)
.unwrap_or_default();
if tool == "ListAgents" {
return None;
}
if tool != "SendMessage" {
return Some(format!("a relay may not use {tool}"));
}
let Some(queued) = queued else {
return Some(
"nothing is queued to send; mail for the host is not answered by the relay".into(),
);
};
let input = hook_input.get("tool_input").cloned().unwrap_or(Value::Null);
let to = input.get("to").and_then(Value::as_str).unwrap_or_default();
let message = input
.get("message")
.and_then(Value::as_str)
.unwrap_or_default();
if to != queued.to {
return Some(format!(
"not the queued outbound message: `to` is {to:?}, the queue says {:?}",
queued.to
));
}
if message != queued.message {
let at = message
.char_indices()
.zip(queued.message.chars())
.find(|((_, sent), queued)| sent != queued)
.map(|((index, _), _)| index)
.unwrap_or_else(|| message.len().min(queued.message.len()));
return Some(format!(
"not the queued outbound message: `message` ({} bytes) differs from the queue ({} \
bytes) at byte {at}: sent {:?}, queued {:?}",
message.len(),
queued.message.len(),
message
.get(at..)
.unwrap_or_default()
.chars()
.take(40)
.collect::<String>(),
queued
.message
.get(at..)
.unwrap_or_default()
.chars()
.take(40)
.collect::<String>(),
));
}
None
}
pub fn gate_denial(reason: &str) -> Value {
json!({
"hookSpecificOutput": {
"hookEventName": "PreToolUse",
"permissionDecision": "deny",
"permissionDecisionReason": reason,
}
})
}
#[derive(Debug, Clone)]
pub struct RelayPaths {
pub directory: PathBuf,
pub queue: PathBuf,
pub receipt: PathBuf,
pub settings: PathBuf,
pub record: PathBuf,
pub sent: PathBuf,
pub lock: PathBuf,
}
impl RelayPaths {
pub fn new(root: &Path, relay_name: &str) -> Self {
let hash = blake3::hash(relay_name.as_bytes()).to_hex();
Self::in_directory(root.join("relays").join(&hash[..24]))
}
pub fn in_directory(directory: PathBuf) -> Self {
Self {
queue: directory.join("queue.json"),
receipt: directory.join("receipt.json"),
settings: directory.join("settings.json"),
record: directory.join("relay.json"),
sent: directory.join("sent.json"),
lock: directory.join("send.lock"),
directory,
}
}
pub fn last_sent(&self) -> HashMap<String, String> {
std::fs::read(&self.sent)
.ok()
.and_then(|bytes| serde_json::from_slice(&bytes).ok())
.unwrap_or_default()
}
}
#[derive(Debug, Clone)]
pub struct RelaySpec {
pub represented: MailAddress,
pub represented_name: String,
pub name: String,
pub paths: RelayPaths,
pub program: PathBuf,
}
impl RelaySpec {
pub fn for_sender(represented: &MailAddress, represented_name: &str) -> std::io::Result<Self> {
let base = represented_name
.split('@')
.next()
.unwrap_or(represented_name);
let name = relay_name(base, &represented.machine);
Ok(Self {
represented: represented.clone(),
represented_name: represented_name.to_string(),
paths: RelayPaths::new(&mail_root(), &name),
name,
program: supercode_program()?,
})
}
}
pub fn relay_arguments(spec: &RelaySpec) -> Vec<String> {
vec![
"--print".into(),
"--input-format".into(),
"stream-json".into(),
"--output-format".into(),
"stream-json".into(),
"--verbose".into(),
"--permission-prompt-tool".into(),
"stdio".into(),
"--model".into(),
RELAY_MODEL.into(),
"--name".into(),
spec.name.clone(),
"--permission-mode".into(),
"bypassPermissions".into(),
"--setting-sources".into(),
"project".into(),
"--settings".into(),
spec.paths.settings.to_string_lossy().into_owned(),
"--tools".into(),
RELAY_TOOLS.into(),
"--no-session-persistence".into(),
]
}
pub fn relay_environment(spec: &RelaySpec, endpoint: &str) -> BTreeMap<String, String> {
let key = spec
.paths
.directory
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_default();
BTreeMap::from([
("ANTHROPIC_BASE_URL".to_string(), endpoint.to_string()),
("ANTHROPIC_API_KEY".to_string(), key),
(
"CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC".to_string(),
"1".to_string(),
),
])
}
pub fn relay_settings(spec: &RelaySpec) -> Value {
let program = shell_quote(&spec.program.to_string_lossy());
let directory = shell_quote(&spec.paths.directory.to_string_lossy());
let represented = shell_quote(&spec.represented.to_string());
let command = |verb: &str| format!("{program} message {verb} {directory}");
json!({
"hooks": {
"PreToolUse": [{
"matcher": ".*",
"hooks": [{"type": "command", "command": command("gate")}],
}],
"PostToolUse": [{
"matcher": "SendMessage",
"hooks": [{"type": "command", "command": command("relay-receipt")}],
}],
"Stop": [{
"hooks": [{"type": "command", "command": command("relay-receipt")}],
}],
"StopFailure": [{
"hooks": [{"type": "command", "command": command("relay-receipt")}],
}],
"UserPromptSubmit": [{
"hooks": [{
"type": "command",
"command": format!("{program} message relay-inbound {represented} {directory}"),
}],
}],
}
})
}
fn shell_quote(value: &str) -> String {
format!("'{}'", value.replace('\'', "'\\''"))
}
pub fn is_inbound_prompt(prompt: &str) -> bool {
let trimmed = prompt.trim_start();
trimmed.starts_with(NATIVE_PEER_OPENING)
|| trimmed.starts_with(NATIVE_IDLE_NOTICE)
|| trimmed.starts_with(NATIVE_DELIVERY_NOTICE)
}
pub fn inbound_block() -> Value {
json!({
"decision": "block",
"reason": "filed in the mailbox of the session this relay speaks for",
})
}
pub fn send_turn(send: &QueuedSend) -> String {
format!(
"Send one message with SendMessage.\n\
to: {}\n\
message: the exact text between the markers, without the markers\n\
---BEGIN MESSAGE---\n{}\n---END MESSAGE---",
send.to, send.message
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NativeEnvelope {
pub from: String,
pub from_name: Option<String>,
pub body: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RelayEvent {
Peer(NativeEnvelope),
IdleNotice(String),
DeliveryNotice(String),
}
pub fn resolve_native_sender(
registry: &[ClaudePeerSession],
native_from: &str,
) -> Option<ClaudePeerSession> {
let socket = native_from.strip_prefix("uds:").unwrap_or(native_from);
registry
.iter()
.find(|session| session.socket_path.to_string_lossy() == socket)
.cloned()
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case", tag = "outcome")]
pub enum RelayReceipt {
Delivered {
detail: String,
native_msg_id: Option<String>,
},
Failed {
detail: String,
},
}
pub fn read_receipt(value: &Value) -> RelayReceipt {
if let Some(reason) = value.get("denied").and_then(Value::as_str) {
return RelayReceipt::Failed {
detail: format!("the relay's gate refused the send: {reason}"),
};
}
if value.get("turn_ended").and_then(Value::as_bool) == Some(true) {
let said = value
.get("said")
.and_then(Value::as_str)
.filter(|said| !said.trim().is_empty())
.map(|said| format!("; it said: {}", said.trim()))
.unwrap_or_default();
return RelayReceipt::Failed {
detail: format!(
"the Claude relay ended its turn without sending; nothing was sent{said}"
),
};
}
let response = value.get("tool_response").cloned().unwrap_or(Value::Null);
let parsed = match &response {
Value::String(text) => serde_json::from_str::<Value>(text).unwrap_or(response.clone()),
Value::Array(blocks) => blocks
.iter()
.find_map(|block| block.get("text").and_then(Value::as_str))
.and_then(|text| serde_json::from_str::<Value>(text).ok())
.unwrap_or(response.clone()),
_ => response.clone(),
};
if parsed.get("success").and_then(Value::as_bool) != Some(true) {
return RelayReceipt::Failed {
detail: format!("Claude did not report the send as successful: {parsed}"),
};
}
RelayReceipt::Delivered {
detail: parsed
.get("message")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
native_msg_id: parsed
.get("msg_id")
.and_then(Value::as_str)
.map(str::to_string),
}
}
pub async fn send_through_relay(
sender: &MailAddress,
sender_name: &str,
to: &str,
message: String,
message_id: &str,
) -> RelayReceipt {
let failed = |detail: String| RelayReceipt::Failed { detail };
let spec = match RelaySpec::for_sender(sender, sender_name) {
Ok(spec) => spec,
Err(error) => return failed(error.to_string()),
};
if let Err(error) = std::fs::create_dir_all(&spec.paths.directory) {
return failed(error.to_string());
}
let lock = match tokio::task::spawn_blocking({
let path = spec.paths.lock.clone();
move || SendLock::acquire(&path)
})
.await
{
Ok(Ok(lock)) => lock,
Ok(Err(error)) => return failed(format!("could not take the relay's send lock: {error}")),
Err(error) => return failed(error.to_string()),
};
let runtime = match ensure_relay_runtime(&spec).await {
Ok(runtime) => runtime,
Err(detail) => return failed(detail),
};
let queued = QueuedSend {
to: to.to_string(),
message,
};
std::fs::remove_file(&spec.paths.receipt).ok();
if let Err(error) = std::fs::write(
&spec.paths.queue,
serde_json::to_vec(&queued).unwrap_or_default(),
) {
return failed(error.to_string());
}
let delivered =
crate::runtime_mail::deliver_to_runtime(&runtime, send_turn(&queued), true).await;
let receipt = match delivered {
Err(detail) => failed(format!("could not reach the Claude relay: {detail}")),
Ok(_) => wait_for_receipt(&spec.paths.receipt).await,
};
std::fs::remove_file(&spec.paths.queue).ok();
if matches!(receipt, RelayReceipt::Delivered { .. }) {
let mut sent = spec.paths.last_sent();
sent.insert(to.to_string(), message_id.to_string());
std::fs::write(
&spec.paths.sent,
serde_json::to_vec(&sent).unwrap_or_default(),
)
.ok();
}
drop(lock);
receipt
}
async fn wait_for_receipt(path: &Path) -> RelayReceipt {
let started = Instant::now();
while started.elapsed() < SEND_TIMEOUT {
if let Some(value) = std::fs::read(path)
.ok()
.and_then(|bytes| serde_json::from_slice::<Value>(&bytes).ok())
{
return read_receipt(&value);
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
RelayReceipt::Failed {
detail: format!(
"the Claude relay did not confirm the send within {} seconds; it may still arrive",
SEND_TIMEOUT.as_secs()
),
}
}
async fn ensure_relay_runtime(
spec: &RelaySpec,
) -> Result<crate::live_runtime::LiveRuntimeRecord, String> {
#[derive(Serialize, Deserialize)]
struct Record {
runtime_id: String,
#[serde(default)]
endpoint: String,
}
let endpoint = relay_endpoint(&spec.program).await?;
if let Some(record) = std::fs::read(&spec.paths.record)
.ok()
.and_then(|bytes| serde_json::from_slice::<Record>(&bytes).ok())
{
if let Some(runtime) =
crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id)
{
if record.endpoint == endpoint {
return Ok(runtime);
}
end_relay_process(&spec.name);
}
}
std::fs::write(
&spec.paths.settings,
serde_json::to_vec_pretty(&relay_settings(spec)).unwrap_or_default(),
)
.map_err(|error| error.to_string())?;
let params = json!({
"harness": "claude-code",
"launch": {
"program": "claude",
"arguments": relay_arguments(spec),
"env": relay_environment(spec, &endpoint),
},
"cwd": spec.paths.directory,
});
let result = machine_rpc(&spec.program, "runtimes.start", ¶ms).await?;
let runtime_id = result
.pointer("/handle/runtime_id")
.and_then(Value::as_str)
.ok_or_else(|| format!("the machine daemon did not start the relay: {result}"))?
.to_string();
std::fs::write(
&spec.paths.record,
serde_json::to_vec(&Record {
runtime_id: runtime_id.clone(),
endpoint,
})
.unwrap_or_default(),
)
.map_err(|error| error.to_string())?;
crate::runtime_mail::controlled_runtime("claude-code", &runtime_id)
.ok_or_else(|| "the relay started but registered no live runtime".to_string())
}
fn end_relay_process(name: &str) {
let registry = registry_dir(&HarnessHomes::default());
for session in read_registry(®istry) {
#[cfg(unix)]
if session.name == name {
unsafe {
libc::kill(session.pid as libc::pid_t, libc::SIGTERM);
}
}
}
}
async fn relay_endpoint(program: &Path) -> Result<String, String> {
if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
return Ok(url);
}
ensure_machine_daemon(program).await?;
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
return Ok(url);
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
Err(
"the relay endpoint is not answering; the machine daemon's `supercode message watch` \
serves it (see mail/machine-daemon.log)"
.into(),
)
}
pub async fn machine_rpc(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
let call = || {
let mut command = tokio::process::Command::new(program);
command
.args(["teams", "rpc", method, ¶ms.to_string()])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped());
command.output()
};
let output = call().await.map_err(|error| error.to_string())?;
if output.status.success() {
return serde_json::from_slice(&output.stdout)
.map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
}
ensure_machine_daemon(program).await?;
let output = call().await.map_err(|error| error.to_string())?;
if output.status.success() {
return serde_json::from_slice(&output.stdout)
.map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
}
Err(error_line(&String::from_utf8_lossy(&output.stderr)))
}
pub async fn ensure_machine_daemon(program: &Path) -> Result<(), String> {
let describe = || {
tokio::process::Command::new(program)
.args(["teams", "describe"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status()
};
if describe().await.is_ok_and(|status| status.success()) {
return Ok(());
}
let root = mail_root();
std::fs::create_dir_all(&root).map_err(|error| error.to_string())?;
let log = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(root.join("machine-daemon.log"))
.map_err(|error| error.to_string())?;
let mut command = std::process::Command::new(program);
command
.args(["teams", "machine", "start", "--supercode"])
.arg(program)
.stdin(std::process::Stdio::null())
.stdout(log.try_clone().map_err(|error| error.to_string())?)
.stderr(log);
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
command.process_group(0);
}
command
.spawn()
.map_err(|error| format!("could not start the machine daemon: {error}"))?;
for _ in 0..50 {
tokio::time::sleep(Duration::from_millis(200)).await;
if describe().await.is_ok_and(|status| status.success()) {
return Ok(());
}
}
Err(format!(
"the machine daemon did not start within 10 seconds; see {}",
root.join("machine-daemon.log").display()
))
}
fn error_line(stderr: &str) -> String {
let lines: Vec<&str> = stderr
.lines()
.map(str::trim)
.filter(|line| !line.is_empty())
.collect();
lines
.iter()
.find(|line| line.starts_with("Error") || line.starts_with("error"))
.or(lines.last())
.copied()
.unwrap_or("the machine daemon failed without saying why")
.chars()
.take(300)
.collect()
}
struct SendLock {
_file: std::fs::File,
}
impl SendLock {
fn acquire(path: &Path) -> std::io::Result<Self> {
let file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(path)?;
#[cfg(unix)]
{
use std::os::unix::io::AsRawFd;
if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
return Err(std::io::Error::last_os_error());
}
}
Ok(Self { _file: file })
}
}
pub fn file_inbound_prompt(
homes: &HarnessHomes,
represented: &MailAddress,
prompt: &str,
last_sent: Option<&HashMap<String, String>>,
) {
match parse_inbound_text(prompt) {
Some(RelayEvent::Peer(native)) => {
let registry = read_registry(®istry_dir(homes));
let sender = resolve_native_sender(®istry, &native.from);
let in_reply_to = sender
.as_ref()
.and_then(|session| last_sent.and_then(|sent| sent.get(&session.name).cloned()));
file_inbound(represented, &native, sender.as_ref(), in_reply_to);
}
Some(RelayEvent::IdleNotice(text) | RelayEvent::DeliveryNotice(text)) => {
file_notice(represented, &text);
}
_ => {}
}
}
pub fn parse_inbound_text(text: &str) -> Option<RelayEvent> {
let trimmed = text.trim_start();
if trimmed.starts_with(NATIVE_IDLE_NOTICE) {
return Some(RelayEvent::IdleNotice(trimmed.to_string()));
}
if trimmed.starts_with(NATIVE_DELIVERY_NOTICE) {
return Some(RelayEvent::DeliveryNotice(trimmed.to_string()));
}
if !trimmed.starts_with(NATIVE_PEER_PREAMBLE) && !trimmed.starts_with(NATIVE_PEER_OPENING) {
return None;
}
let start = text.find(NATIVE_PEER_OPENING)?;
let header_end = start + text[start..].find(">\n")?;
let header = &text[start + NATIVE_PEER_OPENING.len()..header_end];
let body_start = header_end + 2;
let body_end = text.rfind(&format!("\n{NATIVE_PEER_CLOSING}"))?;
if body_end < body_start {
return None;
}
let attributes = parse_attributes(header);
Some(RelayEvent::Peer(NativeEnvelope {
from: attributes.get("from")?.clone(),
from_name: attributes.get("from-name").cloned(),
body: text[body_start..body_end].to_string(),
}))
}
fn parse_attributes(header: &str) -> BTreeMap<String, String> {
let mut attributes = BTreeMap::new();
let mut rest = header;
while let Some(equals) = rest.find("=\"") {
let key = rest[..equals].trim().to_string();
let value_start = equals + 2;
let Some(length) = rest[value_start..].find('"') else {
break;
};
attributes.insert(key, rest[value_start..value_start + length].to_string());
rest = &rest[value_start + length + 1..];
}
attributes
}
fn deliver_home(represented: &MailAddress, envelope: &Envelope) {
if let Err(error) = crate::mailbox::deliver_to(represented, envelope) {
eprintln!(
"could not deliver {} to {represented}: {error}; kept in this machine's mailbox for it",
envelope.id
);
if let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), represented) {
mailbox.deliver(envelope).ok();
}
return;
}
if represented.machine == crate::mailbox::local_machine_name() {
if let Ok(program) = supercode_program() {
std::process::Command::new(program)
.args(["message", "push", &represented.to_string(), &envelope.id])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.ok();
}
}
}
fn file_inbound(
represented: &MailAddress,
native: &NativeEnvelope,
sender: Option<&ClaudePeerSession>,
in_reply_to: Option<String>,
) {
let machine = crate::mailbox::local_machine_name();
let (from, from_name) = match sender {
Some(session) => (
MailAddress::new(&machine, "claude-code", &session.session_id),
format!("{}@{machine}", session.name),
),
None => (
MailAddress::new(&machine, "claude-code", "unknown"),
native
.from_name
.clone()
.unwrap_or_else(|| "an unknown Claude session".into()),
),
};
let Ok(from) = from else { return };
let Ok(mut envelope) = Envelope::new(
from,
from_name,
MailKind::Peer,
ReplyVia::Command,
native.body.clone(),
) else {
return;
};
envelope.native_from = Some(native.from.clone());
if let Some(in_reply_to) = in_reply_to {
envelope.in_reply_to = Some(in_reply_to);
envelope.in_reply_to_inferred = true;
}
deliver_home(represented, &envelope);
}
fn file_notice(represented: &MailAddress, text: &str) {
let Ok(from) = MailAddress::new(
crate::mailbox::local_machine_name(),
"claude-code",
"notice",
) else {
return;
};
let Ok(envelope) = Envelope::new(
from,
"Claude Code",
MailKind::Notice,
ReplyVia::None,
text.to_string(),
) else {
return;
};
deliver_home(represented, &envelope);
}
pub fn supercode_program() -> std::io::Result<PathBuf> {
let current = std::env::current_exe()?;
if current.file_stem().and_then(|stem| stem.to_str()) == Some("supercode") {
return Ok(current);
}
std::env::var_os("PATH")
.iter()
.flat_map(std::env::split_paths)
.map(|directory| directory.join("supercode"))
.find(|candidate| candidate.is_file())
.ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::NotFound,
"the supercode program is not on PATH; Claude relays and the machine daemon need it",
)
})
}