use crate::config::HerdrSettings;
use serde::Serialize;
use std::{
ffi::{OsStr, OsString},
io,
path::{Path, PathBuf},
sync::{
Arc,
atomic::{AtomicU64, Ordering},
},
time::Duration,
};
const SOURCE: &str = "custom:magi-code";
const AGENT: &str = "magi-code";
const DEFAULT_SOCKET_RELATIVE: &str = ".config/herdr/herdr.sock";
static NEXT_REPORT_ID: AtomicU64 = AtomicU64::new(1);
static NEXT_REPORT_SEQ: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum HerdrAgentState {
Idle,
Working,
}
#[derive(Clone)]
pub(crate) struct HerdrReporter {
pane_id: String,
socket_path: PathBuf,
transport: Arc<dyn HerdrTransport>,
}
trait HerdrTransport: Send + Sync {
fn write_line(&self, socket_path: &Path, line: &str) -> io::Result<()>;
}
#[derive(Debug, Clone, Copy)]
struct UnixSocketTransport;
#[cfg(test)]
#[derive(Default)]
struct RecordingHerdrTransport {
lines: Arc<std::sync::Mutex<Vec<String>>>,
}
#[cfg(test)]
impl HerdrTransport for RecordingHerdrTransport {
fn write_line(&self, _socket_path: &Path, line: &str) -> io::Result<()> {
self.lines.lock().unwrap().push(line.to_string());
Ok(())
}
}
impl HerdrReporter {
pub(crate) fn from_env(settings: &HerdrSettings) -> Option<Self> {
Self::from_env_map(settings, RealEnv)
}
fn from_env_map<E: HerdrEnv>(settings: &HerdrSettings, env: E) -> Option<Self> {
if !settings.enabled || env.get_os("HERDR_ENV")? != OsStr::new("1") {
return None;
}
let pane_id = env
.get_os("HERDR_PANE_ID")?
.to_string_lossy()
.trim()
.to_string();
if pane_id.is_empty() {
return None;
}
let socket_path = resolve_socket_path(&env)?;
Some(Self::new_with_transport(
pane_id,
socket_path,
Arc::new(UnixSocketTransport),
))
}
fn new_with_transport(
pane_id: String,
socket_path: PathBuf,
transport: Arc<dyn HerdrTransport>,
) -> Self {
Self {
pane_id,
socket_path,
transport,
}
}
#[cfg(test)]
pub(crate) fn new_for_test() -> (Self, Arc<std::sync::Mutex<Vec<String>>>) {
let lines = Arc::new(std::sync::Mutex::new(Vec::new()));
let transport = Arc::new(RecordingHerdrTransport {
lines: Arc::clone(&lines),
});
(
Self::new_with_transport(
"pane-test".to_string(),
PathBuf::from("/tmp/herdr-test.sock"),
transport,
),
lines,
)
}
pub(crate) fn report(&self, state: HerdrAgentState, custom_status: &str, message: &str) {
let seq = NEXT_REPORT_SEQ.fetch_add(1, Ordering::Relaxed);
if let Ok(line) = report_agent_line(&self.pane_id, state, custom_status, message, seq) {
let _ = self.transport.write_line(&self.socket_path, &line);
}
}
pub(crate) fn report_thinking(&self) {
self.report(HerdrAgentState::Working, "thinking", "magi-code: thinking");
}
pub(crate) fn report_tool(&self, tool_name: &str) {
let tool_name = sanitize_tool_name(tool_name);
self.report(
HerdrAgentState::Working,
"tool",
&format!("magi-code: running {tool_name}"),
);
}
pub(crate) fn report_done(&self) {
self.report(HerdrAgentState::Idle, "done", "magi-code: done");
}
pub(crate) fn report_stopped(&self) {
self.report(HerdrAgentState::Idle, "stopped", "magi-code: stopped");
}
}
pub(crate) struct HerdrTurnReporter {
reporter: Option<HerdrReporter>,
finished: bool,
}
impl HerdrTurnReporter {
pub(crate) fn start(reporter: Option<HerdrReporter>) -> Self {
if let Some(reporter) = &reporter {
reporter.report_thinking();
}
Self {
reporter,
finished: false,
}
}
pub(crate) fn reporter(&self) -> Option<&HerdrReporter> {
self.reporter.as_ref()
}
pub(crate) fn done(&mut self) {
if let Some(reporter) = &self.reporter {
reporter.report_done();
}
self.finished = true;
}
}
impl Drop for HerdrTurnReporter {
fn drop(&mut self) {
if !self.finished
&& let Some(reporter) = &self.reporter
{
reporter.report_stopped();
}
}
}
#[derive(Serialize)]
struct JsonRpcRequest<'a> {
id: String,
method: &'static str,
params: PaneReportAgentParams<'a>,
}
#[derive(Serialize)]
struct PaneReportAgentParams<'a> {
pane_id: &'a str,
source: &'static str,
agent: &'static str,
state: HerdrAgentState,
message: &'a str,
custom_status: &'a str,
seq: u64,
}
fn report_agent_line(
pane_id: &str,
state: HerdrAgentState,
custom_status: &str,
message: &str,
seq: u64,
) -> serde_json::Result<String> {
let id = NEXT_REPORT_ID.fetch_add(1, Ordering::Relaxed);
let request = JsonRpcRequest {
id: format!("magi-code-{id}"),
method: "pane.report_agent",
params: PaneReportAgentParams {
pane_id,
source: SOURCE,
agent: AGENT,
state,
message,
custom_status,
seq,
},
};
let mut line = serde_json::to_string(&request)?;
line.push('\n');
Ok(line)
}
fn sanitize_tool_name(name: &str) -> String {
let mut sanitized = name
.chars()
.filter(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-'))
.take(48)
.collect::<String>();
if sanitized.is_empty() {
sanitized = "tool".to_string();
}
sanitized
}
fn resolve_socket_path<E: HerdrEnv>(env: &E) -> Option<PathBuf> {
if let Some(path) = env.get_os("HERDR_SOCKET_PATH") {
let path = PathBuf::from(path);
if !path.as_os_str().is_empty() {
return Some(path);
}
}
env.get_os("HOME")
.filter(|home| !home.is_empty())
.map(|home| PathBuf::from(home).join(DEFAULT_SOCKET_RELATIVE))
}
trait HerdrEnv {
fn get_os(&self, key: &str) -> Option<OsString>;
}
#[derive(Debug, Clone, Copy)]
struct RealEnv;
impl HerdrEnv for RealEnv {
fn get_os(&self, key: &str) -> Option<OsString> {
std::env::var_os(key)
}
}
impl HerdrTransport for UnixSocketTransport {
fn write_line(&self, socket_path: &Path, line: &str) -> io::Result<()> {
write_line_to_socket(socket_path, line)
}
}
#[cfg(unix)]
fn write_line_to_socket(socket_path: &Path, line: &str) -> io::Result<()> {
use std::{io::Write, os::unix::net::UnixStream};
let mut stream = UnixStream::connect(socket_path)?;
let _ = stream.set_write_timeout(Some(Duration::from_millis(100)));
stream.write_all(line.as_bytes())?;
stream.flush()
}
#[cfg(not(unix))]
fn write_line_to_socket(_socket_path: &Path, _line: &str) -> io::Result<()> {
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
#[derive(Default)]
struct MapEnv(Vec<(&'static str, OsString)>);
impl MapEnv {
fn with(mut self, key: &'static str, value: impl Into<OsString>) -> Self {
self.0.push((key, value.into()));
self
}
}
impl HerdrEnv for MapEnv {
fn get_os(&self, key: &str) -> Option<OsString> {
self.0
.iter()
.find(|(candidate, _)| *candidate == key)
.map(|(_, value)| value.clone())
}
}
#[derive(Default)]
struct RecordingTransport {
lines: Mutex<Vec<String>>,
}
impl HerdrTransport for RecordingTransport {
fn write_line(&self, _socket_path: &Path, line: &str) -> io::Result<()> {
self.lines.lock().unwrap().push(line.to_string());
Ok(())
}
}
struct FailingTransport;
impl HerdrTransport for FailingTransport {
fn write_line(&self, _socket_path: &Path, _line: &str) -> io::Result<()> {
Err(io::Error::new(io::ErrorKind::NotFound, "herdr unavailable"))
}
}
fn report_params(line: &str) -> serde_json::Value {
serde_json::from_str::<serde_json::Value>(line).unwrap()["params"].clone()
}
#[test]
fn herdr_sequence_is_process_monotonic_across_reporters_for_same_source() {
let first_transport = Arc::new(RecordingTransport::default());
let first_reporter = HerdrReporter::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
first_transport.clone(),
);
let second_transport = Arc::new(RecordingTransport::default());
let second_reporter = HerdrReporter::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
second_transport.clone(),
);
first_reporter.report_thinking();
second_reporter.report_thinking();
let first_lines = first_transport.lines.lock().unwrap();
let second_lines = second_transport.lines.lock().unwrap();
let first_params = report_params(&first_lines[0]);
let second_params = report_params(&second_lines[0]);
let first_seq = first_params["seq"].as_u64().unwrap();
let second_seq = second_params["seq"].as_u64().unwrap();
assert_eq!(first_params["source"], SOURCE);
assert_eq!(second_params["source"], SOURCE);
assert!(
second_seq > first_seq,
"second reporter seq {second_seq} must be greater than first reporter seq {first_seq}"
);
}
#[test]
fn herdr_sequence_stays_monotonic_within_reporter() {
let transport = Arc::new(RecordingTransport::default());
let reporter = HerdrReporter::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport.clone(),
);
reporter.report_thinking();
reporter.report_done();
let lines = transport.lines.lock().unwrap();
let first_seq = report_params(&lines[0])["seq"].as_u64().unwrap();
let second_seq = report_params(&lines[1])["seq"].as_u64().unwrap();
assert!(
second_seq > first_seq,
"same reporter seq must increase from {first_seq} to {second_seq}"
);
}
#[test]
fn herdr_reporter_resolver_requires_enabled_env_and_pane() {
let enabled = HerdrSettings {
enabled: true,
..Default::default()
};
assert!(
HerdrReporter::from_env_map(&HerdrSettings::default(), MapEnv::default()).is_none()
);
assert!(HerdrReporter::from_env_map(&enabled, MapEnv::default()).is_none());
assert!(
HerdrReporter::from_env_map(&enabled, MapEnv::default().with("HERDR_ENV", "0"))
.is_none()
);
assert!(
HerdrReporter::from_env_map(&enabled, MapEnv::default().with("HERDR_ENV", "1"))
.is_none()
);
let reporter = HerdrReporter::from_env_map(
&enabled,
MapEnv::default()
.with("HERDR_ENV", "1")
.with("HERDR_PANE_ID", " pane-1 ")
.with("HERDR_SOCKET_PATH", "/tmp/herdr.sock"),
)
.unwrap();
assert_eq!(reporter.pane_id, "pane-1");
assert_eq!(reporter.socket_path, PathBuf::from("/tmp/herdr.sock"));
}
#[test]
fn herdr_report_agent_json_is_newline_framed_and_sanitized() {
let line = report_agent_line(
"pane-1",
HerdrAgentState::Working,
"tool",
"magi-code: running read",
7,
)
.unwrap();
assert!(line.ends_with('\n'));
let value: serde_json::Value = serde_json::from_str(&line).unwrap();
assert_eq!(value["method"], "pane.report_agent");
assert_eq!(value["params"]["pane_id"], "pane-1");
assert_eq!(value["params"]["source"], SOURCE);
assert_eq!(value["params"]["agent"], AGENT);
assert_eq!(value["params"]["state"], "working");
assert_eq!(value["params"]["custom_status"], "tool");
assert_eq!(value["params"]["message"], "magi-code: running read");
assert_eq!(value["params"]["seq"], 7);
assert!(value["params"].get("model").is_none());
assert!(value["params"].get("cwd").is_none());
}
#[test]
fn herdr_completion_reports_clear_working_with_valid_idle_state() {
for (report, expected_status, expected_message) in [
(
HerdrReporter::report_done as fn(&HerdrReporter),
"done",
"magi-code: done",
),
(
HerdrReporter::report_stopped as fn(&HerdrReporter),
"stopped",
"magi-code: stopped",
),
] {
let (reporter, lines) = HerdrReporter::new_for_test();
report(&reporter);
let line = lines.lock().unwrap().pop().unwrap();
let value: serde_json::Value = serde_json::from_str(&line).unwrap();
assert_eq!(value["method"], "pane.report_agent");
assert_eq!(value["params"]["state"], "idle");
assert_eq!(value["params"]["custom_status"], expected_status);
assert_eq!(value["params"]["message"], expected_message);
}
}
#[test]
fn herdr_reporter_writes_best_effort_lines() {
let transport = Arc::new(RecordingTransport::default());
let reporter = HerdrReporter::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport.clone(),
);
reporter.report_tool("bad/tool name with spaces");
let lines = transport.lines.lock().unwrap();
let value: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
assert_eq!(
value["params"]["message"],
"magi-code: running badtoolnamewithspaces"
);
}
#[test]
fn herdr_transport_failure_is_best_effort() {
let reporter = HerdrReporter::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/missing-herdr.sock"),
Arc::new(FailingTransport),
);
reporter.report_thinking();
}
}