Something went wrong. Try again.
A lexicon-driven AppView for ATProto.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989//! Admin surface for dead-lettered events.//!//! Two tables back this://!//! - **`dead_letter_hooks`** (legacy) — written by the pre-trigger-keyed//! indexer when a hook script exhausted retries. Columns are//! per-event-field (lexicon_id, uri, did, collection, rkey, action,//! record). UUID / TEXT primary keys.//! - **`dead_letter_scripts`** (current) — written by the trigger-keyed//! dispatcher in `crate::lua::scripts`. The event-specific fields are//! inside `payload` (JSON). INTEGER primary keys. Carries both record//! and label dead letters via the `host_kind` discriminator.//!//! Both tables are kept readable + manageable through this admin//! surface. Per-id operations route by id format: an id that parses as//! an integer routes to `dead_letter_scripts`, anything else (UUIDs,//! sqlite NULL-stringified primary keys) routes to `dead_letter_hooks`.//! The two id namespaces are disjoint so this dispatch is unambiguous.
use 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::{RecordEventPayload, resolve_record_event, run_record_event_once};use crate::record_handler::RecordEvent;
// ---------------------------------------------------------------------------// Source enum — which table backs a given dead-letter id// ---------------------------------------------------------------------------
#[derive(Clone, Copy, Debug, PartialEq, Eq)]enum DeadLetterSource { /// Pre-trigger-keyed legacy hooks table. LegacyHooks, /// New trigger-keyed scripts table. Scripts,}
impl DeadLetterSource { /// Pick the table by id format. Integer-parseable → Scripts; /// anything else → LegacyHooks. fn from_id(id: &str) -> Self { if id.parse::<i64>().is_ok() { Self::Scripts } else { Self::LegacyHooks } }
fn table(self) -> &'static str { match self { Self::LegacyHooks => "happyview_dead_letter_hooks", Self::Scripts => "happyview_dead_letter_scripts", } }}
// ---------------------------------------------------------------------------// 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 collection: Option<String>, 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./// Populated from either table by `fetch_dead_letter_for_action`.struct DeadLetterRow { id: String, source: DeadLetterSource, /// Discriminator for new-table rows: `"record"` or `"label"`. /// Always `"record"` for legacy rows. Retries are only supported /// for record dead letters. host_kind: String, uri: String, did: String, collection: String, rkey: String, action: String, record: Option<String>,}
// ---------------------------------------------------------------------------// Handlers// ---------------------------------------------------------------------------
/// `GET /admin/dead-letters` — list rows from both tables, merge by/// created_at, paginate via cursor.pub(super) async fn list( auth: UserAuth, State(state): State<AppState>, Query(query): Query<ListQuery>,) -> Result<Json<ListResponse>, AppError> { auth.require(Permission::DeadLettersRead).await?; let limit = query.limit.unwrap_or(50).clamp(1, 100); let resolved = query.resolved.as_deref().unwrap_or("false");
// Fetch up to `limit` rows from each table, then merge + slice. // Two queries instead of a SQL UNION because the schemas differ // (legacy has columns; scripts has payload JSON we parse in Rust). let mut rows = list_legacy( &state, resolved, query.collection.as_deref(), &query.cursor, limit, ) .await?; rows.extend( list_scripts( &state, resolved, query.collection.as_deref(), &query.cursor, limit, ) .await?, );
// Newest first, then truncate. rows.sort_by_key(|r| std::cmp::Reverse(r.created_at)); let truncated = rows.len() as i64 > limit; rows.truncate(limit as usize);
let cursor = if truncated { rows.last().map(|r| r.created_at.to_rfc3339()) } else { None };
Ok(Json(ListResponse { dead_letters: rows, cursor, }))}
/// `GET /admin/dead-letters/count` — sum of unresolved across both tables.pub(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 resolved = query.resolved.as_deref().unwrap_or("false"); let resolved_clause = match resolved { "false" => " AND resolved_at IS NULL", "true" => " AND resolved_at IS NOT NULL", _ => "", };
let collection_clause = if query.collection.is_some() { " AND collection = ?" } else { "" };
let mut total: i64 = 0; for table in [ DeadLetterSource::LegacyHooks.table(), DeadLetterSource::Scripts.table(), ] { let sql = adapt_sql( &format!("SELECT COUNT(*) FROM {table} WHERE 1=1{resolved_clause}{collection_clause}"), backend, ); let mut q = crate::db::query_as::<(i64,)>(&sql); if let Some(ref c) = query.collection { q = q.bind(c); } let (n,) = q .fetch_one(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to count dead letters: {e}")))?; total += n; }
Ok(Json(CountResponse { count: total }))}
/// `GET /admin/dead-letters/{id}` — detail view. Routes by id format.pub(super) async fn detail( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<DeadLetterDetail>, AppError> { auth.require(Permission::DeadLettersRead).await?; match DeadLetterSource::from_id(&id) { DeadLetterSource::LegacyHooks => detail_legacy(&state, &id).await.map(Json), DeadLetterSource::Scripts => detail_scripts(&state, &id).await.map(Json), }}
/// `POST /admin/dead-letters/{id}/dismiss`pub(super) async fn dismiss( auth: UserAuth, State(state): State<AppState>, Path(id): Path<String>,) -> Result<Json<Value>, AppError> { auth.require(Permission::DeadLettersManage).await?; let source = DeadLetterSource::from_id(&id); mark_resolved(&state, &id, source).await?; Ok(Json(serde_json::json!({ "ok": true })))}
/// `POST /admin/dead-letters/{id}/retry`pub(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}/reindex`pub(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/dismiss`pub(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) { for source in [DeadLetterSource::LegacyHooks, DeadLetterSource::Scripts] { let table = source.table(); let mut sql = format!("UPDATE {table} 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 = crate::db::query(&sql).bind(&now); if let Some(ref c) = body.collection { q = q.bind(c); } 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 { mark_resolved(&state, id, DeadLetterSource::from_id(id)).await?; } } else { return Err(AppError::BadRequest( "must provide 'ids' or 'all: true'".into(), )); }
Ok(Json(serde_json::json!({ "ok": true })))}
/// `POST /admin/dead-letters/bulk/retry`pub(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/reindex`pub(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 })))}
// ---------------------------------------------------------------------------// Per-table list / detail// ---------------------------------------------------------------------------
async fn list_legacy( state: &AppState, resolved: &str, collection: Option<&str>, cursor: &Option<String>, limit: i64,) -> Result<Vec<DeadLetterSummary>, AppError> { let backend = state.db_backend; let mut sql = String::from( "SELECT id, lexicon_id, uri, did, collection, rkey, action, error, attempts, created_at, resolved_at FROM happyview_dead_letter_hooks WHERE 1=1", ); match resolved { "false" => sql.push_str(" AND resolved_at IS NULL"), "true" => sql.push_str(" AND resolved_at IS NOT NULL"), _ => {} } if collection.is_some() { sql.push_str(" AND collection = ?"); } if 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 = crate::db::query_as::<( String, String, String, String, String, String, String, String, i64, String, Option<String>, )>(&sql); if let Some(c) = collection { q = q.bind(c); } if let Some(cur) = cursor { q = q.bind(cur); } q = q.bind(limit);
let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to query legacy dead letters: {e}")))?;
Ok(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())}
async fn list_scripts( state: &AppState, resolved: &str, collection: Option<&str>, cursor: &Option<String>, limit: i64,) -> Result<Vec<DeadLetterSummary>, AppError> { let backend = state.db_backend; let mut sql = String::from( "SELECT id, script_ref, host_kind, host_id, payload, error, attempts, created_at, resolved_at FROM happyview_dead_letter_scripts WHERE 1=1", ); match resolved { "false" => sql.push_str(" AND resolved_at IS NULL"), "true" => sql.push_str(" AND resolved_at IS NOT NULL"), _ => {} } if collection.is_some() { sql.push_str(" AND collection = ?"); } if 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 = crate::db::query_as::<( i64, String, String, String, String, String, i64, String, Option<String>, )>(&sql); if let Some(c) = collection { q = q.bind(c); } if let Some(cur) = cursor { q = q.bind(cur); } q = q.bind(limit);
let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to query scripts dead letters: {e}")))?;
Ok(rows .into_iter() .map(|row| { summary_from_scripts_row(&ScriptsDeadLetterRow { id: row.0.to_string(), script_ref: row.1, host_kind: row.2, host_id: row.3, payload: row.4, error: row.5, attempts: row.6, created_at: row.7, resolved_at: row.8, }) }) .collect())}
struct ScriptsDeadLetterRow { id: String, script_ref: String, host_kind: String, host_id: String, payload: String, error: String, attempts: i64, created_at: String, resolved_at: Option<String>,}
fn summary_from_scripts_row(row: &ScriptsDeadLetterRow) -> DeadLetterSummary { let ScriptsDeadLetterRow { id, script_ref, host_kind, host_id, payload, error, attempts, created_at, resolved_at, } = row; let payload_v: Value = serde_json::from_str(payload).unwrap_or(Value::Null); let s = |key: &str| { payload_v .get(key) .and_then(|v| v.as_str()) .unwrap_or_default() .to_string() };
let (lexicon_id, collection, did, rkey, action) = match host_kind.as_str() { "label" => { let collection_from_trigger = script_ref .split_once(':') .map(|(_, suf)| suf.to_string()) .unwrap_or_default(); ( script_ref.clone(), collection_from_trigger, host_id.clone(), s("val"), "label".to_string(), ) } _ => ( s("collection"), s("collection"), s("did"), s("rkey"), s("action"), ), };
DeadLetterSummary { id: id.clone(), lexicon_id, uri: s("uri"), did, collection, rkey, action, error: error.clone(), attempts: *attempts, created_at: parse_dt(created_at), resolved_at: resolved_at.as_deref().map(parse_dt), }}
async fn detail_legacy(state: &AppState, id: &str) -> Result<DeadLetterDetail, AppError> { 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 happyview_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>, ) = crate::db::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(DeadLetterDetail { summary, record })}
async fn detail_scripts(state: &AppState, id: &str) -> Result<DeadLetterDetail, AppError> { let backend = state.db_backend; let sql = adapt_sql( "SELECT id, script_ref, host_kind, host_id, payload, error, attempts, created_at, resolved_at FROM happyview_dead_letter_scripts WHERE id = ?", backend, ); let id_int: i64 = id .parse() .map_err(|_| AppError::NotFound(format!("dead letter {id} not found")))?;
#[allow(clippy::type_complexity)] let row: ( i64, String, String, String, String, String, i64, String, Option<String>, ) = crate::db::query_as(&sql) .bind(id_int) .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 scripts_row = ScriptsDeadLetterRow { id: row.0.to_string(), script_ref: row.1, host_kind: row.2, host_id: row.3, payload: row.4, error: row.5, attempts: row.6, created_at: row.7, resolved_at: row.8, };
let payload_v: Value = serde_json::from_str(&scripts_row.payload).unwrap_or(Value::Null); let record = if scripts_row.host_kind == "label" { Some(payload_v.clone()) } else { payload_v.get("record").cloned() };
let summary = summary_from_scripts_row(&scripts_row);
Ok(DeadLetterDetail { summary, record })}
// ---------------------------------------------------------------------------// Per-row helpers (retry, reindex, mark resolved, update error)// ---------------------------------------------------------------------------
/// Fetch an unresolved dead letter from whichever table holds it.async fn fetch_dead_letter_for_action( state: &AppState, id: &str,) -> Result<DeadLetterRow, AppError> { let source = DeadLetterSource::from_id(id); match source { DeadLetterSource::LegacyHooks => { let backend = state.db_backend; let sql = adapt_sql( "SELECT id, lexicon_id, uri, did, collection, rkey, action, record, error, attempts FROM happyview_dead_letter_hooks WHERE id = ? AND resolved_at IS NULL", backend, ); #[allow(clippy::type_complexity)] let row: ( String, String, String, String, String, String, String, Option<String>, String, i64, ) = crate::db::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, source, host_kind: "record".to_string(), uri: row.2, did: row.3, collection: row.4, rkey: row.5, action: row.6, record: row.7, }) } DeadLetterSource::Scripts => { let backend = state.db_backend; let id_int: i64 = id.parse().unwrap_or_default(); let sql = adapt_sql( "SELECT id, host_kind, payload FROM happyview_dead_letter_scripts WHERE id = ? AND resolved_at IS NULL", backend, ); let row: (i64, String, String) = crate::db::query_as(&sql) .bind(id_int) .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")) })?; let payload_v: Value = serde_json::from_str(&row.2).unwrap_or(Value::Null); let s = |k: &str| { payload_v .get(k) .and_then(|v| v.as_str()) .unwrap_or_default() .to_string() }; // The record body is nested under `payload.record` for // record events; serialize it back to a string for the // retry call site. let record = payload_v .get("record") .filter(|v| !v.is_null()) .map(|v| v.to_string()); Ok(DeadLetterRow { id: row.0.to_string(), source, host_kind: row.1.clone(), uri: s("uri"), did: s("did"), collection: s("collection"), rkey: s("rkey"), action: s("action"), record, }) } }}
async fn mark_resolved( state: &AppState, id: &str, source: DeadLetterSource,) -> Result<(), AppError> { let backend = state.db_backend; let now = now_rfc3339(); let table = source.table(); let sql = adapt_sql( &format!("UPDATE {table} SET resolved_at = ? WHERE id = ?"), backend, ); let q = crate::db::query(&sql).bind(&now); let q = match source { // Scripts table has INTEGER ids; bind as i64 to avoid sqlite's // implicit-conversion quirks. DeadLetterSource::Scripts => q.bind(id.parse::<i64>().unwrap_or(0)), DeadLetterSource::LegacyHooks => q.bind(id), }; let result = q .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to mark dead letter resolved: {e}")))?; if result.rows_affected() == 0 { return Err(AppError::NotFound(format!("dead letter '{id}' not found"))); } Ok(())}
async fn update_error( state: &AppState, id: &str, source: DeadLetterSource, error: &str,) -> Result<(), AppError> { let backend = state.db_backend; let table = source.table(); let sql = adapt_sql( &format!("UPDATE {table} SET error = ?, attempts = attempts + 1 WHERE id = ?"), backend, ); let q = crate::db::query(&sql).bind(error); let q = match source { DeadLetterSource::Scripts => q.bind(id.parse::<i64>().unwrap_or(0)), DeadLetterSource::LegacyHooks => q.bind(id), }; q.execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to update dead letter error: {e}")))?; Ok(())}
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 ids: Vec<String> = Vec::new();
// Legacy table — supports the optional collection filter. let mut sql = String::from("SELECT id FROM happyview_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 = crate::db::query_as::<(String,)>(&sql); if let Some(ref c) = body.collection { q = q.bind(c); } let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to resolve bulk ids: {e}")))?; ids.extend(rows.into_iter().map(|r| r.0));
{ let mut sql = String::from( "SELECT id FROM happyview_dead_letter_scripts WHERE resolved_at IS NULL", ); if body.collection.is_some() { sql.push_str(" AND collection = ?"); } let sql = adapt_sql(&sql, backend); let mut q = crate::db::query_as::<(i64,)>(&sql); if let Some(ref c) = body.collection { q = q.bind(c); } let rows = q .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to resolve bulk ids: {e}")))?; ids.extend(rows.into_iter().map(|r| r.0.to_string())); }
Ok(ids) } else if let Some(ref ids) = body.ids { Ok(ids.clone()) } else { Err(AppError::BadRequest( "must provide 'ids' or 'all: true'".into(), )) }}
/// Retry a single dead letter by re-running its trigger-keyed script.////// Resolves the script via the new dispatcher's cascade/// (`record.<action>:<nsid>` → `record.index:<nsid>`). Label-arrival/// dead letters are not retried (the upstream label is gone; there's/// nothing to feed back into the runner). Caller gets a 400.async fn retry_single(state: &AppState, id: &str) -> Result<(), AppError> { let dl = fetch_dead_letter_for_action(state, id).await?;
if dl.host_kind == "label" { return Err(AppError::BadRequest( "label-arrival dead letters can't be retried — the upstream label \ event is gone. Dismiss this row and let the labeler subscription \ redeliver if needed." .into(), )); }
let resolved = resolve_record_event(state, &dl.collection, &dl.action) .await .ok_or_else(|| { AppError::NotFound(format!( "no script bound for record.{}:{} (or record.index:{})", dl.action, dl.collection, dl.collection )) })?;
let record: Option<Value> = dl .record .as_deref() .and_then(|r| serde_json::from_str(r).ok());
match run_record_event_once( state, &resolved, RecordEventPayload { nsid: &dl.collection, action: &dl.action, uri: &dl.uri, did: &dl.did, rkey: &dl.rkey, record: record.as_ref(), }, ) .await { Ok(_) => { mark_resolved(state, &dl.id, dl.source).await?; Ok(()) } Err(e) => { update_error(state, &dl.id, dl.source, &e).await?; Err(AppError::Internal(format!( "retry failed for dead letter {id}: {e}" ))) } }}
/// Reindex by fetching the record fresh from the PDS. Only applies to/// record-event dead letters.async fn reindex_single(state: &AppState, id: &str) -> Result<(), AppError> { let dl = fetch_dead_letter_for_action(state, id).await?;
if dl.host_kind == "label" { return Err(AppError::BadRequest( "label-arrival dead letters can't be reindexed — they don't have \ a record to fetch." .into(), )); }
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, &dl.id, dl.source).await?;
Ok(())}