pub use self::http::{DashboardResponse, ResponseError};
use crate::{
models::{Job, Stage},
persist::Persist,
storage::{JobIndexChange, JobIndexMetricsBatch, JobIndexRow, JobIndexWaitStats},
JobId, UtcDateTime,
};
use std::{collections::HashMap, sync::Arc};
mod http;
const RECENT_WINDOW_SECONDS: usize = 60;
const WORKER_HEARTBEAT_TTL_SECONDS: usize = 15;
const METRIC_BUCKET_MS: i64 = 60_000;
fn bucket_of(at: UtcDateTime) -> i64 {
at.timestamp_millis().div_euclid(METRIC_BUCKET_MS)
}
#[derive(Default)]
struct PendingMetrics {
transitions: HashMap<(String, i64), (UtcDateTime, i64)>,
waits: HashMap<(String, i64), (UtcDateTime, i64, i64, i64)>,
}
impl PendingMetrics {
fn add_transition(&mut self, stage: String, at: UtcDateTime) {
let entry = self
.transitions
.entry((stage, bucket_of(at)))
.or_insert((at, 0));
entry.1 += 1;
}
fn add_wait(&mut self, mode: &str, at: UtcDateTime, wait_ms: i64) {
let entry = self
.waits
.entry((mode.to_string(), bucket_of(at)))
.or_insert((at, 0, 0, 0));
entry.1 += 1;
entry.2 += wait_ms;
entry.3 = entry.3.max(wait_ms);
}
fn into_batch(self) -> JobIndexMetricsBatch {
JobIndexMetricsBatch {
transitions: self
.transitions
.into_iter()
.map(|((stage, _), (at, count))| (stage, at, count))
.collect(),
waits: self
.waits
.into_iter()
.map(|((mode, _), (at, count, sum, max))| (mode, at, count, sum, max))
.collect(),
}
}
fn restore(&mut self, batch: JobIndexMetricsBatch) {
for (stage, at, count) in batch.transitions {
let entry = self
.transitions
.entry((stage, bucket_of(at)))
.or_insert((at, 0));
entry.1 += count;
}
for (mode, at, count, sum, max) in batch.waits {
let entry = self
.waits
.entry((mode, bucket_of(at)))
.or_insert((at, 0, 0, 0));
entry.1 += count;
entry.2 += sum;
entry.3 = entry.3.max(max);
}
}
}
#[derive(Clone)]
pub struct Stats {
queue_backs_enqueued: bool,
count_cache: Arc<http::CountCache>,
pending_metrics: Arc<std::sync::Mutex<PendingMetrics>>,
storage: Arc<Persist>,
topics: Vec<crate::topic::TopicConfig>,
retained_log_lag: crate::RetainedLogLag,
committer: Arc<dyn crate::backend::JobCommitter>,
}
impl Stats {
pub(crate) fn new(
storage: Arc<Persist>,
topics: Vec<crate::topic::TopicConfig>,
retained_log_lag: crate::RetainedLogLag,
committer: Arc<dyn crate::backend::JobCommitter>,
) -> Self {
Self {
queue_backs_enqueued: storage.inner.queue_backs_enqueued_stage(),
count_cache: Arc::default(),
pending_metrics: Arc::default(),
storage,
topics,
retained_log_lag,
committer,
}
}
pub(crate) async fn flush_metrics(&self) {
let pending = match self.pending_metrics.lock() {
Ok(mut pending) => std::mem::take(&mut *pending),
Err(poisoned) => std::mem::take(&mut *poisoned.into_inner()),
};
let batch = pending.into_batch();
if batch.is_empty() {
return;
}
let namespace = self.storage.key_prefix();
if let Err(error) = self
.storage
.inner
.job_index_record_metrics_batch(namespace, &batch)
.await
{
tracing::warn!(%error, "Could not record dashboard metrics; will retry");
let mut pending = match self.pending_metrics.lock() {
Ok(pending) => pending,
Err(poisoned) => poisoned.into_inner(),
};
pending.restore(batch);
}
}
pub(crate) async fn reconcile_stage_counts(&self) {
let namespace = self.storage.key_prefix();
if let Err(error) = self
.storage
.inner
.job_index_reconcile_stage_counts(namespace)
.await
{
tracing::warn!(%error, "Could not reconcile dashboard stage counts");
}
}
pub(crate) async fn purge_unused_rows(&self) {
const BATCH: usize = 2_000;
const PAUSE: std::time::Duration = std::time::Duration::from_millis(250);
let namespace = self.storage.key_prefix();
let mut total = 0usize;
loop {
match self
.storage
.inner
.purge_unused_dashboard_rows(namespace, BATCH)
.await
{
Ok(0) => break,
Ok(removed) => total += removed,
Err(error) => {
tracing::warn!(%error, "Could not purge unused dashboard rows; will retry next start");
break;
}
}
tokio::time::sleep(PAUSE).await;
}
if total > 0 {
tracing::info!(removed = total, "Purged unused dashboard rows");
}
}
pub(crate) async fn record_queue_sample(&self) {
let namespace = self.storage.key_prefix();
let queued = match self.storage.inner.job_index_stage_counts(namespace).await {
Ok(counts) => counts.get("enqueued").copied().unwrap_or(0),
Err(error) => {
tracing::warn!(%error, "Could not read the queue length for the dashboard graph");
return;
}
};
let partitioned = match self.committer.partition_queue_depth_summary().await {
Ok(summary) => summary.iter().map(|(_, _, depth)| *depth as i64).sum(),
Err(error) => {
tracing::warn!(%error, "Could not read partition backlog for the dashboard graph");
0
}
};
if let Err(error) = self
.storage
.inner
.queue_sample_record(
namespace,
chrono::Utc::now(),
i64::try_from(queued).unwrap_or(i64::MAX),
partitioned,
)
.await
{
tracing::warn!(%error, "Could not record a queue length sample");
}
}
pub(crate) async fn record_transition(&self, job: &Job, date_expire: Option<UtcDateTime>) {
let namespace = self.storage.key_prefix();
let result = match self.index_change(job, date_expire) {
Some(JobIndexChange::Upsert(row)) => {
self.storage.inner.job_index_upsert(namespace, row).await
}
Some(JobIndexChange::Remove(job_id)) => {
self.storage
.inner
.job_index_remove(namespace, &job_id)
.await
}
None => Ok(()),
};
if let Err(error) = result {
tracing::warn!(job_id = %job.id, %error, "Could not update dashboard job index");
}
self.buffer_metrics(job);
}
pub(crate) fn index_change(
&self,
job: &Job,
date_expire: Option<UtcDateTime>,
) -> Option<JobIndexChange> {
if self.queue_backs_enqueued
&& job.topic.is_none()
&& matches!(job.stage, Stage::Enqueued(_))
{
return if job.previous_stages.is_empty() {
None
} else {
Some(JobIndexChange::Remove(job.id.to_string()))
};
}
Some(JobIndexChange::Upsert(job_index_row(job, date_expire)))
}
pub(crate) fn buffer_metrics(&self, job: &Job) {
let at = *job.stage.date();
let wait = match (&job.stage, job.previous_stages.last()) {
(Stage::Running(running), Some(previous)) => Some((
running.date,
running
.date
.signed_duration_since(*previous.date())
.num_milliseconds()
.max(0),
)),
_ => None,
};
let mut pending = match self.pending_metrics.lock() {
Ok(pending) => pending,
Err(poisoned) => poisoned.into_inner(),
};
pending.add_transition(job.stage.get_name(), at);
if let Some((running_at, wait_ms)) = wait {
pending.add_wait(wait_mode(job), running_at, wait_ms);
}
}
pub(crate) async fn record_worker_heartbeat(&self, worker_id: String) -> anyhow::Result<()> {
let key = format!("{}-stats-active-workers", self.storage.key_prefix());
let worker_id = crate::encoder::encode(&worker_id)?;
self.storage
.inner
.apply(vec![
crate::storage::StorageOperation::RangeAdd {
key: key.clone(),
value: worker_id.clone(),
},
crate::storage::StorageOperation::RangeExpire {
key,
value: worker_id,
ttl_seconds: WORKER_HEARTBEAT_TTL_SECONDS,
},
])
.await
}
pub(crate) async fn handle_http(
&self,
prefix: String,
query_string: String,
) -> Result<DashboardResponse, ResponseError> {
self.flush_metrics().await;
http::handle_http_raw(
self.storage.clone(),
&self.topics,
&self.retained_log_lag,
&self.committer,
&self.count_cache,
prefix,
query_string,
)
.await
}
pub(crate) async fn truncate_partition_history(
&self,
topic: &str,
partition: u32,
keep_last: usize,
) -> anyhow::Result<usize> {
self.storage
.inner
.job_index_truncate_partition(self.storage.key_prefix(), topic, partition, keep_last)
.await
}
}
pub(crate) fn job_index_row(job: &Job, date_expire: Option<UtcDateTime>) -> JobIndexRow {
let parent_job_id = match &job.stage {
Stage::Waiting(waiting) => Some(waiting.parent_id.to_string()),
Stage::Delayed(_)
| Stage::Enqueued(_)
| Stage::Running(_)
| Stage::Requeued(_)
| Stage::Success(_)
| Stage::Failed(_) => None,
};
let (wait_ms, wait_mode_value) = if let (Stage::Running(running), Some(previous)) =
(&job.stage, job.previous_stages.last())
{
let wait_ms = running
.date
.signed_duration_since(*previous.date())
.num_milliseconds()
.max(0);
(Some(wait_ms), Some(wait_mode(job).to_string()))
} else {
(None, None)
};
let created_at = *job.previous_stages.first().unwrap_or(&job.stage).date();
JobIndexRow {
job_id: job.id.to_string(),
payload_type: job.payload_type.clone(),
stage: job.stage.get_name(),
stage_date: *job.stage.date(),
revision: i64::try_from(job.previous_stages.len()).unwrap_or(i64::MAX),
created_at,
wait_ms,
wait_mode: wait_mode_value,
topic: job.topic.clone(),
partition: job.partition.map(i64::from),
sequence: job.sequence,
parent_job_id,
date_expire,
}
}
fn wait_mode(job: &Job) -> &'static str {
if job.topic.is_some() && job.partition.is_some() && job.sequence.is_some() {
"sequential"
} else {
"regular"
}
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, serde::Serialize)]
pub(crate) struct WaitStats {
pub count: u64,
pub avg_ms: i64,
pub max_ms: i64,
}
impl From<JobIndexWaitStats> for WaitStats {
fn from(stats: JobIndexWaitStats) -> Self {
Self {
count: stats.count as u64,
avg_ms: stats.avg_ms,
max_ms: stats.max_ms,
}
}
}
#[cfg(all(test, feature = "sqlite"))]
mod tests {
use super::*;
use crate::{
models::{DelayedStage, EnqueuedStage, JobConfig, RunningStage, SuccessStage},
storage::Sqlite,
};
fn job(id: &str, stage: Stage) -> Job {
Job {
id: JobId(id.to_string()),
payload_type: "test".to_string(),
payload: Vec::new(),
config: JobConfig::default(),
stage,
previous_stages: Vec::new(),
recurring_job_id: None,
topic: None,
partition: None,
sequence: None,
}
}
struct NoCommitter;
#[async_trait::async_trait]
impl crate::backend::JobCommitter for NoCommitter {
async fn commit(
&self,
_operations: Vec<crate::storage::StorageOperation>,
_queue_name: &str,
_payload: &[u8],
) -> anyhow::Result<()> {
Err(anyhow::anyhow!("NoCommitter never commits"))
}
async fn claim(
&self,
_job_id: &JobId,
) -> anyhow::Result<Option<Box<dyn crate::backend::JobLease>>> {
Ok(None)
}
}
fn test_stats(persist: Arc<Persist>) -> Stats {
Stats::new(
persist,
Vec::new(),
crate::RetainedLogLag::default(),
Arc::new(NoCommitter),
)
}
#[tokio::test]
async fn truncate_partition_history_keeps_only_the_newest_entries() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-partition-history".to_string(),
));
let stats = test_stats(persist.clone());
for sequence in 0..10 {
let mut j = job(
&format!("job-{sequence}"),
Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
);
j.topic = Some("orders".to_string());
j.partition = Some(0);
j.sequence = Some(sequence);
stats.record_transition(&j, None).await;
}
let namespace = persist.key_prefix();
let all = persist
.inner
.job_index_list_by_partition(namespace, "orders", 0, None, 100)
.await?;
assert_eq!(all.items.len(), 10);
let removed = stats.truncate_partition_history("orders", 0, 4).await?;
assert_eq!(removed, 6);
let remaining = persist
.inner
.job_index_list_by_partition(namespace, "orders", 0, None, 100)
.await?
.items;
assert_eq!(remaining.len(), 4);
assert_eq!(
remaining.iter().map(|row| row.sequence).collect::<Vec<_>>(),
vec![Some(9), Some(8), Some(7), Some(6)]
);
assert_eq!(stats.truncate_partition_history("orders", 0, 4).await?, 0);
assert_eq!(stats.truncate_partition_history("orders", 0, 100).await?, 0);
assert_eq!(stats.truncate_partition_history("orders", 1, 0).await?, 0);
Ok(())
}
#[tokio::test]
async fn wait_time_is_split_by_regular_vs_sequential() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-wait-time".to_string(),
));
let stats = test_stats(persist.clone());
let mut partitioned = job(
"partitioned-job",
Stage::Enqueued(EnqueuedStage {
date: chrono::Utc::now(),
}),
);
partitioned.topic = Some("orders".to_string());
partitioned.partition = Some(0);
partitioned.sequence = Some(1);
let enqueued_stage = partitioned.stage.clone();
partitioned.previous_stages.push(enqueued_stage);
partitioned.stage = Stage::Running(RunningStage {
date: chrono::Utc::now(),
});
stats.record_transition(&partitioned, None).await;
stats.flush_metrics().await;
let now = chrono::Utc::now();
let regular = persist
.inner
.job_index_recent_wait_stats(persist.key_prefix(), "regular", now)
.await?;
let sequential = persist
.inner
.job_index_recent_wait_stats(persist.key_prefix(), "sequential", now)
.await?;
assert_eq!(regular.count, 0);
assert_eq!(sequential.count, 1);
Ok(())
}
#[tokio::test]
async fn recent_throughput_counts_transitions_within_the_window() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-throughput".to_string(),
));
let stats = test_stats(persist.clone());
let enqueued = job(
"throughput-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(&enqueued, None).await;
stats.record_transition(&success, None).await;
stats.flush_metrics().await;
let now = chrono::Utc::now();
let succeeded = persist
.inner
.job_index_recent_transition_count(persist.key_prefix(), "success", now)
.await?;
assert_eq!(succeeded, 1);
Ok(())
}
#[tokio::test]
async fn metrics_are_buffered_and_written_together_on_flush() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-batched".to_string(),
));
let stats = test_stats(persist.clone());
let now = chrono::Utc::now();
for number in 0..25 {
let enqueued = job(
&format!("batched-{number}"),
Stage::Enqueued(EnqueuedStage { date: now }),
);
let mut running = enqueued.clone();
running.previous_stages = vec![Stage::Enqueued(EnqueuedStage {
date: now - chrono::Duration::milliseconds(100 + number),
})];
running.stage = Stage::Running(RunningStage { date: now });
stats.record_transition(&running, None).await;
}
let read = |stage: &'static str| {
let persist = persist.clone();
async move {
persist
.inner
.job_index_recent_transition_count(persist.key_prefix(), stage, now)
.await
}
};
assert_eq!(
read("running").await?,
0,
"nothing is written before the flush"
);
stats.flush_metrics().await;
assert_eq!(read("running").await?, 25);
let wait = persist
.inner
.job_index_recent_wait_stats(persist.key_prefix(), "regular", now)
.await?;
assert_eq!((wait.count, wait.avg_ms, wait.max_ms), (25, 112, 124));
stats.flush_metrics().await;
assert_eq!(
read("running").await?,
25,
"a flush with nothing pending adds nothing"
);
Ok(())
}
#[test]
fn a_failed_flush_puts_its_samples_back_without_losing_any() {
let at = chrono::Utc::now();
let mut pending = PendingMetrics::default();
pending.add_transition("success".to_string(), at);
pending.add_wait("regular", at, 40);
let batch = std::mem::take(&mut pending).into_batch();
pending.add_transition("success".to_string(), at);
pending.add_wait("regular", at, 90);
pending.restore(batch);
let batch = pending.into_batch();
assert_eq!(batch.transitions.iter().map(|t| t.2).sum::<i64>(), 2);
let (_, _, count, sum, max) = batch.waits[0].clone();
assert_eq!((count, sum, max), (2, 130, 90));
}
#[tokio::test]
async fn transitions_move_the_live_stage_counters() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-deltas".to_string(),
));
let stats = test_stats(persist.clone());
let now = chrono::Utc::now();
let counts = || {
let persist = persist.clone();
async move {
let all = persist
.inner
.job_index_stage_counts(persist.key_prefix())
.await?;
anyhow::Ok((
all.get("delayed").copied().unwrap_or(0),
all.get("running").copied().unwrap_or(0),
all.get("success").copied().unwrap_or(0),
))
}
};
let delayed = |id: &str| {
job(
id,
Stage::Delayed(DelayedStage {
date: now,
not_before: now + chrono::Duration::minutes(5),
}),
)
};
for id in ["a", "b", "c"] {
stats.record_transition(&delayed(id), None).await;
}
assert_eq!(counts().await?, (3, 0, 0));
let mut running = job("a", Stage::Running(RunningStage { date: now }));
running.previous_stages = vec![delayed("a").stage];
stats.record_transition(&running, None).await;
assert_eq!(counts().await?, (2, 1, 0));
let mut success = running.clone();
success.previous_stages.push(running.stage.clone());
success.stage = Stage::Success(SuccessStage { date: now });
stats
.record_transition(&success, Some(now + chrono::Duration::hours(1)))
.await;
assert_eq!(counts().await?, (2, 0, 1));
Ok(())
}
#[tokio::test]
async fn saving_a_job_twice_in_one_stage_counts_it_once() -> anyhow::Result<()> {
let persist = Arc::new(Persist::new(
Box::new(Sqlite::new("sqlite::memory:").await?),
"dashboard-resave".to_string(),
));
let stats = test_stats(persist.clone());
let delayed = job(
"occurrence",
Stage::Delayed(DelayedStage {
date: chrono::Utc::now(),
not_before: chrono::Utc::now() + chrono::Duration::minutes(5),
}),
);
for _ in 0..4 {
stats.record_transition(&delayed, None).await;
}
let counts = persist
.inner
.job_index_stage_counts(persist.key_prefix())
.await?;
assert_eq!(counts.get("delayed").copied(), Some(1));
Ok(())
}
#[tokio::test]
async fn queued_jobs_are_not_mirrored_in_the_index() -> anyhow::Result<()> {
let sqlite = Sqlite::new("sqlite::memory:").await?;
let pool = sqlite.pool().clone();
let persist = Arc::new(Persist::new(
Box::new(sqlite),
"dashboard-queued".to_string(),
));
let stats = test_stats(persist.clone());
let now = chrono::Utc::now();
let index_rows = || {
let persist = persist.clone();
async move {
let all = persist
.inner
.job_index_list_by_stage(persist.key_prefix(), "running", None, 100)
.await?;
let requeued = persist
.inner
.job_index_stage_counts(persist.key_prefix())
.await?;
anyhow::Ok((
all.items.len(),
requeued.get("requeued").copied().unwrap_or(0),
))
}
};
let total_rows = || {
let pool = pool.clone();
async move {
anyhow::Ok(
sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM later_jobs_index")
.fetch_one(&pool)
.await?,
)
}
};
let fresh = job("a", Stage::Enqueued(EnqueuedStage { date: now }));
assert!(stats.index_change(&fresh, None).is_none());
stats.record_transition(&fresh, None).await;
assert_eq!(total_rows().await?, 0);
let mut running = job("a", Stage::Running(RunningStage { date: now }));
running.previous_stages = vec![fresh.stage.clone()];
stats.record_transition(&running, None).await;
assert_eq!(index_rows().await?, (1, 0));
let mut queued_again = job("a", Stage::Enqueued(EnqueuedStage { date: now }));
queued_again.previous_stages = vec![running.stage.clone()];
assert!(matches!(
stats.index_change(&queued_again, None),
Some(JobIndexChange::Remove(id)) if id == "a"
));
stats.record_transition(&queued_again, None).await;
assert_eq!(total_rows().await?, 0);
let mut partitioned = job("p", Stage::Enqueued(EnqueuedStage { date: now }));
partitioned.topic = Some("orders".to_string());
assert!(matches!(
stats.index_change(&partitioned, None),
Some(JobIndexChange::Upsert(_))
));
Ok(())
}
}