qrush 2.1.1

Lightweight Job Queue and Task Scheduler for Rust (Actix/Axum + Redis + Cron)
Documentation
// src/routes/metrics_route.rs
//
// Actix adapter for the dashboard. Each handler extracts request data with
// Actix extractors, delegates to the framework-neutral service functions, and
// converts the returned `DashboardResponse` into an Actix `HttpResponse`.

use actix_web::{web, HttpResponse, HttpResponseBuilder};
use actix_web::http::StatusCode;
use actix_web::http::header::{HeaderName, HeaderValue};

use crate::services::http_response::DashboardResponse;
use crate::services::{
    basic_auth_service::BasicAuthMiddleware,
    cron_service::{self, CronActionRequest},
    metrics_service::{self, MetricsQuery},
};
use crate::utils::pagination::PaginationQuery;

/// Convert the neutral dashboard response into an Actix response.
impl From<DashboardResponse> for HttpResponse {
    fn from(resp: DashboardResponse) -> Self {
        let status = StatusCode::from_u16(resp.status).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
        let mut builder = HttpResponseBuilder::new(status);
        builder.content_type(resp.content_type);
        for (name, value) in &resp.headers {
            if let (Ok(n), Ok(v)) = (
                HeaderName::from_bytes(name.as_bytes()),
                HeaderValue::from_str(value),
            ) {
                builder.append_header((n, v));
            }
        }
        builder.body(resp.body)
    }
}

// --- Handlers -------------------------------------------------------------

async fn render_metrics(query: web::Query<MetricsQuery>) -> HttpResponse {
    metrics_service::render_metrics(query.into_inner()).await.into()
}

async fn render_metrics_for_queue(
    path: web::Path<String>,
    query: web::Query<PaginationQuery>,
) -> HttpResponse {
    metrics_service::render_metrics_for_queue(path.into_inner(), query.into_inner())
        .await
        .into()
}

async fn render_dead_jobs(query: web::Query<PaginationQuery>) -> HttpResponse {
    metrics_service::render_dead_jobs(query.into_inner()).await.into()
}

async fn render_delayed_jobs(query: web::Query<PaginationQuery>) -> HttpResponse {
    metrics_service::render_delayed_jobs(query.into_inner()).await.into()
}

async fn render_scheduled_jobs(query: web::Query<PaginationQuery>) -> HttpResponse {
    metrics_service::render_scheduled_jobs(query.into_inner()).await.into()
}

async fn export_queue_csv(path: web::Path<String>) -> HttpResponse {
    metrics_service::export_queue_csv(path.into_inner()).await.into()
}

async fn get_metrics_summary() -> HttpResponse {
    metrics_service::get_metrics_summary().await.into()
}

async fn job_action(payload: web::Json<serde_json::Value>) -> HttpResponse {
    metrics_service::job_action(payload.into_inner()).await.into()
}

async fn render_failed_jobs(query: web::Query<PaginationQuery>) -> HttpResponse {
    metrics_service::render_failed_jobs(query.into_inner()).await.into()
}

async fn render_retry_jobs(query: web::Query<PaginationQuery>) -> HttpResponse {
    metrics_service::render_retry_jobs(query.into_inner()).await.into()
}

async fn render_cron_jobs() -> HttpResponse {
    cron_service::render_cron_jobs().await.into()
}

async fn cron_action(payload: web::Json<CronActionRequest>) -> HttpResponse {
    cron_service::cron_action(payload.into_inner()).await.into()
}

/// Mount the QRush dashboard routes onto an Actix `ServiceConfig`.
pub fn qrush_metrics_routes(cfg: &mut web::ServiceConfig) {
    cfg.service(
        web::scope("metrics")
            .wrap(BasicAuthMiddleware)
            .route("/health", web::get().to(|| async {
                HttpResponse::Ok().body("healthy")
            }))
            .route("", web::get().to(render_metrics))
            .route("/queues/{queue}", web::get().to(render_metrics_for_queue))
            .route("/extras/dead", web::get().to(render_dead_jobs))
            .route("/extras/delayed", web::get().to(render_delayed_jobs))
            .route("/extras/scheduled", web::get().to(render_scheduled_jobs))
            .route("/queues/{queue}/export", web::get().to(export_queue_csv))
            .route("/extras/summary", web::get().to(get_metrics_summary))
            .route("/jobs/action", web::post().to(job_action))
            .route("/extras/cron", web::get().to(render_cron_jobs))
            .route("/cron/action", web::post().to(cron_action))
            .route("/extras/failed", web::get().to(render_failed_jobs))
            .route("/extras/retry", web::get().to(render_retry_jobs))
    );
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn csv_response_maps_status_content_type_and_headers() {
        let resp: HttpResponse =
            DashboardResponse::csv(b"a,b\n".to_vec(), "queue_x.csv").into();
        assert_eq!(resp.status().as_u16(), 200);
        assert_eq!(
            resp.headers().get("content-type").unwrap(),
            "text/csv"
        );
        assert!(resp
            .headers()
            .get("content-disposition")
            .unwrap()
            .to_str()
            .unwrap()
            .contains("queue_x.csv"));
    }

    #[test]
    fn error_status_is_preserved() {
        let resp: HttpResponse =
            DashboardResponse::json_status(404, &serde_json::json!({"error": "nope"})).into();
        assert_eq!(resp.status().as_u16(), 404);
        assert_eq!(
            resp.headers().get("content-type").unwrap(),
            "application/json"
        );
    }
}