use async_trait::async_trait;
use serde_json::{json, Value};
use crate::error::{Error, Result};
use crate::mail_route::{process_ancestry, resolve_caller, LiveSessions};
use crate::mail_send::SendOptions;
use crate::mailbox::{mail_root, Mailbox};
use crate::tools::{Tool, ToolContext, ToolRegistry};
use crate::HarnessHomes;
static STARTED_AS: std::sync::OnceLock<Option<(std::path::PathBuf, (u64, u64))>> =
std::sync::OnceLock::new();
#[cfg(unix)]
fn file_identity(path: &std::path::Path) -> Option<(u64, u64)> {
use std::os::unix::fs::MetadataExt;
std::fs::metadata(path)
.ok()
.map(|metadata| (metadata.dev(), metadata.ino()))
}
#[cfg(not(unix))]
fn file_identity(_path: &std::path::Path) -> Option<(u64, u64)> {
None
}
fn started_as() -> &'static Option<(std::path::PathBuf, (u64, u64))> {
STARTED_AS.get_or_init(|| {
let path = std::env::current_exe().ok()?;
let identity = file_identity(&path)?;
Some((path, identity))
})
}
fn replaced_by() -> Option<std::path::PathBuf> {
let (path, identity) = started_as().as_ref()?;
let now = file_identity(path)?;
(now != *identity).then(|| path.clone())
}
pub fn registry() -> ToolRegistry {
started_as();
let mut registry = ToolRegistry::new();
registry.register(ListAgents);
registry.register(SendMessage);
registry.register(ReadMessages);
registry
}
fn caller(tool: &str, homes: &HarnessHomes) -> Result<crate::mail_route::Caller> {
resolve_caller(homes, &process_ancestry()).map_err(|message| Error::tool(tool, message))
}
struct ListAgents;
#[async_trait]
impl Tool for ListAgents {
fn name(&self) -> &str {
"list_agents"
}
fn description(&self) -> &str {
"List the coding-agent sessions you can message, of every harness (Claude Code, Codex, \
sessions supercode hosts): each one's name, harness, status (busy, idle, hosted), how it \
is reached, and address. `query` keeps the ones whose name or address contains it."
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"query": {"type": "string", "description": "Keep sessions whose name or address contains this."}
},
"additionalProperties": false
})
}
async fn execute(&self, args: Value, _ctx: &ToolContext) -> Result<String> {
let query = args
.get("query")
.and_then(Value::as_str)
.map(str::to_lowercase);
let sessions = LiveSessions::read(&HarnessHomes::default());
let rows: Vec<String> = sessions
.all()
.iter()
.filter(|session| {
query.as_deref().is_none_or(|query| {
session.name.to_lowercase().contains(query)
|| session.address.to_string().to_lowercase().contains(query)
})
})
.map(|session| {
format!(
"{} ({}, {}, {}) {}",
session.name,
session.address.harness,
session.status,
session.door.name(),
session.address
)
})
.collect();
Ok(if rows.is_empty() {
"No session to message matches.".into()
} else {
rows.join("\n")
})
}
}
struct SendMessage;
#[async_trait]
impl Tool for SendMessage {
fn name(&self) -> &str {
"send_message"
}
fn description(&self) -> &str {
"Send a message to another coding-agent session, of any harness, on this machine or \
another of your team's (`name@machine`). `to` is a name from list_agents, `name@machine`, \
or an address. To answer a message, send to its `from` with `reply_to` set to its id. \
The result says whether it was delivered, queued, held for approval, refused or stored."
}
fn parameters(&self) -> Value {
json!({
"type": "object",
"properties": {
"to": {"type": "string", "description": "A session name, name@machine, or address."},
"message": {"type": "string", "description": "The message."},
"subject": {"type": "string", "description": "A short subject shown on the unopened envelope."},
"reply_to": {"type": "string", "description": "Id of the message this answers."},
"notify_when_idle": {"type": "boolean", "description": "Also get one notice when the receiver's next turn ends."}
},
"required": ["to", "message"],
"additionalProperties": false
})
}
async fn execute(&self, args: Value, _ctx: &ToolContext) -> Result<String> {
let text = |key: &str| args.get(key).and_then(Value::as_str).map(str::to_string);
let (Some(to), Some(message)) = (text("to"), text("message")) else {
return Err(Error::tool(self.name(), "`to` and `message` are required"));
};
if message.trim().is_empty() {
return Err(Error::tool(
self.name(),
"Nothing was sent: the message is empty.",
));
}
if let Some(installed) = replaced_by() {
return send_through(&installed, self.name(), &to, &message, &args).await;
}
let homes = HarnessHomes::default();
let caller = caller(self.name(), &homes)?;
let options = SendOptions {
subject: text("subject"),
in_reply_to: text("reply_to"),
notify_when_idle: args
.get("notify_when_idle")
.and_then(Value::as_bool)
.unwrap_or(false),
..Default::default()
};
let outcome = crate::mail_send::send(&homes, &caller, &to, &message, options)
.await
.map_err(|error| Error::tool(self.name(), error.to_string()))?;
if outcome.code == 0 {
Ok(outcome.text)
} else {
Err(Error::tool(self.name(), outcome.text))
}
}
}
async fn send_through(
installed: &std::path::Path,
tool: &str,
to: &str,
message: &str,
args: &Value,
) -> Result<String> {
use tokio::io::AsyncWriteExt;
let mut command = tokio::process::Command::new(installed);
command.args(["message", "send", to]);
if let Some(subject) = args.get("subject").and_then(Value::as_str) {
command.args(["--subject", subject]);
}
if let Some(answered) = args.get("reply_to").and_then(Value::as_str) {
command.args(["--re", answered]);
}
if args.get("notify_when_idle").and_then(Value::as_bool) == Some(true) {
command.arg("--notify-when-idle");
}
command
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped());
let mut child = command.spawn().map_err(|error| {
Error::tool(
tool,
format!("could not start {}: {error}", installed.display()),
)
})?;
if let Some(mut stdin) = child.stdin.take() {
stdin
.write_all(message.as_bytes())
.await
.map_err(|error| Error::tool(tool, error.to_string()))?;
}
let output = child
.wait_with_output()
.await
.map_err(|error| Error::tool(tool, error.to_string()))?;
let said = |bytes: &[u8]| String::from_utf8_lossy(bytes).trim().to_string();
if output.status.success() {
Ok(said(&output.stdout))
} else {
let stderr = said(&output.stderr);
Err(Error::tool(
tool,
if stderr.is_empty() {
said(&output.stdout)
} else {
stderr
},
))
}
}
struct ReadMessages;
#[async_trait]
impl Tool for ReadMessages {
fn name(&self) -> &str {
"read_messages"
}
fn description(&self) -> &str {
"Read your unread messages from other sessions, each with how to answer it. Messages \
usually arrive in your conversation by themselves; this reads any that are waiting."
}
fn parameters(&self) -> Value {
json!({"type": "object", "properties": {}, "additionalProperties": false})
}
async fn execute(&self, _args: Value, _ctx: &ToolContext) -> Result<String> {
let homes = HarnessHomes::default();
let caller = caller(self.name(), &homes)?;
let mailbox = Mailbox::open(&mail_root(), &caller.address)
.map_err(|error| Error::tool(self.name(), error.to_string()))?;
let claimed = mailbox
.claim_unread()
.map_err(|error| Error::tool(self.name(), error.to_string()))?;
let text = if claimed.is_empty() {
"No unread messages.".to_string()
} else {
claimed
.iter()
.map(|stored| stored.envelope.render())
.collect::<Vec<_>>()
.join("\n\n")
};
for stored in &claimed {
mailbox
.acknowledge(stored)
.map_err(|error| Error::tool(self.name(), error.to_string()))?;
}
Ok(text)
}
}