use crate::engine::api::state::ApiState;
use axum::Json;
use axum::extract::{Query, State};
use serde::{Deserialize, Serialize};
#[derive(Deserialize)]
pub struct DlqParams {
pub topic: Option<String>,
pub count: Option<usize>,
}
#[derive(Serialize)]
pub struct DlqMessageResponse {
pub id: String,
pub payload: String,
pub reason: String,
pub original_id: String,
}
pub async fn get_dlq(
State(state): State<ApiState>,
Query(params): Query<DlqParams>,
) -> Json<Vec<DlqMessageResponse>> {
let topic = params.topic.unwrap_or_else(|| "task".to_string());
let count = params.count.unwrap_or(10);
match state.queue_manager.read_dlq(&topic, count).await {
Ok(messages) => {
let response: Vec<DlqMessageResponse> = messages
.into_iter()
.map(|(id, payload, reason, original_id)| {
let payload_str = String::from_utf8(payload)
.unwrap_or_else(|_| "Invalid UTF-8 payload".to_string());
DlqMessageResponse {
id,
payload: payload_str,
reason,
original_id,
}
})
.collect();
Json(response)
}
Err(e) => {
log::error!("Failed to read DLQ: {}", e);
Json(vec![])
}
}
}