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;
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)
}
}
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()
}
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"
);
}
}