later 0.0.48

Distributed Background jobs manager and runner for Rust
//! The dashboard's job-listing/counting index - see `later_jobs_index` in
//! the bundled migrations for the full rationale (one row per job,
//! upserted in place on every stage transition, replacing a family of
//! independently-maintained `later_storage_range` keys that could silently
//! drift out of sync with each other).

use crate::UtcDateTime;

/// One job's current state, as tracked for dashboard listing/counting.
///
/// Not the job's full durable record (payload/config/full stage history) -
/// that still lives in `later_storage`, keyed by job ID, and is fetched
/// separately for a job's detail view. This is only ever a lightweight
/// projection used to answer "which jobs are in stage X" without scanning
/// or trusting a separate list.
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct JobIndexRow {
    /// The job's own ID.
    pub job_id: String,
    /// The job's payload type name.
    pub payload_type: String,
    /// Current stage name (e.g. `"enqueued"`, `"running"`, `"success"`).
    pub stage: String,
    /// When the job entered its current stage.
    pub stage_date: UtcDateTime,
    /// `previous_stages.len()` at write time - a monotonic per-job counter.
    /// A write is only applied if its revision is `>=` whatever is already
    /// stored, so two writes for the same job that arrive out of order can
    /// never let an older one clobber a newer one.
    pub revision: i64,
    /// When the job's very first stage was recorded - fixed at insert,
    /// never updated by a later transition. Backs the listing view's
    /// total-duration column.
    pub created_at: UtcDateTime,
    /// Time this job spent queued before its most recent `Running`
    /// transition, in milliseconds. Set once, on that transition, and left
    /// untouched by every later transition of the same job - a query for
    /// "recent wait samples" filters by when this was recorded, not by
    /// current stage, since a job that has since finished must still count
    /// toward the window it was actually sampled in.
    pub wait_ms: Option<i64>,
    /// "regular" or "sequential" - see `later::stats`'s `wait_mode`.
    pub wait_mode: Option<String>,
    /// Topic this job orders under, when enqueued with partition ordering.
    pub topic: Option<String>,
    /// Partition within `topic` that must process this job in order.
    pub partition: Option<i64>,
    /// Monotonic sequence assigned within `(topic, partition)` at enqueue.
    pub sequence: Option<i64>,
    /// Non-NULL only while this job is `Stage::Waiting` on another job.
    pub parent_job_id: Option<String>,
    /// When this row stops being visible to reads, if it has an expiry.
    pub date_expire: Option<UtcDateTime>,
}

/// A bounded, newest-first page of [`JobIndexRow`]s.
#[derive(Debug, Clone, Default)]
pub struct JobIndexPage {
    /// Rows in the requested order, at most the requested limit.
    pub items: Vec<JobIndexRow>,
    /// Pass back to the next call for the following page. `None` means
    /// there are no older rows left.
    pub next_cursor: Option<i64>,
}

/// Aggregate queue-wait-time stats over a recent window - see
/// [`crate::storage::Storage::job_index_recent_wait_stats`].
#[derive(Debug, Clone, Copy, Default)]
pub struct JobIndexWaitStats {
    /// Number of samples the average/max below are computed from.
    pub count: usize,
    /// Mean wait time across the window, in milliseconds.
    pub avg_ms: i64,
    /// Longest wait time across the window, in milliseconds.
    pub max_ms: i64,
}

/// Dashboard throughput/wait samples already summed per minute bucket by the
/// caller, written in one call - see [`crate::storage::Storage::job_index_record_metrics_batch`].
///
/// Each job transition used to write its own counter row, which made the
/// dashboard's own bookkeeping up to four of a job's nineteen statements,
/// all updating the same few hot rows.
#[derive(Debug, Clone, Default)]
pub struct JobIndexMetricsBatch {
    /// `(stage, an instant inside the bucket, transitions)`.
    pub transitions: Vec<(String, UtcDateTime, i64)>,
    /// `(mode, an instant inside the bucket, samples, sum_ms, max_ms)`.
    pub waits: Vec<(String, UtcDateTime, i64, i64, i64)>,
}

impl JobIndexMetricsBatch {
    /// Whether there is nothing to write.
    pub fn is_empty(&self) -> bool {
        self.transitions.is_empty() && self.waits.is_empty()
    }
}

/// What a job transition does to the dashboard job index.
#[derive(Debug, Clone)]
pub enum JobIndexChange {
    /// Write the job's current row.
    Upsert(JobIndexRow),
    /// Delete the job's row: it moved to a stage the index does not hold.
    Remove(String),
}

/// One per-minute reading of how many jobs were waiting - see
/// [`crate::storage::Storage::queue_sample_record`].
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct QueueSample {
    /// Start of the minute the reading belongs to.
    pub at: UtcDateTime,
    /// Jobs waiting in the shared queue for a worker.
    pub queued: i64,
    /// Jobs waiting across every topic partition.
    pub partitioned: i64,
}