use std::sync::{Arc, OnceLock};
use aion::loader::DeclaredQueues;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum QueueDeclaration {
Declared,
NotDeclared,
Unknown,
}
pub trait QueueDeclarations: Send + Sync {
fn declaration_for(&self, task_queue: &str) -> QueueDeclaration;
}
#[derive(Clone, Default)]
pub struct QueueDeclarationSource {
inner: Arc<OnceLock<Arc<dyn QueueDeclarations>>>,
}
impl std::fmt::Debug for QueueDeclarationSource {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("QueueDeclarationSource")
.field("installed", &self.inner.get().is_some())
.finish()
}
}
impl QueueDeclarationSource {
pub fn install(&self, source: Arc<dyn QueueDeclarations>) {
if self.inner.set(source).is_err() {
tracing::warn!("queue declaration source already installed; ignoring duplicate set");
}
}
#[must_use]
pub fn is_installed(&self) -> bool {
self.inner.get().is_some()
}
#[must_use]
pub fn declaration_for(&self, task_queue: &str) -> QueueDeclaration {
self.inner
.get()
.map_or(QueueDeclaration::Unknown, |source| {
source.declaration_for(task_queue)
})
}
}
pub struct EngineQueueDeclarations {
engine: Arc<aion::Engine>,
}
impl EngineQueueDeclarations {
#[must_use]
pub const fn new(engine: Arc<aion::Engine>) -> Self {
Self { engine }
}
}
impl std::fmt::Debug for EngineQueueDeclarations {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("EngineQueueDeclarations")
}
}
fn declaration_from(read: &DeclaredQueues, task_queue: &str) -> QueueDeclaration {
if read.declares(task_queue) {
QueueDeclaration::Declared
} else if read.found_no_declaration() || !read.covers_every_entry() {
QueueDeclaration::Unknown
} else {
QueueDeclaration::NotDeclared
}
}
impl QueueDeclarations for EngineQueueDeclarations {
fn declaration_for(&self, task_queue: &str) -> QueueDeclaration {
match self.engine.declared_task_queues() {
Ok(read) => declaration_from(&read, task_queue),
Err(error) => {
tracing::warn!(
task_queue,
%error,
"queue declaration lookup failed; treating the declaration as unknown"
);
QueueDeclaration::Unknown
}
}
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use super::super::census::{PoolCensus, classify};
use super::super::taxonomy::QueueServiceReason;
use super::*;
fn queues(names: &[&str]) -> BTreeSet<String> {
names.iter().map(|name| (*name).to_owned()).collect()
}
#[test]
fn a_complete_read_still_reports_an_absent_queue_as_undeclared() {
let read = DeclaredQueues::new(queues(&["orders"]), Vec::new());
assert_eq!(
declaration_from(&read, "checkout"),
QueueDeclaration::NotDeclared
);
assert_eq!(
declaration_from(&read, "orders"),
QueueDeclaration::Declared
);
}
#[test]
fn a_complete_read_that_found_nothing_is_unknown() {
let read = DeclaredQueues::new(BTreeSet::new(), Vec::new());
assert_eq!(declaration_from(&read, "orders"), QueueDeclaration::Unknown);
}
#[test]
fn a_queue_absent_from_a_partial_read_is_never_undeclared() {
let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
assert_eq!(
declaration_from(&read, "checkout"),
QueueDeclaration::Unknown
);
}
#[test]
fn a_positive_find_survives_a_partial_read() {
let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
assert_eq!(
declaration_from(&read, "orders"),
QueueDeclaration::Declared
);
}
#[test]
fn a_mixed_catalog_never_structurally_refuses_the_undecodable_queue() {
let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
let reason = classify(
declaration_from(&read, "checkout"),
&PoolCensus {
workers_in_pool: 1,
workers_serving_activity: 0,
compatible_workers: 0,
eligible_compatible_workers: 0,
compatible_workers_reachability_lost: 0,
last_compatible_poller_age: None,
compatible_workers_at_capacity: 0,
compatible_workers_capacity_unannounced: 0,
},
);
assert_ne!(reason, Some(QueueServiceReason::NoQueueDeclaration));
assert_eq!(reason, Some(QueueServiceReason::PollersIncompatible));
}
struct FixedDeclarations(QueueDeclaration);
impl QueueDeclarations for FixedDeclarations {
fn declaration_for(&self, _task_queue: &str) -> QueueDeclaration {
self.0
}
}
#[test]
fn an_uninstalled_source_answers_unknown() {
let source = QueueDeclarationSource::default();
assert!(!source.is_installed());
assert_eq!(source.declaration_for("general"), QueueDeclaration::Unknown);
}
#[test]
fn an_installed_source_answers_and_is_visible_to_every_clone() {
let source = QueueDeclarationSource::default();
source
.clone()
.install(Arc::new(FixedDeclarations(QueueDeclaration::NotDeclared)));
assert!(source.is_installed());
assert_eq!(
source.declaration_for("general"),
QueueDeclaration::NotDeclared
);
}
#[test]
fn a_duplicate_install_never_changes_the_answer() {
let source = QueueDeclarationSource::default();
source.install(Arc::new(FixedDeclarations(QueueDeclaration::Declared)));
source.install(Arc::new(FixedDeclarations(QueueDeclaration::NotDeclared)));
assert_eq!(
source.declaration_for("general"),
QueueDeclaration::Declared
);
}
}