use axum::{Json, extract::State};
use serde::Serialize;
use super::auth::HttpCaller;
use super::error::HttpWireError;
use crate::ServerState;
#[derive(Debug, Serialize)]
pub(crate) struct UnservedQueueBody {
namespace: String,
task_queue: String,
activity_type: String,
reason: &'static str,
detail: String,
policy: &'static str,
workers_in_pool: u64,
workers_serving_activity: u64,
compatible_workers: u64,
unserved_for_ms: u64,
waiting: Vec<UnservedDispatchBody>,
}
#[derive(Debug, Serialize)]
pub(crate) struct UnservedDispatchBody {
workflow_id: String,
activity_id: String,
node: Option<String>,
waiting_for_ms: u64,
}
pub(crate) async fn list_unserved_queues(
State(state): State<ServerState>,
HttpCaller(caller): HttpCaller,
) -> Result<Json<Vec<UnservedQueueBody>>, HttpWireError> {
let unserved = state
.unserved_queues()
.map_err(|error| HttpWireError(error.to_wire_error()))?;
let body = unserved
.into_iter()
.filter(|queue| caller.can_access(&queue.key.namespace))
.map(project)
.collect();
Ok(Json(body))
}
fn project(queue: crate::worker::UnservedQueue) -> UnservedQueueBody {
let address = crate::worker::ServiceAddress {
namespace: queue.key.namespace.clone(),
task_queue: queue.key.task_queue.clone(),
activity_type: queue.key.activity_type.clone(),
node: None,
};
UnservedQueueBody {
reason: queue.reason.as_str(),
detail: queue.reason.explain(&address),
policy: queue.policy.as_str(),
namespace: queue.key.namespace,
task_queue: queue.key.task_queue,
activity_type: queue.key.activity_type,
workers_in_pool: count(queue.census.workers_in_pool),
workers_serving_activity: count(queue.census.workers_serving_activity),
compatible_workers: count(queue.census.compatible_workers),
unserved_for_ms: millis(queue.unserved_for),
waiting: queue
.waiting
.into_iter()
.map(|dispatch| UnservedDispatchBody {
workflow_id: dispatch.workflow_id.to_string(),
activity_id: dispatch.activity_id.to_string(),
node: dispatch.node,
waiting_for_ms: millis(dispatch.waiting_for),
})
.collect(),
}
}
fn count(value: usize) -> u64 {
u64::try_from(value).unwrap_or(u64::MAX)
}
fn millis(duration: std::time::Duration) -> u64 {
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
}
#[cfg(test)]
#[path = "queues_tests.rs"]
mod tests;