use std::sync::mpsc;
use std::thread;
use crate::client::{Envelope, RequestId};
use crate::message::{GossipEvent, RadioRequest, RadioResponse};
use evy_ecs_storage::ParticleId;
pub struct DaemonHandle {
thread: Option<thread::JoinHandle<()>>,
}
impl DaemonHandle {
pub(crate) fn spawn(
request_rx: mpsc::Receiver<Envelope>,
response_tx: mpsc::Sender<(RequestId, RadioResponse)>,
gossip_tx: mpsc::Sender<GossipEvent>,
) -> Self {
let thread = thread::Builder::new()
.name("evy-radio".into())
.spawn(move || run_daemon(request_rx, response_tx, gossip_tx))
.expect("spawn radio daemon thread");
Self {
thread: Some(thread),
}
}
}
impl Drop for DaemonHandle {
fn drop(&mut self) {
if let Some(handle) = self.thread.take() {
let _ = handle.join();
}
}
}
fn run_daemon(
request_rx: mpsc::Receiver<Envelope>,
response_tx: mpsc::Sender<(RequestId, RadioResponse)>,
_gossip_tx: mpsc::Sender<GossipEvent>,
) {
while let Ok(env) = request_rx.recv() {
let Envelope { id, request } = env;
let response = handle_request(request);
if let Some(resp) = response {
if response_tx.send((id, resp)).is_err() {
return;
}
} else {
return;
}
}
}
fn handle_request(request: RadioRequest) -> Option<RadioResponse> {
match request {
RadioRequest::FetchParticle { particle } => Some(RadioResponse::Fetched {
particle,
bytes: Vec::new(),
}),
RadioRequest::PublishParticle { bytes } => {
let synthetic = ParticleId::from_entity(bytes.len() as u32, 0);
Some(RadioResponse::Published {
particle: synthetic,
})
}
RadioRequest::Subscribe { topic } => Some(RadioResponse::SubscribeAck { topic }),
RadioRequest::Unsubscribe { topic } => Some(RadioResponse::UnsubscribeAck { topic }),
RadioRequest::GossipPublish { .. } => Some(RadioResponse::GossipPublishAck),
RadioRequest::Shutdown => None,
}
}