use super::*;
use crate::job_timing;
use serde::Deserialize;
use serde_json::{json, Value};
use std::collections::HashMap;
const PAGE_SIZE: usize = 25;
const DASHBOARD_BASE_PLACEHOLDER: &str = "__LATER_DASHBOARD_BASE__";
const DASHBOARD_TOPICS_PLACEHOLDER: &str = "__LATER_DASHBOARD_TOPICS__";
const DASHBOARD_LAG_TOPICS_PLACEHOLDER: &str = "__LATER_DASHBOARD_LAG_TOPICS__";
const DASHBOARD_METRICS_ENABLED_PLACEHOLDER: &str = "__LATER_DASHBOARD_METRICS_ENABLED__";
macro_rules! hashmap {
($(($k:literal, $v: literal)),+) => {
{
let mut hm = HashMap::new();
$(hm.insert($k.to_string(), $v.to_string());)+
hm
}
};
}
#[derive(serde::Serialize, serde::Deserialize, Debug)]
#[serde(rename_all = "snake_case", tag = "cmd", deny_unknown_fields)]
pub enum DashboardCmd {
Index,
JobsInStage {
stage: String,
cursor: Option<i64>,
},
Count,
QueueHistory {
minutes: Option<u32>,
},
Job {
id: JobId,
},
Topics,
PartitionJobs {
topic: String,
partition: u32,
cursor: Option<i64>,
},
ConsumerLag,
PartitionQueueDepth,
#[cfg(feature = "prometheus")]
Metrics,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct DashboardQuery {
query: String,
}
#[derive(serde::Serialize, typed_builder::TypedBuilder, Debug)]
pub struct DashboardResponse {
pub status_code: u16,
#[builder(default=HashMap::default())]
pub headers: HashMap<String, String>,
pub body: String,
}
#[derive(thiserror::Error, Debug)]
pub enum ResponseError {
#[error("Internal server error: {0}")]
InternalServer(#[from] anyhow::Error),
#[error("Dashboard JSON error: {0}")]
ParseError(#[from] serde_json::Error),
}
#[tracing::instrument(skip(persist, topics, retained_log_lag, committer, count_cache))]
pub(crate) async fn handle_http_raw(
persist: Arc<Persist>,
topics: &[crate::topic::TopicConfig],
retained_log_lag: &crate::RetainedLogLag,
committer: &Arc<dyn crate::backend::JobCommitter>,
count_cache: &CountCache,
prefix: String,
query_string: String,
) -> Result<DashboardResponse, ResponseError> {
if query_string.is_empty() {
return dashboard_index(&prefix, topics, retained_log_lag);
}
let query = match serde_querystring::from_str::<DashboardQuery>(&query_string) {
Ok(query) => query,
Err(_) => return Ok(DashboardResponse::error(400, "Invalid dashboard query")),
};
let command = match serde_json::from_str::<DashboardCmd>(&query.query) {
Ok(command) => command,
Err(_) => return Ok(DashboardResponse::error(400, "Invalid dashboard command")),
};
handle_http(
persist,
topics,
retained_log_lag,
committer,
count_cache,
&prefix,
command,
)
.await
}
fn topics_json(topics: &[crate::topic::TopicConfig]) -> Value {
json!(topics
.iter()
.map(|topic| json!({
"name": topic.name(),
"partition_count": topic.partition_count(),
}))
.collect::<Vec<_>>())
}
#[cfg(feature = "retained-log")]
async fn consumer_lag_json(retained_log_lag: &crate::RetainedLogLag) -> Value {
let Some(source) = retained_log_lag.source() else {
return json!([]);
};
let mut result = Vec::new();
for topic in retained_log_lag.topics() {
match source.consumer_lag(topic).await {
Ok(lags) => result.extend(lags.into_iter().map(|lag| {
json!({
"topic": topic,
"partition": lag.partition.0,
"group": lag.group,
"end_offset": lag.end_offset,
"committed_offset": lag.committed_offset,
"lag": lag.lag,
})
})),
Err(error) => {
tracing::warn!(%error, topic, "Failed to compute consumer lag");
}
}
}
json!(result)
}
#[cfg(not(feature = "retained-log"))]
async fn consumer_lag_json(_retained_log_lag: &crate::RetainedLogLag) -> Value {
json!([])
}
async fn partition_queue_depth_json(committer: &Arc<dyn crate::backend::JobCommitter>) -> Value {
match committer.partition_queue_depth_summary().await {
Ok(summary) => json!(summary
.into_iter()
.map(|(topic, partition, depth)| json!({
"topic": topic,
"partition": partition.0,
"depth": depth,
}))
.collect::<Vec<_>>()),
Err(error) => {
tracing::warn!(%error, "Failed to summarize partition queue depth");
json!([])
}
}
}
fn dashboard_index(
prefix: &str,
topics: &[crate::topic::TopicConfig],
retained_log_lag: &crate::RetainedLogLag,
) -> Result<DashboardResponse, ResponseError> {
let prefix = serde_json::to_string(prefix)?;
let topics = topics_json(topics).to_string();
let lag_topics = json!(lag_topics(retained_log_lag)).to_string();
Ok(DashboardResponse::html(
include_str!("dash_index.html")
.replace(DASHBOARD_BASE_PLACEHOLDER, &prefix)
.replace(DASHBOARD_TOPICS_PLACEHOLDER, &topics)
.replace(DASHBOARD_LAG_TOPICS_PLACEHOLDER, &lag_topics)
.replace(
DASHBOARD_METRICS_ENABLED_PLACEHOLDER,
&metrics_enabled().to_string(),
),
))
}
#[cfg(feature = "prometheus")]
fn metrics_enabled() -> bool {
true
}
#[cfg(not(feature = "prometheus"))]
fn metrics_enabled() -> bool {
false
}
#[cfg(feature = "retained-log")]
fn lag_topics(retained_log_lag: &crate::RetainedLogLag) -> Vec<String> {
if retained_log_lag.source().is_some() {
retained_log_lag.topics().to_vec()
} else {
Vec::new()
}
}
#[cfg(not(feature = "retained-log"))]
fn lag_topics(_retained_log_lag: &crate::RetainedLogLag) -> Vec<String> {
Vec::new()
}
const COUNT_CACHE_TTL: std::time::Duration = std::time::Duration::from_secs(1);
#[derive(Default)]
pub(crate) struct CountCache {
slot: tokio::sync::Mutex<Option<(std::time::Instant, Value)>>,
}
impl CountCache {
async fn get(&self, persist: &Persist) -> anyhow::Result<Value> {
let mut slot = self.slot.lock().await;
if let Some((computed_at, value)) = slot.as_ref() {
if computed_at.elapsed() < COUNT_CACHE_TTL {
return Ok(value.clone());
}
}
let value = compute_count(persist).await?;
*slot = Some((std::time::Instant::now(), value.clone()));
Ok(value)
}
}
async fn compute_count(persist: &Persist) -> anyhow::Result<Value> {
let namespace = persist.key_prefix();
let jobs_in_stages = persist.inner.job_index_stage_counts(namespace).await?;
let now = chrono::Utc::now();
let mut recent_transitions = HashMap::new();
for stage in Stage::get_all_stage_names() {
let count = persist
.inner
.job_index_recent_transition_count(namespace, &stage, now)
.await?;
recent_transitions.insert(stage, count);
}
let active_workers_key = format!("{namespace}-stats-active-workers");
let active_workers = persist.inner.range_count(&active_workers_key).await?;
let completed = recent_transitions.get("success").copied().unwrap_or(0)
+ recent_transitions.get("failed").copied().unwrap_or(0);
let jobs_per_second = completed as f64 / RECENT_WINDOW_SECONDS as f64;
let running = jobs_in_stages.get("running").copied().unwrap_or(0);
let wait_time_regular: WaitStats = persist
.inner
.job_index_recent_wait_stats(namespace, "regular", now)
.await?
.into();
let wait_time_sequential: WaitStats = persist
.inner
.job_index_recent_wait_stats(namespace, "sequential", now)
.await?
.into();
let db_size_bytes = persist.inner.total_db_size_bytes().await?;
Ok(json!({
"stages": jobs_in_stages,
"db_size_bytes": db_size_bytes,
"performance": {
"window_seconds": RECENT_WINDOW_SECONDS,
"enqueued": recent_transitions.get("enqueued").copied().unwrap_or(0),
"started": recent_transitions.get("running").copied().unwrap_or(0),
"retried": recent_transitions.get("requeued").copied().unwrap_or(0),
"succeeded": recent_transitions.get("success").copied().unwrap_or(0),
"failed": recent_transitions.get("failed").copied().unwrap_or(0),
"completed": completed,
"jobs_per_second": jobs_per_second,
"active_workers": active_workers,
"busy_workers": running.min(active_workers),
"wait_time": {
"regular": wait_time_regular,
"sequential": wait_time_sequential,
},
}
}))
}
async fn handle_http(
persist: Arc<Persist>,
topics: &[crate::topic::TopicConfig],
retained_log_lag: &crate::RetainedLogLag,
committer: &Arc<dyn crate::backend::JobCommitter>,
count_cache: &CountCache,
prefix: &str,
command: DashboardCmd,
) -> Result<DashboardResponse, ResponseError> {
let namespace = persist.key_prefix();
let response = match command {
DashboardCmd::Index => dashboard_index(prefix, topics, retained_log_lag)?,
DashboardCmd::Topics => DashboardResponse::json(topics_json(topics))?,
DashboardCmd::ConsumerLag => {
DashboardResponse::json(consumer_lag_json(retained_log_lag).await)?
}
DashboardCmd::PartitionQueueDepth => {
DashboardResponse::json(partition_queue_depth_json(committer).await)?
}
#[cfg(feature = "prometheus")]
DashboardCmd::Metrics => match crate::metrics::snapshot() {
Ok(snapshot) => DashboardResponse::json(snapshot)?,
Err(error) => {
tracing::warn!(%error, "Failed to snapshot Prometheus metrics");
DashboardResponse::json(json!([]))?
}
},
DashboardCmd::PartitionJobs {
topic,
partition,
cursor,
} => {
let Some(config) = topics.iter().find(|config| config.name() == topic) else {
return Ok(DashboardResponse::error(400, "Unknown topic"));
};
if partition >= config.partition_count() {
return Ok(DashboardResponse::error(400, "Partition out of range"));
}
let page = persist
.inner
.job_index_list_by_partition(namespace, &topic, partition, cursor, PAGE_SIZE)
.await?;
DashboardResponse::json(json!({
"result": page.items.iter().map(job_index_item).collect::<Vec<_>>(),
"paging": {
"next_cursor": page.next_cursor,
"items": page.items.len(),
}
}))?
}
DashboardCmd::JobsInStage { stage, cursor } => {
if !Stage::get_all_stage_names().contains(&stage) {
return Ok(DashboardResponse::error(400, "Unknown job stage"));
}
let page = persist
.inner
.job_index_list_by_stage(namespace, &stage, cursor, PAGE_SIZE)
.await?;
DashboardResponse::json(json!({
"result": page.items.iter().map(job_index_item).collect::<Vec<_>>(),
"paging": {
"next_cursor": page.next_cursor,
"items": page.items.len(),
}
}))?
}
DashboardCmd::Job { id } => match persist.get_job(id).await? {
Some(job) => DashboardResponse::json(get_job_detail(&persist, job).await?)?,
None => DashboardResponse::error(404, "Job not found"),
},
DashboardCmd::Count => DashboardResponse::json(count_cache.get(&persist).await?)?,
DashboardCmd::QueueHistory { minutes } => {
let minutes = minutes.unwrap_or(360).clamp(1, 7 * 24 * 60);
let since = chrono::Utc::now() - chrono::Duration::minutes(i64::from(minutes));
let samples = persist.inner.queue_samples_since(namespace, since).await?;
DashboardResponse::json(json!({
"minutes": minutes,
"points": samples
.iter()
.map(|sample| json!({
"t": sample.at.timestamp_millis(),
"queued": sample.queued,
"partitioned": sample.partitioned,
}))
.collect::<Vec<_>>(),
}))?
}
};
Ok(response)
}
fn job_index_item(row: &crate::storage::JobIndexRow) -> Value {
let now = chrono::Utc::now();
let finished_at = if is_terminal_stage(&row.stage) {
row.stage_date
} else {
now
};
let total_ms = finished_at
.signed_duration_since(row.created_at)
.num_milliseconds()
.max(0);
json!({
"id": row.job_id,
"info": {
"id": row.job_id,
"payload_type": row.payload_type,
"stage": { "type": row.stage, "date": row.stage_date },
"topic": row.topic,
"partition": row.partition,
"sequence": row.sequence,
},
"timing": {
"total_ms": total_ms,
},
})
}
fn is_terminal_stage(stage: &str) -> bool {
stage == "success" || stage == "failed"
}
async fn get_job_detail(persist: &Persist, job: Job) -> anyhow::Result<Value> {
let namespace = persist.key_prefix();
let parent_job_id = job
.previous_stages
.iter()
.chain(std::iter::once(&job.stage))
.find_map(|stage| match stage {
Stage::Waiting(waiting) => Some(waiting.parent_id.clone()),
Stage::Delayed(_)
| Stage::Enqueued(_)
| Stage::Running(_)
| Stage::Requeued(_)
| Stage::Success(_)
| Stage::Failed(_) => None,
});
let continuations = persist
.inner
.job_index_list_continuations(namespace, &job.id.to_string(), PAGE_SIZE + 1)
.await?;
let continuations_truncated = continuations.len() > PAGE_SIZE;
let continuation_job_ids: Vec<String> = continuations
.into_iter()
.take(PAGE_SIZE)
.map(|row| row.job_id)
.collect();
let partition = match (&job.topic, job.partition, job.sequence) {
(Some(topic), Some(partition_id), Some(sequence)) => {
let previous_job_id = persist
.inner
.job_index_partition_neighbor(namespace, topic, partition_id, sequence, true)
.await?
.map(|row| row.job_id);
let next_job_id = persist
.inner
.job_index_partition_neighbor(namespace, topic, partition_id, sequence, false)
.await?
.map(|row| row.job_id);
Some(json!({
"topic": topic,
"partition": partition_id,
"sequence": sequence,
"previous_job_id": previous_job_id,
"next_job_id": next_job_id,
}))
}
_ => None,
};
let timing = job_timing::calculate(&job.previous_stages, &job.stage, chrono::Utc::now());
let payload = match rmp_serde::from_slice::<Value>(&job.payload) {
Ok(value) => json!({ "format": "msgpack", "value": value }),
Err(_) => json!({
"format": "raw",
"value": String::from_utf8_lossy(&job.payload),
}),
};
Ok(json!({
"id": job.id,
"info": {
"id": job.id,
"payload_type": job.payload_type,
"stage": job.stage,
"previous_stages": job.previous_stages,
"topic": job.topic,
"partition": job.partition,
"sequence": job.sequence,
},
"timing": {
"total_ms": timing.total_ms,
"processing_ms": timing.processing_ms,
"waiting_ms": timing.waiting_ms,
},
"relations": {
"parent_job_id": parent_job_id,
"continuation_job_ids": continuation_job_ids,
"continuations_truncated": continuations_truncated,
"partition": partition,
},
"payload": payload,
}))
}
impl DashboardResponse {
pub fn json<T: serde::Serialize>(json: T) -> anyhow::Result<Self> {
Ok(Self {
status_code: 200,
headers: hashmap!(("content-type", "application/json")),
body: serde_json::to_string(&json)?,
})
}
pub fn html(html: String) -> Self {
Self {
status_code: 200,
headers: hashmap!(("content-type", "text/html")),
body: html,
}
}
pub fn error<S: Into<String>>(status_code: u16, error: S) -> Self {
Self {
status_code,
headers: hashmap!(("content-type", "application/json")),
body: json!({ "error": error.into() }).to_string(),
}
}
}
#[cfg(all(test, feature = "sqlite"))]
mod tests {
use super::*;
use crate::{
models::{EnqueuedStage, FailedStage, JobConfig, RunningStage, SuccessStage, WaitingStage},
storage::Sqlite,
};
async fn test_persist() -> Arc<Persist> {
let storage = Sqlite::new("sqlite::memory:")
.await
.expect("create dashboard test storage");
Arc::new(Persist::new(
Box::new(storage),
"later-dashboard-test".to_string(),
))
}
async fn test_committer() -> Arc<dyn crate::backend::JobCommitter> {
use crate::backend::Backend;
let backend = crate::backend::SqliteBackend::connect(
format!("dashboard-test-{}", crate::generate_id()),
"sqlite::memory:",
)
.await
.expect("connect dashboard test committer");
Box::new(backend)
.into_parts()
.expect("split dashboard test backend")
.committer
}
async fn test_stats(persist: Arc<Persist>) -> Stats {
Stats::new(
persist,
Vec::new(),
crate::RetainedLogLag::default(),
test_committer().await,
)
}
async fn save_and_record(persist: &Persist, stats: &Stats, job: &Job) {
persist
.inner
.apply(
Persist::job_operations(persist.key_prefix(), job).expect("build job operations"),
)
.await
.expect("save durable job record");
stats.record_transition(job, None).await;
stats.flush_metrics().await;
}
fn job(id: &str, stage: Stage) -> Job {
Job {
id: JobId(id.to_string()),
payload_type: "test_job".to_string(),
payload: Vec::new(),
config: JobConfig::default(),
stage,
previous_stages: Vec::new(),
recurring_job_id: None,
topic: None,
partition: None,
sequence: None,
}
}
async fn add_job(persist: &Arc<Persist>, job: Job) {
persist
.inner
.apply(
Persist::job_operations(persist.key_prefix(), &job).expect("build job operations"),
)
.await
.expect("save durable job record");
let stats = test_stats(persist.clone()).await;
stats.record_transition(&job, None).await;
stats.flush_metrics().await;
}
#[tokio::test]
async fn job_response_includes_total_duration_since_creation() -> anyhow::Result<()> {
let persist = test_persist().await;
let mut success = job(
"timed-job",
Stage::Success(SuccessStage {
date: chrono::DateTime::from_timestamp(10, 0).unwrap_or_default(),
}),
);
success.previous_stages = vec![
Stage::Enqueued(EnqueuedStage {
date: chrono::DateTime::from_timestamp(1, 0).unwrap_or_default(),
}),
Stage::Running(RunningStage {
date: chrono::DateTime::from_timestamp(4, 0).unwrap_or_default(),
}),
];
add_job(&persist, success).await;
let page = persist
.inner
.job_index_list_by_stage(persist.key_prefix(), "success", None, 10)
.await?;
let row = page.items.first().expect("indexed job row");
assert_eq!(job_index_item(row)["timing"], json!({ "total_ms": 9_000 }));
Ok(())
}
#[tokio::test]
async fn job_detail_includes_a_decoded_payload() -> anyhow::Result<()> {
let persist = test_persist().await;
let full_job = Job {
id: JobId("payload-job".to_string()),
payload_type: "test_job".to_string(),
payload: crate::encoder::encode(("hello@example.com", "hi"))?,
config: JobConfig::default(),
stage: Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
previous_stages: Vec::new(),
recurring_job_id: None,
topic: None,
partition: None,
sequence: None,
};
persist
.inner
.apply(Persist::job_operations(persist.key_prefix(), &full_job)?)
.await?;
let detail = get_job_detail(&persist, full_job).await?;
assert_eq!(
detail["payload"],
json!({ "format": "msgpack", "value": ["hello@example.com", "hi"] })
);
Ok(())
}
#[tokio::test]
async fn job_detail_falls_back_to_raw_text_for_non_msgpack_payloads() -> anyhow::Result<()> {
let persist = test_persist().await;
let full_job = Job {
id: JobId("raw-payload-job".to_string()),
payload_type: "test_job".to_string(),
payload: vec![0xd9, 0xff],
config: JobConfig::default(),
stage: Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
previous_stages: Vec::new(),
recurring_job_id: None,
topic: None,
partition: None,
sequence: None,
};
persist
.inner
.apply(Persist::job_operations(persist.key_prefix(), &full_job)?)
.await?;
let detail = get_job_detail(&persist, full_job).await?;
assert_eq!(detail["payload"]["format"], "raw");
Ok(())
}
#[tokio::test]
async fn job_details_link_both_sides_of_a_continuation() -> anyhow::Result<()> {
let persist = test_persist().await;
let stats = test_stats(persist.clone()).await;
let parent = job(
"parent-job",
Stage::Success(SuccessStage {
date: chrono::Utc::now(),
}),
);
let mut child = job(
"child-job",
Stage::Waiting(WaitingStage {
date: chrono::Utc::now(),
parent_id: parent.id.clone(),
}),
);
save_and_record(&persist, &stats, &parent).await;
save_and_record(&persist, &stats, &child).await;
child.previous_stages.push(child.stage.clone());
child.stage = Stage::Running(RunningStage {
date: chrono::Utc::now(),
});
save_and_record(&persist, &stats, &child).await;
let parent_response = request(
persist.clone(),
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22parent-job%22%7D",
)
.await;
let child_response = request(
persist,
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22child-job%22%7D",
)
.await;
let parent_body = serde_json::from_str::<Value>(&parent_response.body)?;
let child_body = serde_json::from_str::<Value>(&child_response.body)?;
assert_eq!(
(
parent_body["relations"].clone(),
child_body["relations"].clone()
),
(
json!({
"parent_job_id": null,
"continuation_job_ids": ["child-job"],
"continuations_truncated": false,
"partition": null,
}),
json!({
"parent_job_id": "parent-job",
"continuation_job_ids": [],
"continuations_truncated": false,
"partition": null,
}),
)
);
Ok(())
}
#[tokio::test]
async fn job_details_link_previous_and_next_in_the_same_partition() -> anyhow::Result<()> {
let persist = test_persist().await;
let stats = test_stats(persist.clone()).await;
let make = |id: &str, sequence: i64| {
let mut j = job(
id,
Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
);
j.topic = Some("orders".to_string());
j.partition = Some(2);
j.sequence = Some(sequence);
j
};
for (id, sequence) in [("seq-0", 0), ("seq-1", 1), ("seq-2", 2)] {
save_and_record(&persist, &stats, &make(id, sequence)).await;
}
let first = request(
persist.clone(),
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22seq-0%22%7D",
)
.await;
let middle = request(
persist.clone(),
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22seq-1%22%7D",
)
.await;
let last = request(
persist,
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22seq-2%22%7D",
)
.await;
let partition = |response: &DashboardResponse| -> Value {
serde_json::from_str::<Value>(&response.body).expect("decode dashboard response")
["relations"]["partition"]
.clone()
};
assert_eq!(
(partition(&first), partition(&middle), partition(&last)),
(
json!({
"topic": "orders",
"partition": 2,
"sequence": 0,
"previous_job_id": null,
"next_job_id": "seq-1",
}),
json!({
"topic": "orders",
"partition": 2,
"sequence": 1,
"previous_job_id": "seq-0",
"next_job_id": "seq-2",
}),
json!({
"topic": "orders",
"partition": 2,
"sequence": 2,
"previous_job_id": "seq-1",
"next_job_id": null,
}),
)
);
Ok(())
}
#[tokio::test]
async fn topics_command_lists_configured_topics_with_partition_counts() -> anyhow::Result<()> {
let persist = test_persist().await;
let topics = vec![
crate::topic::TopicConfig::new("orders", 4)?,
crate::topic::TopicConfig::new("shipments", 2)?,
];
let response =
request_with_topics(persist, &topics, "query=%7B%22cmd%22%3A%22topics%22%7D").await;
let body: Value = serde_json::from_str(&response.body)?;
assert_eq!(
body,
json!([
{"name": "orders", "partition_count": 4},
{"name": "shipments", "partition_count": 2},
])
);
Ok(())
}
#[tokio::test]
async fn partition_jobs_command_lists_newest_sequence_first_and_validates_topic_and_partition(
) -> anyhow::Result<()> {
let persist = test_persist().await;
let stats = test_stats(persist.clone()).await;
let topics = vec![crate::topic::TopicConfig::new("orders", 4)?];
let make = |id: &str, partition: u32, sequence: i64| {
let mut j = job(
id,
Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
);
j.topic = Some("orders".to_string());
j.partition = Some(partition);
j.sequence = Some(sequence);
j
};
for (id, sequence) in [("first", 0), ("second", 1), ("third", 2)] {
stats.record_transition(&make(id, 1, sequence), None).await;
}
stats
.record_transition(&make("other-partition", 0, 0), None)
.await;
let response = request_with_topics(
persist.clone(),
&topics,
"query=%7B%22cmd%22%3A%22partition_jobs%22%2C%22topic%22%3A%22orders%22%2C%22partition%22%3A1%2C%22cursor%22%3Anull%7D",
)
.await;
let body: Value = serde_json::from_str(&response.body)?;
let ids: Vec<&str> = body["result"]
.as_array()
.expect("job list")
.iter()
.map(|job| job["id"].as_str().expect("job id"))
.collect();
assert_eq!(ids, vec!["third", "second", "first"]);
assert_eq!(body["paging"]["items"], json!(3));
let unknown_topic = request_with_topics(
persist.clone(),
&topics,
"query=%7B%22cmd%22%3A%22partition_jobs%22%2C%22topic%22%3A%22missing%22%2C%22partition%22%3A0%2C%22cursor%22%3Anull%7D",
)
.await;
assert_eq!(unknown_topic.status_code, 400);
let out_of_range = request_with_topics(
persist,
&topics,
"query=%7B%22cmd%22%3A%22partition_jobs%22%2C%22topic%22%3A%22orders%22%2C%22partition%22%3A99%2C%22cursor%22%3Anull%7D",
)
.await;
assert_eq!(out_of_range.status_code, 400);
Ok(())
}
#[cfg(feature = "retained-log")]
#[tokio::test]
async fn consumer_lag_command_reports_lag_from_the_configured_source() -> anyhow::Result<()> {
use crate::retained_log::SqliteLog;
use crate::topic::PartitionKey;
let log = SqliteLog::connect("dashboard-lag-test", "sqlite::memory:").await?;
log.register_topic(crate::topic::TopicConfig::new("events", 1)?)
.await?;
log.append("events", PartitionKey::from("a"), b"1".to_vec())
.await?;
log.append("events", PartitionKey::from("b"), b"2".to_vec())
.await?;
log.commit_offset("events", "consumer", crate::topic::PartitionId(0), 2)
.await?;
let lag_config =
crate::RetainedLogLag::new(vec!["events".to_string()], std::sync::Arc::new(log));
let persist = test_persist().await;
let response = handle_http_raw(
persist,
&[],
&lag_config,
&test_committer().await,
&CountCache::default(),
"/dashboard".to_string(),
"query=%7B%22cmd%22%3A%22consumer_lag%22%7D".to_string(),
)
.await?;
let body: Vec<Value> = serde_json::from_str(&response.body)?;
assert_eq!(
body,
vec![json!({
"topic": "events",
"partition": 0,
"group": "consumer",
"end_offset": 3,
"committed_offset": 2,
"lag": 1,
})]
);
Ok(())
}
#[tokio::test]
async fn partition_queue_depth_command_counts_not_yet_completed_jobs() -> anyhow::Result<()> {
use crate::backend::Backend;
let backend =
crate::backend::SqliteBackend::connect("dashboard-depth-test", "sqlite::memory:")
.await?;
let committer = Box::new(backend).into_parts()?.committer;
committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 2)?])
.await?;
for (id, partition) in [("a", 0), ("b", 0), ("c", 1)] {
committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(partition),
&JobId(id.to_string()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
}
let persist = test_persist().await;
let topics = vec![crate::topic::TopicConfig::new("orders", 2)?];
let response = handle_http_raw(
persist,
&topics,
&crate::RetainedLogLag::default(),
&committer,
&CountCache::default(),
"/dashboard".to_string(),
"query=%7B%22cmd%22%3A%22partition_queue_depth%22%7D".to_string(),
)
.await?;
let mut body: Vec<Value> = serde_json::from_str(&response.body)?;
body.sort_by_key(|row| row["partition"].as_u64().unwrap());
assert_eq!(
body,
vec![
json!({"topic": "orders", "partition": 0, "depth": 2}),
json!({"topic": "orders", "partition": 1, "depth": 1}),
]
);
Ok(())
}
async fn request(persist: Arc<Persist>, query: &str) -> DashboardResponse {
request_with_topics(persist, &[], query).await
}
async fn request_with_topics(
persist: Arc<Persist>,
topics: &[crate::topic::TopicConfig],
query: &str,
) -> DashboardResponse {
handle_http_raw(
persist,
topics,
&crate::RetainedLogLag::default(),
&test_committer().await,
&CountCache::default(),
"/jobs/dashboard".to_string(),
query.to_string(),
)
.await
.expect("handle dashboard request")
}
#[tokio::test]
async fn validates_query_stage_and_cursor() {
let persist = test_persist().await;
let responses = [
request(persist.clone(), "query=not-json").await,
request(
persist.clone(),
"query=%7B%22cmd%22%3A%22count%22%7D&query=%7B%22cmd%22%3A%22count%22%7D",
)
.await,
request(
persist.clone(),
"query=%7B%22cmd%22%3A%22count%22%7D&extra=value",
)
.await,
request(
persist.clone(),
"query=%7B%22cmd%22%3A%22jobs_in_stage%22%2C%22stage%22%3A%22unknown%22%2C%22cursor%22%3Anull%7D",
)
.await,
request(
persist,
"query=%7B%22cmd%22%3A%22jobs_in_stage%22%2C%22stage%22%3A%22success%22%2C%22page%22%3A1%7D",
)
.await,
request(
test_persist().await,
"query=%7B%22cmd%22%3A%22jobs_in_stage%22%2C%22stage%22%3A%22success%22%2C%22cursor%22%3A%22bad%22%7D",
)
.await,
];
assert_eq!(
responses
.iter()
.map(|response| (response.status_code, response.body.as_str()))
.collect::<Vec<_>>(),
vec![
(400, r#"{"error":"Invalid dashboard command"}"#),
(400, r#"{"error":"Invalid dashboard query"}"#),
(400, r#"{"error":"Invalid dashboard query"}"#),
(400, r#"{"error":"Unknown job stage"}"#),
(400, r#"{"error":"Invalid dashboard command"}"#),
(400, r#"{"error":"Invalid dashboard command"}"#),
]
);
}
#[tokio::test]
async fn lists_newest_jobs_first_and_pages_correctly() {
let persist = test_persist().await;
for index in 0..27 {
add_job(
&persist,
job(
&format!("job-{index:02}"),
Stage::Running(RunningStage {
date: chrono::Utc::now(),
}),
),
)
.await;
}
let response = request(
persist.clone(),
"query=%7B%22cmd%22%3A%22jobs_in_stage%22%2C%22stage%22%3A%22running%22%2C%22cursor%22%3Anull%7D",
)
.await;
let body: Value = serde_json::from_str(&response.body).expect("decode dashboard response");
let cursor = body["paging"]["next_cursor"]
.as_i64()
.expect("next dashboard cursor");
let second_page = request(
persist,
&format!(
"query=%7B%22cmd%22%3A%22jobs_in_stage%22%2C%22stage%22%3A%22running%22%2C%22cursor%22%3A{cursor}%7D"
),
)
.await;
let ids = body["result"]
.as_array()
.expect("job list")
.iter()
.map(|job| job["id"].as_str().expect("job id"))
.collect::<Vec<_>>();
assert_eq!(response.status_code, 200);
assert_eq!(
ids,
(2..=26)
.rev()
.map(|index| format!("job-{index:02}"))
.collect::<Vec<_>>()
);
assert_eq!(body["paging"]["items"], json!(25));
let second_body =
serde_json::from_str::<Value>(&second_page.body).expect("decode second dashboard page");
assert_eq!(
(
second_body["result"]
.as_array()
.expect("second job list")
.iter()
.map(|job| job["id"].as_str().expect("second page job id"))
.collect::<Vec<_>>(),
second_body["paging"]["items"].clone(),
),
(vec!["job-01", "job-00"], json!(2))
);
}
#[tokio::test]
async fn an_out_of_order_transition_cannot_overwrite_a_newer_one() -> anyhow::Result<()> {
let persist = test_persist().await;
let stats = test_stats(persist.clone()).await;
let enqueued = job(
"racy-job",
Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
);
let mut success = enqueued.clone();
success.previous_stages = vec![enqueued.stage.clone()];
success.stage = Stage::Success(SuccessStage {
date: chrono::Utc::now(),
});
stats.record_transition(&success, None).await;
stats.record_transition(&enqueued, None).await;
let counts = persist
.inner
.job_index_stage_counts(persist.key_prefix())
.await?;
assert_eq!(counts.get("success").copied(), Some(1));
assert_eq!(counts.get("enqueued").copied(), None);
Ok(())
}
#[tokio::test]
async fn counts_failed_jobs_separately_and_returns_not_found() {
let persist = test_persist().await;
add_job(
&persist,
job(
"failed-job",
Stage::Failed(FailedStage {
date: chrono::Utc::now(),
reason: "handler error".to_string(),
}),
),
)
.await;
let count = request(persist.clone(), "query=%7B%22cmd%22%3A%22count%22%7D").await;
let missing = request(
persist,
"query=%7B%22cmd%22%3A%22job%22%2C%22id%22%3A%22missing%22%7D",
)
.await;
let mut count_body =
serde_json::from_str::<Value>(&count.body).expect("decode count response");
let db_size_bytes = count_body
.as_object_mut()
.expect("count response is a JSON object")
.remove("db_size_bytes")
.and_then(|value| value.as_u64());
assert!(
db_size_bytes.is_some_and(|size| size > 0),
"expected a positive db_size_bytes, got {db_size_bytes:?}"
);
assert_eq!(
count_body,
json!({
"stages": {
"failed": 1
},
"performance": {
"window_seconds": 60,
"enqueued": 0,
"started": 0,
"retried": 0,
"succeeded": 0,
"failed": 1,
"completed": 1,
"jobs_per_second": 1.0 / 60.0,
"active_workers": 0,
"busy_workers": 0,
"wait_time": {
"regular": {"count": 0, "avg_ms": 0, "max_ms": 0},
"sequential": {"count": 0, "avg_ms": 0, "max_ms": 0}
}
}
})
);
assert_eq!(
(missing.status_code, missing.body),
(404, r#"{"error":"Job not found"}"#.to_string())
);
}
#[tokio::test]
async fn index_uses_a_safe_configured_path() {
let persist = test_persist().await;
let response = handle_http_raw(
persist.clone(),
&[],
&crate::RetainedLogLag::default(),
&test_committer().await,
&CountCache::default(),
"/dash\";alert(1)//".to_string(),
String::new(),
)
.await
.expect("render dashboard index");
assert!(response
.body
.contains(r#"const BASE_APP = "/dash\";alert(1)//";"#));
assert!(!response.body.contains(DASHBOARD_BASE_PLACEHOLDER));
assert!(!response.body.contains("https://"));
assert!(response.body.contains("id=\"previous\""));
assert!(response.body.contains("id=\"next\""));
assert!(response.body.contains("id=\"jobs-per-second\""));
assert!(response.body.contains("id=\"detail-total-duration\""));
assert!(response.body.contains("id=\"relations\""));
assert!(response.body.contains("relationRow(\"Parent\""));
assert!(response.body.contains("relationRow(\"Continues with\""));
assert!(response
.body
.contains("relationRow(\"Previous in partition\""));
assert!(response.body.contains("relationRow(\"Next in partition\""));
assert!(response.body.contains("id=\"detail-partition-row\""));
assert!(response
.body
.contains("formatDuration(job.timing.total_ms)"));
assert!(response.body.contains("id=\"active-workers\""));
assert!(response
.body
.contains("loadStage(state.stage, state.cursor, false, true)"));
assert!(response.body.contains("silent = false"));
assert!(response.body.contains(r#"const TOPICS = [];"#));
assert!(!response.body.contains(DASHBOARD_TOPICS_PLACEHOLDER));
assert!(response.body.contains("id=\"tabs\""));
assert!(response.body.contains("id=\"all-jobs\""));
assert!(response.body.contains("id=\"topics-menu\""));
assert!(response.body.contains("id=\"topic-list\""));
assert!(response.body.contains("id=\"partition-picker\""));
assert!(response.body.contains("id=\"partition-select\""));
assert!(response.body.contains("id=\"job-table\""));
assert!(response.body.contains("class=\"table-scroll\""));
assert!(response.body.contains("th class=\"seq\""));
assert!(response.body.contains("id=\"detail-previous\""));
assert!(response.body.contains("id=\"detail-next\""));
assert!(response.body.contains("id=\"detail-position\""));
assert!(response.body.contains("function loadAdjacentJob(offset)"));
assert!(response
.body
.contains("if (!elements.details.open) elements.details.showModal();"));
assert!(response.body.contains("cmd: \"partition_jobs\""));
assert!(response.body.contains("id=\"monitoring-tab\""));
assert!(response.body.contains("id=\"monitoring\""));
assert!(response.body.contains("id=\"monitoring-grid\""));
assert!(!response
.body
.contains(DASHBOARD_METRICS_ENABLED_PLACEHOLDER));
let index_command = request(persist, "query=%7B%22cmd%22%3A%22index%22%7D").await;
assert_eq!(index_command.status_code, 200);
assert!(index_command
.body
.contains("<title>Later Dashboard</title>"));
}
#[cfg(feature = "prometheus")]
#[tokio::test]
async fn metrics_command_returns_the_prometheus_snapshot_as_json() -> anyhow::Result<()> {
crate::metrics::record_job_transition(
"metrics-command-test",
&job(
"metrics-job",
Stage::Success(SuccessStage {
date: chrono::Utc::now(),
}),
),
);
let persist = test_persist().await;
let response = request(persist, "query=%7B%22cmd%22%3A%22metrics%22%7D").await;
assert_eq!(response.status_code, 200);
let body: Value = serde_json::from_str(&response.body)?;
let families = body.as_array().expect("metrics families array");
let transitions = families
.iter()
.find(|family| family["name"] == "later_job_transitions_total")
.expect("later_job_transitions_total family present");
assert_eq!(transitions["type"], "counter");
assert!(transitions["metrics"]
.as_array()
.expect("metrics array")
.iter()
.any(|metric| metric["labels"]["queue"] == "metrics-command-test"
&& metric["labels"]["stage"] == "success"
&& metric["value"].as_f64().unwrap_or(0.0) >= 1.0));
Ok(())
}
#[tokio::test]
async fn index_embeds_configured_topics_for_the_dashboard_tabs() {
let persist = test_persist().await;
let topics = vec![crate::topic::TopicConfig::new("orders", 4).expect("build topic config")];
let response = handle_http_raw(
persist,
&topics,
&crate::RetainedLogLag::default(),
&test_committer().await,
&CountCache::default(),
"/dashboard".to_string(),
String::new(),
)
.await
.expect("render dashboard index");
assert!(response
.body
.contains(r#"const TOPICS = [{"name":"orders","partition_count":4}];"#));
}
#[tokio::test]
async fn count_snapshot_is_shared_until_it_expires() -> anyhow::Result<()> {
let persist = test_persist().await;
let stats = test_stats(persist.clone()).await;
let cache = CountCache::default();
let enqueued = |id: &str| {
job(
id,
Stage::Running(RunningStage {
date: chrono::Utc::now(),
}),
)
};
let stages = |value: &Value| value["stages"]["running"].as_u64().unwrap_or(0);
stats.record_transition(&enqueued("one"), None).await;
stats.flush_metrics().await;
let first = cache.get(&persist).await?;
assert_eq!(stages(&first), 1);
stats.record_transition(&enqueued("two"), None).await;
stats.flush_metrics().await;
let within_ttl = cache.get(&persist).await?;
assert_eq!(
stages(&within_ttl),
1,
"a second viewer inside the TTL must reuse the first computation"
);
tokio::time::sleep(COUNT_CACHE_TTL + std::time::Duration::from_millis(100)).await;
assert_eq!(stages(&cache.get(&persist).await?), 2);
Ok(())
}
#[tokio::test]
async fn queue_history_returns_stored_minutes_oldest_first_within_the_range(
) -> anyhow::Result<()> {
let persist = test_persist().await;
let now = chrono::Utc::now();
for (minutes_ago, queued, partitioned) in [(200, 900, 50), (90, 700, 40), (5, 300, 10)] {
persist
.inner
.queue_sample_record(
persist.key_prefix(),
now - chrono::Duration::minutes(minutes_ago),
queued,
partitioned,
)
.await?;
}
let history = |query: &'static str| {
let persist = persist.clone();
async move {
let response = request(persist, query).await;
serde_json::from_str::<Value>(&response.body).expect("decode history")
}
};
let all = history("query=%7B%22cmd%22%3A%22queue_history%22%7D").await;
let queued: Vec<_> = all["points"]
.as_array()
.expect("points")
.iter()
.map(|point| point["queued"].as_i64().expect("queued"))
.collect();
assert_eq!(queued, vec![900, 700, 300]);
assert_eq!(all["points"][2]["partitioned"], json!(10));
let recent =
history("query=%7B%22cmd%22%3A%22queue_history%22%2C%22minutes%22%3A60%7D").await;
assert_eq!(recent["points"].as_array().map(Vec::len), Some(1));
assert_eq!(recent["minutes"], json!(60));
let clamped =
history("query=%7B%22cmd%22%3A%22queue_history%22%2C%22minutes%22%3A0%7D").await;
assert_eq!(clamped["minutes"], json!(1));
Ok(())
}
#[tokio::test]
async fn the_sampler_records_the_current_queue_length() -> anyhow::Result<()> {
use crate::mq::MqClient;
let sqlite = crate::storage::Sqlite::new("sqlite::memory:").await?;
let queue = crate::mq::sql::SqliteQueue::new(sqlite.pool().clone(), "mq")?;
let publisher = queue.new_publisher("sampled").await?;
for number in 0..7 {
publisher
.publish(&crate::encoder::encode(
&crate::models::AmqpCommand::ExecuteJob(crate::models::JobAmqp {
payload_type: "work".to_string(),
id: JobId(format!("job-{number}")),
}),
)?)
.await?;
}
let persist = Arc::new(Persist::new(Box::new(sqlite), "sampled".to_string()));
let stats = test_stats(persist.clone()).await;
stats.record_queue_sample().await;
let samples = persist
.inner
.queue_samples_since("sampled", chrono::Utc::now() - chrono::Duration::minutes(5))
.await?;
assert_eq!(samples.len(), 1);
assert_eq!((samples[0].queued, samples[0].partitioned), (7, 0));
Ok(())
}
}