mirror of
https://github.com/herdrdev/herdr.git
synced 2026-09-22 16:01:07 +00:00
@@ -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
|
||||
|
||||
|
||||
@@ -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::<ApiRequestMessage>();
|
||||
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::<ApiRequestMessage>();
|
||||
|
||||
+39
-16
@@ -310,6 +310,19 @@ impl ActiveSubscription {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn poll_for_wait(
|
||||
&mut self,
|
||||
api_tx: &ApiRequestSender,
|
||||
event_hub: &EventHub,
|
||||
) -> Result<Option<serde_json::Value>, 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<SubscriptionEventEnvelope> {
|
||||
self.poll_result(api_tx, event_hub).ok().flatten()
|
||||
}
|
||||
|
||||
fn poll_result(
|
||||
&mut self,
|
||||
api_tx: &ApiRequestSender,
|
||||
event_hub: &EventHub,
|
||||
) -> Result<Option<SubscriptionEventEnvelope>, 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::<ErrorResponse>(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,
|
||||
|
||||
+10
-2
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user