use flume::Receiver;
use std::time::Duration;
use crate::{HookID, MemberState, TickCommand, TickManagerHandle, TickStateReply};
#[derive(Debug, Clone)]
pub struct TickMember {
pub id: usize,
manager_handle: TickManagerHandle,
receiver: Receiver<TickStateReply>,
}
impl TickMember {
pub fn new(manager_handle: TickManagerHandle, speed_factor: usize) -> Self {
let (sender, receiver) = flume::bounded(10);
manager_handle
.send(TickCommand::Register(sender, speed_factor))
.unwrap();
let id = expect_id(&receiver);
Self {
id,
manager_handle,
receiver,
}
}
pub fn set_state(&self, state: MemberState) {
self.manager_handle
.send(TickCommand::ChangeMemberState(self.id, state))
.unwrap();
}
pub fn wait_for_tick(&self) {
self.set_state(MemberState::Finished);
loop {
match expect_reply(&self.receiver) {
Ok(TickStateReply::Tick) => break,
_ => continue,
}
}
}
}
fn expect_reply(
receiver: &Receiver<TickStateReply>,
) -> Result<TickStateReply, flume::RecvTimeoutError> {
receiver.recv_timeout(Duration::from_secs(1))
}
impl Drop for TickMember {
fn drop(&mut self) {
let _ = self.manager_handle.send(TickCommand::Unregister(self.id));
}
}
fn expect_id(receiver: &Receiver<TickStateReply>) -> HookID {
let reply = match expect_reply(receiver) {
Ok(reply) => reply,
Err(e) => panic!(
"Did not receive TickStateReply in time while waiting for HookID: {}",
e
),
};
match reply {
TickStateReply::SelfID(id) => id,
unexpected => panic!("Expected SelfID, got {:?}", unexpected),
}
}