aion_server/worker/queue_service/
declarations.rs1use std::sync::{Arc, OnceLock};
21
22use aion::loader::DeclaredQueues;
23
24#[derive(Clone, Copy, Debug, Eq, PartialEq)]
26pub enum QueueDeclaration {
27 Declared,
29 NotDeclared,
32 Unknown,
36}
37
38pub trait QueueDeclarations: Send + Sync {
40 fn declaration_for(&self, task_queue: &str) -> QueueDeclaration;
43}
44
45#[derive(Clone, Default)]
52pub struct QueueDeclarationSource {
53 inner: Arc<OnceLock<Arc<dyn QueueDeclarations>>>,
54}
55
56impl std::fmt::Debug for QueueDeclarationSource {
57 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
58 formatter
59 .debug_struct("QueueDeclarationSource")
60 .field("installed", &self.inner.get().is_some())
61 .finish()
62 }
63}
64
65impl QueueDeclarationSource {
66 pub fn install(&self, source: Arc<dyn QueueDeclarations>) {
69 if self.inner.set(source).is_err() {
70 tracing::warn!("queue declaration source already installed; ignoring duplicate set");
71 }
72 }
73
74 #[must_use]
76 pub fn is_installed(&self) -> bool {
77 self.inner.get().is_some()
78 }
79
80 #[must_use]
83 pub fn declaration_for(&self, task_queue: &str) -> QueueDeclaration {
84 self.inner
85 .get()
86 .map_or(QueueDeclaration::Unknown, |source| {
87 source.declaration_for(task_queue)
88 })
89 }
90}
91
92pub struct EngineQueueDeclarations {
94 engine: Arc<aion::Engine>,
95}
96
97impl EngineQueueDeclarations {
98 #[must_use]
100 pub const fn new(engine: Arc<aion::Engine>) -> Self {
101 Self { engine }
102 }
103}
104
105impl std::fmt::Debug for EngineQueueDeclarations {
106 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
107 formatter.write_str("EngineQueueDeclarations")
108 }
109}
110
111fn declaration_from(read: &DeclaredQueues, task_queue: &str) -> QueueDeclaration {
120 if read.declares(task_queue) {
121 QueueDeclaration::Declared
122 } else if read.found_no_declaration() || !read.covers_every_entry() {
123 QueueDeclaration::Unknown
124 } else {
125 QueueDeclaration::NotDeclared
126 }
127}
128
129impl QueueDeclarations for EngineQueueDeclarations {
130 fn declaration_for(&self, task_queue: &str) -> QueueDeclaration {
131 match self.engine.declared_task_queues() {
132 Ok(read) => declaration_from(&read, task_queue),
133 Err(error) => {
134 tracing::warn!(
135 task_queue,
136 %error,
137 "queue declaration lookup failed; treating the declaration as unknown"
138 );
139 QueueDeclaration::Unknown
140 }
141 }
142 }
143}
144
145#[cfg(test)]
146mod tests {
147 use std::collections::BTreeSet;
148
149 use super::super::census::{PoolCensus, classify};
150 use super::super::taxonomy::QueueServiceReason;
151 use super::*;
152
153 fn queues(names: &[&str]) -> BTreeSet<String> {
154 names.iter().map(|name| (*name).to_owned()).collect()
155 }
156
157 #[test]
160 fn a_complete_read_still_reports_an_absent_queue_as_undeclared() {
161 let read = DeclaredQueues::new(queues(&["orders"]), Vec::new());
162 assert_eq!(
163 declaration_from(&read, "checkout"),
164 QueueDeclaration::NotDeclared
165 );
166 assert_eq!(
167 declaration_from(&read, "orders"),
168 QueueDeclaration::Declared
169 );
170 }
171
172 #[test]
173 fn a_complete_read_that_found_nothing_is_unknown() {
174 let read = DeclaredQueues::new(BTreeSet::new(), Vec::new());
175 assert_eq!(declaration_from(&read, "orders"), QueueDeclaration::Unknown);
176 }
177
178 #[test]
182 fn a_queue_absent_from_a_partial_read_is_never_undeclared() {
183 let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
184 assert_eq!(
185 declaration_from(&read, "checkout"),
186 QueueDeclaration::Unknown
187 );
188 }
189
190 #[test]
191 fn a_positive_find_survives_a_partial_read() {
192 let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
193 assert_eq!(
194 declaration_from(&read, "orders"),
195 QueueDeclaration::Declared
196 );
197 }
198
199 #[test]
204 fn a_mixed_catalog_never_structurally_refuses_the_undecodable_queue() {
205 let read = DeclaredQueues::new(queues(&["orders"]), vec!["sha256:legacy".to_owned()]);
206 let reason = classify(
207 declaration_from(&read, "checkout"),
208 &PoolCensus {
209 workers_in_pool: 1,
210 workers_serving_activity: 0,
211 compatible_workers: 0,
212 eligible_compatible_workers: 0,
213 compatible_workers_reachability_lost: 0,
214 last_compatible_poller_age: None,
215 compatible_workers_at_capacity: 0,
216 compatible_workers_capacity_unannounced: 0,
217 },
218 );
219 assert_ne!(reason, Some(QueueServiceReason::NoQueueDeclaration));
220 assert_eq!(reason, Some(QueueServiceReason::PollersIncompatible));
221 }
222
223 struct FixedDeclarations(QueueDeclaration);
224
225 impl QueueDeclarations for FixedDeclarations {
226 fn declaration_for(&self, _task_queue: &str) -> QueueDeclaration {
227 self.0
228 }
229 }
230
231 #[test]
232 fn an_uninstalled_source_answers_unknown() {
233 let source = QueueDeclarationSource::default();
234 assert!(!source.is_installed());
235 assert_eq!(source.declaration_for("general"), QueueDeclaration::Unknown);
236 }
237
238 #[test]
239 fn an_installed_source_answers_and_is_visible_to_every_clone() {
240 let source = QueueDeclarationSource::default();
241 source
242 .clone()
243 .install(Arc::new(FixedDeclarations(QueueDeclaration::NotDeclared)));
244 assert!(source.is_installed());
245 assert_eq!(
246 source.declaration_for("general"),
247 QueueDeclaration::NotDeclared
248 );
249 }
250
251 #[test]
252 fn a_duplicate_install_never_changes_the_answer() {
253 let source = QueueDeclarationSource::default();
254 source.install(Arc::new(FixedDeclarations(QueueDeclaration::Declared)));
255 source.install(Arc::new(FixedDeclarations(QueueDeclaration::NotDeclared)));
256 assert_eq!(
257 source.declaration_for("general"),
258 QueueDeclaration::Declared
259 );
260 }
261}