use aion_core::{ActivityId, Event, StepState, UnservedActivity, WorkflowId};
use chrono::{DateTime, Utc};
use super::census::classify;
use super::declarations::QueueDeclarationSource;
use super::state::QueueServiceState;
use super::taxonomy::ServiceAddress;
use crate::error::ServerError;
use crate::worker::registry::ConnectedWorkerRegistry;
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OpenActivity {
pub activity_id: ActivityId,
pub activity_type: String,
pub task_queue: String,
pub node: Option<String>,
pub attempt: u32,
pub dispatched_at: DateTime<Utc>,
}
#[must_use]
pub fn open_activities_in_active_segment(history: &[Event]) -> Vec<OpenActivity> {
aion_core::open_steps(history)
.into_iter()
.filter_map(|step| match step.state {
StepState::Dispatched {
attempt,
dispatched_at,
} => Some(OpenActivity {
activity_id: step.activity_id,
activity_type: step.activity_type,
task_queue: step.task_queue,
node: step.node,
attempt,
dispatched_at,
}),
StepState::Scheduled | StepState::Reopened { .. } => None,
})
.collect()
}
pub struct ActivityReachability<'a> {
pub registry: &'a ConnectedWorkerRegistry,
pub declarations: &'a QueueDeclarationSource,
pub state: &'a QueueServiceState,
}
impl ActivityReachability<'_> {
pub fn unserved(
&self,
namespace: &str,
workflow_id: &WorkflowId,
history: &[Event],
) -> Result<Vec<UnservedActivity>, ServerError> {
let mut unserved = Vec::new();
for open in open_activities_in_active_segment(history) {
let address = ServiceAddress {
namespace: namespace.to_owned(),
task_queue: open.task_queue,
activity_type: open.activity_type,
node: open.node,
};
let census = self.registry.pool_census(
&address.namespace,
&address.task_queue,
&address.activity_type,
address.node.as_deref(),
)?;
let declaration = self.declarations.declaration_for(&address.task_queue);
let Some(reason) = classify(declaration, &census) else {
continue;
};
let dispatch_parked = self.state.is_unserved(workflow_id, &open.activity_id)?;
unserved.push(UnservedActivity {
activity_id: open.activity_id,
reason: reason.as_str().to_owned(),
detail: reason.explain(&address),
activity_type: address.activity_type,
task_queue: address.task_queue,
node: address.node,
attempt: open.attempt,
dispatched_at: open.dispatched_at,
workers_in_pool: count(census.workers_in_pool),
workers_serving_activity: count(census.workers_serving_activity),
compatible_workers: count(census.compatible_workers),
dispatch_parked,
});
}
Ok(unserved)
}
}
fn count(value: usize) -> u64 {
u64::try_from(value).unwrap_or(u64::MAX)
}
#[cfg(test)]
#[path = "reachability_tests.rs"]
mod tests;