use crate::config::HerdrSettings;
use serde::Serialize;
#[cfg(unix)]
use std::time::Duration;
use std::{
collections::BTreeMap,
ffi::{OsStr, OsString},
io,
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicU64, Ordering},
},
time::{SystemTime, UNIX_EPOCH},
};
const SOURCE: &str = "custom:magi-code";
const METADATA_SOURCE: &str = "custom:magi-code:metadata";
const AGENT: &str = "magi-code";
const DEFAULT_SOCKET_RELATIVE: &str = ".config/herdr/herdr.sock";
const MAX_LOCAL_TEXT_CHARS: usize = 64;
const MAX_LOCAL_IDENTIFIER_CHARS: usize = 64;
const MAX_TOOL_NAME_CHARS: usize = 48;
const MAX_SESSION_ID_BYTES: usize = 512;
#[cfg(unix)]
const IPC_TIMEOUT: Duration = Duration::from_millis(100);
static NEXT_REPORT_ID: AtomicU64 = AtomicU64::new(1);
static NEXT_REPORT_SEQ: AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum HerdrAgentState {
Idle,
Working,
Blocked,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum HerdrSessionStartSource {
Startup,
Resume,
New,
Select,
}
pub(crate) struct HerdrOwner {
shared: Arc<HerdrShared>,
}
#[derive(Clone)]
pub(crate) struct HerdrReporter {
shared: Arc<HerdrShared>,
}
impl std::fmt::Debug for HerdrReporter {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("HerdrReporter")
.finish_non_exhaustive()
}
}
struct HerdrShared {
pane_id: String,
socket_path: PathBuf,
transport: Arc<dyn HerdrTransport>,
clock: Arc<dyn HerdrClock>,
io: Mutex<HerdrIoState>,
}
struct HerdrIoState {
released: bool,
}
trait HerdrTransport: Send + Sync {
fn write_line(&self, socket_path: &Path, line: &str) -> io::Result<()>;
}
trait HerdrClock: Send + Sync {
fn epoch_ms(&self) -> u64;
}
mod transport;
use transport::UnixSocketTransport;
#[derive(Debug, Clone, Copy)]
struct SystemHerdrClock;
impl HerdrClock for SystemHerdrClock {
fn epoch_ms(&self) -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.min(u64::MAX as u128) as u64
}
}
impl HerdrOwner {
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 = sanitize_identifier(&env.get_os("HERDR_PANE_ID")?.to_string_lossy())?;
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::new_with_transport_and_clock(
pane_id,
socket_path,
transport,
Arc::new(SystemHerdrClock),
)
}
fn new_with_transport_and_clock(
pane_id: String,
socket_path: PathBuf,
transport: Arc<dyn HerdrTransport>,
clock: Arc<dyn HerdrClock>,
) -> Self {
Self {
shared: Arc::new(HerdrShared {
pane_id,
socket_path,
transport,
clock,
io: Mutex::new(HerdrIoState { released: false }),
}),
}
}
pub(crate) fn reporter(&self) -> HerdrReporter {
HerdrReporter {
shared: Arc::clone(&self.shared),
}
}
#[cfg(test)]
pub(crate) fn new_for_test() -> (Self, Arc<Mutex<Vec<String>>>) {
let lines = Arc::new(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 release(self) {
let mut io = self
.shared
.io
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if io.released {
return;
}
io.released = true;
let Some((clear_seq, release_seq)) =
reserve_report_seq_range(self.shared.clock.as_ref(), 2)
else {
return;
};
let Ok(clear_line) =
report_metadata_line(&self.shared.pane_id, MetadataUpdate::clear_all(), clear_seq)
else {
return;
};
let Ok(release_line) = release_agent_line(&self.shared.pane_id, release_seq) else {
return;
};
let _ = self
.shared
.transport
.write_line(&self.shared.socket_path, &clear_line);
let _ = self
.shared
.transport
.write_line(&self.shared.socket_path, &release_line);
}
}
impl HerdrReporter {
fn with_open_request<F>(&self, build: F)
where
F: FnOnce(u64) -> Option<String>,
{
let io = self
.shared
.io
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if io.released {
return;
}
let Some(seq) = next_report_seq(self.shared.clock.as_ref()) else {
return;
};
if let Some(line) = build(seq) {
let _ = self
.shared
.transport
.write_line(&self.shared.socket_path, &line);
}
}
fn send_agent_report(
&self,
state: HerdrAgentState,
message: Option<&str>,
session_id: Option<&str>,
) {
let message = message.and_then(sanitize_message);
let session_id = session_id.and_then(validate_session_id);
self.with_open_request(|seq| {
report_agent_line(
&self.shared.pane_id,
state,
message.as_deref(),
seq,
session_id,
)
.ok()
});
}
fn report(&self, state: HerdrAgentState, message: Option<&str>) {
self.send_agent_report(state, message, None);
}
pub(crate) fn report_ready(&self) {
self.report(HerdrAgentState::Idle, Some("ready"));
}
pub(crate) fn report_thinking(&self) {
self.report(HerdrAgentState::Working, Some("thinking"));
}
pub(crate) fn report_tool(&self, tool_name: &'static str) {
let tool_name = sanitize_tool_name(tool_name);
let message = format!("running {tool_name}");
self.report(HerdrAgentState::Working, Some(&message));
}
pub(crate) fn report_bash(&self) {
self.report(HerdrAgentState::Working, Some("running bash"));
}
pub(crate) fn report_done(&self) {
self.report(HerdrAgentState::Idle, Some("done"));
}
pub(crate) fn report_cancelled(&self) {
self.report(HerdrAgentState::Idle, Some("cancelled"));
}
pub(crate) fn report_blocked(&self) {
self.report(HerdrAgentState::Blocked, Some("needs attention"));
}
fn report_unfinished_blocked(&self) {
self.report(HerdrAgentState::Blocked, Some("turn incomplete"));
}
#[cfg(test)]
pub(crate) fn report_agent_session(&self, session_id: &str) {
self.report_agent_session_with_optional_source(session_id, None);
}
pub(crate) fn report_agent_session_with_source(
&self,
session_id: &str,
source: HerdrSessionStartSource,
) {
self.report_agent_session_with_optional_source(session_id, Some(source));
}
fn report_agent_session_with_optional_source(
&self,
session_id: &str,
session_start_source: Option<HerdrSessionStartSource>,
) {
let Some(session_id) = validate_session_id(session_id) else {
return;
};
self.with_open_request(|seq| {
report_agent_session_line(&self.shared.pane_id, session_id, session_start_source, seq)
.ok()
});
}
pub(crate) fn report_metadata(&self, title: Option<&str>, summary: Option<&str>) {
let Some(title) = metadata_field(title) else {
return;
};
let Some(summary) = metadata_field(summary) else {
return;
};
self.with_open_request(|seq| {
report_metadata_line(
&self.shared.pane_id,
MetadataUpdate {
title,
summary,
presentation: PresentationUpdate::Replace,
},
seq,
)
.ok()
});
}
pub(crate) fn report_title(&self, title: Option<&str>) {
let Some(title) = metadata_field(title) else {
return;
};
self.with_open_request(|seq| {
report_metadata_line(
&self.shared.pane_id,
MetadataUpdate {
title,
summary: MetadataField::Keep,
presentation: PresentationUpdate::Replace,
},
seq,
)
.ok()
});
}
#[cfg(test)]
pub(crate) fn report_summary(&self, summary: Option<&str>) {
let Some(summary) = metadata_field(summary) else {
return;
};
self.with_open_request(|seq| {
report_metadata_line(
&self.shared.pane_id,
MetadataUpdate {
title: MetadataField::Keep,
summary,
presentation: PresentationUpdate::Keep,
},
seq,
)
.ok()
});
}
#[cfg(test)]
pub(crate) fn clear_metadata(&self) {
self.with_open_request(|seq| {
report_metadata_line(&self.shared.pane_id, MetadataUpdate::clear_all(), seq).ok()
});
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum HerdrTurnOutcome {
Done,
Cancelled,
Blocked,
}
pub(crate) struct HerdrTurnReporter {
reporter: Option<HerdrReporter>,
outcome: Option<HerdrTurnOutcome>,
}
impl HerdrTurnReporter {
pub(crate) fn start(reporter: Option<HerdrReporter>) -> Self {
if let Some(reporter) = &reporter {
reporter.report_thinking();
}
Self {
reporter,
outcome: None,
}
}
pub(crate) fn pending(reporter: Option<HerdrReporter>) -> Self {
Self {
reporter,
outcome: None,
}
}
pub(crate) fn finish_done(&mut self) {
self.finish(HerdrTurnOutcome::Done);
}
pub(crate) fn finish_cancelled(&mut self) {
self.finish(HerdrTurnOutcome::Cancelled);
}
pub(crate) fn finish_blocked(&mut self) {
self.finish(HerdrTurnOutcome::Blocked);
}
pub(crate) fn finish_result<T>(&mut self, result: &anyhow::Result<T>) {
match result {
Ok(_) => self.finish_done(),
Err(error) if crate::cancellation::is_run_canceled(error) => self.finish_cancelled(),
Err(_) => self.finish_blocked(),
}
}
pub(crate) fn finish_bash_result(
&mut self,
result: &anyhow::Result<crate::agent::AgentRunOutput>,
) {
match result {
Ok(output) if output.tool_results.len() == 1 && output.tool_results[0].success => {
self.finish_done()
}
Ok(_) => self.finish_blocked(),
Err(error) if crate::cancellation::is_run_canceled(error) => self.finish_cancelled(),
Err(_) => self.finish_blocked(),
}
}
fn finish(&mut self, outcome: HerdrTurnOutcome) {
if self.outcome.is_some() {
return;
}
self.outcome = Some(outcome);
if let Some(reporter) = &self.reporter {
match outcome {
HerdrTurnOutcome::Done => reporter.report_done(),
HerdrTurnOutcome::Cancelled => reporter.report_cancelled(),
HerdrTurnOutcome::Blocked => reporter.report_blocked(),
}
}
}
#[cfg(test)]
pub(crate) fn done(&mut self) {
self.finish(HerdrTurnOutcome::Done);
}
#[cfg(test)]
pub(crate) fn cancelled(&mut self) {
self.finish(HerdrTurnOutcome::Cancelled);
}
#[cfg(test)]
pub(crate) fn blocked(&mut self) {
self.finish(HerdrTurnOutcome::Blocked);
}
}
impl Drop for HerdrTurnReporter {
fn drop(&mut self) {
if self.outcome.is_none() {
self.outcome = Some(HerdrTurnOutcome::Blocked);
if let Some(reporter) = &self.reporter {
reporter.report_unfinished_blocked();
}
}
}
}
#[derive(Serialize)]
struct JsonRpcRequest<P> {
id: String,
method: &'static str,
params: P,
}
#[derive(Serialize)]
struct PaneReportAgentParams<'a> {
pane_id: &'a str,
source: &'static str,
agent: &'static str,
state: HerdrAgentState,
#[serde(skip_serializing_if = "Option::is_none")]
message: Option<&'a str>,
seq: u64,
#[serde(skip_serializing_if = "Option::is_none")]
agent_session_id: Option<&'a str>,
}
#[derive(Serialize)]
struct PaneReportAgentSessionParams<'a> {
pane_id: &'a str,
source: &'static str,
agent: &'static str,
seq: u64,
agent_session_id: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
session_start_source: Option<HerdrSessionStartSource>,
}
#[derive(Serialize)]
struct PaneReportMetadataParams<'a> {
pane_id: &'a str,
source: &'static str,
agent: &'static str,
applies_to_source: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
title: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
display_agent: Option<&'static str>,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
state_labels: BTreeMap<&'static str, &'static str>,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
tokens: BTreeMap<&'static str, Option<&'a str>>,
clear_title: bool,
clear_display_agent: bool,
clear_state_labels: bool,
seq: u64,
}
#[derive(Serialize)]
struct PaneReleaseAgentParams<'a> {
pane_id: &'a str,
source: &'static str,
agent: &'static str,
seq: u64,
}
#[derive(Debug, PartialEq, Eq)]
enum MetadataField {
Keep,
Set(String),
Clear,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PresentationUpdate {
#[cfg(test)]
Keep,
Replace,
Clear,
}
struct MetadataUpdate {
title: MetadataField,
summary: MetadataField,
presentation: PresentationUpdate,
}
impl MetadataUpdate {
fn clear_all() -> Self {
Self {
title: MetadataField::Clear,
summary: MetadataField::Clear,
presentation: PresentationUpdate::Clear,
}
}
}
fn json_rpc_line<P: Serialize>(method: &'static str, params: P) -> serde_json::Result<String> {
let id = NEXT_REPORT_ID.fetch_add(1, Ordering::Relaxed);
let request = JsonRpcRequest {
id: format!("magi-code-{id}"),
method,
params,
};
let mut line = serde_json::to_string(&request)?;
line.push('\n');
Ok(line)
}
fn report_agent_line(
pane_id: &str,
state: HerdrAgentState,
message: Option<&str>,
seq: u64,
session_id: Option<&str>,
) -> serde_json::Result<String> {
json_rpc_line(
"pane.report_agent",
PaneReportAgentParams {
pane_id,
source: SOURCE,
agent: AGENT,
state,
message,
seq,
agent_session_id: session_id,
},
)
}
fn report_agent_session_line(
pane_id: &str,
session_id: &str,
session_start_source: Option<HerdrSessionStartSource>,
seq: u64,
) -> serde_json::Result<String> {
json_rpc_line(
"pane.report_agent_session",
PaneReportAgentSessionParams {
pane_id,
source: SOURCE,
agent: AGENT,
seq,
agent_session_id: session_id,
session_start_source,
},
)
}
fn report_metadata_line(
pane_id: &str,
update: MetadataUpdate,
seq: u64,
) -> serde_json::Result<String> {
let title = match &update.title {
MetadataField::Set(value) => Some(value.as_str()),
MetadataField::Keep | MetadataField::Clear => None,
};
let clear_title = matches!(update.title, MetadataField::Clear);
let (display_agent, state_labels, clear_display_agent, clear_state_labels) =
match update.presentation {
#[cfg(test)]
PresentationUpdate::Keep => (None, BTreeMap::new(), false, false),
PresentationUpdate::Replace => {
let mut state_labels = BTreeMap::new();
state_labels.insert("blocked", "needs attention");
state_labels.insert("done", "done");
state_labels.insert("idle", "ready");
state_labels.insert("working", "working");
(Some(AGENT), state_labels, false, false)
}
PresentationUpdate::Clear => (None, BTreeMap::new(), true, true),
};
let mut tokens = BTreeMap::new();
match &update.summary {
MetadataField::Set(value) => {
tokens.insert("summary", Some(value.as_str()));
}
MetadataField::Clear => {
tokens.insert("summary", None);
}
MetadataField::Keep => {}
}
json_rpc_line(
"pane.report_metadata",
PaneReportMetadataParams {
pane_id,
source: METADATA_SOURCE,
agent: AGENT,
applies_to_source: SOURCE,
title,
display_agent,
state_labels,
tokens,
clear_title,
clear_display_agent,
clear_state_labels,
seq,
},
)
}
fn release_agent_line(pane_id: &str, seq: u64) -> serde_json::Result<String> {
json_rpc_line(
"pane.release_agent",
PaneReleaseAgentParams {
pane_id,
source: SOURCE,
agent: AGENT,
seq,
},
)
}
fn next_report_seq(clock: &dyn HerdrClock) -> Option<u64> {
reserve_report_seq_range(clock, 1).map(|(first, _last)| first)
}
fn reserve_report_seq_range(clock: &dyn HerdrClock, count: u64) -> Option<(u64, u64)> {
let epoch_ms = clock.epoch_ms();
loop {
let previous = NEXT_REPORT_SEQ.load(Ordering::Relaxed);
let (first, last) = allocate_report_seq_range_at(previous, epoch_ms, count)?;
if NEXT_REPORT_SEQ
.compare_exchange(previous, last, Ordering::SeqCst, Ordering::Relaxed)
.is_ok()
{
return Some((first, last));
}
}
}
fn allocate_report_seq_range_at(previous: u64, epoch_ms: u64, count: u64) -> Option<(u64, u64)> {
if count == 0 {
return None;
}
let first = previous.checked_add(1)?.max(epoch_ms.saturating_mul(1_000));
let last = first.checked_add(count - 1)?;
Some((first, last))
}
fn metadata_field(value: Option<&str>) -> Option<MetadataField> {
match value {
None => Some(MetadataField::Clear),
Some(value) if value.chars().any(char::is_control) => None,
Some(value) => sanitize_metadata_text(value).map(MetadataField::Set),
}
}
fn validate_session_id(value: &str) -> Option<&str> {
(!value.is_empty()
&& value.len() <= MAX_SESSION_ID_BYTES
&& !value.chars().any(char::is_control))
.then_some(value)
}
fn sanitize_message(message: &str) -> Option<String> {
sanitize_visible_text(message, MAX_LOCAL_TEXT_CHARS)
}
fn sanitize_metadata_text(value: &str) -> Option<String> {
sanitize_visible_text(value, MAX_LOCAL_TEXT_CHARS)
}
fn sanitize_visible_text(value: &str, max_chars: usize) -> Option<String> {
let mut sanitized = String::new();
for ch in value.chars() {
if ch.is_control() {
continue;
}
sanitized.push(ch);
if sanitized.chars().count() == max_chars {
break;
}
}
let sanitized = sanitized.trim().to_string();
(!sanitized.is_empty()).then_some(sanitized)
}
fn sanitize_identifier(value: &str) -> Option<String> {
let mut sanitized = String::new();
for ch in value.chars() {
if ch.is_control() {
continue;
}
sanitized.push(ch);
if sanitized.chars().count() == MAX_LOCAL_IDENTIFIER_CHARS {
break;
}
}
let sanitized = sanitized.trim().to_string();
(!sanitized.is_empty()).then_some(sanitized)
}
fn sanitize_tool_name(name: &str) -> String {
let mut sanitized = name
.chars()
.filter(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-'))
.take(MAX_TOOL_NAME_CHARS)
.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)
}
}
#[cfg(test)]
struct RecordingHerdrTransport {
lines: Arc<Mutex<Vec<String>>>,
}
#[cfg(test)]
impl HerdrTransport for RecordingHerdrTransport {
fn write_line(&self, _socket_path: &Path, line: &str) -> io::Result<()> {
self.lines
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.push(line.to_string());
Ok(())
}
}
#[cfg(test)]
mod tests;