use std::{
collections::HashMap,
env,
ffi::{OsStr, OsString},
io::{BufRead, BufReader, Write},
path::{Path, PathBuf},
process::{Command, Stdio},
sync::Mutex,
sync::atomic::{AtomicU64, Ordering},
sync::mpsc,
thread,
time::{Duration, Instant},
};
use anyhow::{Context, Result, bail};
use interprocess::local_socket::Stream;
use interprocess::local_socket::traits::Stream as _StreamExt;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use super::{
Backend, Capabilities, HerdrExt, PaneSpec, ProcessInfo, SessionState, SessionStop, Split,
TabLayout,
};
use crate::model::SplitDirection;
const SESSION_START_TIMEOUT: Duration = Duration::from_secs(10);
static REQUEST_ID: AtomicU64 = AtomicU64::new(1);
#[derive(Debug)]
pub struct HerdrClient {
socket_path: PathBuf,
created_root_tabs: Mutex<HashMap<String, String>>,
}
impl HerdrClient {
pub fn new(socket_path: PathBuf) -> Self {
Self {
socket_path,
created_root_tabs: Mutex::new(HashMap::new()),
}
}
pub fn discover(explicit_socket: Option<&Path>, session: Option<&str>) -> Self {
Self::new(resolve_socket_path(explicit_socket, session))
}
pub fn socket_path(&self) -> &Path {
&self.socket_path
}
pub fn request(&self, method: &str, params: Value) -> Result<Value> {
self.request_with_timeout(method, params, None)
}
pub fn request_with_timeout(
&self,
method: &str,
params: Value,
timeout: Option<Duration>,
) -> Result<Value> {
let id = format!("drove:{}", REQUEST_ID.fetch_add(1, Ordering::Relaxed));
let request = json!({"id": id, "method": method, "params": params});
let stream = connect(&self.socket_path).with_context(|| {
format!("cannot connect to Herdr at {}", self.socket_path.display())
})?;
let _ = stream.set_recv_timeout(timeout);
let mut stream = BufReader::new(stream);
serde_json::to_writer(stream.get_mut(), &request).context("cannot encode Herdr request")?;
stream
.get_mut()
.write_all(b"\n")
.context("cannot send Herdr request")?;
stream
.get_mut()
.flush()
.context("cannot flush Herdr request")?;
let line = read_line_with_timeout(stream, timeout, "Herdr")?;
let response: ApiResponse =
serde_json::from_str(&line).context("Herdr returned invalid JSON")?;
if response.id != request["id"] {
bail!("Herdr response id did not match the request");
}
if let Some(error) = response.error {
bail!("Herdr API error {}: {}", error.code, error.message);
}
response.result.context("Herdr response omitted result")
}
pub fn ping(&self) -> Result<Value> {
self.request("ping", json!({}))
}
pub fn snapshot(&self) -> Result<SessionSnapshot> {
let result = self.request("session.snapshot", json!({}))?;
let snapshot = result
.get("snapshot")
.cloned()
.context("session.snapshot response omitted snapshot")?;
let mut snapshot: SessionSnapshot =
serde_json::from_value(snapshot).context("invalid Herdr session snapshot")?;
for pane in &mut snapshot.panes {
pane.process_info = self.pane_process_info(&pane.pane_id)?;
}
snapshot.caller_pane_id = caller_pane_id_from_env();
Ok(snapshot)
}
pub fn export_layout(&self, tab_id: &str) -> Result<ExportedLayout> {
let result = self.request("layout.export", json!({"tab_id": tab_id}))?;
let layout = result
.get("layout")
.cloned()
.context("layout.export response omitted layout")?;
serde_json::from_value(layout).context("invalid Herdr layout export")
}
pub fn create_workspace(&self, label: &str, cwd: &Path) -> Result<String> {
let result = self.request(
"workspace.create",
json!({"label": label, "cwd": cwd, "focus": false}),
)?;
let workspace_id = find_string(&result, "workspace_id")
.map(ToOwned::to_owned)
.context("workspace.create response omitted workspace_id")?;
if let Some(tab_id) = find_string(&result, "tab_id") {
self.created_root_tabs
.lock()
.expect("root tab lock poisoned")
.insert(workspace_id.clone(), tab_id.to_owned());
}
Ok(workspace_id)
}
pub fn take_root_tab(&self, workspace_id: &str) -> Option<String> {
self.created_root_tabs
.lock()
.expect("root tab lock poisoned")
.remove(workspace_id)
}
pub fn apply_layout(
&self,
workspace_id: &str,
tab_id: Option<&str>,
tab_label: &str,
root: Value,
) -> Result<ExportedLayout> {
let mut params = json!({
"tab_label": tab_label,
"focus": false,
"root": root,
});
if let Some(tab_id) = tab_id {
params["tab_id"] = json!(tab_id);
} else {
params["workspace_id"] = json!(workspace_id);
}
let result = self.request("layout.apply", params)?;
let layout = result
.get("layout")
.cloned()
.or_else(|| result.get("created_layout").cloned())
.context("layout.apply response omitted layout")?;
serde_json::from_value(layout).context("invalid applied Herdr layout")
}
pub fn rename_workspace(&self, workspace_id: &str, label: &str) -> Result<()> {
self.request(
"workspace.rename",
json!({"workspace_id": workspace_id, "label": label}),
)?;
Ok(())
}
pub fn rename_tab(&self, tab_id: &str, label: &str) -> Result<()> {
self.request("tab.rename", json!({"tab_id": tab_id, "label": label}))?;
Ok(())
}
pub fn start_agent(
&self,
pane_id: &str,
name: &str,
kind: &str,
args: &[String],
) -> Result<()> {
self.request(
"agent.start",
json!({
"pane_id": pane_id,
"name": name,
"kind": kind,
"args": args,
}),
)?;
Ok(())
}
pub fn report_workspace_status(&self, workspace_id: &str, status: &str) -> Result<()> {
self.request(
"workspace.report_metadata",
json!({
"workspace_id": workspace_id,
"source": "drove",
"tokens": {"drove_status": status},
}),
)?;
Ok(())
}
pub fn report_metadata(
&self,
address: &str,
tokens: &std::collections::BTreeMap<String, String>,
) -> Result<()> {
if is_pane_id(address) {
self.report_pane_metadata(address, tokens)
} else {
self.report_workspace_metadata(address, tokens)
}
}
pub fn report_workspace_metadata(
&self,
workspace_id: &str,
tokens: &std::collections::BTreeMap<String, String>,
) -> Result<()> {
self.request(
"workspace.report_metadata",
json!({
"workspace_id": workspace_id,
"source": "drove",
"tokens": tokens,
}),
)?;
Ok(())
}
pub fn report_pane_metadata(
&self,
pane_id: &str,
tokens: &std::collections::BTreeMap<String, String>,
) -> Result<()> {
self.request(
"pane.report_metadata",
json!({
"pane_id": pane_id,
"source": "drove",
"tokens": tokens,
}),
)?;
Ok(())
}
pub fn split_pane(
&self,
target_pane_id: &str,
direction: SplitDirection,
ratio: f64,
command: Option<&[String]>,
cwd: Option<&Path>,
) -> Result<String> {
let mut params = json!({
"target_pane_id": target_pane_id,
"direction": direction,
"ratio": ratio,
"focus": false,
});
if let Some(cwd) = cwd {
params["cwd"] = json!(cwd);
}
let result = self.request("pane.split", params)?;
let pane = result
.get("pane")
.context("pane.split response omitted pane")?;
let pane_id = pane
.get("pane_id")
.and_then(Value::as_str)
.context("pane.split response omitted pane_id")?
.to_owned();
if let Some(command) = command {
self.run_command(&pane_id, command)?;
}
Ok(pane_id)
}
pub fn run_command(&self, pane_id: &str, command: &[String]) -> Result<()> {
if command.is_empty() {
return Ok(());
}
self.request(
"pane.send_text",
json!({"pane_id": pane_id, "text": shell_join(command)}),
)?;
self.request(
"pane.send_keys",
json!({"pane_id": pane_id, "keys": ["Enter"]}),
)?;
Ok(())
}
pub fn close_pane(&self, pane_id: &str) -> Result<()> {
self.request("pane.close", json!({"pane_id": pane_id}))?;
Ok(())
}
pub fn close_workspace(&self, workspace_id: &str) -> Result<()> {
self.request("workspace.close", json!({"workspace_id": workspace_id}))?;
Ok(())
}
pub fn set_ratio(&self, tab_id: &str, ratios: &[f64]) -> Result<()> {
for (index, ratio) in ratios.iter().enumerate() {
let path: Vec<bool> = std::iter::repeat_n(true, index).collect();
self.request(
"layout.set_split_ratio",
json!({"tab_id": tab_id, "path": path, "ratio": ratio}),
)?;
}
Ok(())
}
pub fn rename_pane(&self, pane_id: &str, label: &str) -> Result<()> {
self.request("pane.rename", json!({"pane_id": pane_id, "label": label}))?;
Ok(())
}
pub fn prompt_agent(&self, target: &str, text: &str) -> Result<()> {
self.request("agent.prompt", json!({"target": target, "text": text}))?;
Ok(())
}
pub fn pane_process_info(&self, pane_id: &str) -> Result<Option<ProcessInfo>> {
let result = self.request("pane.process_info", json!({"pane_id": pane_id}))?;
let info = result
.get("process_info")
.cloned()
.context("pane.process_info response omitted process_info")?;
let info: RawPaneProcessInfo =
serde_json::from_value(info).context("invalid Herdr process info")?;
Ok(info.into_process_info())
}
pub fn output(&self, pane_id: &str, timeout: Duration) -> Result<String> {
let result = self.request_with_timeout(
"pane.read",
json!({
"pane_id": pane_id,
"source": "recent",
"format": "text",
"strip_ansi": true,
}),
Some(timeout),
)?;
result
.get("read")
.and_then(|read| read.get("text"))
.and_then(Value::as_str)
.map(ToOwned::to_owned)
.context("pane.read response omitted text")
}
}
fn is_pane_id(address: &str) -> bool {
address
.rsplit_once(':')
.is_some_and(|(_, segment)| segment.starts_with('p'))
}
fn shell_join(argv: &[String]) -> String {
argv.iter()
.map(|arg| format!("'{}'", arg.replace('\'', r"'\''")))
.collect::<Vec<_>>()
.join(" ")
}
#[derive(Debug, Clone, Deserialize)]
struct RawPaneProcessInfo {
#[serde(default)]
shell_pid: Option<u32>,
#[serde(default)]
foreground_processes: Vec<RawPaneProcess>,
}
#[derive(Debug, Clone, Deserialize)]
struct RawPaneProcess {
pid: u32,
#[serde(default)]
argv: Option<Vec<String>>,
name: String,
}
impl RawPaneProcessInfo {
fn into_process_info(self) -> Option<ProcessInfo> {
if let Some(process) = self.foreground_processes.into_iter().next() {
let command = process.argv.unwrap_or_else(|| vec![process.name]);
return Some(ProcessInfo {
command,
pid: Some(process.pid),
});
}
self.shell_pid.map(|pid| ProcessInfo {
command: Vec::new(),
pid: Some(pid),
})
}
}
impl Backend for HerdrClient {
fn capabilities(&self) -> Capabilities {
Capabilities {
workspace_env: true,
pane_command_at_create: true,
metadata_tokens: true,
process_info: true,
events: true,
readiness_output: true,
}
}
fn caller_pane_id(&self) -> Option<String> {
caller_pane_id_from_env()
}
fn snapshot(&self) -> Result<SessionSnapshot> {
HerdrClient::snapshot(self)
}
fn create_workspace(&self, label: &str, cwd: &Path) -> Result<String> {
HerdrClient::create_workspace(self, label, cwd)
}
fn rename_workspace(&self, workspace_id: &str, label: &str) -> Result<()> {
HerdrClient::rename_workspace(self, workspace_id, label)
}
fn create_pane(&self, workspace_id: &str, spec: &PaneSpec) -> Result<String> {
let label = spec.label.as_deref().unwrap_or("pane");
let mut leaf = json!({"type": "pane", "label": label});
if let Some(command) = &spec.command {
leaf["command"] = json!(command);
}
if let Some(cwd) = &spec.cwd {
leaf["cwd"] = json!(cwd);
}
let layout = HerdrClient::apply_layout(self, workspace_id, None, label, leaf)?;
layout
.pane_ids_preorder()
.into_iter()
.next()
.context("layout.apply for create_pane returned no pane")
}
fn close_pane(&self, pane_id: &str) -> Result<()> {
HerdrClient::close_pane(self, pane_id)
}
fn rename_pane(&self, pane_id: &str, label: &str) -> Result<()> {
HerdrClient::rename_pane(self, pane_id, label)
}
fn restart_command(&self, pane_id: &str, argv: &[String]) -> Result<()> {
HerdrClient::run_command(self, pane_id, argv)
}
fn prompt_agent(&self, pane_id: &str, prompt: &str) -> Result<()> {
HerdrClient::prompt_agent(self, pane_id, prompt)
}
fn process_info(&self, pane_id: &str) -> Result<Option<ProcessInfo>> {
HerdrClient::pane_process_info(self, pane_id)
}
fn report_tokens(
&self,
address: &str,
tokens: &std::collections::BTreeMap<String, String>,
) -> Result<()> {
HerdrClient::report_metadata(self, address, tokens)
}
fn output(&self, pane_id: &str, timeout: Duration) -> Result<String> {
HerdrClient::output(self, pane_id, timeout)
}
fn herdr(&self) -> Option<&dyn HerdrExt> {
Some(self)
}
}
impl HerdrExt for HerdrClient {
fn create_tab(
&self,
workspace_id: &str,
label: &str,
split: Split,
ratios: &[f64],
panes: &[PaneSpec],
existing_tab: Option<&str>,
) -> Result<TabLayout> {
let first = panes.first();
let first_label = first
.and_then(|spec| spec.label.clone())
.unwrap_or_else(|| label.to_owned());
let mut leaf = json!({"type": "pane", "label": first_label});
if let Some(command) = first.and_then(|spec| spec.command.as_ref()) {
leaf["command"] = json!(command);
}
if let Some(cwd) = first.and_then(|spec| spec.cwd.as_ref()) {
leaf["cwd"] = json!(cwd);
}
let layout = HerdrClient::apply_layout(self, workspace_id, existing_tab, label, leaf)?;
let mut pane_ids = layout.pane_ids_preorder();
let tab_id = layout.tab_id;
if existing_tab.is_some() {
HerdrClient::rename_tab(self, &tab_id, label)?;
}
for spec in panes.iter().skip(1) {
let pane_id = <Self as HerdrExt>::split_pane(self, &tab_id, spec, split)?;
pane_ids.push(pane_id);
}
if !ratios.is_empty() {
HerdrClient::set_ratio(self, &tab_id, ratios)?;
}
Ok(TabLayout { tab_id, pane_ids })
}
fn split_pane(&self, tab_id: &str, spec: &PaneSpec, split: Split) -> Result<String> {
let target = HerdrClient::export_layout(self, tab_id)?
.pane_ids_preorder()
.into_iter()
.next_back()
.with_context(|| format!("tab `{tab_id}` has no pane to split"))?;
let command = spec.command.as_deref();
let pane_id =
HerdrClient::split_pane(self, &target, split, 0.5, command, spec.cwd.as_deref())?;
if let Some(label) = &spec.label {
HerdrClient::rename_pane(self, &pane_id, label)?;
}
Ok(pane_id)
}
fn set_ratio(&self, tab_id: &str, ratios: &[f64]) -> Result<()> {
HerdrClient::set_ratio(self, tab_id, ratios)
}
fn rename_tab(&self, tab_id: &str, label: &str) -> Result<()> {
HerdrClient::rename_tab(self, tab_id, label)
}
fn take_root_tab(&self, workspace_id: &str) -> Option<String> {
HerdrClient::take_root_tab(self, workspace_id)
}
fn start_agent(&self, pane_id: &str, name: &str, kind: &str, args: &[String]) -> Result<()> {
HerdrClient::start_agent(self, pane_id, name, kind, args)
}
fn focus_workspace(&self, id: &str) -> Result<()> {
self.request("workspace.focus", json!({"workspace_id": id}))?;
Ok(())
}
fn ensure_session(&self, name: &str) -> Result<SessionState> {
if self.ping().is_ok() {
return Ok(SessionState::Running);
}
let hint = format!("herdr --session {name}");
if start_session_server(name).is_err() {
return Ok(SessionState::CannotStart { hint });
}
if self.wait_for_ping(SESSION_START_TIMEOUT) {
Ok(SessionState::Started)
} else {
Ok(SessionState::CannotStart { hint })
}
}
fn stop_session(&self, name: &str) -> Result<SessionStop> {
run_stop_session(&herdr_bin_path(), name)
}
}
impl HerdrClient {
fn wait_for_ping(&self, timeout: Duration) -> bool {
let deadline = Instant::now() + timeout;
loop {
if self.ping().is_ok() {
return true;
}
if Instant::now() >= deadline {
return false;
}
thread::sleep(Duration::from_millis(100));
}
}
}
fn herdr_bin_path() -> OsString {
env::var_os("HERDR_BIN_PATH").unwrap_or_else(|| OsString::from("herdr"))
}
fn start_session_server(name: &str) -> Result<()> {
Command::new(herdr_bin_path())
.arg("server")
.arg("--session")
.arg(name)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.context("cannot start the Herdr session server")?;
Ok(())
}
fn run_stop_session(bin: &OsStr, name: &str) -> Result<SessionStop> {
let stop = Command::new(bin)
.args(["session", "stop", name, "--json"])
.output()
.context("cannot run the herdr binary")?;
let stopped = if stop.status.success() {
true
} else if stop_failed_because_not_running(&stop.stdout, &stop.stderr) {
false
} else {
bail!(
"herdr session stop {name} failed: {}",
String::from_utf8_lossy(&stop.stderr).trim()
);
};
let delete = Command::new(bin)
.args(["session", "delete", name, "--json"])
.output()
.context("cannot run the herdr binary")?;
if !delete.status.success() {
bail!(
"herdr session delete {name} failed: {}",
String::from_utf8_lossy(&delete.stderr).trim()
);
}
Ok(SessionStop {
stopped,
deleted: true,
})
}
fn stop_failed_because_not_running(stdout: &[u8], stderr: &[u8]) -> bool {
[stdout, stderr].into_iter().any(|bytes| {
let Ok(text) = std::str::from_utf8(bytes) else {
return false;
};
let Ok(value) = serde_json::from_str::<Value>(text.trim()) else {
return false;
};
let code = value
.get("code")
.or_else(|| value.get("error").and_then(|error| error.get("code")));
code.and_then(Value::as_str) == Some("session_stop_failed")
})
}
#[derive(Debug, Deserialize)]
struct ApiResponse {
id: Value,
#[serde(default)]
result: Option<Value>,
#[serde(default)]
error: Option<ApiError>,
}
#[derive(Debug, Deserialize)]
struct ApiError {
code: String,
message: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SessionSnapshot {
#[serde(default)]
pub version: String,
#[serde(default)]
pub protocol: u32,
#[serde(default)]
pub workspaces: Vec<WorkspaceInfo>,
#[serde(default)]
pub tabs: Vec<TabInfo>,
#[serde(default)]
pub panes: Vec<PaneInfo>,
#[serde(default)]
pub agents: Vec<AgentInfo>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub caller_pane_id: Option<String>,
}
impl SessionSnapshot {
pub fn workspace(&self, id: &str) -> Option<&WorkspaceInfo> {
self.workspaces
.iter()
.find(|workspace| workspace.workspace_id == id)
}
pub fn tab(&self, id: &str) -> Option<&TabInfo> {
self.tabs.iter().find(|tab| tab.tab_id == id)
}
pub fn pane(&self, id: &str) -> Option<&PaneInfo> {
self.panes.iter().find(|pane| pane.pane_id == id)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkspaceInfo {
pub workspace_id: String,
#[serde(default)]
pub label: String,
#[serde(default)]
pub tokens: std::collections::BTreeMap<String, String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TabInfo {
pub tab_id: String,
pub workspace_id: String,
#[serde(default)]
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PaneInfo {
pub pane_id: String,
pub tab_id: String,
pub workspace_id: String,
#[serde(default)]
pub cwd: Option<PathBuf>,
#[serde(default)]
pub tokens: std::collections::BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub process_info: Option<ProcessInfo>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentInfo {
pub pane_id: String,
#[serde(default)]
pub agent: String,
#[serde(default)]
pub agent_status: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ExportedLayout {
pub workspace_id: String,
pub tab_id: String,
#[serde(default)]
pub root: Value,
}
impl ExportedLayout {
pub fn pane_ids_preorder(&self) -> Vec<String> {
fn visit(node: &Value, ids: &mut Vec<String>) {
if node.get("type").and_then(Value::as_str) == Some("pane") {
if let Some(id) = node.get("pane_id").and_then(Value::as_str) {
ids.push(id.to_owned());
}
return;
}
if let Some(first) = node.get("first") {
visit(first, ids);
}
if let Some(second) = node.get("second") {
visit(second, ids);
}
}
let mut ids = Vec::new();
visit(&self.root, &mut ids);
ids
}
}
fn caller_pane_id_from_env() -> Option<String> {
env::var("HERDR_PANE_ID").ok().filter(|id| !id.is_empty())
}
fn find_string<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
value.get(key).and_then(Value::as_str).or_else(|| {
value
.as_object()
.and_then(|object| object.values().find_map(|child| find_string(child, key)))
})
}
pub fn resolve_socket_path(explicit_socket: Option<&Path>, session: Option<&str>) -> PathBuf {
if let Some(path) = explicit_socket {
return path.to_owned();
}
if let Some(session) = session {
return herdr_config_dir()
.join("sessions")
.join(session)
.join("herdr.sock");
}
if let Ok(path) = env::var("HERDR_SOCKET_PATH") {
return PathBuf::from(path);
}
if let Ok(session) = env::var("HERDR_SESSION")
&& !session.is_empty()
&& session != "default"
{
return herdr_config_dir()
.join("sessions")
.join(session)
.join("herdr.sock");
}
herdr_config_dir().join("herdr.sock")
}
fn herdr_config_dir() -> PathBuf {
if let Ok(root) = env::var("XDG_CONFIG_HOME") {
return PathBuf::from(root).join("herdr");
}
#[cfg(windows)]
{
if let Ok(root) = env::var("APPDATA") {
return PathBuf::from(root).join("herdr");
}
if let Ok(root) = env::var("USERPROFILE") {
return PathBuf::from(root)
.join("AppData")
.join("Roaming")
.join("herdr");
}
}
if let Ok(home) = env::var("HOME") {
return PathBuf::from(home).join(".config").join("herdr");
}
env::temp_dir().join("herdr")
}
fn read_line_with_timeout(
mut stream: BufReader<Stream>,
timeout: Option<Duration>,
label: &str,
) -> Result<String> {
let (sender, receiver) = mpsc::channel();
thread::spawn(move || {
let mut line = String::new();
let outcome = stream.read_line(&mut line).map(|read| (read, line));
let _ = sender.send(outcome);
});
let (read, line) = match timeout {
Some(duration) => receiver
.recv_timeout(duration)
.map_err(|_| anyhow::anyhow!("{label} response timed out after {duration:?}"))?
.with_context(|| format!("cannot read {label} response"))?,
None => receiver
.recv()
.context("response reader thread disconnected without a result")?
.with_context(|| format!("cannot read {label} response"))?,
};
if read == 0 {
bail!("{label} closed the socket without a response");
}
if !line.ends_with('\n') {
bail!("{label} returned a truncated response without an NDJSON newline");
}
Ok(line)
}
#[cfg(unix)]
fn connect(path: &Path) -> std::io::Result<Stream> {
use interprocess::local_socket::{GenericFilePath, prelude::*};
Stream::connect(path.to_fs_name::<GenericFilePath>()?)
}
#[cfg(windows)]
fn connect(path: &Path) -> std::io::Result<Stream> {
use interprocess::local_socket::{GenericNamespaced, prelude::*};
Stream::connect(
path.to_string_lossy()
.to_string()
.to_ns_name::<GenericNamespaced>()?,
)
}
#[cfg(test)]
mod tests {
use std::{fs, thread};
use interprocess::local_socket::{Listener, ListenerOptions, traits::Listener as _};
use super::*;
#[test]
fn resolves_named_session_before_environment_override() {
let path = resolve_socket_path(Some(Path::new("/tmp/custom.sock")), Some("work"));
assert_eq!(path, PathBuf::from("/tmp/custom.sock"));
}
#[test]
fn extracts_layout_panes_in_preorder() {
let layout = ExportedLayout {
workspace_id: "w1".into(),
tab_id: "w1:t1".into(),
root: json!({
"type": "split",
"first": {"type": "pane", "pane_id": "w1:p1"},
"second": {"type": "pane", "pane_id": "w1:p2"}
}),
};
assert_eq!(layout.pane_ids_preorder(), ["w1:p1", "w1:p2"]);
}
#[test]
fn shell_join_quotes_every_argument() {
let joined = shell_join(&[
"echo".into(),
"$(rm -rf /)".into(),
"a;b".into(),
"`whoami`".into(),
"*.rs".into(),
"a>b".into(),
"it's".into(),
]);
assert_eq!(
joined,
r"'echo' '$(rm -rf /)' 'a;b' '`whoami`' '*.rs' 'a>b' 'it'\''s'"
);
}
#[test]
fn exchanges_one_ndjson_request_with_fake_server() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {"type": "pong"}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let result = HerdrClient::new(path).ping().expect("ping");
assert_eq!(result["type"], "pong");
server.join().expect("server thread");
}
#[test]
fn agent_start_targets_the_resolved_pane() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-agent.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
assert_eq!(request["method"], "agent.start");
assert_eq!(request["params"]["pane_id"], "w1:p2");
assert_eq!(request["params"]["kind"], "cursor");
let response = json!({
"id": request["id"],
"result": {"type": "agent_started"}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
HerdrClient::new(path)
.start_agent("w1:p2", "review", "cursor", &["--model".into()])
.expect("start agent");
server.join().expect("server thread");
}
#[test]
fn split_pane_types_and_submits_the_command() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-split.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let mut requests = Vec::new();
for _ in 0..3 {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let result = match request["method"].as_str().expect("method") {
"pane.split" => json!({"pane": {"pane_id": "w1:p3"}}),
"pane.send_text" | "pane.send_keys" => json!({"type": "ok"}),
other => panic!("unexpected method {other}"),
};
let response = json!({"id": request["id"], "result": result});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
requests.push(request);
}
requests
});
let pane_id = HerdrClient::new(path)
.split_pane(
"w1:p1",
SplitDirection::Down,
0.5,
Some(&["echo".into(), "hi there".into()]),
None,
)
.expect("split pane");
assert_eq!(pane_id, "w1:p3");
let requests = server.join().expect("server thread");
assert_eq!(requests[0]["method"], "pane.split");
assert_eq!(requests[0]["params"]["target_pane_id"], "w1:p1");
assert_eq!(requests[0]["params"]["direction"], "down");
assert_eq!(requests[1]["method"], "pane.send_text");
assert_eq!(requests[1]["params"]["pane_id"], "w1:p3");
assert_eq!(requests[1]["params"]["text"], "'echo' 'hi there'");
assert_eq!(requests[2]["method"], "pane.send_keys");
assert_eq!(requests[2]["params"]["keys"], json!(["Enter"]));
}
#[test]
fn set_ratio_sends_one_right_leaning_path_per_gap() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-ratio.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let mut requests = Vec::new();
for _ in 0..2 {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {"type": "layout_split_ratio_set"}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
requests.push(request);
}
requests
});
HerdrClient::new(path)
.set_ratio("w1:t1", &[0.5, 0.3])
.expect("set ratio");
let requests = server.join().expect("server thread");
assert_eq!(requests[0]["params"]["path"], json!([]));
assert_eq!(requests[0]["params"]["ratio"], 0.5);
assert_eq!(requests[1]["params"]["path"], json!([true]));
assert_eq!(requests[1]["params"]["ratio"], 0.3);
}
#[test]
fn create_tab_opens_the_first_pane_splits_the_rest_then_sets_ratios_last() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-createtab.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let mut requests = Vec::new();
for _ in 0..5 {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let result = match request["method"].as_str().expect("method") {
"layout.apply" => json!({
"layout": {
"workspace_id": "w1",
"tab_id": "w1:t2",
"root": {"type": "pane", "pane_id": "w1:p2"},
}
}),
"layout.export" => json!({
"layout": {
"workspace_id": "w1",
"tab_id": "w1:t2",
"root": {"type": "pane", "pane_id": "w1:p2"},
}
}),
"pane.split" => json!({"pane": {"pane_id": "w1:p3"}}),
"pane.rename" => json!({"type": "ok"}),
"layout.set_split_ratio" => json!({"type": "layout_split_ratio_set"}),
other => panic!("unexpected method {other}"),
};
let response = json!({"id": request["id"], "result": result});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
requests.push(request);
}
requests
});
let panes = [
PaneSpec {
label: Some("editor".into()),
..PaneSpec::default()
},
PaneSpec {
label: Some("tests".into()),
..PaneSpec::default()
},
];
let layout = HerdrExt::create_tab(
&HerdrClient::new(path),
"w1",
"main",
SplitDirection::Right,
&[0.67],
&panes,
None,
)
.expect("create tab");
assert_eq!(layout.tab_id, "w1:t2");
assert_eq!(layout.pane_ids, ["w1:p2", "w1:p3"]);
let requests = server.join().expect("server thread");
let methods: Vec<&str> = requests
.iter()
.map(|request| request["method"].as_str().expect("method"))
.collect();
assert_eq!(methods[0], "layout.apply");
assert_eq!(requests[0]["params"]["root"]["label"], "editor");
let split_at = methods
.iter()
.position(|method| *method == "pane.split")
.expect("a split happened");
let ratio_at = methods
.iter()
.position(|method| *method == "layout.set_split_ratio")
.expect("a ratio was set");
assert!(
ratio_at > split_at,
"ratios must be applied after the split exists: {methods:?}"
);
}
#[test]
fn create_workspace_captures_the_root_tab_id_from_root_pane() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory
.path()
.join("herdr-create-workspace-root-pane.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {
"workspace_id": "w9",
"root_pane": {"pane_id": "w9:p1", "tab_id": "w9:t1"},
}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let client = HerdrClient::new(path);
let workspace_id = client
.create_workspace("dev", Path::new("."))
.expect("create workspace");
assert_eq!(workspace_id, "w9");
assert_eq!(client.take_root_tab("w9"), Some("w9:t1".to_owned()));
server.join().expect("server thread");
}
#[test]
fn create_workspace_captures_the_root_tab_id_from_tab() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-create-workspace-tab.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {
"workspace_id": "w9",
"tab": {"tab_id": "w9:t1", "label": "1"},
}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let client = HerdrClient::new(path);
let workspace_id = client
.create_workspace("dev", Path::new("."))
.expect("create workspace");
assert_eq!(
client.take_root_tab(&workspace_id),
Some("w9:t1".to_owned())
);
server.join().expect("server thread");
}
#[test]
fn take_root_tab_returns_it_only_once() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-take-root-tab-once.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {
"workspace_id": "w9",
"root_pane": {"tab_id": "w9:t1"},
}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let client = HerdrClient::new(path);
client
.create_workspace("dev", Path::new("."))
.expect("create workspace");
assert_eq!(client.take_root_tab("w9"), Some("w9:t1".to_owned()));
assert_eq!(client.take_root_tab("w9"), None);
server.join().expect("server thread");
}
#[test]
fn create_tab_applies_onto_an_existing_tab_and_renames_it_instead_of_opening_a_new_one() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-createtab-existing.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let mut requests = Vec::new();
for _ in 0..2 {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let result = match request["method"].as_str().expect("method") {
"layout.apply" => json!({
"layout": {
"workspace_id": "w1",
"tab_id": "w1:t1",
"root": {"type": "pane", "pane_id": "w1:p1"},
}
}),
"tab.rename" => json!({"type": "ok"}),
other => panic!("unexpected method {other}"),
};
let response = json!({"id": request["id"], "result": result});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
requests.push(request);
}
requests
});
let panes = [PaneSpec {
label: Some("editor".into()),
..PaneSpec::default()
}];
let layout = HerdrExt::create_tab(
&HerdrClient::new(path),
"w1",
"main",
SplitDirection::Right,
&[],
&panes,
Some("w1:t1"),
)
.expect("create tab onto existing tab");
assert_eq!(layout.tab_id, "w1:t1");
assert_eq!(layout.pane_ids, ["w1:p1"]);
let requests = server.join().expect("server thread");
assert_eq!(requests[0]["method"], "layout.apply");
assert_eq!(requests[0]["params"]["tab_id"], "w1:t1");
assert!(
requests[0]["params"].get("workspace_id").is_none(),
"reusing an existing tab must not also address the workspace: {:?}",
requests[0]
);
assert_eq!(requests[1]["method"], "tab.rename");
assert_eq!(requests[1]["params"]["tab_id"], "w1:t1");
assert_eq!(requests[1]["params"]["label"], "main");
}
#[test]
fn report_metadata_addresses_a_pane_id_at_pane_report_metadata() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-tokens-pane.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "ok"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
let mut tokens = std::collections::BTreeMap::new();
tokens.insert("drove_digest".to_owned(), "abc123".to_owned());
HerdrClient::new(path)
.report_metadata("w1:p2", &tokens)
.expect("report metadata");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "pane.report_metadata");
assert_eq!(request["params"]["pane_id"], "w1:p2");
}
#[test]
fn report_metadata_addresses_a_workspace_id_at_workspace_report_metadata() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-tokens-workspace.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "ok"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
let mut tokens = std::collections::BTreeMap::new();
tokens.insert("drove_digest".to_owned(), "abc123".to_owned());
HerdrClient::new(path)
.report_metadata("w1", &tokens)
.expect("report metadata");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "workspace.report_metadata");
assert_eq!(request["params"]["workspace_id"], "w1");
}
#[test]
fn process_info_prefers_the_foreground_process_over_the_shell() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-process-info.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {"process_info": {
"pane_id": "w1:p1",
"shell_pid": 100,
"foreground_processes": [
{"pid": 200, "name": "cargo", "argv": ["cargo", "test"]}
]
}}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let info = HerdrClient::new(path)
.pane_process_info("w1:p1")
.expect("process info")
.expect("some process");
assert_eq!(info.command, vec!["cargo".to_owned(), "test".to_owned()]);
assert_eq!(info.pid, Some(200));
server.join().expect("server thread");
}
#[test]
fn process_info_falls_back_to_the_shell_when_nothing_is_foreground() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-process-info-shell.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {"process_info": {
"pane_id": "w1:p1",
"shell_pid": 100,
"foreground_processes": []
}}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let info = HerdrClient::new(path)
.pane_process_info("w1:p1")
.expect("process info")
.expect("some process");
assert!(info.command.is_empty());
assert_eq!(info.pid, Some(100));
server.join().expect("server thread");
}
#[test]
fn output_reads_recent_pane_text() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-output.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
assert_eq!(request["method"], "pane.read");
assert_eq!(request["params"]["source"], "recent");
let response = json!({
"id": request["id"],
"result": {"read": {"text": "scaffold: watching for changes"}}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let text = HerdrClient::new(path)
.output("w1:p1", Duration::from_secs(1))
.expect("pane output");
assert_eq!(text, "scaffold: watching for changes");
server.join().expect("server thread");
}
#[test]
fn output_times_out_on_a_withheld_response() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-output-timeout.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
thread::sleep(Duration::from_millis(500));
});
let error = HerdrClient::new(path)
.output("w1:p1", Duration::from_millis(100))
.expect_err("withheld response times out");
assert!(error.to_string().contains("Herdr response"));
server.join().expect("server thread");
}
#[test]
fn close_pane_sends_the_target_pane_id() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-close.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "ok"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
HerdrClient::new(path)
.close_pane("w1:p2")
.expect("close pane");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "pane.close");
assert_eq!(request["params"]["pane_id"], "w1:p2");
}
#[test]
fn close_workspace_sends_the_target_workspace_id() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-close-workspace.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "ok"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
HerdrClient::new(path)
.close_workspace("w1")
.expect("close workspace");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "workspace.close");
assert_eq!(request["params"]["workspace_id"], "w1");
}
#[test]
fn rename_pane_sends_the_pane_id_and_label() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-rename-pane.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "ok"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
HerdrClient::new(path)
.rename_pane("w1:p2", "renamed")
.expect("rename pane");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "pane.rename");
assert_eq!(request["params"]["pane_id"], "w1:p2");
assert_eq!(request["params"]["label"], "renamed");
}
#[test]
fn prompt_agent_sends_the_target_and_text() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-prompt.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "agent_prompted"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
HerdrClient::new(path)
.prompt_agent("review", "fix the bug")
.expect("prompt agent");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "agent.prompt");
assert_eq!(request["params"]["target"], "review");
assert_eq!(request["params"]["text"], "fix the bug");
}
#[test]
fn focus_workspace_sends_the_workspace_id() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-focus.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({"id": request["id"], "result": {"type": "workspace_focused"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
request
});
HerdrExt::focus_workspace(&HerdrClient::new(path), "w1").expect("focus workspace");
let request = server.join().expect("server thread");
assert_eq!(request["method"], "workspace.focus");
assert_eq!(request["params"]["workspace_id"], "w1");
}
#[test]
fn ensure_session_returns_running_on_a_reachable_socket_without_starting_anything() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-ensure.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
assert_eq!(
request["method"], "ping",
"a reachable socket is detected by ping alone, never by shelling out"
);
let response = json!({"id": request["id"], "result": {"type": "pong"}});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
stream.get_mut().write_all(b"\n").expect("newline");
});
let state =
HerdrExt::ensure_session(&HerdrClient::new(path), "unused").expect("ensure session");
assert_eq!(state, SessionState::Running);
server.join().expect("server thread");
}
#[cfg(unix)]
#[test]
fn stop_session_stops_then_deletes_in_order() {
let directory = tempfile::tempdir().expect("tempdir");
let log = directory.path().join("argv.log");
let bin = write_fake_herdr(&directory, &log, "exit 0");
let stop = run_stop_session(bin.as_os_str(), "x").expect("stop session");
assert_eq!(
stop,
SessionStop {
stopped: true,
deleted: true
}
);
let calls = fs::read_to_string(&log).expect("argv log");
let mut lines = calls.lines();
assert_eq!(lines.next(), Some("session stop x --json"));
assert_eq!(lines.next(), Some("session delete x --json"));
assert_eq!(lines.next(), None);
}
#[cfg(unix)]
#[test]
fn stop_session_still_deletes_a_session_that_was_already_stopped() {
let directory = tempfile::tempdir().expect("tempdir");
let log = directory.path().join("argv.log");
let bin = write_fake_herdr(
&directory,
&log,
r#"if [ "$2" = "stop" ]; then echo '{"code":"session_stop_failed","message":"not running"}' >&2; exit 1; fi
exit 0"#,
);
let stop = run_stop_session(bin.as_os_str(), "x").expect("stop session");
assert_eq!(
stop,
SessionStop {
stopped: false,
deleted: true
}
);
let calls = fs::read_to_string(&log).expect("argv log");
let mut lines = calls.lines();
assert_eq!(lines.next(), Some("session stop x --json"));
assert_eq!(lines.next(), Some("session delete x --json"));
assert_eq!(lines.next(), None);
}
#[cfg(unix)]
#[test]
fn stop_session_with_an_undocumented_stop_failure_is_an_error_and_never_deletes() {
let directory = tempfile::tempdir().expect("tempdir");
let log = directory.path().join("argv.log");
let bin = write_fake_herdr(
&directory,
&log,
r#"if [ "$2" = "stop" ]; then echo '{"code":"internal_error","message":"disk full"}' >&2; exit 1; fi
exit 0"#,
);
let error = run_stop_session(bin.as_os_str(), "x").expect_err("undocumented stop failure");
assert!(error.to_string().contains("disk full"));
let calls = fs::read_to_string(&log).expect("argv log");
let mut lines = calls.lines();
assert_eq!(lines.next(), Some("session stop x --json"));
assert_eq!(
lines.next(),
None,
"a stop failure for an undocumented reason must never reach delete"
);
}
#[test]
fn stop_session_with_a_missing_binary_is_an_error() {
let directory = tempfile::tempdir().expect("tempdir");
let bin = directory.path().join("no-such-herdr-binary");
let error = run_stop_session(bin.as_os_str(), "x").expect_err("missing binary");
assert!(error.to_string().contains("cannot run the herdr binary"));
}
#[cfg(unix)]
#[test]
fn stop_session_with_a_failed_delete_is_an_error() {
let directory = tempfile::tempdir().expect("tempdir");
let log = directory.path().join("argv.log");
let bin = write_fake_herdr(
&directory,
&log,
r#"if [ "$2" = "delete" ]; then echo "boom" >&2; exit 1; fi
exit 0"#,
);
let error = run_stop_session(bin.as_os_str(), "x").expect_err("delete failure");
assert!(error.to_string().contains("boom"));
let calls = fs::read_to_string(&log).expect("argv log");
let mut lines = calls.lines();
assert_eq!(lines.next(), Some("session stop x --json"));
assert_eq!(lines.next(), Some("session delete x --json"));
assert_eq!(lines.next(), None);
}
#[cfg(unix)]
fn write_fake_herdr(directory: &tempfile::TempDir, log: &Path, body: &str) -> PathBuf {
use std::os::unix::fs::PermissionsExt;
let script = directory.path().join("herdr");
fs::write(
&script,
format!("#!/bin/sh\necho \"$*\" >> '{}'\n{body}\n", log.display()),
)
.expect("write fake herdr script");
fs::set_permissions(&script, fs::Permissions::from_mode(0o755))
.expect("make fake herdr script executable");
warm_up(&script);
let _ = fs::write(log, "");
script
}
#[cfg(unix)]
fn warm_up(script: &Path) {
for _ in 0..50 {
match Command::new(script).arg("--warmup").output() {
Ok(_) => return,
Err(error) if error.raw_os_error() == Some(26) => {
thread::sleep(Duration::from_millis(20));
}
Err(error) => panic!("warm up fake herdr script: {error}"),
}
}
panic!("fake herdr script stayed text-busy after 50 retries");
}
#[test]
fn rejects_truncated_ndjson_response() {
let directory = tempfile::tempdir().expect("tempdir");
let path = directory.path().join("herdr-truncated.sock");
let listener = bind(&path).expect("bind fake Herdr");
let server = thread::spawn(move || {
let stream = listener.accept().expect("accept");
let mut stream = BufReader::new(stream);
let mut line = String::new();
stream.read_line(&mut line).expect("read");
let request: Value = serde_json::from_str(&line).expect("request JSON");
let response = json!({
"id": request["id"],
"result": {"type": "pong"}
});
serde_json::to_writer(stream.get_mut(), &response).expect("write JSON");
});
let error = HerdrClient::new(path)
.ping()
.expect_err("truncated response");
assert!(error.to_string().contains("truncated response"));
server.join().expect("server thread");
}
#[cfg(unix)]
fn bind(path: &Path) -> std::io::Result<Listener> {
use interprocess::local_socket::{GenericFilePath, prelude::*};
ListenerOptions::new()
.name(path.to_fs_name::<GenericFilePath>()?)
.create_sync()
}
#[cfg(windows)]
fn bind(path: &Path) -> std::io::Result<Listener> {
use interprocess::local_socket::{GenericNamespaced, prelude::*};
ListenerOptions::new()
.name(
path.to_string_lossy()
.to_string()
.to_ns_name::<GenericNamespaced>()?,
)
.create_sync()
}
}