use super::Orchestrator;
use crate::app::events::AppEvent;
#[derive(Default)]
pub(super) struct SlotsDiscovery {
epoch: u64,
pending: bool,
known: Option<u32>,
}
impl SlotsDiscovery {
fn invalidate(&mut self) {
self.epoch += 1;
self.pending = false;
self.known = None;
}
pub(super) fn known(&self) -> Option<u32> {
self.known
}
fn apply(&mut self, epoch: u64, slots: Option<u32>) -> bool {
if epoch != self.epoch {
return false;
}
self.pending = false;
self.known = slots;
true
}
}
impl Orchestrator {
pub(super) fn refresh_engine_slots(&mut self) {
self.slots.invalidate();
self.emit_engine_slots();
let Some(backend) = self.engines.backend.clone() else {
return;
};
self.slots.pending = true;
let epoch = self.slots.epoch;
let tx = self.slots_tx.clone();
tokio::spawn(async move {
let _ = tx.send((epoch, backend.parallel_slots().await));
});
}
pub(super) fn handle_slots_result(&mut self, epoch: u64, slots: Option<u32>) {
if !self.slots.apply(epoch, slots) {
return;
}
if let Some(n) = self.slots.known {
tracing::info!(slots = n, "engine reported its slot count");
}
self.emit_engine_slots();
}
fn emit_engine_slots(&self) {
let _ = self.evt_tx.send(AppEvent::EngineSlots(self.slots.known));
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_stale_answer_is_dropped_and_a_current_one_lands() {
let mut d = SlotsDiscovery::default();
let old = d.epoch;
d.pending = true;
d.invalidate();
assert!(!d.pending);
assert!(!d.apply(old, Some(4)), "the previous engine's answer");
assert_eq!(d.known, None);
assert!(d.apply(d.epoch, Some(4)));
assert_eq!(d.known, Some(4));
assert!(d.apply(d.epoch, None));
assert_eq!(d.known, None);
}
}