use std::collections::{BTreeMap, VecDeque};
use std::fs::{self, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
use serde::{Deserialize, Serialize};
use crate::error::{Error, Result};
const TRANSCRIPT_PAGE_BYTES: usize = 32 * 1024;
const REVIVE_TRANSCRIPT_BUDGET: usize = 16 * 1024;
const MAX_QUEUE_PER_CHILD: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub enum ChildStatus {
Starting,
Running,
Done,
Failed,
Cancelled,
Killed,
}
impl ChildStatus {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Starting => "starting",
Self::Running => "running",
Self::Done => "done",
Self::Failed => "failed",
Self::Cancelled => "cancelled",
Self::Killed => "killed",
}
}
#[must_use]
pub const fn settled(self) -> bool {
matches!(
self,
Self::Done | Self::Failed | Self::Cancelled | Self::Killed
)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ChildEntry {
pub id: String,
pub name: String,
pub task: String,
pub pid: Option<u32>,
pub status: ChildStatus,
pub started_ms: u64,
pub finished_ms: Option<u64>,
pub output_bytes: usize,
pub transcript_path: PathBuf,
pub steer_path: PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub revived_from: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct BusMessage {
pub seq: u64,
pub from: String,
pub to: String,
pub body: String,
pub sent_ms: u64,
}
#[derive(Debug, Default)]
pub struct AgentHubRegistry {
entries: BTreeMap<String, ChildEntry>,
seq: u64,
bus: BTreeMap<String, VecDeque<BusMessage>>,
bus_seq: u64,
dir: Option<PathBuf>,
}
static REGISTRY: OnceLock<Mutex<AgentHubRegistry>> = OnceLock::new();
pub fn registry() -> &'static Mutex<AgentHubRegistry> {
REGISTRY.get_or_init(|| Mutex::new(AgentHubRegistry::default()))
}
fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| u64::try_from(d.as_millis()).unwrap_or(u64::MAX))
}
impl AgentHubRegistry {
#[doc(hidden)]
pub fn set_dir_for_tests(&mut self, dir: PathBuf) {
self.dir = Some(dir);
}
fn dir(&mut self) -> Result<PathBuf> {
let dir = self.dir.clone().unwrap_or_else(|| {
crate::config::Config::global_dir()
.join("agent-hub")
.join(std::process::id().to_string())
});
fs::create_dir_all(&dir).map_err(|e| {
Error::tool(
"hub",
format!("create agent-hub dir {}: {e}", dir.display()),
)
})?;
self.dir = Some(dir.clone());
Ok(dir)
}
pub fn register(&mut self, name: &str, task: &str) -> Result<ChildEntry> {
self.seq = self.seq.saturating_add(1);
let id = format!("{}-{}", sanitize_id(name), self.seq);
let dir = self.dir()?;
let entry = ChildEntry {
id: id.clone(),
name: name.to_string(),
task: truncate_chars(task, 500),
pid: None,
status: ChildStatus::Starting,
started_ms: now_ms(),
finished_ms: None,
output_bytes: 0,
transcript_path: dir.join(format!("{id}.transcript.jsonl")),
steer_path: dir.join(format!("{id}.steer")),
revived_from: None,
};
self.entries.insert(id, entry.clone());
Ok(entry)
}
pub fn mark_running(&mut self, id: &str, pid: u32) {
if let Some(entry) = self.entries.get_mut(id) {
entry.pid = Some(pid);
entry.status = ChildStatus::Running;
}
}
pub fn settle(&mut self, id: &str, status: ChildStatus) {
if let Some(entry) = self.entries.get_mut(id) {
entry.status = status;
entry.finished_ms = Some(now_ms());
}
}
pub fn append_transcript(&mut self, id: &str, line: &str) {
let Some(entry) = self.entries.get_mut(id) else {
return;
};
entry.output_bytes = entry.output_bytes.saturating_add(line.len());
if let Ok(mut file) = OpenOptions::new()
.create(true)
.append(true)
.open(&entry.transcript_path)
{
let _ = writeln!(file, "{line}");
}
}
#[must_use]
pub fn roster(&self) -> Vec<ChildEntry> {
self.entries.values().cloned().collect()
}
#[must_use]
pub fn get(&self, id: &str) -> Option<ChildEntry> {
self.entries.get(id).cloned()
}
pub fn transcript_page(&self, id: &str) -> Result<String> {
let entry = self
.entries
.get(id)
.ok_or_else(|| Error::validation(format!("hub: unknown child '{id}'")))?;
let raw = match fs::read_to_string(&entry.transcript_path) {
Ok(text) => text,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => String::new(),
Err(err) => {
return Err(Error::tool(
"hub",
format!("read transcript {}: {err}", entry.transcript_path.display()),
));
}
};
let tail = tail_bytes(&raw, TRANSCRIPT_PAGE_BYTES);
let mut vault = crate::secrets::SecretVault::default();
let (masked, _audit) = crate::secrets::obfuscate(&tail, &mut vault, &[]);
Ok(masked)
}
pub fn steer(&mut self, id: &str, from: &str, body: &str) -> Result<BusMessage> {
let entry = self
.entries
.get(id)
.ok_or_else(|| Error::validation(format!("hub: unknown child '{id}'")))?
.clone();
if entry.status.settled() {
return Err(Error::validation(format!(
"hub: cannot steer '{}' — status {}",
id,
entry.status.as_str()
)));
}
let message = self.enqueue_bus(id, from, body);
append_steer_line(&entry.steer_path, &message)?;
Ok(message)
}
pub fn bus_send(&mut self, to: &str, from: &str, body: &str) -> Result<BusMessage> {
self.steer(to, from, body)
}
#[must_use]
pub fn inbox(&self, id: &str) -> Vec<BusMessage> {
self.bus
.get(id)
.map(|q| q.iter().cloned().collect())
.unwrap_or_default()
}
fn enqueue_bus(&mut self, to: &str, from: &str, body: &str) -> BusMessage {
self.bus_seq = self.bus_seq.saturating_add(1);
let message = BusMessage {
seq: self.bus_seq,
from: from.to_string(),
to: to.to_string(),
body: body.to_string(),
sent_ms: now_ms(),
};
let queue = self.bus.entry(to.to_string()).or_default();
if queue.len() >= MAX_QUEUE_PER_CHILD {
queue.pop_front();
}
queue.push_back(message.clone());
message
}
pub fn mark_killed(&mut self, id: &str) {
self.settle(id, ChildStatus::Killed);
}
pub fn revive(&mut self, from_id: &str) -> Result<(ChildEntry, String)> {
let prior = self
.entries
.get(from_id)
.ok_or_else(|| Error::validation(format!("hub: unknown child '{from_id}'")))?
.clone();
if !prior.status.settled() {
return Err(Error::validation(format!(
"hub: cannot revive '{from_id}' — still {}",
prior.status.as_str()
)));
}
let transcript = fs::read_to_string(&prior.transcript_path).unwrap_or_default();
let tail = tail_bytes(&transcript, REVIVE_TRANSCRIPT_BUDGET);
let task = format!(
"{}\n\n[Continuation of a prior run ({}). Its transcript tail follows; \
pick up where it left off and finish the task.]\n{}",
prior.task,
prior.status.as_str(),
tail
);
let mut entry = self.register(&prior.name, &task)?;
entry.revived_from = Some(from_id.to_string());
self.entries.insert(entry.id.clone(), entry.clone());
Ok((entry, task))
}
pub fn cleanup_session_files() {
let dir = crate::config::Config::global_dir()
.join("agent-hub")
.join(std::process::id().to_string());
let _ = fs::remove_dir_all(dir);
}
}
fn append_steer_line(path: &Path, message: &BusMessage) -> Result<()> {
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(path)
.map_err(|e| Error::tool("hub", format!("append steer queue {}: {e}", path.display())))?;
let line = serde_json::to_string(message)
.map_err(|e| Error::validation(format!("serialize bus message: {e}")))?;
writeln!(file, "{line}")
.map_err(|e| Error::tool("hub", format!("write steer queue {}: {e}", path.display())))
}
pub fn drain_steer_file(path: &Path) -> Vec<String> {
let Ok(raw) = fs::read_to_string(path) else {
return Vec::new();
};
if raw.is_empty() {
return Vec::new();
}
let _ = fs::write(path, "");
raw.lines()
.filter_map(|line| {
serde_json::from_str::<BusMessage>(line)
.ok()
.map(|m| format!("[hub:{}] {}", m.from, m.body))
})
.collect()
}
fn sanitize_id(name: &str) -> String {
name.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect()
}
fn truncate_chars(text: &str, max: usize) -> String {
if text.chars().count() <= max {
return text.to_string();
}
text.chars().take(max).collect()
}
fn tail_bytes(text: &str, max: usize) -> String {
if text.len() <= max {
return text.to_string();
}
let start = text.len() - max;
let boundary = text[start..].find('\n').map_or(start, |i| start + i + 1);
text[boundary.min(text.len())..].to_string()
}
#[cfg(test)]
mod tests {
use super::*;
fn fresh_registry() -> AgentHubRegistry {
AgentHubRegistry::default()
}
#[test]
fn registry_fsm_running_to_done() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("scout", "inspect the code").expect("register");
assert_eq!(entry.status, ChildStatus::Starting);
reg.mark_running(&entry.id, 4242);
assert_eq!(
reg.get(&entry.id).expect("get").status,
ChildStatus::Running
);
reg.settle(&entry.id, ChildStatus::Done);
let settled = reg.get(&entry.id).expect("get");
assert_eq!(settled.status, ChildStatus::Done);
assert!(settled.finished_ms.is_some());
let _ = fs::remove_dir_all(&temp);
}
#[test]
fn steer_refuses_settled_child() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test2-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("scout", "task").expect("register");
reg.settle(&entry.id, ChildStatus::Failed);
let err = reg.steer(&entry.id, "parent", "hello").unwrap_err();
assert!(err.to_string().contains("cannot steer"));
let _ = fs::remove_dir_all(&temp);
}
#[test]
fn bus_preserves_delivery_order() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test3-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("worker", "task").expect("register");
reg.mark_running(&entry.id, 1);
reg.bus_send(&entry.id, "a", "first").expect("send 1");
reg.bus_send(&entry.id, "b", "second").expect("send 2");
let inbox = reg.inbox(&entry.id);
assert_eq!(inbox.len(), 2);
assert_eq!(inbox[0].body, "first");
assert_eq!(inbox[1].body, "second");
assert!(inbox[0].seq < inbox[1].seq);
let drained = drain_steer_file(&entry.steer_path);
assert_eq!(
drained,
vec!["[hub:a] first".to_string(), "[hub:b] second".to_string()]
);
assert!(drain_steer_file(&entry.steer_path).is_empty());
let _ = fs::remove_dir_all(&temp);
}
#[test]
fn transcript_page_redacts_secret_shapes() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test4-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("scout", "task").expect("register");
reg.append_transcript(&entry.id, "{\"note\":\"key is sk-ant-api03-AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\"}");
let page = reg.transcript_page(&entry.id).expect("page");
assert!(
!page.contains("sk-ant-api03-AAAA"),
"raw secret leaked: {page}"
);
assert!(page.contains("<pi-secret:"));
let _ = fs::remove_dir_all(&temp);
}
#[test]
fn revive_carries_transcript_context() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test5-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("scout", "original task").expect("register");
reg.append_transcript(
&entry.id,
"{\"type\":\"message_end\",\"text\":\"half-done\"}",
);
reg.settle(&entry.id, ChildStatus::Failed);
let (revived, task) = reg.revive(&entry.id).expect("revive");
assert_eq!(revived.revived_from.as_deref(), Some(entry.id.as_str()));
assert!(task.contains("original task"));
assert!(task.contains("half-done"));
assert!(task.contains("Continuation"));
let _ = fs::remove_dir_all(&temp);
}
#[test]
fn revive_refuses_running_child() {
let mut reg = fresh_registry();
let temp = std::env::temp_dir().join(format!("pi-agent-hub-test6-{}", std::process::id()));
reg.dir = Some(temp.clone());
let entry = reg.register("scout", "task").expect("register");
reg.mark_running(&entry.id, 9);
let err = reg.revive(&entry.id).unwrap_err();
assert!(err.to_string().contains("still running"));
let _ = fs::remove_dir_all(&temp);
}
}