From d0beea46d4153896c9de668ee8204375cac6bfa8 Mon Sep 17 00:00:00 2001 From: akbash-bot <300245827+akbash-bot@users.noreply.github.com> Date: Thu, 16 Jul 2026 15:27:44 +0000 Subject: [PATCH] fix: stop agent wait when pane closes refs #1439 --- docs/next/CHANGELOG.md | 1 + src/api/server.rs | 58 ++++++++++++++++++++++++++++++++++++++++ src/api/subscriptions.rs | 55 ++++++++++++++++++++++++++----------- src/api/wait.rs | 12 +++++++-- 4 files changed, 108 insertions(+), 18 deletions(-) diff --git a/docs/next/CHANGELOG.md b/docs/next/CHANGELOG.md index 9b0c452f..3f75e127 100644 --- a/docs/next/CHANGELOG.md +++ b/docs/next/CHANGELOG.md @@ -8,6 +8,7 @@ ### Fixed - Live handoff now preserves installed plugins and no longer lets the next plugin installation overwrite the existing registry. (#893) +- `herdr wait agent-status` now returns `pane_not_found` promptly when its target pane closes instead of waiting for the full timeout. (#1439) ## [0.7.4] - 2026-07-15 diff --git a/src/api/server.rs b/src/api/server.rs index c4c76703..2378e1bd 100644 --- a/src/api/server.rs +++ b/src/api/server.rs @@ -1124,6 +1124,64 @@ mod tests { responder.join().unwrap(); } + #[test] + fn events_wait_agent_status_returns_not_found_when_pane_closes() { + let event_hub = EventHub::default(); + let responder_event_hub = event_hub.clone(); + let (api_tx, mut api_rx) = mpsc::unbounded_channel::(); + let responder = std::thread::spawn(move || { + let mut pane_get_count = 0; + while let Some(msg) = api_rx.blocking_recv() { + let Method::PaneGet(_) = msg.request.method else { + panic!("unexpected request: {:?}", msg.request.method); + }; + pane_get_count += 1; + let response = if pane_get_count == 1 { + serde_json::to_string(&SuccessResponse { + id: msg.request.id, + result: ResponseResult::PaneInfo { + pane: pane_info("pane_1", crate::api::schema::AgentStatus::Unknown), + }, + }) + .unwrap() + } else { + if pane_get_count == 2 { + responder_event_hub.push(crate::api::schema::EventEnvelope { + event: crate::api::schema::EventKind::PaneClosed, + data: crate::api::schema::EventData::PaneClosed { + pane_id: "pane_1".into(), + workspace_id: "ws_1".into(), + }, + }); + } + error_response_json( + msg.request.id, + "pane_not_found", + "pane pane_1 not found".into(), + ) + }; + msg.respond_to.send(response).unwrap(); + } + }); + + let (mut client, server, _path) = local_stream_pair("api-events-wait-pane-close"); + client + .write_all(br#"{"id":"wait_close","method":"events.wait","params":{"match_event":{"event":"pane_agent_status_changed","pane_id":"pane_1","agent_status":"done"},"timeout_ms":500}}"#) + .unwrap(); + client.write_all(b"\n").unwrap(); + client.flush().unwrap(); + + let running = Arc::new(AtomicBool::new(true)); + handle_connection(server, &api_tx, &event_hub, &running, None).unwrap(); + + let response: serde_json::Value = serde_json::from_str(&read_line(&mut client)).unwrap(); + assert_eq!(response["id"], "wait_close"); + assert_eq!(response["error"]["code"], "pane_not_found"); + assert_eq!(response["error"]["message"], "pane pane_1 not found"); + drop(api_tx); + responder.join().unwrap(); + } + #[test] fn wait_for_output_stops_when_client_disconnects() { let (api_tx, mut api_rx) = mpsc::unbounded_channel::(); diff --git a/src/api/subscriptions.rs b/src/api/subscriptions.rs index 531fb7bb..13a340b6 100644 --- a/src/api/subscriptions.rs +++ b/src/api/subscriptions.rs @@ -310,6 +310,19 @@ impl ActiveSubscription { } } } + + pub(super) fn poll_for_wait( + &mut self, + api_tx: &ApiRequestSender, + event_hub: &EventHub, + ) -> Result, ErrorResponse> { + match self { + Self::AgentStatusChanged(subscription) => Ok(subscription + .poll_result(api_tx, event_hub)? + .and_then(|event| serde_json::to_value(event).ok())), + _ => Ok(self.poll(api_tx, event_hub)), + } + } } impl ActiveEventSubscription { @@ -366,6 +379,14 @@ impl ActiveAgentStatusChangedSubscription { api_tx: &ApiRequestSender, event_hub: &EventHub, ) -> Option { + self.poll_result(api_tx, event_hub).ok().flatten() + } + + fn poll_result( + &mut self, + api_tx: &ApiRequestSender, + event_hub: &EventHub, + ) -> Result, ErrorResponse> { let mut saw_status_event = false; for (sequence, event) in event_hub.events_after(self.last_sequence) { self.last_sequence = sequence; @@ -401,7 +422,7 @@ impl ActiveAgentStatusChangedSubscription { } self.initial_event = None; - return Some(SubscriptionEventEnvelope { + return Ok(Some(SubscriptionEventEnvelope { event: SubscriptionEventKind::PaneAgentStatusChanged, data: SubscriptionEventData::PaneAgentStatusChanged(PaneAgentStatusChangedEvent { pane_id, @@ -412,18 +433,18 @@ impl ActiveAgentStatusChangedSubscription { display_agent, state_labels, }), - }); + })); } if saw_status_event { self.initial_event = None; } else if event_hub.current_sequence() != self.last_sequence { - return None; + return Ok(None); } else if let Some(event) = self.initial_event.take() { - return Some(SubscriptionEventEnvelope { + return Ok(Some(SubscriptionEventEnvelope { event: SubscriptionEventKind::PaneAgentStatusChanged, data: SubscriptionEventData::PaneAgentStatusChanged(event), - }); + })); } let before_snapshot_sequence = self.last_sequence; @@ -431,18 +452,18 @@ impl ActiveAgentStatusChangedSubscription { format!("{}:pane", self.request_prefix), &self.pane_id, api_tx, - ) - .ok()?; + ); let after_snapshot_sequence = event_hub.current_sequence(); if after_snapshot_sequence != before_snapshot_sequence { - return None; + return Ok(None); } + let pane = pane?; let event = self.event_from_snapshot(pane); if event.is_some() { self.last_sequence = after_snapshot_sequence; } - event + Ok(event) } fn event_from_snapshot( @@ -585,13 +606,15 @@ fn pane_get( }, })?; 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(), - }, - }); + let response = + 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(), + }, + })?; + return Err(response); } serde_json::from_value(value["result"]["pane"].clone()).map_err(|_| ErrorResponse { id: request_id, diff --git a/src/api/wait.rs b/src/api/wait.rs index 17ffb992..d556fa96 100644 --- a/src/api/wait.rs +++ b/src/api/wait.rs @@ -153,8 +153,16 @@ pub(super) fn wait_for_event( return Ok(None); } - if let Some(event) = active.poll(api_tx, event_hub) { - return Ok(Some(wait_matched_response(&request_id, event))); + match active.poll_for_wait(api_tx, event_hub) { + Ok(Some(event)) => return Ok(Some(wait_matched_response(&request_id, event))), + Ok(None) => {} + Err(mut response) if response.error.code == "pane_not_found" => { + response.id = request_id; + return serde_json::to_string(&response) + .map(Some) + .map_err(std::io::Error::other); + } + Err(_) => {} } if deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) {