diff --git a/Cargo.lock b/Cargo.lock index aa88d227..7f4fe64f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -459,6 +459,7 @@ dependencies = [ "libc", "portable-pty", "ratatui", + "regex", "serde", "serde_json", "tokio", diff --git a/Cargo.toml b/Cargo.toml index 3119832f..a73f53b1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,6 +16,7 @@ bytes = "1" crossterm = "0.29" portable-pty = "0.9" ratatui = "0.30" +regex = "1" serde = { version = "1", features = ["derive"] } serde_json = "1" tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "time"] } diff --git a/src/api/mod.rs b/src/api/mod.rs new file mode 100644 index 00000000..642dc70d --- /dev/null +++ b/src/api/mod.rs @@ -0,0 +1,348 @@ +pub mod schema; + +use std::fs; +use std::io::{BufRead, BufReader, Write}; +use std::os::unix::net::{UnixListener, UnixStream}; +use std::path::{Path, PathBuf}; + +use tracing::{debug, error, info, warn}; + +use regex::Regex; + +use crate::api::schema::{ + ErrorBody, ErrorResponse, Method, Request, ResponseResult, SuccessResponse, +}; + +pub const SOCKET_PATH_ENV_VAR: &str = "HERDR_SOCKET_PATH"; + +pub struct ApiRequestMessage { + pub request: Request, + pub respond_to: std::sync::mpsc::Sender, +} + +pub fn socket_path() -> PathBuf { + if let Ok(path) = std::env::var(SOCKET_PATH_ENV_VAR) { + return PathBuf::from(path); + } + + if let Ok(dir) = std::env::var("XDG_RUNTIME_DIR") { + return PathBuf::from(dir).join("herdr.sock"); + } + + if let Ok(dir) = std::env::var("XDG_CONFIG_HOME") { + return PathBuf::from(dir).join("herdr/herdr.sock"); + } + + if let Ok(home) = std::env::var("HOME") { + return PathBuf::from(home).join(".config/herdr/herdr.sock"); + } + + PathBuf::from("/tmp/herdr.sock") +} + +pub struct ServerHandle { + _thread: std::thread::JoinHandle<()>, + path: PathBuf, +} + +impl Drop for ServerHandle { + fn drop(&mut self) { + if let Err(err) = fs::remove_file(&self.path) { + if err.kind() != std::io::ErrorKind::NotFound { + warn!(path = %self.path.display(), err = %err, "failed to remove api socket on shutdown"); + } + } + } +} + +pub fn start_server( + api_tx: std::sync::mpsc::Sender, +) -> std::io::Result { + let path = socket_path(); + prepare_socket_path(&path)?; + + let listener = UnixListener::bind(&path)?; + info!(path = %path.display(), "api server listening"); + + let thread = std::thread::spawn(move || { + for stream in listener.incoming() { + match stream { + Ok(stream) => { + if let Err(err) = handle_connection(stream, &api_tx) { + warn!(err = %err, "api connection failed"); + } + } + Err(err) => { + error!(err = %err, "api listener accept failed"); + break; + } + } + } + debug!("api server thread exiting"); + }); + + Ok(ServerHandle { + _thread: thread, + path, + }) +} + +fn prepare_socket_path(path: &Path) -> std::io::Result<()> { + if let Some(parent) = path.parent() { + fs::create_dir_all(parent)?; + } + + if let Err(err) = fs::remove_file(path) { + if err.kind() != std::io::ErrorKind::NotFound { + return Err(err); + } + } + + Ok(()) +} + +fn handle_connection( + mut stream: UnixStream, + api_tx: &std::sync::mpsc::Sender, +) -> std::io::Result<()> { + let mut line = String::new(); + { + let mut reader = BufReader::new(&stream); + let read = reader.read_line(&mut line)?; + if read == 0 { + return Ok(()); + } + } + + let line = line.trim(); + if line.is_empty() { + return Ok(()); + } + + let response = match serde_json::from_str::(line) { + Ok(request) => handle_request(request, api_tx), + Err(err) => serde_json::to_string(&ErrorResponse { + id: String::new(), + error: ErrorBody { + code: "invalid_request".into(), + message: format!("invalid request: {err}"), + }, + })?, + }; + + stream.write_all(response.as_bytes())?; + stream.write_all(b"\n")?; + stream.flush()?; + Ok(()) +} + +fn handle_request(request: Request, api_tx: &std::sync::mpsc::Sender) -> String { + let request_id = request.id.clone(); + match request.method { + Method::Ping(_) => serde_json::to_string(&SuccessResponse { + id: request.id, + result: ResponseResult::Pong { + version: env!("CARGO_PKG_VERSION").into(), + }, + }) + .unwrap_or_else(|_| { + r#"{"id":"","error":{"code":"internal_error","message":"failed to encode response"}}"# + .to_string() + }), + Method::PaneWaitForOutput(params) => wait_for_output(request_id, params, api_tx), + _ => dispatch_to_app(request, api_tx), + } +} + +fn wait_for_output( + request_id: String, + params: crate::api::schema::PaneWaitForOutputParams, + api_tx: &std::sync::mpsc::Sender, +) -> String { + let deadline = params + .timeout_ms + .map(|ms| std::time::Instant::now() + std::time::Duration::from_millis(ms)); + + let regex = match ¶ms.r#match { + crate::api::schema::OutputMatch::Regex { value } => match Regex::new(value) { + Ok(regex) => Some(regex), + Err(err) => { + return serde_json::to_string(&ErrorResponse { + id: request_id, + error: ErrorBody { + code: "invalid_regex".into(), + message: err.to_string(), + }, + }) + .unwrap(); + } + }, + crate::api::schema::OutputMatch::Substring { .. } => None, + }; + + loop { + let read_request = Request { + id: format!("{request_id}:read"), + method: Method::PaneRead(crate::api::schema::PaneReadParams { + pane_id: params.pane_id.clone(), + source: params.source.clone(), + lines: params.lines, + strip_ansi: params.strip_ansi, + }), + }; + let response = dispatch_to_app(read_request, api_tx); + let Ok(value) = serde_json::from_str::(&response) else { + return response; + }; + if value.get("error").is_some() { + let mut value = value; + value["id"] = serde_json::Value::String(request_id); + return serde_json::to_string(&value).unwrap(); + } + + let read_value = value["result"]["read"].clone(); + let Ok(read) = serde_json::from_value::(read_value) + else { + return serde_json::to_string(&ErrorResponse { + id: request_id, + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane read result".into(), + }, + }) + .unwrap(); + }; + + let matched_line = match ¶ms.r#match { + crate::api::schema::OutputMatch::Substring { value } => read + .text + .lines() + .find(|line| line.contains(value)) + .map(|line| line.to_string()), + crate::api::schema::OutputMatch::Regex { .. } => regex.as_ref().and_then(|re| { + read.text + .lines() + .find(|line| re.is_match(line)) + .map(|line| line.to_string()) + }), + }; + if matched_line.is_some() { + let revision = read.revision; + return serde_json::to_string(&SuccessResponse { + id: request_id, + result: ResponseResult::OutputMatched { + pane_id: params.pane_id, + revision, + matched_line, + read, + }, + }) + .unwrap(); + } + + if deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) { + return serde_json::to_string(&ErrorResponse { + id: request_id, + error: ErrorBody { + code: "timeout".into(), + message: "timed out waiting for output match".into(), + }, + }) + .unwrap(); + } + + std::thread::sleep(std::time::Duration::from_millis(100)); + } +} + +fn dispatch_to_app( + request: Request, + api_tx: &std::sync::mpsc::Sender, +) -> String { + let (respond_to, response_rx) = std::sync::mpsc::channel(); + if let Err(err) = api_tx.send(ApiRequestMessage { + request, + respond_to, + }) { + return serde_json::to_string(&ErrorResponse { + id: String::new(), + error: ErrorBody { + code: "server_unavailable".into(), + message: format!("failed to dispatch request: {err}"), + }, + }) + .unwrap_or_else(|_| { + r#"{"id":"","error":{"code":"internal_error","message":"failed to encode error response"}}"#.to_string() + }); + } + + response_rx.recv().unwrap_or_else(|err| { + serde_json::to_string(&ErrorResponse { + id: String::new(), + error: ErrorBody { + code: "server_unavailable".into(), + message: format!("request handling failed: {err}"), + }, + }) + .unwrap_or_else(|_| { + r#"{"id":"","error":{"code":"internal_error","message":"failed to encode error response"}}"#.to_string() + }) + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn socket_path_prefers_explicit_env_override() { + let unique = format!("/tmp/herdr-test-{}.sock", std::process::id()); + std::env::set_var(SOCKET_PATH_ENV_VAR, &unique); + assert_eq!(socket_path(), PathBuf::from(&unique)); + std::env::remove_var(SOCKET_PATH_ENV_VAR); + } + + #[test] + fn ping_request_returns_pong() { + let (tx, _rx) = std::sync::mpsc::channel(); + let response = handle_request( + Request { + id: "req_1".into(), + method: Method::Ping(crate::api::schema::PingParams::default()), + }, + &tx, + ); + + let parsed: SuccessResponse = serde_json::from_str(&response).unwrap(); + assert_eq!(parsed.id, "req_1"); + assert!(matches!(parsed.result, ResponseResult::Pong { .. })); + } + + #[test] + fn request_dispatches_to_app_channel() { + let (tx, rx) = std::sync::mpsc::channel(); + let request = Request { + id: "req_2".into(), + method: Method::WorkspaceList(crate::api::schema::EmptyParams::default()), + }; + + let request_for_thread = request.clone(); + let thread = std::thread::spawn(move || handle_request(request_for_thread, &tx)); + + let msg = rx.recv().unwrap(); + assert_eq!(msg.request.id, "req_2"); + msg.respond_to + .send( + serde_json::to_string(&SuccessResponse { + id: "req_2".into(), + result: ResponseResult::Ok {}, + }) + .unwrap(), + ) + .unwrap(); + + let response = thread.join().unwrap(); + let parsed: SuccessResponse = serde_json::from_str(&response).unwrap(); + assert_eq!(parsed.id, "req_2"); + } +} diff --git a/src/api/schema.rs b/src/api/schema.rs new file mode 100644 index 00000000..29509211 --- /dev/null +++ b/src/api/schema.rs @@ -0,0 +1,537 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Request { + pub id: String, + #[serde(flatten)] + pub method: Method, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "method", content = "params")] +pub enum Method { + #[serde(rename = "ping")] + Ping(PingParams), + #[serde(rename = "workspace.create")] + WorkspaceCreate(WorkspaceCreateParams), + #[serde(rename = "workspace.list")] + WorkspaceList(EmptyParams), + #[serde(rename = "workspace.get")] + WorkspaceGet(WorkspaceTarget), + #[serde(rename = "workspace.focus")] + WorkspaceFocus(WorkspaceTarget), + #[serde(rename = "workspace.rename")] + WorkspaceRename(WorkspaceRenameParams), + #[serde(rename = "workspace.close")] + WorkspaceClose(WorkspaceTarget), + #[serde(rename = "pane.split")] + PaneSplit(PaneSplitParams), + #[serde(rename = "pane.list")] + PaneList(PaneListParams), + #[serde(rename = "pane.get")] + PaneGet(PaneTarget), + #[serde(rename = "pane.send_text")] + PaneSendText(PaneSendTextParams), + #[serde(rename = "pane.send_keys")] + PaneSendKeys(PaneSendKeysParams), + #[serde(rename = "pane.read")] + PaneRead(PaneReadParams), + #[serde(rename = "pane.close")] + PaneClose(PaneTarget), + #[serde(rename = "events.subscribe")] + EventsSubscribe(EventsSubscribeParams), + #[serde(rename = "events.wait")] + EventsWait(EventsWaitParams), + #[serde(rename = "pane.wait_for_output")] + PaneWaitForOutput(PaneWaitForOutputParams), +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)] +pub struct EmptyParams {} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)] +pub struct PingParams {} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct WorkspaceTarget { + pub workspace_id: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneTarget { + pub pane_id: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct WorkspaceCreateParams { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, + #[serde(default)] + pub focus: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct WorkspaceRenameParams { + pub workspace_id: String, + pub label: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneSplitParams { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace_id: Option, + pub target_pane_id: String, + pub direction: SplitDirection, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, + #[serde(default)] + pub focus: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum SplitDirection { + Right, + Down, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)] +pub struct PaneListParams { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace_id: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneSendTextParams { + pub pane_id: String, + pub text: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneSendKeysParams { + pub pane_id: String, + pub keys: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneReadParams { + pub pane_id: String, + pub source: ReadSource, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub lines: Option, + #[serde(default = "default_true")] + pub strip_ansi: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ReadSource { + Visible, + Recent, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct EventsSubscribeParams { + pub events: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub workspace_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub pane_id: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct EventsWaitParams { + pub match_event: EventMatch, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub timeout_ms: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneWaitForOutputParams { + pub pane_id: String, + pub source: ReadSource, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub lines: Option, + pub r#match: OutputMatch, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub timeout_ms: Option, + #[serde(default = "default_true")] + pub strip_ansi: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum OutputMatch { + Substring { value: String }, + Regex { value: String }, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "event", rename_all = "snake_case")] +pub enum EventMatch { + WorkspaceCreated { + #[serde(default, skip_serializing_if = "Option::is_none")] + workspace_id: Option, + }, + WorkspaceClosed { + workspace_id: String, + }, + WorkspaceRenamed { + workspace_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + label: Option, + }, + WorkspaceFocused { + workspace_id: String, + }, + PaneCreated { + #[serde(default, skip_serializing_if = "Option::is_none")] + pane_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + workspace_id: Option, + }, + PaneClosed { + pane_id: String, + }, + PaneFocused { + pane_id: String, + }, + PaneOutputChanged { + pane_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + min_revision: Option, + }, + PaneExited { + pane_id: String, + }, + PaneAgentDetected { + pane_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + agent: Option, + }, + PaneAgentStateChanged { + pane_id: String, + state: PaneAgentState, + }, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum EventKind { + WorkspaceCreated, + WorkspaceClosed, + WorkspaceRenamed, + WorkspaceFocused, + PaneCreated, + PaneClosed, + PaneFocused, + PaneOutputChanged, + PaneExited, + PaneAgentDetected, + PaneAgentStateChanged, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SuccessResponse { + pub id: String, + pub result: ResponseResult, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ErrorResponse { + pub id: String, + pub error: ErrorBody, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ErrorBody { + pub code: String, + pub message: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum ResponseResult { + Pong { + version: String, + }, + WorkspaceInfo { + workspace: WorkspaceInfo, + }, + WorkspaceList { + workspaces: Vec, + }, + PaneInfo { + pane: PaneInfo, + }, + PaneList { + panes: Vec, + }, + PaneRead { + read: PaneReadResult, + }, + SubscriptionStarted { + events: Vec, + }, + WaitMatched { + event: EventEnvelope, + }, + OutputMatched { + pane_id: String, + revision: u64, + matched_line: Option, + read: PaneReadResult, + }, + Ok {}, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct WorkspaceInfo { + pub workspace_id: String, + pub number: usize, + pub label: String, + pub focused: bool, + pub pane_count: usize, + pub agent_state: PaneAgentState, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneInfo { + pub pane_id: String, + pub workspace_id: String, + pub focused: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub agent: Option, + pub agent_state: PaneAgentState, + pub revision: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneReadResult { + pub pane_id: String, + pub workspace_id: String, + pub source: ReadSource, + pub text: String, + pub revision: u64, + pub truncated: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct EventEnvelope { + pub event: EventKind, + pub data: EventData, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum EventData { + WorkspaceCreated { + workspace: WorkspaceInfo, + }, + WorkspaceClosed { + workspace_id: String, + }, + WorkspaceRenamed { + workspace_id: String, + label: String, + }, + WorkspaceFocused { + workspace_id: String, + }, + PaneCreated { + pane: PaneInfo, + }, + PaneClosed { + pane_id: String, + workspace_id: String, + }, + PaneFocused { + pane_id: String, + workspace_id: String, + }, + PaneOutputChanged { + pane_id: String, + workspace_id: String, + revision: u64, + }, + PaneExited { + pane_id: String, + workspace_id: String, + }, + PaneAgentDetected { + pane_id: String, + workspace_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + agent: Option, + }, + PaneAgentStateChanged { + pane_id: String, + workspace_id: String, + state: PaneAgentState, + }, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum PaneAgentState { + Idle, + Busy, + Waiting, + Unknown, +} + +fn default_true() -> bool { + true +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn request_round_trips_for_pane_read() { + let request = Request { + id: "req_1".into(), + method: Method::PaneRead(PaneReadParams { + pane_id: "p_1".into(), + source: ReadSource::Recent, + lines: Some(80), + strip_ansi: true, + }), + }; + + let json = serde_json::to_string(&request).unwrap(); + let restored: Request = serde_json::from_str(&json).unwrap(); + assert_eq!(restored, request); + } + + #[test] + fn request_uses_dot_method_names() { + let request = Request { + id: "req_1".into(), + method: Method::WorkspaceCreate(WorkspaceCreateParams { + cwd: Some("/tmp".into()), + focus: true, + }), + }; + + let json = serde_json::to_value(&request).unwrap(); + assert_eq!(json["method"], "workspace.create"); + } + + #[test] + fn unknown_method_is_rejected() { + let json = r#"{"id":"req_1","method":"nope","params":{}}"#; + let err = serde_json::from_str::(json) + .unwrap_err() + .to_string(); + assert!(err.contains("unknown variant")); + } + + #[test] + fn missing_required_params_are_rejected() { + let json = r#"{"id":"req_1","method":"pane.send_text","params":{"pane_id":"p_1"}}"#; + let err = serde_json::from_str::(json) + .unwrap_err() + .to_string(); + assert!(err.contains("text")); + } + + #[test] + fn pane_wait_for_output_defaults_strip_ansi_to_true() { + let json = r#" + { + "id": "req_1", + "method": "pane.wait_for_output", + "params": { + "pane_id": "p_1", + "source": "recent", + "match": { "type": "substring", "value": "ready" } + } + } + "#; + + let request: Request = serde_json::from_str(json).unwrap(); + let Method::PaneWaitForOutput(params) = request.method else { + panic!("wrong method parsed"); + }; + assert!(params.strip_ansi); + } + + #[test] + fn event_envelope_round_trips() { + let event = EventEnvelope { + event: EventKind::PaneOutputChanged, + data: EventData::PaneOutputChanged { + pane_id: "p_1".into(), + workspace_id: "w_1".into(), + revision: 42, + }, + }; + + let json = serde_json::to_string(&event).unwrap(); + let restored: EventEnvelope = serde_json::from_str(&json).unwrap(); + assert_eq!(restored, event); + } + + #[test] + fn success_response_round_trips() { + let response = SuccessResponse { + id: "req_1".into(), + result: ResponseResult::Pong { + version: "0.1.2".into(), + }, + }; + + let json = serde_json::to_string(&response).unwrap(); + let restored: SuccessResponse = serde_json::from_str(&json).unwrap(); + assert_eq!(restored, response); + } + + #[test] + fn error_response_round_trips() { + let response = ErrorResponse { + id: "req_1".into(), + error: ErrorBody { + code: "pane_not_found".into(), + message: "pane p_1 not found".into(), + }, + }; + + let json = serde_json::to_string(&response).unwrap(); + let restored: ErrorResponse = serde_json::from_str(&json).unwrap(); + assert_eq!(restored, response); + } + + #[test] + fn event_wait_parses_typed_match() { + let json = r#" + { + "id": "req_9", + "method": "events.wait", + "params": { + "match_event": { + "event": "pane_agent_state_changed", + "pane_id": "p_1", + "state": "waiting" + }, + "timeout_ms": 30000 + } + } + "#; + + let request: Request = serde_json::from_str(json).unwrap(); + let Method::EventsWait(params) = request.method else { + panic!("wrong method parsed"); + }; + assert_eq!( + params.match_event, + EventMatch::PaneAgentStateChanged { + pane_id: "p_1".into(), + state: PaneAgentState::Waiting, + } + ); + } +} diff --git a/src/app/mod.rs b/src/app/mod.rs index f8036315..48700577 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -28,13 +28,19 @@ pub struct App { pub state: AppState, pub event_tx: mpsc::Sender, event_rx: mpsc::Receiver, + api_rx: std::sync::mpsc::Receiver, no_session: bool, config_diagnostic_deadline: Option, toast_deadline: Option, } impl App { - pub fn new(config: &Config, no_session: bool, config_diagnostic: Option) -> Self { + pub fn new( + config: &Config, + no_session: bool, + config_diagnostic: Option, + api_rx: std::sync::mpsc::Receiver, + ) -> Self { let (prefix_code, prefix_mods) = config.prefix_key(); let (event_tx, event_rx) = mpsc::channel::(64); @@ -114,6 +120,7 @@ impl App { state, event_tx, event_rx, + api_rx, no_session, } } @@ -141,6 +148,11 @@ impl App { crate::ui::render(&self.state, frame); })?; + while let Ok(msg) = self.api_rx.try_recv() { + let response = self.handle_api_request(msg.request); + let _ = msg.respond_to.send(response); + } + // Drain internal events while let Ok(ev) = self.event_rx.try_recv() { let previous_toast = self.state.toast.clone(); @@ -192,6 +204,265 @@ impl App { Ok(()) } + fn handle_api_request(&mut self, request: crate::api::schema::Request) -> String { + use bytes::Bytes; + + use crate::api::schema::{ + ErrorBody, ErrorResponse, Method, PaneListParams, PaneReadResult, ReadSource, + ResponseResult, SuccessResponse, + }; + + let response = match request.method { + Method::WorkspaceList(_) => SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceList { + workspaces: self + .state + .workspaces + .iter() + .enumerate() + .map(|(idx, _)| workspace_info(&self.state, idx)) + .collect(), + }, + }, + Method::WorkspaceGet(target) => { + let Some(index) = parse_workspace_id(&target.workspace_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "workspace_not_found".into(), + message: format!("workspace {} not found", target.workspace_id), + }, + }) + .unwrap(); + }; + let Some(_) = self.state.workspaces.get(index) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "workspace_not_found".into(), + message: format!("workspace {} not found", target.workspace_id), + }, + }) + .unwrap(); + }; + SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceInfo { + workspace: workspace_info(&self.state, index), + }, + } + } + Method::WorkspaceCreate(params) => { + let cwd = params + .cwd + .map(std::path::PathBuf::from) + .or_else(|| std::env::current_dir().ok()) + .unwrap_or_else(|| std::path::PathBuf::from("/")); + match self.create_workspace_with_options(cwd, params.focus) { + Ok(index) => SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceInfo { + workspace: workspace_info(&self.state, index), + }, + }, + Err(err) => { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "workspace_create_failed".into(), + message: err.to_string(), + }, + }) + .unwrap(); + } + } + } + Method::PaneList(PaneListParams { workspace_id }) => { + match self.collect_panes_for_workspace(workspace_id.as_deref()) { + Ok(panes) => SuccessResponse { + id: request.id, + result: ResponseResult::PaneList { panes }, + }, + Err((code, message)) => { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { code, message }, + }) + .unwrap(); + } + } + } + Method::PaneGet(target) => { + let Some((ws_idx, pane_id)) = parse_pane_id(&target.pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", target.pane_id), + }, + }) + .unwrap(); + }; + let Some(pane) = self.pane_info(ws_idx, pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", target.pane_id), + }, + }) + .unwrap(); + }; + SuccessResponse { + id: request.id, + result: ResponseResult::PaneInfo { pane }, + } + } + Method::PaneRead(params) => { + let Some((ws_idx, pane_id)) = parse_pane_id(¶ms.pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + let Some((pane, workspace_id)) = self.lookup_runtime(ws_idx, pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + let requested_lines = params.lines.unwrap_or(80).min(1000) as usize; + let text = match params.source { + ReadSource::Visible => pane.visible_text(), + ReadSource::Recent => pane.recent_text(requested_lines), + }; + SuccessResponse { + id: request.id, + result: ResponseResult::PaneRead { + read: PaneReadResult { + pane_id: params.pane_id, + workspace_id, + source: params.source, + text, + revision: 0, + truncated: false, + }, + }, + } + } + Method::PaneSendText(params) => { + let Some((ws_idx, pane_id)) = parse_pane_id(¶ms.pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + let Some(runtime) = self.lookup_runtime_sender(ws_idx, pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + if let Err(err) = runtime.0.try_send(Bytes::from(params.text)) { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_send_failed".into(), + message: err.to_string(), + }, + }) + .unwrap(); + } + SuccessResponse { + id: request.id, + result: ResponseResult::Ok {}, + } + } + Method::PaneSendKeys(params) => { + let Some((ws_idx, pane_id)) = parse_pane_id(¶ms.pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + let Some(runtime) = self.lookup_runtime_sender(ws_idx, pane_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.pane_id), + }, + }) + .unwrap(); + }; + for key in params.keys { + let Some(key_event) = parse_api_key(&key) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "invalid_key".into(), + message: format!("unsupported key {}", key), + }, + }) + .unwrap(); + }; + let kitty = runtime + .1 + .kitty_keyboard + .load(std::sync::atomic::Ordering::Relaxed); + let bytes = crate::input::encode_key(key_event, kitty); + if let Err(err) = runtime.0.try_send(Bytes::from(bytes)) { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_send_failed".into(), + message: err.to_string(), + }, + }) + .unwrap(); + } + } + SuccessResponse { + id: request.id, + result: ResponseResult::Ok {}, + } + } + _ => { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "not_implemented".into(), + message: "method not implemented yet".into(), + }, + }) + .unwrap(); + } + }; + + serde_json::to_string(&response).unwrap() + } + pub(crate) fn complete_onboarding(&mut self) { let (sound_enabled, toast_enabled) = match self.state.onboarding_selected { 0 => (false, false), @@ -220,7 +491,6 @@ impl App { /// Create a workspace with a real PTY (needs event_tx). fn create_workspace(&mut self) { - let (rows, cols) = self.state.estimate_pane_size(); let initial_cwd = self .state .active @@ -229,17 +499,181 @@ impl App { .and_then(|rt| rt.cwd()) .or_else(|| std::env::current_dir().ok()) .unwrap_or_else(|| std::path::PathBuf::from("/")); - match Workspace::new(initial_cwd, rows, cols, self.event_tx.clone()) { - Ok(ws) => { - self.state.workspaces.push(ws); - let idx = self.state.workspaces.len() - 1; - self.state.switch_workspace(idx); - self.state.mode = Mode::Terminal; - } - Err(e) => { - error!(err = %e, "failed to create workspace"); - self.state.mode = Mode::Navigate; - } + if let Err(e) = self.create_workspace_with_options(initial_cwd, true) { + error!(err = %e, "failed to create workspace"); + self.state.mode = Mode::Navigate; } } + + fn create_workspace_with_options( + &mut self, + initial_cwd: std::path::PathBuf, + focus: bool, + ) -> std::io::Result { + let (rows, cols) = self.state.estimate_pane_size(); + let ws = Workspace::new(initial_cwd, rows, cols, self.event_tx.clone())?; + self.state.workspaces.push(ws); + let idx = self.state.workspaces.len() - 1; + if focus || self.state.active.is_none() { + self.state.switch_workspace(idx); + self.state.mode = Mode::Terminal; + } + Ok(idx) + } + + fn collect_panes_for_workspace( + &self, + workspace_id: Option<&str>, + ) -> Result, (String, String)> { + if let Some(workspace_id) = workspace_id { + let Some(ws_idx) = parse_workspace_id(workspace_id) else { + return Err(( + "workspace_not_found".into(), + format!("workspace {workspace_id} not found"), + )); + }; + let Some(ws) = self.state.workspaces.get(ws_idx) else { + return Err(( + "workspace_not_found".into(), + format!("workspace {workspace_id} not found"), + )); + }; + Ok(ws + .layout + .pane_ids() + .into_iter() + .filter_map(|pane_id| self.pane_info(ws_idx, pane_id)) + .collect()) + } else { + Ok(self + .state + .workspaces + .iter() + .enumerate() + .flat_map(|(ws_idx, ws)| { + ws.layout + .pane_ids() + .into_iter() + .filter_map(move |pane_id| self.pane_info(ws_idx, pane_id)) + }) + .collect()) + } + } + + fn pane_info( + &self, + ws_idx: usize, + pane_id: crate::layout::PaneId, + ) -> Option { + let ws = self.state.workspaces.get(ws_idx)?; + let pane = ws.panes.get(&pane_id)?; + let runtime = ws.runtimes.get(&pane_id); + Some(crate::api::schema::PaneInfo { + pane_id: format!("p_{}_{}", ws_idx + 1, pane_id.raw()), + workspace_id: format!("w_{}", ws_idx + 1), + focused: self.state.active == Some(ws_idx) && ws.layout.focused() == pane_id, + cwd: runtime + .and_then(|rt| rt.cwd()) + .map(|cwd| cwd.display().to_string()), + agent: pane.detected_agent.map(agent_name), + agent_state: pane_agent_state(pane.state), + revision: 0, + }) + } + + fn lookup_runtime( + &self, + ws_idx: usize, + pane_id: crate::layout::PaneId, + ) -> Option<(&crate::pane::PaneRuntime, String)> { + let ws = self.state.workspaces.get(ws_idx)?; + let runtime = ws.runtimes.get(&pane_id)?; + Some((runtime, format!("w_{}", ws_idx + 1))) + } + + fn lookup_runtime_sender( + &self, + ws_idx: usize, + pane_id: crate::layout::PaneId, + ) -> Option<( + &tokio::sync::mpsc::Sender, + &crate::pane::PaneRuntime, + )> { + let ws = self.state.workspaces.get(ws_idx)?; + let runtime = ws.runtimes.get(&pane_id)?; + Some((&runtime.sender, runtime)) + } +} + +fn parse_api_key(key: &str) -> Option { + use crossterm::event::{KeyCode, KeyEvent, KeyModifiers}; + + let normalized = key.trim(); + match normalized { + "Enter" | "enter" => Some(KeyEvent::new(KeyCode::Enter, KeyModifiers::empty())), + "Tab" | "tab" => Some(KeyEvent::new(KeyCode::Tab, KeyModifiers::empty())), + "Esc" | "esc" => Some(KeyEvent::new(KeyCode::Esc, KeyModifiers::empty())), + "Backspace" | "backspace" => Some(KeyEvent::new(KeyCode::Backspace, KeyModifiers::empty())), + "Up" | "up" => Some(KeyEvent::new(KeyCode::Up, KeyModifiers::empty())), + "Down" | "down" => Some(KeyEvent::new(KeyCode::Down, KeyModifiers::empty())), + "Left" | "left" => Some(KeyEvent::new(KeyCode::Left, KeyModifiers::empty())), + "Right" | "right" => Some(KeyEvent::new(KeyCode::Right, KeyModifiers::empty())), + "C-c" | "c-c" | "ctrl+c" => Some(KeyEvent::new(KeyCode::Char('c'), KeyModifiers::CONTROL)), + _ if normalized.len() == 1 => normalized + .chars() + .next() + .map(|ch| KeyEvent::new(KeyCode::Char(ch), KeyModifiers::empty())), + _ => None, + } +} + +fn parse_workspace_id(id: &str) -> Option { + id.strip_prefix("w_")?.parse::().ok()?.checked_sub(1) +} + +fn parse_pane_id(id: &str) -> Option<(usize, crate::layout::PaneId)> { + let rest = id.strip_prefix("p_")?; + let (ws_raw, pane_raw) = rest.split_once('_')?; + let ws_idx = ws_raw.parse::().ok()?.checked_sub(1)?; + let pane_id = crate::layout::PaneId::from_raw(pane_raw.parse::().ok()?); + Some((ws_idx, pane_id)) +} + +fn pane_agent_state(state: crate::detect::AgentState) -> crate::api::schema::PaneAgentState { + match state { + crate::detect::AgentState::Idle => crate::api::schema::PaneAgentState::Idle, + crate::detect::AgentState::Busy => crate::api::schema::PaneAgentState::Busy, + crate::detect::AgentState::Waiting => crate::api::schema::PaneAgentState::Waiting, + crate::detect::AgentState::Unknown => crate::api::schema::PaneAgentState::Unknown, + } +} + +fn workspace_info(state: &AppState, index: usize) -> crate::api::schema::WorkspaceInfo { + let ws = &state.workspaces[index]; + let (agg_state, _) = ws.aggregate_state(); + crate::api::schema::WorkspaceInfo { + workspace_id: format!("w_{}", index + 1), + number: index + 1, + label: ws.display_name(), + focused: state.active == Some(index), + pane_count: ws.panes.len(), + agent_state: pane_agent_state(agg_state), + } +} + +fn agent_name(agent: crate::detect::Agent) -> String { + match agent { + crate::detect::Agent::Pi => "pi", + crate::detect::Agent::Claude => "claude", + crate::detect::Agent::Codex => "codex", + crate::detect::Agent::Gemini => "gemini", + crate::detect::Agent::Cursor => "cursor", + crate::detect::Agent::Cline => "cline", + crate::detect::Agent::OpenCode => "opencode", + crate::detect::Agent::GithubCopilot => "copilot", + crate::detect::Agent::Kimi => "kimi", + crate::detect::Agent::Droid => "droid", + crate::detect::Agent::Amp => "amp", + } + .to_string() } diff --git a/src/main.rs b/src/main.rs index 6bcaa111..e3c54541 100644 --- a/src/main.rs +++ b/src/main.rs @@ -18,6 +18,7 @@ const NESTED_HERDR_MESSAGES: [&str; 6] = [ "recursion detected. base case not found. aborting.", ]; +mod api; mod app; mod config; mod detect; @@ -223,6 +224,9 @@ fn main() -> io::Result<()> { init_logging(); + let (api_tx, api_rx) = std::sync::mpsc::channel(); + let _api_server = api::start_server(api_tx)?; + let no_session = std::env::args().any(|a| a == "--no-session"); let in_tmux = std::env::var("TMUX").is_ok(); @@ -284,7 +288,7 @@ fn main() -> io::Result<()> { std::io::stdout().flush()?; } - let mut app = app::App::new(config, no_session, config_diagnostic); + let mut app = app::App::new(config, no_session, config_diagnostic, api_rx); let result = app.run(&mut terminal).await; // Reset modifyOtherKeys if we enabled it diff --git a/src/pane.rs b/src/pane.rs index fa08bf68..9e8d6b1e 100644 --- a/src/pane.rs +++ b/src/pane.rs @@ -102,6 +102,50 @@ impl Drop for PaneRuntime { } } +fn trim_trailing_blank_rows(rows: &mut Vec) { + while rows.last().is_some_and(|row| row.trim().is_empty()) { + rows.pop(); + } +} + +fn recent_text_from_parser(parser: &mut vt100::Parser, lines: usize) -> String { + let screen = parser.screen_mut(); + let original_scrollback = screen.scrollback(); + + screen.set_scrollback(usize::MAX); + let max_scrollback = screen.scrollback(); + + let (_, cols) = screen.size(); + screen.set_scrollback(0); + let visible_rows: Vec = screen + .rows(0, cols) + .map(|row| row.trim_end().to_string()) + .collect(); + let extra_rows = lines.saturating_sub(visible_rows.len()).min(max_scrollback); + + let mut rows = Vec::with_capacity(extra_rows + visible_rows.len()); + if extra_rows > 0 { + for offset in (1..=extra_rows).rev() { + screen.set_scrollback(offset); + if let Some(row) = screen.rows(0, cols).next() { + rows.push(row.trim_end().to_string()); + } + } + } + + screen.set_scrollback(original_scrollback); + rows.extend(visible_rows); + trim_trailing_blank_rows(&mut rows); + + let start = rows.len().saturating_sub(lines); + let text = rows[start..].join("\n"); + if text.is_empty() { + text + } else { + format!("{text}\n") + } +} + impl PaneRuntime { pub fn spawn( pane_id: PaneId, @@ -418,6 +462,30 @@ impl PaneRuntime { } } + pub fn visible_text(&self) -> String { + let Ok(content) = self.screen_content.read() else { + return String::new(); + }; + let mut rows: Vec = content + .lines() + .map(|line| line.trim_end().to_string()) + .collect(); + trim_trailing_blank_rows(&mut rows); + let text = rows.join("\n"); + if text.is_empty() { + text + } else { + format!("{text}\n") + } + } + + pub fn recent_text(&self, lines: usize) -> String { + self.parser + .write() + .map(|mut parser| recent_text_from_parser(&mut parser, lines)) + .unwrap_or_default() + } + /// Get the current working directory of the child shell process. pub fn cwd(&self) -> Option { let pid = self.child_pid.load(Ordering::Relaxed); @@ -429,6 +497,23 @@ impl PaneRuntime { mod tests { use super::*; + #[test] + fn recent_text_reconstructs_scrollback_tail() { + let responses = PtyResponses::new(); + let mut parser = vt100::Parser::new_with_callbacks(3, 10, 100, responses); + parser.process(b"a\r\nb\r\nc\r\nd\r\ne"); + + let recent = recent_text_from_parser(&mut parser, 4); + assert_eq!(recent, "b\nc\nd\ne\n"); + } + + #[test] + fn trim_trailing_blank_rows_drops_empty_viewport_tail() { + let mut rows = vec!["hello".to_string(), "".to_string(), " ".to_string()]; + trim_trailing_blank_rows(&mut rows); + assert_eq!(rows, vec!["hello".to_string()]); + } + #[test] fn claude_busy_is_sticky_for_short_gap() { let now = std::time::Instant::now(); diff --git a/tests/api_ping.rs b/tests/api_ping.rs new file mode 100644 index 00000000..810b6ffe --- /dev/null +++ b/tests/api_ping.rs @@ -0,0 +1,246 @@ +use std::fs; +use std::io::{BufRead, BufReader, Write}; +use std::os::unix::net::UnixStream; +use std::path::{Path, PathBuf}; +use std::thread; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +use portable_pty::{native_pty_system, Child, CommandBuilder, MasterPty, PtySize}; + +fn unique_test_dir() -> PathBuf { + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_nanos()) + .unwrap_or(0); + std::env::temp_dir().join(format!("herdr-api-test-{}-{nanos}", std::process::id())) +} + +struct SpawnedHerdr { + _master: Box, + child: Box, +} + +fn wait_for_socket(path: &Path, timeout: Duration) { + let deadline = Instant::now() + timeout; + while Instant::now() < deadline { + if path.exists() && UnixStream::connect(path).is_ok() { + return; + } + thread::sleep(Duration::from_millis(25)); + } + panic!("socket did not appear at {}", path.display()); +} + +fn spawn_herdr(config_home: &Path, runtime_dir: &Path, socket_path: &Path) -> SpawnedHerdr { + fs::create_dir_all(config_home.join("herdr")).unwrap(); + fs::create_dir_all(runtime_dir).unwrap(); + fs::write( + config_home.join("herdr/config.toml"), + "onboarding = false\n", + ) + .unwrap(); + + let pair = native_pty_system() + .openpty(PtySize { + rows: 24, + cols: 80, + pixel_width: 0, + pixel_height: 0, + }) + .unwrap(); + + let mut cmd = CommandBuilder::new(env!("CARGO_BIN_EXE_herdr")); + cmd.arg("--no-session"); + cmd.env("XDG_CONFIG_HOME", config_home); + cmd.env("XDG_RUNTIME_DIR", runtime_dir); + cmd.env("HERDR_SOCKET_PATH", socket_path); + + let child = pair.slave.spawn_command(cmd).unwrap(); + + SpawnedHerdr { + _master: pair.master, + child, + } +} + +fn send_request(socket_path: &Path, json: &str) -> serde_json::Value { + let mut stream = UnixStream::connect(socket_path).unwrap(); + stream.write_all(json.as_bytes()).unwrap(); + stream.write_all(b"\n").unwrap(); + stream.flush().unwrap(); + + let mut line = String::new(); + let mut reader = BufReader::new(stream); + reader.read_line(&mut line).unwrap(); + serde_json::from_str(&line).unwrap() +} + +#[test] +fn ping_over_socket_returns_version() { + let base = unique_test_dir(); + let config_home = base.join("config"); + let runtime_dir = base.join("runtime"); + let socket_path = runtime_dir.join("herdr.sock"); + + let mut child = spawn_herdr(&config_home, &runtime_dir, &socket_path); + wait_for_socket(&socket_path, Duration::from_secs(5)); + + let value = send_request( + &socket_path, + r#"{"id":"req_1","method":"ping","params":{}}"#, + ); + assert_eq!(value["id"], "req_1"); + assert_eq!(value["result"]["type"], "pong"); + assert_eq!(value["result"]["version"], env!("CARGO_PKG_VERSION")); + + let _ = child.child.kill(); + let _ = child.child.wait(); + let _ = fs::remove_dir_all(base); +} + +#[test] +fn workspace_list_and_create_round_trip() { + let base = unique_test_dir(); + let config_home = base.join("config"); + let runtime_dir = base.join("runtime"); + let socket_path = runtime_dir.join("herdr.sock"); + + let mut child = spawn_herdr(&config_home, &runtime_dir, &socket_path); + wait_for_socket(&socket_path, Duration::from_secs(5)); + + let empty = send_request( + &socket_path, + r#"{"id":"req_2","method":"workspace.list","params":{}}"#, + ); + assert_eq!(empty["id"], "req_2"); + assert_eq!(empty["result"]["type"], "workspace_list"); + assert_eq!(empty["result"]["workspaces"].as_array().unwrap().len(), 0); + + let created = send_request( + &socket_path, + &format!( + r#"{{"id":"req_3","method":"workspace.create","params":{{"cwd":"{}","focus":true}}}}"#, + base.display() + ), + ); + assert_eq!(created["id"], "req_3"); + assert_eq!(created["result"]["type"], "workspace_info"); + assert_eq!(created["result"]["workspace"]["workspace_id"], "w_1"); + assert_eq!(created["result"]["workspace"]["number"], 1); + assert_eq!(created["result"]["workspace"]["focused"], true); + + let listed = send_request( + &socket_path, + r#"{"id":"req_4","method":"workspace.list","params":{}}"#, + ); + let workspaces = listed["result"]["workspaces"].as_array().unwrap(); + assert_eq!(workspaces.len(), 1); + assert_eq!(workspaces[0]["workspace_id"], "w_1"); + + let fetched = send_request( + &socket_path, + r#"{"id":"req_5","method":"workspace.get","params":{"workspace_id":"w_1"}}"#, + ); + assert_eq!(fetched["result"]["workspace"]["workspace_id"], "w_1"); + + let panes = send_request( + &socket_path, + r#"{"id":"req_6","method":"pane.list","params":{}}"#, + ); + let panes = panes["result"]["panes"].as_array().unwrap(); + assert_eq!(panes.len(), 1); + assert_eq!(panes[0]["workspace_id"], "w_1"); + let pane_id = panes[0]["pane_id"].as_str().unwrap().to_string(); + + let pane = send_request( + &socket_path, + &format!( + r#"{{"id":"req_7","method":"pane.get","params":{{"pane_id":"{}"}}}}"#, + pane_id + ), + ); + assert_eq!(pane["result"]["pane"]["pane_id"], pane_id); + + let read = send_request( + &socket_path, + &format!( + r#"{{"id":"req_8","method":"pane.read","params":{{"pane_id":"{}","source":"visible"}}}}"#, + pane_id + ), + ); + assert_eq!(read["result"]["read"]["pane_id"], pane_id); + assert!(read["result"]["read"]["text"].is_string()); + + let send_text = send_request( + &socket_path, + &format!( + r#"{{"id":"req_9","method":"pane.send_text","params":{{"pane_id":"{}","text":"echo alpha; echo beta; echo gamma"}}}}"#, + pane_id + ), + ); + assert_eq!(send_text["result"]["type"], "ok"); + + let send_enter = send_request( + &socket_path, + &format!( + r#"{{"id":"req_10","method":"pane.send_keys","params":{{"pane_id":"{}","keys":["Enter"]}}}}"#, + pane_id + ), + ); + assert_eq!(send_enter["result"]["type"], "ok"); + + std::thread::sleep(Duration::from_millis(300)); + + let recent = send_request( + &socket_path, + &format!( + r#"{{"id":"req_11","method":"pane.read","params":{{"pane_id":"{}","source":"recent","lines":20}}}}"#, + pane_id + ), + ); + let recent_text = recent["result"]["read"]["text"].as_str().unwrap(); + assert!(recent_text.contains("beta") || recent_text.contains("gamma")); + + let waited = send_request( + &socket_path, + &format!( + r#"{{"id":"req_12","method":"pane.wait_for_output","params":{{"pane_id":"{}","source":"recent","lines":40,"match":{{"type":"substring","value":"gamma"}},"timeout_ms":2000}}}}"#, + pane_id + ), + ); + assert_eq!(waited["result"]["type"], "output_matched"); + assert!(waited["result"]["matched_line"] + .as_str() + .unwrap() + .contains("gamma")); + assert!(waited["result"]["read"]["text"] + .as_str() + .unwrap() + .contains("gamma")); + + let waited_regex = send_request( + &socket_path, + &format!( + r#"{{"id":"req_13","method":"pane.wait_for_output","params":{{"pane_id":"{}","source":"recent","lines":40,"match":{{"type":"regex","value":"alp.*gamma"}},"timeout_ms":2000}}}}"#, + pane_id + ), + ); + assert_eq!(waited_regex["result"]["type"], "output_matched"); + assert!(waited_regex["result"]["matched_line"] + .as_str() + .unwrap() + .contains("alpha")); + + let timeout = send_request( + &socket_path, + &format!( + r#"{{"id":"req_14","method":"pane.wait_for_output","params":{{"pane_id":"{}","source":"recent","lines":10,"match":{{"type":"substring","value":"definitely-not-there"}},"timeout_ms":200}}}}"#, + pane_id + ), + ); + assert_eq!(timeout["error"]["code"], "timeout"); + + let _ = child.child.kill(); + let _ = child.child.wait(); + let _ = fs::remove_dir_all(base); +}