use crate::error::RpcResult;
use std::sync::{Arc, Mutex};
use workflow_core::channel::Multiplexer;
#[derive(Default, Debug, Clone, Copy, PartialEq, Eq)]
pub enum RpcState {
Connected,
#[default]
Disconnected,
}
#[derive(Default)]
struct Inner {
state: Mutex<RpcState>,
multiplexer: Multiplexer<RpcState>,
descriptor: Mutex<Option<String>>,
}
#[derive(Default, Clone)]
pub struct RpcCtl {
inner: Arc<Inner>,
}
impl RpcCtl {
pub fn new() -> Self {
Self { inner: Arc::new(Inner::default()) }
}
pub fn with_descriptor<Str: ToString>(descriptor: Option<Str>) -> Self {
if let Some(descriptor) = descriptor {
Self { inner: Arc::new(Inner { descriptor: Mutex::new(Some(descriptor.to_string())), ..Inner::default() }) }
} else {
Self::default()
}
}
pub fn multiplexer(&self) -> &Multiplexer<RpcState> {
&self.inner.multiplexer
}
pub fn is_connected(&self) -> bool {
*self.inner.state.lock().unwrap() == RpcState::Connected
}
pub fn state(&self) -> RpcState {
*self.inner.state.lock().unwrap()
}
pub async fn signal_open(&self) -> RpcResult<()> {
*self.inner.state.lock().unwrap() = RpcState::Connected;
Ok(self.inner.multiplexer.broadcast(RpcState::Connected).await?)
}
pub async fn signal_close(&self) -> RpcResult<()> {
*self.inner.state.lock().unwrap() = RpcState::Disconnected;
Ok(self.inner.multiplexer.broadcast(RpcState::Disconnected).await?)
}
pub fn try_signal_open(&self) -> RpcResult<()> {
*self.inner.state.lock().unwrap() = RpcState::Connected;
Ok(self.inner.multiplexer.try_broadcast(RpcState::Connected)?)
}
pub fn try_signal_close(&self) -> RpcResult<()> {
*self.inner.state.lock().unwrap() = RpcState::Disconnected;
Ok(self.inner.multiplexer.try_broadcast(RpcState::Disconnected)?)
}
pub fn set_descriptor(&self, descriptor: Option<String>) {
*self.inner.descriptor.lock().unwrap() = descriptor;
}
pub fn descriptor(&self) -> Option<String> {
self.inner.descriptor.lock().unwrap().clone()
}
}