use std::time::{Duration, Instant};
use aion_core::{ActivityId, WorkflowId};
use super::census::{PoolCensus, classify};
use super::declarations::QueueDeclarationSource;
use super::policy::{QueueServiceConfig, QueueServicePolicy};
use super::state::{Parked, QueueServiceState};
use super::taxonomy::{
ExpiredClock, QueueServiceReason, ServiceAddress, WorkerUnavailable, millis,
};
use crate::worker::registry::{ConnectedWorkerRegistry, WorkerHandle};
const NEVER: &str = "never";
#[derive(Clone, Debug)]
pub enum SelectionRefusal {
Unavailable(Box<WorkerUnavailable>),
NotAccepting {
reason: String,
},
Registry {
reason: String,
},
}
impl SelectionRefusal {
#[must_use]
pub fn reason_string(&self) -> String {
match self {
Self::Unavailable(unavailable) => unavailable.reason_string(),
Self::NotAccepting { reason } | Self::Registry { reason } => reason.clone(),
}
}
}
pub struct ServiceWait<'a> {
pub registry: &'a ConnectedWorkerRegistry,
pub declarations: &'a QueueDeclarationSource,
pub config: &'a QueueServiceConfig,
pub state: &'a QueueServiceState,
pub address: &'a ServiceAddress,
pub workflow_id: &'a WorkflowId,
pub activity_id: &'a ActivityId,
}
pub fn select_worker_or_refuse(
wait: &ServiceWait<'_>,
accepting: &mut dyn FnMut() -> Result<(), String>,
park: &mut dyn FnMut(Option<Duration>),
) -> Result<WorkerHandle, SelectionRefusal> {
let started_at = Instant::now();
let address = wait.address;
let policy = wait
.config
.policy_for(&address.namespace, &address.task_queue);
let deadline = wait.config.availability_deadline_for(policy);
let mut reported: Option<QueueServiceReason> = None;
let outcome = loop {
match wait.registry.select_worker(
&address.namespace,
&address.task_queue,
&address.activity_type,
address.node.as_deref(),
) {
Ok(Some(worker)) => {
if let Some(reason) = reported {
tracing::info!(
namespace = %address.namespace,
task_queue = %address.task_queue,
activity_type = %address.activity_type,
workflow_id = %wait.workflow_id,
activity_id = %wait.activity_id,
queue_service_reason = reason.as_str(),
waited_ms = millis(started_at.elapsed()),
"queue service restored; the parked dispatch has a worker"
);
}
break Ok(worker);
}
Ok(None) => {
if let Err(reason) = accepting() {
break Err(SelectionRefusal::NotAccepting { reason });
}
let waited = started_at.elapsed();
let observed =
match observe_selection_miss(wait, policy, deadline, waited, reported) {
Ok(Some(observed)) => observed,
Ok(None) => continue,
Err(refusal) => break Err(refusal),
};
let (reason, census) = (observed.reason, observed.census);
reported = Some(reason);
if reason.is_structural() {
break Err(unavailable(address, reason, None, waited, census));
}
match deadline {
Some(deadline) => {
let Some(remaining) = deadline
.checked_sub(waited)
.filter(|remaining| !remaining.is_zero())
else {
break Err(unavailable(
address,
reason,
Some(ExpiredClock::ServiceAvailability),
waited,
census,
));
};
park(Some(remaining));
}
None => park(None),
}
}
Err(error) => {
break Err(SelectionRefusal::Registry {
reason: format!("registry error: {error}"),
});
}
}
};
clear_unserved(wait);
outcome
}
pub struct ObservedMiss {
pub reason: QueueServiceReason,
pub census: PoolCensus,
}
pub fn observe_selection_miss(
wait: &ServiceWait<'_>,
policy: QueueServicePolicy,
deadline: Option<Duration>,
waited: Duration,
reported: Option<QueueServiceReason>,
) -> Result<Option<ObservedMiss>, SelectionRefusal> {
let address = wait.address;
let census = wait
.registry
.pool_census(
&address.namespace,
&address.task_queue,
&address.activity_type,
address.node.as_deref(),
)
.map_err(|error| SelectionRefusal::Registry {
reason: format!("registry error: {error}"),
})?;
let declaration = wait.declarations.declaration_for(&address.task_queue);
let Some(reason) = classify(declaration, &census) else {
return Ok(None);
};
mark_unserved(wait, reason, policy, census);
report(&Report {
wait,
reason,
policy,
census,
waited,
deadline,
repeat: reported == Some(reason),
});
Ok(Some(ObservedMiss { reason, census }))
}
pub fn clear_selection_miss(wait: &ServiceWait<'_>) {
clear_unserved(wait);
}
fn unavailable(
address: &ServiceAddress,
reason: QueueServiceReason,
clock: Option<ExpiredClock>,
waited: Duration,
census: PoolCensus,
) -> SelectionRefusal {
SelectionRefusal::Unavailable(Box::new(WorkerUnavailable {
reason,
clock,
waited,
address: address.clone(),
census,
}))
}
fn mark_unserved(
wait: &ServiceWait<'_>,
reason: QueueServiceReason,
policy: QueueServicePolicy,
census: PoolCensus,
) {
if let Err(error) = wait.state.mark(Parked {
address: wait.address,
reason,
policy,
census,
workflow_id: wait.workflow_id,
activity_id: wait.activity_id,
}) {
tracing::error!(%error, "failed to publish the unserved queue state");
}
}
fn clear_unserved(wait: &ServiceWait<'_>) {
if let Err(error) = wait
.state
.clear(wait.address, wait.workflow_id, wait.activity_id)
{
tracing::error!(%error, "failed to clear the unserved queue state");
}
}
struct Report<'a> {
wait: &'a ServiceWait<'a>,
reason: QueueServiceReason,
policy: QueueServicePolicy,
census: PoolCensus,
waited: Duration,
deadline: Option<Duration>,
repeat: bool,
}
fn report(report: &Report<'_>) {
let Report {
wait,
reason,
policy,
census,
waited,
deadline,
repeat,
} = report;
let address = wait.address;
let age = census
.last_compatible_poller_age
.map_or_else(|| NEVER.to_owned(), |age| millis(age).to_string());
let deadline_ms = deadline.map_or_else(|| NEVER.to_owned(), |value| millis(value).to_string());
if *repeat {
tracing::debug!(
namespace = %address.namespace,
task_queue = %address.task_queue,
activity_type = %address.activity_type,
node = address.node.as_deref(),
workflow_id = %wait.workflow_id,
activity_id = %wait.activity_id,
queue_service_reason = reason.as_str(),
queue_service_policy = policy.as_str(),
last_compatible_poller_age_ms = %age,
waited_ms = millis(*waited),
"queue still unserved"
);
return;
}
tracing::warn!(
namespace = %address.namespace,
task_queue = %address.task_queue,
activity_type = %address.activity_type,
node = address.node.as_deref(),
workflow_id = %wait.workflow_id,
activity_id = %wait.activity_id,
queue_service_reason = reason.as_str(),
queue_service_policy = policy.as_str(),
last_compatible_poller_age_ms = %age,
service_availability_deadline_ms = %deadline_ms,
waited_ms = millis(*waited),
workers_in_pool = census.workers_in_pool,
workers_serving_activity = census.workers_serving_activity,
compatible_workers = census.compatible_workers,
"dispatch is parked on a queue that is not being served"
);
}