use parking_lot::Mutex;
use std::sync::{
mpsc::Receiver,
Arc
};
use super::{ShardManager, ShardManagerMessage};
#[derive(Debug)]
pub struct ShardManagerMonitor {
pub manager: Arc<Mutex<ShardManager>>,
pub rx: Receiver<ShardManagerMessage>,
}
impl ShardManagerMonitor {
pub fn run(&mut self) {
debug!("Starting shard manager worker");
while let Ok(value) = self.rx.recv() {
match value {
ShardManagerMessage::Restart(shard_id) => {
self.manager.lock().restart(shard_id);
},
ShardManagerMessage::ShardUpdate { id, latency, stage } => {
let manager = self.manager.lock();
let mut runners = manager.runners.lock();
if let Some(runner) = runners.get_mut(&id) {
runner.latency = latency;
runner.stage = stage;
}
}
ShardManagerMessage::Shutdown(shard_id) => {
self.manager.lock().shutdown(shard_id);
},
ShardManagerMessage::ShutdownAll => {
self.manager.lock().shutdown_all();
break;
},
ShardManagerMessage::ShutdownInitiated => break,
}
}
}
}