use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use aion::{ActivityDispatch, ActivityDispatcher as _};
use aion_core::{ActivityId, RunId, WorkflowId};
use aion_server::worker::{ConnectedWorkerRegistry, HeartbeatTracker, WorkerActivityDispatcher};
use tracing::Level;
#[path = "test_support/capture.rs"]
mod capture;
use capture::{CaptureLayer, CapturedEvent, capture_log, render};
use tracing_subscriber::layer::SubscriberExt as _;
use tracing_subscriber::registry::Registry;
type TestError = Box<dyn std::error::Error>;
const OBSERVATION_WINDOW: Duration = Duration::from_millis(2_500);
fn dispatch_to_unserved_queue() -> ActivityDispatch {
ActivityDispatch {
namespace: "default".to_owned(),
task_queue: "nobody-serves-this".to_owned(),
node: None,
workflow_id: WorkflowId::new_v4(),
run_id: RunId::new_v4(),
activity_id: ActivityId::from_sequence_position(0),
name: "greet".to_owned(),
input: "{}".to_owned(),
config: "{}".to_owned(),
attempt: 1,
advisory: false,
labels: BTreeMap::new(),
}
}
#[test]
fn parked_dispatch_on_an_unserved_queue_states_its_reason_at_warn() -> Result<(), TestError> {
let events = capture_log();
let registry = ConnectedWorkerRegistry::default();
let dispatcher = WorkerActivityDispatcher::new(
registry.clone(),
"default",
HeartbeatTracker::new(Duration::from_secs(5)),
);
let sink = Arc::clone(&events);
let parked = std::thread::spawn(move || {
let subscriber = Registry::default().with(CaptureLayer { events: sink });
tracing::subscriber::with_default(subscriber, || {
dispatcher.dispatch(dispatch_to_unserved_queue())
})
});
std::thread::sleep(OBSERVATION_WINDOW);
let captured = events.lock().map_err(|_| "capture log poisoned")?.clone();
let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
drop(worker_rx);
let registration = registry.register_namespaces(
[String::from("default")],
"nobody-serves-this",
None,
[String::from("greet")].iter(),
worker_tx,
aion_server::worker::UNBOUNDED_SENDER_WORKER_CONCURRENCY,
)?;
let outcome = parked
.join()
.map_err(|_| "parked dispatch thread panicked")?;
assert!(
outcome.is_err(),
"the released dispatch must resolve, not return a result: {outcome:?}"
);
registration.deregister()?;
let reasoned: Vec<&CapturedEvent> = captured
.iter()
.filter(|event| event.fields.contains_key("queue_service_reason"))
.collect();
assert!(
!reasoned.is_empty(),
"a dispatch parked on an unserved queue emitted no reasoned event; \
captured instead: {}",
render(&captured)
);
assert!(
reasoned
.iter()
.any(|event| event.level == Level::WARN && !event.message.is_empty()),
"the unserved park was never escalated to a WARN carrying a message: {reasoned:?}"
);
assert!(
reasoned.iter().any(|event| event
.fields
.get("queue_service_reason")
.is_some_and(|reason| reason == "NO_LIVE_POLLERS")),
"the park did not classify the queue as NO_LIVE_POLLERS: {reasoned:?}"
);
assert!(
reasoned
.iter()
.any(|event| event.fields.contains_key("last_compatible_poller_age_ms")),
"the park carried no last-compatible-poller age: {reasoned:?}"
);
Ok(())
}