From 8b7848aac33ec67dee09e8d0b35fd70384aed1bb Mon Sep 17 00:00:00 2001 From: Ogulcan Celik Date: Sun, 29 Mar 2026 17:20:53 +0300 Subject: [PATCH] feat(api): add socket subscriptions and lifecycle events --- README.md | 28 ++- SOCKET_API.md | 360 +++++++++++++++++++++++++++++++++ src/api/mod.rs | 504 +++++++++++++++++++++++++++++++++++++++++++--- src/api/schema.rs | 157 ++++++++++++++- src/app/mod.rs | 375 ++++++++++++++++++++++++++++++++-- src/main.rs | 5 +- tests/api_ping.rs | 297 +++++++++++++++++++++++++++ 7 files changed, 1662 insertions(+), 64 deletions(-) create mode 100644 SOCKET_API.md diff --git a/README.md b/README.md index 1c6fa003..c0ac109f 100644 --- a/README.md +++ b/README.md @@ -7,7 +7,7 @@

herd your agents.

- herdr.dev · install · usage · configuration + herdr.dev · install · usage · configuration · socket api

--- @@ -174,12 +174,25 @@ this means detection works with any supported agent, installed any way, with zer the heuristics are pattern-matched against each agent's actual terminal output: prompt boxes, spinners, "waiting for input" messages, tool execution indicators. detection runs on a separate async task per pane, polled every 300-500ms, decoupled from terminal rendering. +## socket api + +herdr now has a local unix socket API for scripts, tools, and coding agents. + +you can: +- create, focus, rename, and close workspaces +- list, inspect, read, split, and close panes +- send text / keys into panes +- wait for output matches +- subscribe to lifecycle, agent, and output-match events over a single long-lived connection + +see [`SOCKET_API.md`](./SOCKET_API.md) for request shapes, examples, and subscription behavior. + ## what's coming -- **notification hooks**: a socket API so any agent or script can report its state to herdr. for agents without built-in detection, wire up a simple hook. +- **notification hooks**: richer agent/script-side state reporting on top of the socket foundation, so unsupported tools can report status directly to herdr. - **in-app preferences**: rerun onboarding and adjust things like sound and toast notifications without editing config by hand. - **native notifications**: OS-level notifications when an agent needs attention and herdr isn't in focus. -- **agent API**: `herdr create`, `herdr split`, `herdr send`, so agents and scripts can manage herdr workspaces programmatically. +- **agent cli wrapper**: `herdr pane ...`, `herdr wait ...`, and similar shell-friendly commands layered on top of the socket API. ## built with agents @@ -191,6 +204,8 @@ there will be rough edges. if you hit one, [open an issue](https://github.com/og ## cli +current built-in commands: + ``` herdr launch herdr herdr update download and install the latest version @@ -200,6 +215,11 @@ herdr --no-session start without restoring or saving sessions herdr --help show help ``` +programmatic control today lives in the socket API: +- [`SOCKET_API.md`](./SOCKET_API.md) + +shell-friendly wrapper commands like `herdr pane ...` and `herdr wait ...` are planned on top of that API. + ## building from source ```bash @@ -212,7 +232,7 @@ cargo build --release ## testing ```bash -just test # unit tests (157 tests) +just test # unit tests just test-integration # LLM-based integration tests just test-all # both ``` diff --git a/SOCKET_API.md b/SOCKET_API.md new file mode 100644 index 00000000..48c74741 --- /dev/null +++ b/SOCKET_API.md @@ -0,0 +1,360 @@ +# herdr socket api + +herdr exposes a local unix socket API for scripts, tools, and coding agents that want to control a running herdr instance or subscribe to pane/workspace events. + +this is the low-level integration surface. +a CLI wrapper on top of it is planned, but the socket API is the foundation. + +## transport + +- transport: unix domain socket +- encoding: newline-delimited JSON +- request/response: one JSON request per line, one JSON response per line +- subscriptions: send `events.subscribe`, receive an ack, then keep the same connection open for pushed events + +socket path resolution: + +1. `HERDR_SOCKET_PATH` +2. `$XDG_RUNTIME_DIR/herdr.sock` +3. `$XDG_CONFIG_HOME/herdr/herdr.sock` +4. `$HOME/.config/herdr/herdr.sock` +5. `/tmp/herdr.sock` + +## request shape + +all requests use this envelope: + +```json +{ + "id": "req_1", + "method": "ping", + "params": {} +} +``` + +success responses: + +```json +{ + "id": "req_1", + "result": { + "type": "pong", + "version": "0.1.2" + } +} +``` + +error responses: + +```json +{ + "id": "req_1", + "error": { + "code": "pane_not_found", + "message": "pane p_1_99 not found" + } +} +``` + +## ids + +workspace ids look like: + +- `w_1` +- `w_2` + +pane ids look like: + +- `p_1_1` +- `p_1_2` +- `p_2_1` + +that means: +- first number = workspace number +- second number = pane id inside that workspace + +## core request methods + +currently useful methods include: + +### basic +- `ping` + +### workspace +- `workspace.list` +- `workspace.get` +- `workspace.create` +- `workspace.focus` +- `workspace.rename` +- `workspace.close` + +### pane +- `pane.list` +- `pane.get` +- `pane.read` +- `pane.send_text` +- `pane.send_keys` +- `pane.split` +- `pane.close` + +### waits / events +- `pane.wait_for_output` +- `events.subscribe` + +## example: create a workspace + +```json +{ + "id": "req_create", + "method": "workspace.create", + "params": { + "cwd": "/home/can/Projects/herdr", + "focus": true + } +} +``` + +example response: + +```json +{ + "id": "req_create", + "result": { + "type": "workspace_info", + "workspace": { + "workspace_id": "w_1", + "number": 1, + "label": "herdr", + "focused": true, + "pane_count": 1, + "agent_state": "unknown" + } + } +} +``` + +## example: read pane output + +```json +{ + "id": "req_read", + "method": "pane.read", + "params": { + "pane_id": "p_1_1", + "source": "recent", + "lines": 80 + } +} +``` + +`source` can be: +- `visible` +- `recent` + +## example: send text and press enter + +low-level input is intentionally explicit: + +```json +{ + "id": "req_send_text", + "method": "pane.send_text", + "params": { + "pane_id": "p_1_1", + "text": "bun run dev" + } +} +``` + +then: + +```json +{ + "id": "req_send_keys", + "method": "pane.send_keys", + "params": { + "pane_id": "p_1_1", + "keys": ["Enter"] + } +} +``` + +this is kept separate on purpose. sending text is not always the same thing as submitting it. + +a future CLI wrapper will likely offer a more ergonomic `pane run` style command on top of this. + +## example: one-shot wait for output + +```json +{ + "id": "req_wait", + "method": "pane.wait_for_output", + "params": { + "pane_id": "p_1_1", + "source": "recent", + "lines": 200, + "match": { "type": "substring", "value": "ready" }, + "timeout_ms": 30000 + } +} +``` + +regex matching is also supported: + +```json +{ + "type": "regex", + "value": "server.*ready" +} +``` + +## subscriptions + +`events.subscribe` is the long-lived pubsub entrypoint. + +you send a subscribe request once, get an ack on the same connection, and then keep reading newline-delimited JSON events from that same socket. + +### subscription ack + +```json +{ + "id": "sub_1", + "result": { + "type": "subscription_started" + } +} +``` + +## supported subscriptions + +### lifecycle / base events +- `workspace.created` +- `workspace.closed` +- `workspace.focused` +- `pane.created` +- `pane.closed` +- `pane.focused` +- `pane.exited` +- `pane.agent_detected` +- `pane.agent_state_changed` + +### parameterized event +- `pane.output_matched` + +## example: subscribe to lifecycle events + +```json +{ + "id": "sub_life", + "method": "events.subscribe", + "params": { + "subscriptions": [ + { "type": "workspace.created" }, + { "type": "workspace.focused" }, + { "type": "pane.created" }, + { "type": "pane.focused" }, + { "type": "pane.agent_detected" }, + { "type": "pane.closed" }, + { "type": "workspace.closed" } + ] + } +} +``` + +example pushed event: + +```json +{ + "event": "workspace_created", + "data": { + "workspace": { + "workspace_id": "w_1", + "number": 1, + "label": "herdr", + "focused": true, + "pane_count": 1, + "agent_state": "unknown" + } + } +} +``` + +## example: subscribe to output matches and agent state changes + +```json +{ + "id": "sub_1", + "method": "events.subscribe", + "params": { + "subscriptions": [ + { + "type": "pane.output_matched", + "pane_id": "p_1_1", + "source": "recent", + "lines": 200, + "match": { "type": "substring", "value": "ready" } + }, + { + "type": "pane.agent_state_changed", + "pane_id": "p_1_1", + "state": "idle" + } + ] + } +} +``` + +example pushed `pane.output_matched` event: + +```json +{ + "event": "pane.output_matched", + "data": { + "pane_id": "p_1_1", + "matched_line": "server ready", + "read": { + "pane_id": "p_1_1", + "workspace_id": "w_1", + "source": "recent", + "text": "...server ready...", + "revision": 0, + "truncated": false + } + } +} +``` + +example pushed `pane.agent_state_changed` event: + +```json +{ + "event": "pane.agent_state_changed", + "data": { + "pane_id": "p_1_1", + "workspace_id": "w_1", + "state": "idle", + "agent": "pi" + } +} +``` + +## behavior notes + +- `pane.output_matched` emits when a subscription transitions into a matching state. it does not repeatedly spam the same visible match on every poll. +- closing the socket connection ends the subscription. +- there is no separate transport for events. +- the same herdr process can serve regular request/response calls and long-lived subscription connections at the same time. + +## intended layering + +recommended architecture: +- socket api = foundational integration protocol +- future `herdr ...` commands = ergonomic wrapper for humans and coding agents + +for agent workflows, the future CLI should likely expose blocking commands like: +- `herdr pane read ...` +- `herdr pane run ...` +- `herdr wait output ...` +- `herdr wait agent-state ...` + +but those should sit on top of this socket surface rather than replacing it. diff --git a/src/api/mod.rs b/src/api/mod.rs index 642dc70d..e2c29490 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -4,13 +4,16 @@ use std::fs; use std::io::{BufRead, BufReader, Write}; use std::os::unix::net::{UnixListener, UnixStream}; use std::path::{Path, PathBuf}; +use std::time::Duration; use tracing::{debug, error, info, warn}; use regex::Regex; use crate::api::schema::{ - ErrorBody, ErrorResponse, Method, Request, ResponseResult, SuccessResponse, + ErrorBody, ErrorResponse, Method, PaneAgentStateChangedEvent, PaneOutputMatchedEvent, Request, + ResponseResult, Subscription, SubscriptionEventData, SubscriptionEventEnvelope, + SubscriptionEventKind, SuccessResponse, }; pub const SOCKET_PATH_ENV_VAR: &str = "HERDR_SOCKET_PATH"; @@ -20,6 +23,46 @@ pub struct ApiRequestMessage { pub respond_to: std::sync::mpsc::Sender, } +#[derive(Clone, Default)] +pub struct EventHub { + inner: std::sync::Arc>, +} + +#[derive(Default)] +struct EventHubState { + next_sequence: u64, + events: Vec<(u64, crate::api::schema::EventEnvelope)>, +} + +impl EventHub { + const MAX_EVENTS: usize = 512; + + pub fn push(&self, event: crate::api::schema::EventEnvelope) { + let Ok(mut state) = self.inner.lock() else { + return; + }; + state.next_sequence += 1; + let sequence = state.next_sequence; + state.events.push((sequence, event)); + let overflow = state.events.len().saturating_sub(Self::MAX_EVENTS); + if overflow > 0 { + state.events.drain(0..overflow); + } + } + + pub fn events_after(&self, sequence: u64) -> Vec<(u64, crate::api::schema::EventEnvelope)> { + let Ok(state) = self.inner.lock() else { + return Vec::new(); + }; + state + .events + .iter() + .filter(|(event_sequence, _)| *event_sequence > sequence) + .cloned() + .collect() + } +} + pub fn socket_path() -> PathBuf { if let Ok(path) = std::env::var(SOCKET_PATH_ENV_VAR) { return PathBuf::from(path); @@ -57,6 +100,7 @@ impl Drop for ServerHandle { pub fn start_server( api_tx: std::sync::mpsc::Sender, + event_hub: EventHub, ) -> std::io::Result { let path = socket_path(); prepare_socket_path(&path)?; @@ -68,9 +112,13 @@ pub fn start_server( for stream in listener.incoming() { match stream { Ok(stream) => { - if let Err(err) = handle_connection(stream, &api_tx) { - warn!(err = %err, "api connection failed"); - } + let api_tx = api_tx.clone(); + let event_hub = event_hub.clone(); + std::thread::spawn(move || { + if let Err(err) = handle_connection(stream, &api_tx, &event_hub) { + warn!(err = %err, "api connection failed"); + } + }); } Err(err) => { error!(err = %err, "api listener accept failed"); @@ -104,6 +152,7 @@ fn prepare_socket_path(path: &Path) -> std::io::Result<()> { fn handle_connection( mut stream: UnixStream, api_tx: &std::sync::mpsc::Sender, + event_hub: &EventHub, ) -> std::io::Result<()> { let mut line = String::new(); { @@ -119,21 +168,41 @@ fn handle_connection( 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}"), - }, - })?, + let request = match serde_json::from_str::(line) { + Ok(request) => request, + Err(err) => { + write_json_line( + &mut stream, + &ErrorResponse { + id: String::new(), + error: ErrorBody { + code: "invalid_request".into(), + message: format!("invalid request: {err}"), + }, + }, + )?; + return Ok(()); + } }; - stream.write_all(response.as_bytes())?; - stream.write_all(b"\n")?; - stream.flush()?; - Ok(()) + match request.method { + Method::EventsSubscribe(params) => { + stream_subscriptions(stream, request.id, params, api_tx, event_hub) + } + method => { + let response = handle_request( + Request { + id: request.id, + method, + }, + api_tx, + ); + 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 { @@ -213,19 +282,7 @@ fn wait_for_output( .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()) - }), - }; + let matched_line = match_output(&read.text, ¶ms.r#match, regex.as_ref()); if matched_line.is_some() { let revision = read.revision; return serde_json::to_string(&SuccessResponse { @@ -255,6 +312,393 @@ fn wait_for_output( } } +fn stream_subscriptions( + mut stream: UnixStream, + request_id: String, + params: crate::api::schema::EventsSubscribeParams, + api_tx: &std::sync::mpsc::Sender, + event_hub: &EventHub, +) -> std::io::Result<()> { + let mut subscriptions = Vec::with_capacity(params.subscriptions.len()); + for (index, subscription) in params.subscriptions.into_iter().enumerate() { + let active = + match ActiveSubscription::new(subscription, &request_id, index, api_tx, event_hub) { + Ok(active) => active, + Err(response) => { + write_json_line(&mut stream, &response)?; + return Ok(()); + } + }; + subscriptions.push(active); + } + + write_json_line( + &mut stream, + &SuccessResponse { + id: request_id, + result: ResponseResult::SubscriptionStarted {}, + }, + )?; + + loop { + for subscription in &mut subscriptions { + if let Some(event) = subscription.poll(api_tx, event_hub) { + write_json_line(&mut stream, &event)?; + } + } + std::thread::sleep(Duration::from_millis(100)); + } +} + +fn write_json_line(stream: &mut UnixStream, value: &T) -> std::io::Result<()> { + let encoded = serde_json::to_string(value) + .map_err(|err| std::io::Error::other(format!("failed to encode json: {err}")))?; + stream.write_all(encoded.as_bytes())?; + stream.write_all(b"\n")?; + stream.flush()?; + Ok(()) +} + +fn match_output( + text: &str, + matcher: &crate::api::schema::OutputMatch, + regex: Option<&Regex>, +) -> Option { + match matcher { + crate::api::schema::OutputMatch::Substring { value } => text + .lines() + .find(|line| line.contains(value)) + .map(|line| line.to_string()), + crate::api::schema::OutputMatch::Regex { .. } => regex.and_then(|re| { + text.lines() + .find(|line| re.is_match(line)) + .map(|line| line.to_string()) + }), + } +} + +struct ActiveOutputMatchedSubscription { + pane_id: String, + source: crate::api::schema::ReadSource, + lines: Option, + matcher: crate::api::schema::OutputMatch, + regex: Option, + strip_ansi: bool, + currently_matching: bool, + request_prefix: String, +} + +struct ActiveAgentStateChangedSubscription { + pane_id: String, + state_filter: Option, + last_state: Option, + request_prefix: String, +} + +struct ActiveEventSubscription { + event_kind: crate::api::schema::EventKind, + last_sequence: u64, +} + +enum ActiveSubscription { + Event(ActiveEventSubscription), + OutputMatched(ActiveOutputMatchedSubscription), + AgentStateChanged(ActiveAgentStateChangedSubscription), +} + +impl ActiveSubscription { + fn new( + subscription: Subscription, + request_id: &str, + index: usize, + api_tx: &std::sync::mpsc::Sender, + _event_hub: &EventHub, + ) -> Result { + match subscription { + Subscription::WorkspaceCreated {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::WorkspaceCreated, + last_sequence: 0, + })), + Subscription::WorkspaceClosed {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::WorkspaceClosed, + last_sequence: 0, + })), + Subscription::WorkspaceFocused {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::WorkspaceFocused, + last_sequence: 0, + })), + Subscription::PaneCreated {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::PaneCreated, + last_sequence: 0, + })), + Subscription::PaneClosed {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::PaneClosed, + last_sequence: 0, + })), + Subscription::PaneFocused {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::PaneFocused, + last_sequence: 0, + })), + Subscription::PaneExited {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::PaneExited, + last_sequence: 0, + })), + Subscription::PaneAgentDetected {} => Ok(Self::Event(ActiveEventSubscription { + event_kind: crate::api::schema::EventKind::PaneAgentDetected, + last_sequence: 0, + })), + Subscription::PaneOutputMatched { + pane_id, + source, + lines, + r#match, + strip_ansi, + } => { + let regex = match &r#match { + crate::api::schema::OutputMatch::Regex { value } => match Regex::new(value) { + Ok(regex) => Some(regex), + Err(err) => { + return Err(ErrorResponse { + id: request_id.to_string(), + error: ErrorBody { + code: "invalid_regex".into(), + message: err.to_string(), + }, + }); + } + }, + crate::api::schema::OutputMatch::Substring { .. } => None, + }; + + let probe = pane_read( + format!("{request_id}:sub:{index}:probe"), + &pane_id, + source.clone(), + lines, + strip_ansi, + api_tx, + ); + if let Err(error) = probe { + return Err(error); + } + + Ok(Self::OutputMatched(ActiveOutputMatchedSubscription { + pane_id, + source, + lines, + matcher: r#match, + regex, + strip_ansi, + currently_matching: false, + request_prefix: format!("{request_id}:sub:{index}"), + })) + } + Subscription::PaneAgentStateChanged { pane_id, state } => { + let probe = + match pane_get(format!("{request_id}:sub:{index}:probe"), &pane_id, api_tx) { + Ok(probe) => probe, + Err(error) => return Err(error), + }; + + Ok(Self::AgentStateChanged( + ActiveAgentStateChangedSubscription { + pane_id, + state_filter: state, + last_state: Some(probe.agent_state), + request_prefix: format!("{request_id}:sub:{index}"), + }, + )) + } + } + } + + fn poll( + &mut self, + api_tx: &std::sync::mpsc::Sender, + event_hub: &EventHub, + ) -> Option { + match self { + Self::Event(subscription) => subscription.poll(event_hub), + Self::OutputMatched(subscription) => { + serde_json::to_value(subscription.poll(api_tx)?).ok() + } + Self::AgentStateChanged(subscription) => { + serde_json::to_value(subscription.poll(api_tx)?).ok() + } + } + } +} + +impl ActiveEventSubscription { + fn poll(&mut self, event_hub: &EventHub) -> Option { + for (sequence, event) in event_hub.events_after(self.last_sequence) { + self.last_sequence = sequence; + if event.event == self.event_kind { + return serde_json::to_value(event).ok(); + } + } + None + } +} + +impl ActiveOutputMatchedSubscription { + fn poll( + &mut self, + api_tx: &std::sync::mpsc::Sender, + ) -> Option { + let read = pane_read( + format!("{}:read", self.request_prefix), + &self.pane_id, + self.source.clone(), + self.lines, + self.strip_ansi, + api_tx, + ) + .ok()?; + + let matched_line = match_output(&read.text, &self.matcher, self.regex.as_ref()); + match matched_line { + Some(matched_line) => { + if self.currently_matching { + return None; + } + self.currently_matching = true; + Some(SubscriptionEventEnvelope { + event: SubscriptionEventKind::PaneOutputMatched, + data: SubscriptionEventData::PaneOutputMatched(PaneOutputMatchedEvent { + pane_id: self.pane_id.clone(), + matched_line, + read, + }), + }) + } + None => { + self.currently_matching = false; + None + } + } + } +} + +impl ActiveAgentStateChangedSubscription { + fn poll( + &mut self, + api_tx: &std::sync::mpsc::Sender, + ) -> Option { + let pane = pane_get( + format!("{}:pane", self.request_prefix), + &self.pane_id, + api_tx, + ) + .ok()?; + let current_state = pane.agent_state; + let previous_state = self.last_state.replace(current_state); + if previous_state.is_none() || previous_state == Some(current_state) { + return None; + } + if self + .state_filter + .is_some_and(|wanted| wanted != current_state) + { + return None; + } + + Some(SubscriptionEventEnvelope { + event: SubscriptionEventKind::PaneAgentStateChanged, + data: SubscriptionEventData::PaneAgentStateChanged(PaneAgentStateChangedEvent { + pane_id: pane.pane_id, + workspace_id: pane.workspace_id, + state: current_state, + agent: pane.agent, + }), + }) + } +} + +fn pane_read( + request_id: String, + pane_id: &str, + source: crate::api::schema::ReadSource, + lines: Option, + strip_ansi: bool, + api_tx: &std::sync::mpsc::Sender, +) -> Result { + let response = dispatch_to_app( + Request { + id: request_id.clone(), + method: Method::PaneRead(crate::api::schema::PaneReadParams { + pane_id: pane_id.to_string(), + source, + lines, + strip_ansi, + }), + }, + api_tx, + ); + let value: serde_json::Value = serde_json::from_str(&response).map_err(|_| ErrorResponse { + id: request_id.clone(), + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane read response".into(), + }, + })?; + if value.get("error").is_some() { + return serde_json::from_value(value).map_err(|_| ErrorResponse { + id: request_id, + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane read error".into(), + }, + }); + } + serde_json::from_value(value["result"]["read"].clone()).map_err(|_| ErrorResponse { + id: request_id, + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane read result".into(), + }, + }) +} + +fn pane_get( + request_id: String, + pane_id: &str, + api_tx: &std::sync::mpsc::Sender, +) -> Result { + let response = dispatch_to_app( + Request { + id: request_id.clone(), + method: Method::PaneGet(crate::api::schema::PaneTarget { + pane_id: pane_id.to_string(), + }), + }, + api_tx, + ); + let value: serde_json::Value = serde_json::from_str(&response).map_err(|_| ErrorResponse { + id: request_id.clone(), + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane get response".into(), + }, + })?; + if value.get("error").is_some() { + return serde_json::from_value(value).map_err(|_| ErrorResponse { + id: request_id, + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane get error".into(), + }, + }); + } + serde_json::from_value(value["result"]["pane"].clone()).map_err(|_| ErrorResponse { + id: request_id, + error: ErrorBody { + code: "internal_error".into(), + message: "failed to decode pane get result".into(), + }, + }) +} + fn dispatch_to_app( request: Request, api_tx: &std::sync::mpsc::Sender, diff --git a/src/api/schema.rs b/src/api/schema.rs index 29509211..7de0bbed 100644 --- a/src/api/schema.rs +++ b/src/api/schema.rs @@ -132,11 +132,44 @@ pub enum ReadSource { #[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, + pub subscriptions: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "type")] +pub enum Subscription { + #[serde(rename = "workspace.created")] + WorkspaceCreated {}, + #[serde(rename = "workspace.closed")] + WorkspaceClosed {}, + #[serde(rename = "workspace.focused")] + WorkspaceFocused {}, + #[serde(rename = "pane.created")] + PaneCreated {}, + #[serde(rename = "pane.closed")] + PaneClosed {}, + #[serde(rename = "pane.focused")] + PaneFocused {}, + #[serde(rename = "pane.exited")] + PaneExited {}, + #[serde(rename = "pane.agent_detected")] + PaneAgentDetected {}, + #[serde(rename = "pane.output_matched")] + PaneOutputMatched { + pane_id: String, + source: ReadSource, + #[serde(default, skip_serializing_if = "Option::is_none")] + lines: Option, + r#match: OutputMatch, + #[serde(default = "default_true")] + strip_ansi: bool, + }, + #[serde(rename = "pane.agent_state_changed")] + PaneAgentStateChanged { + pane_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + state: Option, + }, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -270,9 +303,7 @@ pub enum ResponseResult { PaneRead { read: PaneReadResult, }, - SubscriptionStarted { - events: Vec, - }, + SubscriptionStarted {}, WaitMatched { event: EventEnvelope, }, @@ -324,6 +355,43 @@ pub struct EventEnvelope { pub data: EventData, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +pub enum SubscriptionEventKind { + #[serde(rename = "pane.output_matched")] + PaneOutputMatched, + #[serde(rename = "pane.agent_state_changed")] + PaneAgentStateChanged, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SubscriptionEventEnvelope { + pub event: SubscriptionEventKind, + pub data: SubscriptionEventData, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(untagged)] +pub enum SubscriptionEventData { + PaneOutputMatched(PaneOutputMatchedEvent), + PaneAgentStateChanged(PaneAgentStateChangedEvent), +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneOutputMatchedEvent { + pub pane_id: String, + pub matched_line: String, + pub read: PaneReadResult, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct PaneAgentStateChangedEvent { + pub pane_id: String, + pub workspace_id: String, + pub state: PaneAgentState, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub agent: Option, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "type", rename_all = "snake_case")] pub enum EventData { @@ -476,6 +544,79 @@ mod tests { assert_eq!(restored, event); } + #[test] + fn subscribe_request_parses_parameterized_subscriptions() { + let json = r#" + { + "id": "sub_1", + "method": "events.subscribe", + "params": { + "subscriptions": [ + { + "type": "pane.output_matched", + "pane_id": "p_1_1", + "source": "recent", + "lines": 200, + "match": { "type": "substring", "value": "auth: received" } + }, + { + "type": "pane.agent_state_changed", + "pane_id": "p_1_1", + "state": "waiting" + } + ] + } + } + "#; + + let request: Request = serde_json::from_str(json).unwrap(); + let Method::EventsSubscribe(params) = request.method else { + panic!("wrong method parsed"); + }; + assert_eq!(params.subscriptions.len(), 2); + assert!(matches!( + ¶ms.subscriptions[0], + Subscription::PaneOutputMatched { + pane_id, + source: ReadSource::Recent, + lines: Some(200), + r#match: OutputMatch::Substring { value }, + strip_ansi: true, + } if pane_id == "p_1_1" && value == "auth: received" + )); + assert!(matches!( + ¶ms.subscriptions[1], + Subscription::PaneAgentStateChanged { + pane_id, + state: Some(PaneAgentState::Waiting), + } if pane_id == "p_1_1" + )); + } + + #[test] + fn subscription_event_envelope_round_trips() { + let event = SubscriptionEventEnvelope { + event: SubscriptionEventKind::PaneOutputMatched, + data: SubscriptionEventData::PaneOutputMatched(PaneOutputMatchedEvent { + pane_id: "p_1_1".into(), + matched_line: "auth: received".into(), + read: PaneReadResult { + pane_id: "p_1_1".into(), + workspace_id: "w_1".into(), + source: ReadSource::Recent, + text: "auth: received\n".into(), + revision: 0, + truncated: false, + }, + }), + }; + + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"event\":\"pane.output_matched\"")); + let restored: SubscriptionEventEnvelope = serde_json::from_str(&json).unwrap(); + assert_eq!(restored, event); + } + #[test] fn success_response_round_trips() { let response = SuccessResponse { diff --git a/src/app/mod.rs b/src/app/mod.rs index 48700577..d5fa0ef2 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -29,6 +29,8 @@ pub struct App { pub event_tx: mpsc::Sender, event_rx: mpsc::Receiver, api_rx: std::sync::mpsc::Receiver, + event_hub: crate::api::EventHub, + last_focus: Option<(usize, crate::layout::PaneId)>, no_session: bool, config_diagnostic_deadline: Option, toast_deadline: Option, @@ -40,6 +42,7 @@ impl App { no_session: bool, config_diagnostic: Option, api_rx: std::sync::mpsc::Receiver, + event_hub: crate::api::EventHub, ) -> Self { let (prefix_code, prefix_mods) = config.prefix_key(); let (event_tx, event_rx) = mpsc::channel::(64); @@ -111,6 +114,13 @@ impl App { std::thread::spawn(move || crate::update::auto_update(update_tx)); } + let last_focus = state.active.and_then(|idx| { + state + .workspaces + .get(idx) + .map(|ws| (idx, ws.layout.focused())) + }); + Self { config_diagnostic_deadline: state .config_diagnostic @@ -121,6 +131,8 @@ impl App { event_tx, event_rx, api_rx, + event_hub, + last_focus, no_session, } } @@ -148,25 +160,15 @@ impl App { crate::ui::render(&self.state, frame); })?; + // Drain internal events first so API reads observe fresh pane state. + self.drain_internal_events(); + 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(); - self.state.handle_app_event(ev); - if self.state.toast != previous_toast { - self.toast_deadline = self.state.toast.as_ref().map(|toast| { - let duration = match toast.kind { - ToastKind::NeedsAttention => Duration::from_secs(8), - ToastKind::Finished => Duration::from_secs(5), - }; - Instant::now() + duration - }); - } - } + self.sync_focus_events(); if event::poll(Duration::from_millis(16))? { match event::read()? { @@ -204,7 +206,114 @@ impl App { Ok(()) } + fn drain_internal_events(&mut self) { + while let Ok(ev) = self.event_rx.try_recv() { + match &ev { + AppEvent::PaneDied { pane_id } => { + if let Some((ws_idx, _)) = self.find_pane(*pane_id) { + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneExited, + data: crate::api::schema::EventData::PaneExited { + pane_id: format!("p_{}_{}", ws_idx + 1, pane_id.raw()), + workspace_id: format!("w_{}", ws_idx + 1), + }, + }); + } + } + AppEvent::StateChanged { + pane_id, + agent, + state, + } => { + if let Some((ws_idx, pane)) = self.find_pane(*pane_id) { + let pane_id = format!("p_{}_{}", ws_idx + 1, pane_id.raw()); + let workspace_id = format!("w_{}", ws_idx + 1); + if pane.detected_agent != *agent { + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneAgentDetected, + data: crate::api::schema::EventData::PaneAgentDetected { + pane_id: pane_id.clone(), + workspace_id: workspace_id.clone(), + agent: agent.map(agent_name), + }, + }); + } + if pane.state != *state { + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneAgentStateChanged, + data: crate::api::schema::EventData::PaneAgentStateChanged { + pane_id, + workspace_id, + state: pane_agent_state(*state), + }, + }); + } + } + } + AppEvent::UpdateReady { .. } => {} + } + + let previous_toast = self.state.toast.clone(); + self.state.handle_app_event(ev); + if self.state.toast != previous_toast { + self.toast_deadline = self.state.toast.as_ref().map(|toast| { + let duration = match toast.kind { + ToastKind::NeedsAttention => Duration::from_secs(8), + ToastKind::Finished => Duration::from_secs(5), + }; + Instant::now() + duration + }); + } + } + } + + fn emit_event(&self, event: crate::api::schema::EventEnvelope) { + self.event_hub.push(event); + } + + fn sync_focus_events(&mut self) { + let current_focus = self.state.active.and_then(|idx| { + self.state + .workspaces + .get(idx) + .map(|ws| (idx, ws.layout.focused())) + }); + if current_focus == self.last_focus { + return; + } + + if let Some((ws_idx, pane_id)) = current_focus { + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::WorkspaceFocused, + data: crate::api::schema::EventData::WorkspaceFocused { + workspace_id: format!("w_{}", ws_idx + 1), + }, + }); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneFocused, + data: crate::api::schema::EventData::PaneFocused { + pane_id: format!("p_{}_{}", ws_idx + 1, pane_id.raw()), + workspace_id: format!("w_{}", ws_idx + 1), + }, + }); + } + + self.last_focus = current_focus; + } + + fn find_pane( + &self, + pane_id: crate::layout::PaneId, + ) -> Option<(usize, &crate::pane::PaneState)> { + self.state + .workspaces + .iter() + .enumerate() + .find_map(|(ws_idx, ws)| ws.panes.get(&pane_id).map(|pane| (ws_idx, pane))) + } + fn handle_api_request(&mut self, request: crate::api::schema::Request) -> String { + self.drain_internal_events(); use bytes::Bytes; use crate::api::schema::{ @@ -260,12 +369,32 @@ impl App { .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), - }, - }, + Ok(index) => { + let workspace = workspace_info(&self.state, index); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::WorkspaceCreated, + data: crate::api::schema::EventData::WorkspaceCreated { + workspace: workspace.clone(), + }, + }); + if let Some(pane_id) = self.state.workspaces[index] + .layout + .pane_ids() + .first() + .copied() + { + if let Some(pane) = self.pane_info(index, pane_id) { + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneCreated, + data: crate::api::schema::EventData::PaneCreated { pane }, + }); + } + } + SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceInfo { workspace }, + } + } Err(err) => { return serde_json::to_string(&ErrorResponse { id: request.id, @@ -278,6 +407,160 @@ impl App { } } } + Method::WorkspaceFocus(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(); + }; + if self.state.workspaces.get(index).is_none() { + 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(); + } + self.state.switch_workspace(index); + SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceInfo { + workspace: workspace_info(&self.state, index), + }, + } + } + Method::WorkspaceRename(params) => { + let Some(index) = parse_workspace_id(¶ms.workspace_id) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "workspace_not_found".into(), + message: format!("workspace {} not found", params.workspace_id), + }, + }) + .unwrap(); + }; + let Some(ws) = self.state.workspaces.get_mut(index) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "workspace_not_found".into(), + message: format!("workspace {} not found", params.workspace_id), + }, + }) + .unwrap(); + }; + ws.set_custom_name(params.label.clone()); + SuccessResponse { + id: request.id, + result: ResponseResult::WorkspaceInfo { + workspace: workspace_info(&self.state, index), + }, + } + } + Method::WorkspaceClose(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(); + }; + if self.state.workspaces.get(index).is_none() { + 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(); + } + self.state.selected = index; + self.state.close_selected_workspace(); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::WorkspaceClosed, + data: crate::api::schema::EventData::WorkspaceClosed { + workspace_id: target.workspace_id, + }, + }); + SuccessResponse { + id: request.id, + result: ResponseResult::Ok {}, + } + } + Method::PaneSplit(params) => { + let Some((ws_idx, target_pane_id)) = parse_pane_id(¶ms.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", params.target_pane_id), + }, + }) + .unwrap(); + }; + let (rows, cols) = self.state.estimate_pane_size(); + let Some(ws) = self.state.workspaces.get_mut(ws_idx) else { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_not_found".into(), + message: format!("pane {} not found", params.target_pane_id), + }, + }) + .unwrap(); + }; + ws.layout.focus_pane(target_pane_id); + let direction = match params.direction { + crate::api::schema::SplitDirection::Right => { + ratatui::layout::Direction::Horizontal + } + crate::api::schema::SplitDirection::Down => { + ratatui::layout::Direction::Vertical + } + }; + let new_pane_id = match ws.split_focused( + direction, + rows, + cols, + params.cwd.map(std::path::PathBuf::from), + ) { + Ok(new_pane_id) => new_pane_id, + Err(err) => { + return serde_json::to_string(&ErrorResponse { + id: request.id, + error: ErrorBody { + code: "pane_split_failed".into(), + message: err.to_string(), + }, + }) + .unwrap(); + } + }; + if !params.focus { + ws.layout.focus_pane(target_pane_id); + } + let pane = self.pane_info(ws_idx, new_pane_id).unwrap(); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneCreated, + data: crate::api::schema::EventData::PaneCreated { pane: pane.clone() }, + }); + SuccessResponse { + id: request.id, + result: ResponseResult::PaneInfo { pane }, + } + } Method::PaneList(PaneListParams { workspace_id }) => { match self.collect_panes_for_workspace(workspace_id.as_deref()) { Ok(panes) => SuccessResponse { @@ -395,6 +678,58 @@ impl App { result: ResponseResult::Ok {}, } } + Method::PaneClose(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(ws) = self.state.workspaces.get_mut(ws_idx) 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 workspace_id = format!("w_{}", ws_idx + 1); + let pane_count = ws.layout.pane_count(); + if pane_count <= 1 { + self.state.selected = ws_idx; + self.state.close_selected_workspace(); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneClosed, + data: crate::api::schema::EventData::PaneClosed { + pane_id: target.pane_id.clone(), + workspace_id: workspace_id.clone(), + }, + }); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::WorkspaceClosed, + data: crate::api::schema::EventData::WorkspaceClosed { workspace_id }, + }); + } else { + ws.remove_pane(pane_id); + self.emit_event(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneClosed, + data: crate::api::schema::EventData::PaneClosed { + pane_id: target.pane_id, + workspace_id, + }, + }); + } + 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 { diff --git a/src/main.rs b/src/main.rs index e3c54541..8b31cdc7 100644 --- a/src/main.rs +++ b/src/main.rs @@ -225,7 +225,8 @@ fn main() -> io::Result<()> { init_logging(); let (api_tx, api_rx) = std::sync::mpsc::channel(); - let _api_server = api::start_server(api_tx)?; + let event_hub = api::EventHub::default(); + let _api_server = api::start_server(api_tx, event_hub.clone())?; let no_session = std::env::args().any(|a| a == "--no-session"); let in_tmux = std::env::var("TMUX").is_ok(); @@ -288,7 +289,7 @@ fn main() -> io::Result<()> { std::io::stdout().flush()?; } - let mut app = app::App::new(config, no_session, config_diagnostic, api_rx); + let mut app = app::App::new(config, no_session, config_diagnostic, api_rx, event_hub); let result = app.run(&mut terminal).await; // Reset modifyOtherKeys if we enabled it diff --git a/tests/api_ping.rs b/tests/api_ping.rs index 810b6ffe..a102eb27 100644 --- a/tests/api_ping.rs +++ b/tests/api_ping.rs @@ -32,6 +32,15 @@ fn wait_for_socket(path: &Path, timeout: Duration) { } fn spawn_herdr(config_home: &Path, runtime_dir: &Path, socket_path: &Path) -> SpawnedHerdr { + spawn_herdr_with_path(config_home, runtime_dir, socket_path, None) +} + +fn spawn_herdr_with_path( + config_home: &Path, + runtime_dir: &Path, + socket_path: &Path, + path_override: Option<&Path>, +) -> SpawnedHerdr { fs::create_dir_all(config_home.join("herdr")).unwrap(); fs::create_dir_all(runtime_dir).unwrap(); fs::write( @@ -54,6 +63,9 @@ fn spawn_herdr(config_home: &Path, runtime_dir: &Path, socket_path: &Path) -> Sp cmd.env("XDG_CONFIG_HOME", config_home); cmd.env("XDG_RUNTIME_DIR", runtime_dir); cmd.env("HERDR_SOCKET_PATH", socket_path); + if let Some(path) = path_override { + cmd.env("PATH", path); + } let child = pair.slave.spawn_command(cmd).unwrap(); @@ -75,6 +87,38 @@ fn send_request(socket_path: &Path, json: &str) -> serde_json::Value { serde_json::from_str(&line).unwrap() } +fn open_subscription(socket_path: &Path, json: &str) -> (UnixStream, BufReader) { + 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 reader = BufReader::new(stream.try_clone().unwrap()); + (stream, reader) +} + +fn read_json_line(reader: &mut BufReader, timeout: Duration) -> serde_json::Value { + reader.get_ref().set_read_timeout(Some(timeout)).unwrap(); + let mut line = String::new(); + reader.read_line(&mut line).unwrap(); + serde_json::from_str(&line).unwrap() +} + +fn wait_for_event( + reader: &mut BufReader, + expected: &str, + timeout: Duration, +) -> serde_json::Value { + let deadline = Instant::now() + timeout; + loop { + let remaining = deadline.saturating_duration_since(Instant::now()); + let value = read_json_line(reader, remaining.max(Duration::from_millis(1))); + if value["event"] == expected { + return value; + } + } +} + #[test] fn ping_over_socket_returns_version() { let base = unique_test_dir(); @@ -244,3 +288,256 @@ fn workspace_list_and_create_round_trip() { let _ = child.child.wait(); let _ = fs::remove_dir_all(base); } + +#[test] +fn events_subscribe_streams_lifecycle_and_agent_events() { + 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 bin_dir = base.join("bin"); + + fs::create_dir_all(&bin_dir).unwrap(); + let fake_pi = bin_dir.join("pi"); + fs::write( + &fake_pi, + "#!/bin/sh\nprintf 'Working...\\n'\nsleep 1\nprintf '\\033[2J\\033[Hdone\\n'\n", + ) + .unwrap(); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mut perms = fs::metadata(&fake_pi).unwrap().permissions(); + perms.set_mode(0o755); + fs::set_permissions(&fake_pi, perms).unwrap(); + } + + let inherited_path = std::env::var("PATH").unwrap_or_default(); + let path_override = format!("{}:{}", bin_dir.display(), inherited_path); + let mut child = spawn_herdr_with_path( + &config_home, + &runtime_dir, + &socket_path, + Some(Path::new(&path_override)), + ); + wait_for_socket(&socket_path, Duration::from_secs(5)); + + let (_stream, mut reader) = open_subscription( + &socket_path, + r#"{"id":"sub_life","method":"events.subscribe","params":{"subscriptions":[{"type":"workspace.created"},{"type":"workspace.focused"},{"type":"pane.created"},{"type":"pane.focused"},{"type":"pane.agent_detected"},{"type":"pane.closed"},{"type":"workspace.closed"}]}}"#, + ); + + let ack = read_json_line(&mut reader, Duration::from_secs(2)); + assert_eq!(ack["id"], "sub_life"); + assert_eq!(ack["result"]["type"], "subscription_started"); + + let created = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l1","method":"workspace.create","params":{{"cwd":"{}","focus":true}}}}"#, + base.display() + ), + ); + let workspace_id = created["result"]["workspace"]["workspace_id"] + .as_str() + .unwrap() + .to_string(); + + let workspace_created = + wait_for_event(&mut reader, "workspace_created", Duration::from_secs(2)); + assert_eq!( + workspace_created["data"]["workspace"]["workspace_id"], + workspace_id + ); + let workspace_focused = + wait_for_event(&mut reader, "workspace_focused", Duration::from_secs(2)); + assert_eq!(workspace_focused["data"]["workspace_id"], workspace_id); + let pane_created = wait_for_event(&mut reader, "pane_created", Duration::from_secs(2)); + let pane_id = pane_created["data"]["pane"]["pane_id"] + .as_str() + .unwrap() + .to_string(); + let pane_focused = wait_for_event(&mut reader, "pane_focused", Duration::from_secs(2)); + assert_eq!(pane_focused["data"]["pane_id"], pane_id); + + let send_pi = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l2","method":"pane.send_text","params":{{"pane_id":"{}","text":"pi"}}}}"#, + pane_id + ), + ); + assert_eq!(send_pi["result"]["type"], "ok"); + let send_enter = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l3","method":"pane.send_keys","params":{{"pane_id":"{}","keys":["Enter"]}}}}"#, + pane_id + ), + ); + assert_eq!(send_enter["result"]["type"], "ok"); + + let agent_detected = wait_for_event(&mut reader, "pane_agent_detected", Duration::from_secs(3)); + assert_eq!(agent_detected["data"]["pane_id"], pane_id); + assert_eq!(agent_detected["data"]["agent"], "pi"); + + let split = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l4","method":"pane.split","params":{{"target_pane_id":"{}","direction":"right","focus":true}}}}"#, + pane_id + ), + ); + let split_pane_id = split["result"]["pane"]["pane_id"] + .as_str() + .unwrap() + .to_string(); + let split_created = wait_for_event(&mut reader, "pane_created", Duration::from_secs(2)); + assert_eq!(split_created["data"]["pane"]["pane_id"], split_pane_id); + + let closed = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l5","method":"pane.close","params":{{"pane_id":"{}"}}}}"#, + split_pane_id + ), + ); + assert_eq!(closed["result"]["type"], "ok"); + let pane_closed = wait_for_event(&mut reader, "pane_closed", Duration::from_secs(2)); + assert_eq!(pane_closed["data"]["pane_id"], split_pane_id); + + let closed_ws = send_request( + &socket_path, + &format!( + r#"{{"id":"req_l6","method":"workspace.close","params":{{"workspace_id":"{}"}}}}"#, + workspace_id + ), + ); + assert_eq!(closed_ws["result"]["type"], "ok"); + let workspace_closed = wait_for_event(&mut reader, "workspace_closed", Duration::from_secs(2)); + assert_eq!(workspace_closed["data"]["workspace_id"], workspace_id); + + let _ = child.child.kill(); + let _ = child.child.wait(); + let _ = fs::remove_dir_all(base); +} + +#[test] +fn events_subscribe_streams_output_and_agent_state_events() { + 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 bin_dir = base.join("bin"); + + fs::create_dir_all(&bin_dir).unwrap(); + let fake_pi = bin_dir.join("pi"); + fs::write( + &fake_pi, + "#!/bin/sh\nprintf 'Working...\\n'\nsleep 1\nprintf '\\033[2J\\033[Hdone\\n'\n", + ) + .unwrap(); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mut perms = fs::metadata(&fake_pi).unwrap().permissions(); + perms.set_mode(0o755); + fs::set_permissions(&fake_pi, perms).unwrap(); + } + + let inherited_path = std::env::var("PATH").unwrap_or_default(); + let path_override = format!("{}:{}", bin_dir.display(), inherited_path); + let mut child = spawn_herdr_with_path( + &config_home, + &runtime_dir, + &socket_path, + Some(Path::new(&path_override)), + ); + wait_for_socket(&socket_path, Duration::from_secs(5)); + + let created = send_request( + &socket_path, + &format!( + r#"{{"id":"req_20","method":"workspace.create","params":{{"cwd":"{}","focus":true}}}}"#, + base.display() + ), + ); + assert_eq!(created["result"]["workspace"]["workspace_id"], "w_1"); + + let panes = send_request( + &socket_path, + r#"{"id":"req_21","method":"pane.list","params":{}}"#, + ); + let pane_id = panes["result"]["panes"][0]["pane_id"] + .as_str() + .unwrap() + .to_string(); + + let (_stream, mut reader) = open_subscription( + &socket_path, + &format!( + r#"{{"id":"sub_1","method":"events.subscribe","params":{{"subscriptions":[{{"type":"pane.output_matched","pane_id":"{}","source":"recent","lines":40,"match":{{"type":"substring","value":"hello from socket"}}}},{{"type":"pane.agent_state_changed","pane_id":"{}","state":"idle"}}]}}}}"#, + pane_id, pane_id, + ), + ); + + let ack = read_json_line(&mut reader, Duration::from_secs(2)); + assert_eq!(ack["id"], "sub_1"); + assert_eq!(ack["result"]["type"], "subscription_started"); + + let send_text = send_request( + &socket_path, + &format!( + r#"{{"id":"req_22","method":"pane.send_text","params":{{"pane_id":"{}","text":"echo hello from socket"}}}}"#, + pane_id + ), + ); + assert_eq!(send_text["result"]["type"], "ok"); + let send_enter = send_request( + &socket_path, + &format!( + r#"{{"id":"req_23","method":"pane.send_keys","params":{{"pane_id":"{}","keys":["Enter"]}}}}"#, + pane_id + ), + ); + assert_eq!(send_enter["result"]["type"], "ok"); + + let output_event = read_json_line(&mut reader, Duration::from_secs(3)); + assert_eq!(output_event["event"], "pane.output_matched"); + assert_eq!(output_event["data"]["pane_id"], pane_id); + assert!(output_event["data"]["matched_line"] + .as_str() + .unwrap() + .contains("hello from socket")); + assert!(output_event["data"]["read"]["text"] + .as_str() + .unwrap() + .contains("hello from socket")); + + let send_pi = send_request( + &socket_path, + &format!( + r#"{{"id":"req_24","method":"pane.send_text","params":{{"pane_id":"{}","text":"pi"}}}}"#, + pane_id + ), + ); + assert_eq!(send_pi["result"]["type"], "ok"); + let send_enter = send_request( + &socket_path, + &format!( + r#"{{"id":"req_25","method":"pane.send_keys","params":{{"pane_id":"{}","keys":["Enter"]}}}}"#, + pane_id + ), + ); + assert_eq!(send_enter["result"]["type"], "ok"); + + let agent_idle = read_json_line(&mut reader, Duration::from_secs(8)); + assert_eq!(agent_idle["event"], "pane.agent_state_changed"); + assert_eq!(agent_idle["data"]["pane_id"], pane_id); + assert_eq!(agent_idle["data"]["state"], "idle"); + assert_eq!(agent_idle["data"]["agent"], "pi"); + + let _ = child.child.kill(); + let _ = child.child.wait(); + let _ = fs::remove_dir_all(base); +}