clusterflux-public/crates/clusterflux-node/src/coordinator_session.rs
Clusterflux release 3996d91ee0 Public release release-ffc6f14a3723
Source commit: ffc6f14a3723ea6aa513613a5406d737f21b201d

Public tree identity: sha256:a77575c31dda8de868bc8de91d94df55c46f0b028e755e56a8cce8e27f72a42b
2026-07-19 16:41:50 +02:00

38 lines
1.2 KiB
Rust

#[cfg(test)]
use clusterflux_control::endpoint_identity;
use clusterflux_control::ControlSession;
use clusterflux_core::coordinator_wire_request;
use serde_json::Value;
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 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)?)
}