Something went wrong. Try again.
A lexicon-driven AppView for ATProto.
Something went wrong. Try again.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598use axum::{ Json, extract::{Path, Query, State},};use serde::{Deserialize, Serialize};use serde_json::Value;
use super::auth::UserAuth;use super::permissions::Permission;use crate::AppState;use crate::db::{adapt_sql, now_rfc3339, parse_dt};use crate::error::AppError;use crate::lua::{HookEvent, run_hook_once};use crate::record_handler::RecordEvent;
// ---------------------------------------------------------------------------// Query / request / response types// ---------------------------------------------------------------------------
#[derive(Deserialize)]pub struct ListQuery { pub collection: Option<String>, pub resolved: Option<String>, pub cursor: Option<String>, pub limit: Option<i64>,}
#[derive(Deserialize)]pub struct CountQuery { pub resolved: Option<String>,}
#[derive(Serialize)]pub struct DeadLetterSummary { pub id: String, pub lexicon_id: String, pub uri: String, pub did: String, pub collection: String, pub rkey: String, pub action: String, pub error: String, pub attempts: i64, pub created_at: chrono::DateTime<chrono::Utc>, #[serde(skip_serializing_if = "Option::is_none")] pub resolved_at: Option<chrono::DateTime<chrono::Utc>>,}
#[derive(Serialize)]pub struct DeadLetterDetail { #[serde(flatten)] pub summary: DeadLetterSummary, pub record: Option<Value>,}
#[derive(Serialize)]pub struct ListResponse { pub dead_letters: Vec<DeadLetterSummary>, #[serde(skip_serializing_if = "Option::is_none")] pub cursor: Option<String>,}
#[derive(Serialize)]pub struct CountResponse { pub count: i64,}
#[derive(Deserialize)]pub struct BulkRequest { pub ids: Option<Vec<String>>, pub all: Option<bool>, pub collection: Option<String>,}
/// Internal row type for fetching action data needed by retry/reindex.#[allow(dead_code)]struct DeadLetterRow { id: String, lexicon_id: String, uri: String, did: String, collection: String, rkey: String, action: String, record: Option<String>, error: String, attempts: i64,}
// ---------------------------------------------------------------------------// Handlers// ---------------------------------------------------------------------------
/// GET /admin/dead-letterspub(super) async fn list( auth: UserAuth, State(state): State<AppState>, Query(query): Query<ListQuery>,) -> Result<Json<ListResponse>, AppError> { auth.require(Permission::DeadLettersRead).await?; let backend = state.db_backend; let limit = query.limit.unwrap_or(50).clamp(1, 100);
let mut sql = String::from( "SELECT id, lexicon_id, uri, did, collection, rkey, action, error, attempts, created_at, resolved_at FROM dead_letter_hooks WHERE 1=1", );
let resolved_filter = query.resolved.as_deref().unwrap_or("false"); match resolved_filter { "false" => sql.push_str(" AND resolved_at IS NULL"), "true" => sql.push_str(" AND resolved_at IS NOT NULL"), _ => {} // no filter }
if query.collection.is_some() { sql.push_str(" AND collection = ?"); } if query.cursor.is_some() { sql.push_str(" AND created_at < ?"); }
sql.push_str(" ORDER BY created_at DESC LIMIT ?");
let sql = adapt_sql(&sql, backend);
#[allow(clippy::type_complexity)] let mut q = sqlx::query_as::< _, ( String, String, String, String, String, String, String, String, i64, String, Option<String>, ), >(&sql);
if let Some(ref collection) = query.collection { q = q.bind(collection); } if let Some(ref cursor) = query.cursor { q = q.bind(cursor); } q = q.bind(limit);
let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to query dead letters: {e}")))?;
let dead_letters: Vec<DeadLetterSummary> = rows .into_iter() .map(|row| DeadLetterSummary { id: row.0, lexicon_id: row.1, uri: row.2, did: row.3, collection: row.4, rkey: row.5, action: row.6, error: row.7, attempts: row.8, created_at: parse_dt(&row.9), resolved_at: row.10.as_deref().map(parse_dt), }) .collect();
let cursor = if dead_letters.len() as i64 >= limit { dead_letters.last().map(|dl| dl.created_at.to_rfc3339()) } else { None };
Ok(Json(ListResponse { dead_letters, cursor, }))}
/// GET /admin/dead-letters/countpub(super) async fn count( auth: UserAuth, State(state): State<AppState>, Query(query): Query<CountQuery>,) -> Result<Json<CountResponse>, AppError> { auth.require(Permission::DeadLettersRead).await?; let backend = state.db_backend;
let mut sql = String::from("SELECT COUNT(*) FROM dead_letter_hooks WHERE 1=1");
let resolved_filter = query.resolved.as_deref().unwrap_or("false"); match resolved_filter { "false" => sql.push_str(" AND resolved_at IS NULL"), "true" => sql.push_str(" AND resolved_at IS NOT NULL"), _ => {} }
let sql = adapt_sql(&sql, backend); let (count,): (i64,) = sqlx::query_as(&sql) .fetch_one(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to count dead letters: {e}")))?;
Ok(Json(CountResponse { count }))}
/// GET /admin/dead-letters/{id}pub(super) async fn detail( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<DeadLetterDetail>, AppError> { auth.require(Permission::DeadLettersRead).await?; let backend = state.db_backend;
let sql = adapt_sql( "SELECT id, lexicon_id, uri, did, collection, rkey, action, error, attempts, created_at, resolved_at, record FROM dead_letter_hooks WHERE id = ?", backend, );
#[allow(clippy::type_complexity)] let row: ( String, String, String, String, String, String, String, String, i64, String, Option<String>, Option<String>, ) = sqlx::query_as(&sql) .bind(&id) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to fetch dead letter: {e}")))? .ok_or_else(|| AppError::NotFound(format!("dead letter {id} not found")))?;
let summary = DeadLetterSummary { id: row.0, lexicon_id: row.1, uri: row.2, did: row.3, collection: row.4, rkey: row.5, action: row.6, error: row.7, attempts: row.8, created_at: parse_dt(&row.9), resolved_at: row.10.as_deref().map(parse_dt), };
let record = row.11.as_deref().and_then(|r| serde_json::from_str(r).ok());
Ok(Json(DeadLetterDetail { summary, record }))}
/// POST /admin/dead-letters/{id}/dismisspub(super) async fn dismiss( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; let dl = fetch_dead_letter_for_action(&state, &id).await?; if dl.id.is_empty() { return Err(AppError::NotFound(format!("dead letter {id} not found"))); } mark_resolved(&state, &id).await?; Ok(Json(serde_json::json!({ "ok": true })))}
/// POST /admin/dead-letters/{id}/retrypub(super) async fn retry( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; retry_single(&state, &id).await?; Ok(Json(serde_json::json!({ "ok": true })))}
/// POST /admin/dead-letters/{id}/reindexpub(super) async fn reindex( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; reindex_single(&state, &id).await?; Ok(Json(serde_json::json!({ "ok": true })))}
/// POST /admin/dead-letters/bulk/dismisspub(super) async fn bulk_dismiss( auth: UserAuth, State(state): State<AppState>, Json(body): Json<BulkRequest>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; let backend = state.db_backend; let now = now_rfc3339();
if body.all == Some(true) { let mut sql = String::from("UPDATE dead_letter_hooks SET resolved_at = ? WHERE resolved_at IS NULL"); if body.collection.is_some() { sql.push_str(" AND collection = ?"); } let sql = adapt_sql(&sql, backend); let mut q = sqlx::query(&sql).bind(&now); if let Some(ref collection) = body.collection { q = q.bind(collection); } q.execute(&state.db) .await .map_err(|e| AppError::Internal(format!("bulk dismiss failed: {e}")))?; } else if let Some(ref ids) = body.ids { for id in ids { let sql = adapt_sql( "UPDATE dead_letter_hooks SET resolved_at = ? WHERE id = ? AND resolved_at IS NULL", backend, ); sqlx::query(&sql) .bind(&now) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("bulk dismiss failed for {id}: {e}")))?; } } else { return Err(AppError::BadRequest( "must provide 'ids' or 'all: true'".into(), )); }
Ok(Json(serde_json::json!({ "ok": true })))}
/// POST /admin/dead-letters/bulk/retrypub(super) async fn bulk_retry( auth: UserAuth, State(state): State<AppState>, Json(body): Json<BulkRequest>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; let ids = resolve_bulk_ids(&state, &body).await?; for id in &ids { retry_single(&state, id).await?; } Ok(Json(serde_json::json!({ "ok": true })))}
/// POST /admin/dead-letters/bulk/reindexpub(super) async fn bulk_reindex( auth: UserAuth, State(state): State<AppState>, Json(body): Json<BulkRequest>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; let ids = resolve_bulk_ids(&state, &body).await?; for id in &ids { reindex_single(&state, id).await?; } Ok(Json(serde_json::json!({ "ok": true })))}
// ---------------------------------------------------------------------------// Helper functions// ---------------------------------------------------------------------------
/// Fetch an unresolved dead letter by ID, returning an error if not found or already resolved.async fn fetch_dead_letter_for_action( state: &AppState, id: &str,) -> Result<DeadLetterRow, AppError> { let backend = state.db_backend; let sql = adapt_sql( "SELECT id, lexicon_id, uri, did, collection, rkey, action, record, error, attempts FROM dead_letter_hooks WHERE id = ? AND resolved_at IS NULL", backend, );
let row: ( String, String, String, String, String, String, String, Option<String>, String, i64, ) = sqlx::query_as(&sql) .bind(id) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to fetch dead letter: {e}")))? .ok_or_else(|| { AppError::NotFound(format!("dead letter {id} not found or already resolved")) })?;
Ok(DeadLetterRow { id: row.0, lexicon_id: row.1, uri: row.2, did: row.3, collection: row.4, rkey: row.5, action: row.6, record: row.7, error: row.8, attempts: row.9, })}
/// Mark a dead letter as resolved.async fn mark_resolved(state: &AppState, id: &str) -> Result<(), AppError> { let backend = state.db_backend; let now = now_rfc3339(); let sql = adapt_sql( "UPDATE dead_letter_hooks SET resolved_at = ? WHERE id = ?", backend, ); sqlx::query(&sql) .bind(&now) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to mark dead letter resolved: {e}")))?; Ok(())}
/// Update the error message and increment attempts.async fn update_error(state: &AppState, id: &str, error: &str) -> Result<(), AppError> { let backend = state.db_backend; let sql = adapt_sql( "UPDATE dead_letter_hooks SET error = ?, attempts = attempts + 1 WHERE id = ?", backend, ); sqlx::query(&sql) .bind(error) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to update dead letter error: {e}")))?; Ok(())}
/// Resolve a BulkRequest into a list of dead letter IDs.async fn resolve_bulk_ids(state: &AppState, body: &BulkRequest) -> Result<Vec<String>, AppError> { if body.all == Some(true) { let backend = state.db_backend; let mut sql = String::from("SELECT id FROM dead_letter_hooks WHERE resolved_at IS NULL"); if body.collection.is_some() { sql.push_str(" AND collection = ?"); } let sql = adapt_sql(&sql, backend); let mut q = sqlx::query_as::<_, (String,)>(&sql); if let Some(ref collection) = body.collection { q = q.bind(collection); } let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to resolve bulk ids: {e}")))?; Ok(rows.into_iter().map(|r| r.0).collect()) } else if let Some(ref ids) = body.ids { Ok(ids.clone()) } else { Err(AppError::BadRequest( "must provide 'ids' or 'all: true'".into(), )) }}
/// Fetch the index_hook script directly from the lexicons table, bypassing the in-memory registry.async fn get_index_hook_from_db( state: &AppState, lexicon_id: &str,) -> Result<Option<String>, AppError> { let backend = state.db_backend; let sql = adapt_sql("SELECT index_hook FROM lexicons WHERE id = ?", backend); let row: Option<(Option<String>,)> = sqlx::query_as(&sql) .bind(lexicon_id) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to fetch index hook: {e}")))?; Ok(row.and_then(|r| r.0))}
/// Retry a single dead letter by re-running its hook script.async fn retry_single(state: &AppState, id: &str) -> Result<(), AppError> { let dl = fetch_dead_letter_for_action(state, id).await?;
let script = get_index_hook_from_db(state, &dl.lexicon_id) .await? .ok_or_else(|| { AppError::NotFound(format!("no index hook found for lexicon {}", dl.lexicon_id)) })?;
let record: Option<Value> = dl .record .as_deref() .and_then(|r| serde_json::from_str(r).ok());
let event = HookEvent { state, lexicon_id: &dl.lexicon_id, script: &script, action: &dl.action, uri: &dl.uri, did: &dl.did, collection: &dl.collection, rkey: &dl.rkey, record: record.as_ref(), };
match run_hook_once(&event).await { Ok(_) => { mark_resolved(state, id).await?; Ok(()) } Err(e) => { update_error(state, id, &e).await?; Err(AppError::Internal(format!( "retry failed for dead letter {id}: {e}" ))) } }}
/// Reindex a single dead letter by fetching the record fresh from the PDS.async fn reindex_single(state: &AppState, id: &str) -> Result<(), AppError> { let dl = fetch_dead_letter_for_action(state, id).await?;
let pds_endpoint = crate::profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, &dl.did).await?;
let url = format!( "{}/xrpc/com.atproto.repo.getRecord?repo={}&collection={}&rkey={}", pds_endpoint, dl.did, dl.collection, dl.rkey );
let resp = state .http .get(&url) .send() .await .map_err(|e| AppError::Internal(format!("failed to fetch record from PDS: {e}")))?;
if !resp.status().is_success() { let status = resp.status(); let body = resp.text().await.unwrap_or_else(|_| "unknown error".into()); return Err(AppError::Internal(format!( "PDS returned {status} fetching record: {body}" ))); }
let body: Value = resp .json() .await .map_err(|e| AppError::Internal(format!("failed to parse PDS response: {e}")))?;
let record = body.get("value").cloned(); let cid = body .get("cid") .and_then(|v| v.as_str()) .map(|s| s.to_string());
let event = RecordEvent { did: dl.did.clone(), collection: dl.collection.clone(), rkey: dl.rkey.clone(), action: dl.action.clone(), record, cid, };
crate::record_handler::handle_record_event(state, &event).await; mark_resolved(state, id).await?;
Ok(())}