clusterflux-public/crates/clusterflux-node/src/coordinator_session.rs
Clusterflux Release 26fdcb9d84 Update public backend API surface
Private source commit: ba3f7ce2b6d9
2026-07-26 17:22:33 +02:00

49 lines
1.5 KiB
Rust

#[cfg(test)]
use clusterflux_control::endpoint_identity;
use clusterflux_control::ControlSession;
use clusterflux_core::coordinator_wire_request;
use serde_json::Value;
use std::time::Duration;
pub(crate) struct CoordinatorSession {
inner: ControlSession,
}
impl CoordinatorSession {
pub(crate) fn connect(addr: &str) -> Result<Self, Box<dyn std::error::Error>> {
Ok(Self {
inner: ControlSession::connect(addr)?,
})
}
pub(crate) fn connect_with_timeouts(
addr: &str,
connect_timeout: Duration,
io_timeout: Duration,
) -> Result<Self, Box<dyn std::error::Error>> {
Ok(Self {
inner: ControlSession::connect_with_timeouts(addr, connect_timeout, io_timeout)?,
})
}
pub(crate) fn request(&mut self, value: Value) -> Result<Value, Box<dyn std::error::Error>> {
let request_id = format!("node-{}", self.inner.requests() + 1);
let wire_request = coordinator_wire_request(request_id, value);
let response = self.inner.request(&wire_request)?;
if response.get("type").and_then(Value::as_str) == Some("error") {
return Err(format!("coordinator error: {response}").into());
}
Ok(response)
}
pub(crate) fn requests(&self) -> usize {
self.inner.requests() as usize
}
}
#[cfg(test)]
pub(crate) fn control_endpoint_identity(
endpoint: &str,
) -> Result<String, Box<dyn std::error::Error>> {
Ok(endpoint_identity(endpoint)?)
}