use std::collections::HashMap;
use std::time::Duration;
use crate::error::ServerError;
use crate::worker::heartbeat::DispatchExclusion;
use crate::worker::queue_service::PoolCensus;
use super::capacity::worker_is_at_capacity;
use super::reservation::worker_matches_node;
use super::{ActivityKey, ConnectedWorkerRegistry, PoolAddress};
impl ConnectedWorkerRegistry {
pub fn ineligible_workers_over_tiers(
&self,
namespace: &str,
task_queue: &str,
activity_type: &str,
tiers: &[Option<String>],
) -> Result<usize, ServerError> {
let state = self.state()?;
let key = ActivityKey::new(PoolAddress::new(namespace, task_queue), activity_type);
Ok(state.by_activity.get(&key).map_or(0, |workers| {
workers
.values()
.filter(|worker| state.dispatch_ineligible.contains_key(&worker.id))
.filter(|worker| {
tiers
.iter()
.any(|tier| worker_matches_node(worker, tier.as_deref()))
})
.count()
}))
}
pub fn pool_census(
&self,
namespace: &str,
task_queue: &str,
activity_type: &str,
node: Option<&str>,
) -> Result<PoolCensus, ServerError> {
let state = self.state()?;
let key = ActivityKey::new(PoolAddress::new(namespace, task_queue), activity_type);
let workers_in_pool = state
.workers
.values()
.filter(|worker| {
worker.task_queue == task_queue && worker.namespaces.contains(namespace)
})
.count();
let serving = state.by_activity.get(&key);
let workers_serving_activity = serving.map_or(0, HashMap::len);
let compatible: Vec<_> = serving.map_or_else(Vec::new, |workers| {
workers
.values()
.filter(|worker| worker_matches_node(worker, node))
.collect()
});
let compatible_workers = compatible.len();
let eligible_compatible_workers = compatible
.iter()
.filter(|worker| !state.dispatch_ineligible.contains_key(&worker.id))
.filter(|worker| !worker_is_at_capacity(&state, worker))
.count();
let compatible_workers_reachability_lost = compatible
.iter()
.filter(|worker| {
matches!(
state.dispatch_ineligible.get(&worker.id),
Some(&DispatchExclusion::ReachabilityLost)
)
})
.count();
let compatible_workers_capacity_unannounced = compatible
.iter()
.filter(|worker| worker.max_concurrency.is_none())
.count();
let compatible_workers_at_capacity = compatible
.iter()
.filter(|worker| worker.max_concurrency.is_some())
.filter(|worker| worker_is_at_capacity(&state, worker))
.count();
let last_compatible_poller_age = if compatible_workers > 0 {
Some(Duration::ZERO)
} else {
state
.last_departure
.get(&key)
.and_then(|by_node| {
by_node
.iter()
.filter(|(departed_node, _)| match node {
None => true,
Some(node) => departed_node.as_deref() == Some(node),
})
.map(|(_, departed_at)| *departed_at)
.max()
})
.map(|departed_at| departed_at.elapsed())
};
Ok(PoolCensus {
workers_in_pool,
workers_serving_activity,
compatible_workers,
eligible_compatible_workers,
compatible_workers_reachability_lost,
compatible_workers_at_capacity,
compatible_workers_capacity_unannounced,
last_compatible_poller_age,
})
}
}