mirror of
https://github.com/herdrdev/herdr.git
synced 2026-10-02 08:00:23 +00:00
125 lines
4.5 KiB
Rust
125 lines
4.5 KiB
Rust
use std::sync::atomic::AtomicBool;
|
|
use std::sync::Arc;
|
|
|
|
use regex::Regex;
|
|
|
|
use crate::api::schema::{
|
|
ErrorBody, ErrorResponse, Method, Request, ResponseResult, SuccessResponse,
|
|
};
|
|
use crate::api::server::{
|
|
dispatch_to_app_with_timeout, should_stop_connection, APP_RESPONSE_TIMEOUT,
|
|
CONNECTION_POLL_INTERVAL,
|
|
};
|
|
use crate::api::subscriptions::{match_output, output_match_read_source};
|
|
use crate::api::ApiRequestSender;
|
|
use crate::ipc::LocalStream;
|
|
|
|
pub(super) fn wait_for_output(
|
|
request_id: String,
|
|
params: crate::api::schema::PaneWaitForOutputParams,
|
|
stream: &mut LocalStream,
|
|
api_tx: &ApiRequestSender,
|
|
running: &Arc<AtomicBool>,
|
|
) -> std::io::Result<Option<String>> {
|
|
crate::logging::api_wait_started(&request_id, ¶ms.pane_id, params.timeout_ms);
|
|
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 Ok(Some(
|
|
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 {
|
|
if should_stop_connection(stream, running)? {
|
|
crate::logging::api_wait_completed(&request_id, ¶ms.pane_id, "client_disconnected");
|
|
return Ok(None);
|
|
}
|
|
|
|
let read_request = Request {
|
|
id: format!("{request_id}:read"),
|
|
method: Method::PaneRead(crate::api::schema::PaneReadParams {
|
|
pane_id: params.pane_id.clone(),
|
|
source: output_match_read_source(¶ms.source),
|
|
lines: params.lines,
|
|
format: crate::api::schema::ReadFormat::Text,
|
|
strip_ansi: params.strip_ansi,
|
|
}),
|
|
};
|
|
let response =
|
|
dispatch_to_app_with_timeout(read_request, api_tx, Some(APP_RESPONSE_TIMEOUT));
|
|
let Ok(value) = serde_json::from_str::<serde_json::Value>(&response) else {
|
|
return Ok(Some(response));
|
|
};
|
|
if value.get("error").is_some() {
|
|
let mut value = value;
|
|
value["id"] = serde_json::Value::String(request_id.clone());
|
|
return Ok(Some(serde_json::to_string(&value).unwrap()));
|
|
}
|
|
|
|
let read_value = value["result"]["read"].clone();
|
|
let Ok(read) = serde_json::from_value::<crate::api::schema::PaneReadResult>(read_value)
|
|
else {
|
|
return Ok(Some(
|
|
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_output(&read.text, ¶ms.r#match, regex.as_ref());
|
|
if matched_line.is_some() {
|
|
let revision = read.revision;
|
|
crate::logging::api_wait_completed(&request_id, ¶ms.pane_id, "matched");
|
|
return Ok(Some(
|
|
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) {
|
|
crate::logging::api_wait_timed_out(&request_id, ¶ms.pane_id);
|
|
return Ok(Some(
|
|
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(CONNECTION_POLL_INTERVAL);
|
|
}
|
|
}
|