use std::sync::atomic::{AtomicBool, Ordering};
use serde_json::{Value, json};
use super::ask::{self, Budget};
use super::{Site, loaded};
use crate::boundary::reply::{self, Reply};
use crate::registry::mailbox::Capture;
pub fn invoke(
site: &Site,
entry: &loaded::Entry,
input: &Value,
stop: &AtomicBool,
) -> Result<Capture, String> {
let invocation = routed(
site,
&json!({ "op": "invoke", "client": entry.client,
"tool": entry.tool.name, "input": input }),
stop,
)?
.0;
let poll = json!({ "op": "capture", "invocation": invocation });
for _ in 0..site.patience.waits {
if let (_, Some(capture)) = routed(site, &poll, stop)? {
return Ok(capture);
}
if stop.load(Ordering::Relaxed) {
return Err(format!("stopped while {} was running it", entry.client));
}
std::thread::sleep(site.patience.tick);
}
Err(format!(
"client {:?} did not answer invocation {invocation} in time; \
it may be offline, or the tool is still running there",
entry.client
))
}
fn routed(
site: &Site,
request: &Value,
stop: &AtomicBool,
) -> Result<(String, Option<Capture>), String> {
let envelope = ask::ask(&site.state_root, request, site.budget, stop)?;
match reply::decode(&envelope) {
Ok(Ok(Reply::Routed {
invocation,
capture,
})) => Ok((invocation, capture)),
Ok(Ok(other)) => Err(format!(
"engine answered {other:?}, not a routed invocation"
)),
Ok(Err(refusal)) => Err(refusal),
Err(e) => Err(format!("undecodable engine reply: {e}")),
}
}
pub fn patience() -> Budget {
Budget {
waits: 240,
tick: std::time::Duration::from_millis(500),
}
}
#[cfg(test)]
mod tests;