use crate::auth_catalog::build_auth_catalog;
use crate::dlq_replay::{self, ReplayInputs, reader::DlqDecryptor};
use crate::error::CliError;
use crate::serve::error::ServeError;
use crate::serve::load::load_submission;
use crate::serve::rbac::AuthContext;
use crate::serve::runner::ConfigFormatWire;
use crate::serve::state::ServerState;
use axum::Json;
use axum::extract::{Extension, State};
use chrono::Utc;
use serde::Deserialize;
use serde_json::Value;
fn cli_to_serve(e: CliError) -> ServeError {
match e {
CliError::Config(m) => ServeError::BadConfig(m),
other => ServeError::Internal(other.to_string()),
}
}
fn default_sample_limit() -> usize {
5
}
#[derive(Debug, Deserialize)]
pub struct DlqInspectRequest {
pub location: String,
#[serde(default)]
pub reason: Option<String>,
#[serde(default = "default_sample_limit")]
pub limit: usize,
#[serde(default)]
pub encryption_keys: Vec<String>,
}
pub async fn inspect(
State(_state): State<ServerState>,
Json(req): Json<DlqInspectRequest>,
) -> Result<Json<Value>, ServeError> {
let location = req.location;
let reason = req.reason;
let limit = req.limit;
let dec = DlqDecryptor::from_keys(&req.encryption_keys)
.map_err(CliError::from)
.map_err(cli_to_serve)?;
let summary = tokio::task::spawn_blocking(move || {
dlq_replay::inspect(&location, reason.as_deref(), limit, &dec)
})
.await
.map_err(|e| ServeError::Internal(format!("dlq inspect task: {e}")))?
.map_err(cli_to_serve)?;
Ok(Json(serde_json::to_value(summary).unwrap_or(Value::Null)))
}
#[derive(Debug, Deserialize)]
pub struct DlqReplayRequest {
pub config: String,
#[serde(default)]
pub config_format: ConfigFormatWire,
pub from: String,
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub failed_dlq: Option<String>,
#[serde(default)]
pub row: Option<String>,
#[serde(default)]
pub dry_run: bool,
#[serde(default)]
pub encryption_keys: Vec<String>,
}
pub async fn replay(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Json(req): Json<DlqReplayRequest>,
) -> Result<Json<Value>, ServeError> {
let loaded = load_submission(
&req.config,
req.config_format.into(),
state.default_base().as_ref(),
)
.await?;
let auth = build_auth_catalog(loaded.cfg.auth.as_ref()).map_err(cli_to_serve)?;
let pipeline_name = loaded
.cfg
.name
.clone()
.unwrap_or_else(|| "dlq-replay".to_string());
let outcome = dlq_replay::replay(
&loaded.cfg,
&req.from,
ReplayInputs {
reason: req.reason.as_deref(),
failed_dlq: req.failed_dlq.as_deref(),
row: req.row.as_deref(),
dry_run: req.dry_run,
pipeline_name,
execution: loaded.cfg.execution.clone(),
auth,
clock: Utc::now().fixed_offset(),
decryptor: DlqDecryptor::from_keys(&req.encryption_keys)
.map_err(CliError::from)
.map_err(cli_to_serve)?,
},
)
.await;
let result_label = if outcome.is_ok() { "ok" } else { "error" };
crate::serve::audit::write(&state, &actor, "dlq.replay", None, None, result_label).await;
let outcome = outcome.map_err(cli_to_serve)?;
Ok(Json(serde_json::to_value(outcome).unwrap_or(Value::Null)))
}
#[derive(Debug, Deserialize)]
pub struct DlqDiscardRequest {
pub location: String,
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub before_ms: Option<i64>,
#[serde(default)]
pub delete: bool,
#[serde(default)]
pub encryption_keys: Vec<String>,
}
pub async fn discard(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Json(req): Json<DlqDiscardRequest>,
) -> Result<Json<Value>, ServeError> {
let location = req.location;
let reason = req.reason;
let before_ms = req.before_ms;
let delete = req.delete;
let dec = DlqDecryptor::from_keys(&req.encryption_keys)
.map_err(CliError::from)
.map_err(cli_to_serve)?;
let outcome = tokio::task::spawn_blocking(move || {
dlq_replay::discard(&location, reason.as_deref(), before_ms, delete, &dec)
})
.await
.map_err(|e| ServeError::Internal(format!("dlq discard task: {e}")));
let result_label = if matches!(&outcome, Ok(Ok(_))) {
"ok"
} else {
"error"
};
crate::serve::audit::write(&state, &actor, "dlq.discard", None, None, result_label).await;
let outcome = outcome?.map_err(cli_to_serve)?;
Ok(Json(serde_json::to_value(outcome).unwrap_or(Value::Null)))
}