use aion_core::WorkflowId;
use aion_proto::WireError;
use axum::{Json, extract::State};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use super::auth::HttpCaller;
use super::error::HttpWireError;
use crate::ServerState;
#[derive(Clone, Debug, Serialize, Deserialize)]
pub(crate) struct DeadLettersRequest {
pub workflow_id: String,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub(crate) struct DeadLetterRow {
pub dispatch_key: String,
pub ordinal: u64,
pub activity_type: String,
pub namespace: String,
pub task_queue: String,
pub node: Option<String>,
pub attempt: u32,
pub failure_delivered: bool,
pub redrivable: bool,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub(crate) struct DeadLettersResponse {
pub dead_letters: Vec<DeadLetterRow>,
}
pub(crate) async fn list_dead_letters(
State(state): State<ServerState>,
HttpCaller(caller): HttpCaller,
Json(request): Json<DeadLettersRequest>,
) -> Result<Json<DeadLettersResponse>, HttpWireError> {
super::cluster_command::deploy_gate(&caller)?;
let Some(outbox_store) = state.outbox_store() else {
return Err(HttpWireError(WireError::invalid_state(
"the durable outbox is not commissioned on this server, so it has no dead letters",
)));
};
let workflow_id = Uuid::parse_str(&request.workflow_id)
.map(WorkflowId::new)
.map_err(|_error| HttpWireError(WireError::invalid_input("workflow_id must be a UUID")))?;
let rows = crate::worker::list_dead_letters(outbox_store.as_ref(), &workflow_id)
.await
.map_err(|error| HttpWireError(crate::ServerError::from(error).to_wire_error()))?;
Ok(Json(DeadLettersResponse {
dead_letters: rows
.into_iter()
.map(|row| DeadLetterRow {
dispatch_key: row.dispatch_key,
ordinal: row.ordinal,
activity_type: row.activity_type,
namespace: row.namespace,
task_queue: row.task_queue,
node: row.node,
attempt: row.attempt,
failure_delivered: row.failure_delivered,
redrivable: !row.failure_delivered,
})
.collect(),
}))
}