use crate::connection::WorkerConnectionKey;
use crate::tenant::{WorkerIndexEntry, WorkerTarget};
use appcore_types::CapabilityName;
use std::collections::{HashMap, HashSet};
use std::sync::atomic::AtomicU64;
use std::sync::Arc;
use std::time::Instant;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PendingRequestKind {
Peer,
Mesh,
}
#[derive(Debug, Clone)]
pub(crate) struct PendingRequest {
pub(crate) kind: PendingRequestKind,
pub(crate) worker_generation: u64,
pub(crate) deadline: Instant,
pub(crate) response_limit: Option<usize>,
}
#[derive(Debug, Default, Clone)]
pub struct CapabilityRegistry {
capability_to_workers: HashMap<CapabilityName, HashSet<WorkerConnectionKey>>,
worker_to_capabilities: HashMap<WorkerConnectionKey, HashSet<CapabilityName>>,
pub(crate) worker_by_core: HashMap<appcore_types::CoreId, WorkerIndexEntry>,
pub(crate) worker_by_target: HashMap<WorkerTarget, WorkerIndexEntry>,
pub(crate) worker_index_rebuilds: u64,
pub(crate) worker_index_inconsistencies: Arc<AtomicU64>,
pub(crate) pending_requests: HashMap<String, PendingRequest>,
}
impl CapabilityRegistry {
pub fn new() -> Self {
Self::default()
}
pub fn register(&mut self, worker: WorkerConnectionKey, capabilities: Vec<CapabilityName>) {
self.deregister(&worker);
let mut caps_set = HashSet::new();
for cap in capabilities {
self.capability_to_workers
.entry(cap.clone())
.or_default()
.insert(worker.clone());
caps_set.insert(cap);
}
self.worker_to_capabilities.insert(worker, caps_set);
}
pub fn deregister(&mut self, worker: &WorkerConnectionKey) {
if let Some(caps) = self.worker_to_capabilities.remove(worker) {
for cap in caps {
if let Some(workers) = self.capability_to_workers.get_mut(&cap) {
workers.remove(worker);
if workers.is_empty() {
self.capability_to_workers.remove(&cap);
}
}
}
}
}
pub fn resolve(&self, capability: &CapabilityName) -> Option<&HashSet<WorkerConnectionKey>> {
self.capability_to_workers.get(capability)
}
pub fn all_capabilities(&self) -> Vec<CapabilityName> {
self.capability_to_workers.keys().cloned().collect()
}
pub fn capabilities_for(&self, worker: &WorkerConnectionKey) -> Vec<CapabilityName> {
let mut capabilities = self
.worker_to_capabilities
.get(worker)
.into_iter()
.flatten()
.cloned()
.collect::<Vec<_>>();
capabilities.sort_by(|left, right| left.as_str().cmp(right.as_str()));
capabilities
}
}