use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use serde_json::{Value, json};
use crate::boundary::deposit;
use crate::boundary::reply::{self, Reply};
use crate::registry::roster::ClientRow;
#[derive(Debug, Clone, Copy)]
pub struct Budget {
pub waits: u32,
pub tick: Duration,
}
impl Budget {
pub fn span(&self) -> Duration {
self.tick * self.waits
}
}
impl Default for Budget {
fn default() -> Self {
Self {
waits: 80,
tick: Duration::from_millis(125),
}
}
}
const SEED: &str = "toolhost";
pub fn ask(
state_root: &Path,
request: &Value,
budget: Budget,
stop: &AtomicBool,
) -> Result<Value, String> {
let id = deposit::mint(state_root, SEED).map_err(|e| format!("gesture id: {e}"))?;
deposit::deposit(state_root, &id, request).map_err(|e| format!("deposit: {e}"))?;
for _ in 0..budget.waits {
if let Some(envelope) = deposit::read_reply(state_root, &id) {
return Ok(envelope);
}
if stop.load(Ordering::Relaxed) {
return Err("stopped while waiting for the engine".to_owned());
}
std::thread::sleep(budget.tick);
}
Err("no engine answered; is yog running on this world?".to_owned())
}
pub fn roster(
state_root: &Path,
workspace: &str,
budget: Budget,
stop: &AtomicBool,
) -> Result<Vec<ClientRow>, String> {
let envelope = ask(
state_root,
&json!({ "op": "clients", "workspace": workspace }),
budget,
stop,
)?;
match reply::decode(&envelope) {
Ok(Ok(Reply::Clients(rows))) => Ok(rows),
Ok(Ok(other)) => Err(format!("engine answered {other:?}, not a client roster")),
Ok(Err(refusal)) => Err(refusal),
Err(e) => Err(format!("undecodable engine reply: {e}")),
}
}
#[cfg(test)]
mod tests;