diff --git a/migrations/postgres/20260421000000_dead_letters_resolved_at.sql b/migrations/postgres/20260421000000_dead_letters_resolved_at.sql new file mode 100644 index 0000000..d7cade8 --- /dev/null +++ b/migrations/postgres/20260421000000_dead_letters_resolved_at.sql @@ -0,0 +1,2 @@ +ALTER TABLE dead_letter_hooks ADD COLUMN resolved_at TIMESTAMPTZ; +CREATE INDEX idx_dead_letter_hooks_resolved_at ON dead_letter_hooks (resolved_at); diff --git a/migrations/sqlite/20260421000000_dead_letters_resolved_at.sql b/migrations/sqlite/20260421000000_dead_letters_resolved_at.sql new file mode 100644 index 0000000..07018a8 --- /dev/null +++ b/migrations/sqlite/20260421000000_dead_letters_resolved_at.sql @@ -0,0 +1,2 @@ +ALTER TABLE dead_letter_hooks ADD COLUMN resolved_at TEXT; +CREATE INDEX idx_dead_letter_hooks_resolved_at ON dead_letter_hooks (resolved_at); diff --git a/packages/docs/docs/guides/event-logs.md b/packages/docs/docs/guides/event-logs.md index 88ae61e..9f723a0 100644 --- a/packages/docs/docs/guides/event-logs.md +++ b/packages/docs/docs/guides/event-logs.md @@ -79,7 +79,7 @@ Logged when a user attempts to access an endpoint they don't have permission for | `hook.executed` | info | Record AT URI | `lexicon_id` | | `hook.dead_lettered` | error | Record AT URI | `lexicon_id`, `error` | -Logged when [index hooks](index-hooks.md) run. Dead-lettered events indicate a hook failed all retry attempts. +Logged when [index hooks](index-hooks.md) run. Dead-lettered events indicate a hook failed all retry attempts. You can manage dead letters from the **Data > Dead Letters** page in the dashboard — see [Dead Letters](#dead-letters) below. ### Backfill events @@ -128,6 +128,18 @@ Set `EVENT_LOG_RETENTION_DAYS=0` to disable automatic cleanup and keep logs inde See [Configuration](../getting-started/configuration.md) for all environment variables. +## Dead Letters + +When an index hook fails after all retry attempts, the event is stored in the dead letters queue. You can manage dead letters from the **Data > Dead Letters** page in the dashboard. + +From the dead letters page you can: + +- **Retry Hook** — replay the stored record through the index hook (use after fixing a hook script) +- **Re-index** — fetch the record fresh from the PDS and run it through the full indexing pipeline (use when the record may have changed) +- **Dismiss** — mark the dead letter as resolved without retrying + +Bulk actions are available for selected rows or all entries matching the current filters. + ## Next steps - [Admin API — Event Logs](../reference/admin/events.md) — full query parameters and response format diff --git a/src/admin/dead_letters.rs b/src/admin/dead_letters.rs new file mode 100644 index 0000000..0bcfd4b --- /dev/null +++ b/src/admin/dead_letters.rs @@ -0,0 +1,600 @@ +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::{HookEvent, run_hook_once}; +use crate::record_handler::RecordEvent; + +// --------------------------------------------------------------------------- +// Query / request / response types +// --------------------------------------------------------------------------- + +#[derive(Deserialize)] +pub struct ListQuery { + pub collection: Option, + pub resolved: Option, + pub cursor: Option, + pub limit: Option, +} + +#[derive(Deserialize)] +pub struct CountQuery { + pub resolved: Option, +} + +#[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, + #[serde(skip_serializing_if = "Option::is_none")] + pub resolved_at: Option>, +} + +#[derive(Serialize)] +pub struct DeadLetterDetail { + #[serde(flatten)] + pub summary: DeadLetterSummary, + pub record: Option, +} + +#[derive(Serialize)] +pub struct ListResponse { + pub dead_letters: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + pub cursor: Option, +} + +#[derive(Serialize)] +pub struct CountResponse { + pub count: i64, +} + +#[derive(Deserialize)] +pub struct BulkRequest { + pub ids: Option>, + pub all: Option, + pub collection: Option, +} + +/// 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, + error: String, + attempts: i64, +} + +// --------------------------------------------------------------------------- +// Handlers +// --------------------------------------------------------------------------- + +/// GET /admin/dead-letters +pub(super) async fn list( + auth: UserAuth, + State(state): State, + Query(query): Query, +) -> Result, 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, + ), + >(&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 = 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/count +pub(super) async fn count( + auth: UserAuth, + State(state): State, + Query(query): Query, +) -> Result, 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, + Path(id): Path, +) -> Result, 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, + Option, + ) = 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}/dismiss +pub(super) async fn dismiss( + auth: UserAuth, + State(state): State, + Path(id): Path, +) -> Result, 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}/retry +pub(super) async fn retry( + auth: UserAuth, + State(state): State, + Path(id): Path, +) -> Result, 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, + Path(id): Path, +) -> Result, 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, + Json(body): Json, +) -> Result, 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/retry +pub(super) async fn bulk_retry( + auth: UserAuth, + State(state): State, + Json(body): Json, +) -> Result, 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, + Json(body): Json, +) -> Result, 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 { + 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, + 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, 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, AppError> { + let backend = state.db_backend; + let sql = adapt_sql( + "SELECT index_hook FROM lexicons WHERE id = ?", + backend, + ); + let row: Option<(Option,)> = 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 = 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(()) +} diff --git a/src/admin/mod.rs b/src/admin/mod.rs index 4df9974..b4a243f 100644 --- a/src/admin/mod.rs +++ b/src/admin/mod.rs @@ -2,6 +2,7 @@ mod api_clients; mod api_keys; pub(crate) mod auth; mod backfill; +mod dead_letters; mod domains; mod events; mod labelers; @@ -102,4 +103,19 @@ pub fn admin_routes(_state: AppState) -> Router { .route("/domains", post(domains::create).get(domains::list)) .route("/domains/{id}", delete(domains::delete)) .route("/domains/{id}/primary", post(domains::set_primary)) + .route("/dead-letters", get(dead_letters::list)) + .route("/dead-letters/count", get(dead_letters::count)) + .route( + "/dead-letters/bulk/dismiss", + post(dead_letters::bulk_dismiss), + ) + .route("/dead-letters/bulk/retry", post(dead_letters::bulk_retry)) + .route( + "/dead-letters/bulk/reindex", + post(dead_letters::bulk_reindex), + ) + .route("/dead-letters/{id}", get(dead_letters::detail)) + .route("/dead-letters/{id}/dismiss", post(dead_letters::dismiss)) + .route("/dead-letters/{id}/retry", post(dead_letters::retry)) + .route("/dead-letters/{id}/reindex", post(dead_letters::reindex)) } diff --git a/src/admin/permissions.rs b/src/admin/permissions.rs index 129ae62..0674ca0 100644 --- a/src/admin/permissions.rs +++ b/src/admin/permissions.rs @@ -2,7 +2,7 @@ use std::collections::HashSet; use serde::{Deserialize, Serialize}; -/// All 27 permissions in the system. +/// All 29 permissions in the system. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] pub enum Permission { #[serde(rename = "lexicons:create")] @@ -78,6 +78,11 @@ pub enum Permission { ApiClientsEdit, #[serde(rename = "api-clients:delete")] ApiClientsDelete, + + #[serde(rename = "dead-letters:read")] + DeadLettersRead, + #[serde(rename = "dead-letters:manage")] + DeadLettersManage, } impl Permission { @@ -115,6 +120,8 @@ impl Permission { Self::ApiClientsCreate => "api-clients:create", Self::ApiClientsEdit => "api-clients:edit", Self::ApiClientsDelete => "api-clients:delete", + Self::DeadLettersRead => "dead-letters:read", + Self::DeadLettersManage => "dead-letters:manage", } } @@ -152,6 +159,8 @@ impl Permission { Self::ApiClientsCreate, Self::ApiClientsEdit, Self::ApiClientsDelete, + Self::DeadLettersRead, + Self::DeadLettersManage, ]) } } @@ -178,12 +187,14 @@ impl Template { Permission::BackfillRead, Permission::StatsRead, Permission::EventsRead, + Permission::DeadLettersRead, ]), Self::Operator => { let mut perms = Self::Viewer.permissions(); perms.insert(Permission::BackfillCreate); perms.insert(Permission::ApiKeysCreate); perms.insert(Permission::ApiKeysDelete); + perms.insert(Permission::DeadLettersManage); perms } Self::Manager => { diff --git a/src/lua/execute.rs b/src/lua/execute.rs index 1fdf106..ed4d97f 100644 --- a/src/lua/execute.rs +++ b/src/lua/execute.rs @@ -969,7 +969,7 @@ pub async fn execute_hook_script(event: &HookEvent<'_>) -> Option { /// Returns `Ok(None)` when `handle()` returns nil (meaning "skip indexing"), /// `Ok(Some(value))` when it returns a table (use that as the record), or /// `Ok(Some(original))` for other non-nil types. -async fn run_hook_once(event: &HookEvent<'_>) -> Result, String> { +pub async fn run_hook_once(event: &HookEvent<'_>) -> Result, String> { let lua = sandbox::create_sandbox().map_err(|e| format!("failed to create Lua VM: {e}"))?; let backend = event.state.db_backend; diff --git a/src/lua/mod.rs b/src/lua/mod.rs index 56f5f3d..bbf7ddb 100644 --- a/src/lua/mod.rs +++ b/src/lua/mod.rs @@ -9,6 +9,6 @@ mod tid; mod xrpc_api; pub(crate) use execute::{ - HookEvent, execute_hook_script, execute_procedure_script, execute_query_script, + HookEvent, execute_hook_script, execute_procedure_script, execute_query_script, run_hook_once, }; pub(crate) use sandbox::validate_script; diff --git a/web/src/app/dashboard/dead-letters/page.tsx b/web/src/app/dashboard/dead-letters/page.tsx new file mode 100644 index 0000000..24ff16f --- /dev/null +++ b/web/src/app/dashboard/dead-letters/page.tsx @@ -0,0 +1,616 @@ +"use client"; + +import { useCallback, useEffect, useMemo, useRef, useState } from "react"; +import { + type ColumnDef, + type ColumnFiltersState, + type RowSelectionState, + type VisibilityState, + getCoreRowModel, + useReactTable, +} from "@tanstack/react-table"; + +import { + getDeadLetters, + getDeadLetter, + retryDeadLetter, + reindexDeadLetter, + dismissDeadLetter, + bulkRetryDeadLetters, + bulkReindexDeadLetters, + bulkDismissDeadLetters, +} from "@/lib/api"; +import type { DeadLetterSummary, DeadLetterDetail } from "@/types/dead-letters"; +import { DataTable } from "@/components/data-table/data-table"; +import { DataTableColumnHeader } from "@/components/data-table/data-table-column-header"; +import { DataTableToolbar } from "@/components/data-table/data-table-toolbar"; +import { CodeBlock } from "@/components/code-block"; +import { MonacoEditor } from "@/components/monaco-editor"; +import { SiteHeader } from "@/components/site-header"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { + Sheet, + SheetContent, + SheetFooter, + SheetHeader, + SheetTitle, +} from "@/components/ui/sheet"; +import { + DropdownMenu, + DropdownMenuContent, + DropdownMenuItem, + DropdownMenuTrigger, +} from "@/components/ui/dropdown-menu"; +import { + Select, + SelectContent, + SelectItem, + SelectTrigger, + SelectValue, +} from "@/components/ui/select"; +import { + ChevronLeft, + ChevronRight, + RotateCcw, + RefreshCw, + XCircle, + MoreHorizontal, +} from "lucide-react"; +import { Checkbox } from "@/components/ui/checkbox"; + +function actionBadge(action: string) { + switch (action) { + case "delete": + return delete; + case "update": + return ( + + update + + ); + default: + return {action}; + } +} + +function timeAgo(dateStr: string): string { + const now = Date.now(); + const then = new Date(dateStr).getTime(); + const seconds = Math.floor((now - then) / 1000); + if (seconds < 60) return `${seconds}s ago`; + const minutes = Math.floor(seconds / 60); + if (minutes < 60) return `${minutes}m ago`; + const hours = Math.floor(minutes / 60); + if (hours < 24) return `${hours}h ago`; + const days = Math.floor(hours / 24); + return `${days}d ago`; +} + +function DetailBody({ detail }: { detail: DeadLetterDetail }) { + return ( +
+
+
+ URI +

{detail.uri}

+
+
+ DID +

{detail.did}

+
+
+ Collection +

{detail.collection}

+
+
+ Lexicon +

{detail.lexicon_id}

+
+
+ Action +

{actionBadge(detail.action)}

+
+
+ Attempts +

{detail.attempts}

+
+
+ Created +

+ {new Date(detail.created_at).toLocaleString()} +

+
+ {detail.resolved_at && ( +
+ Resolved +

+ {new Date(detail.resolved_at).toLocaleString()} +

+
+ )} +
+ +
+ Error +
+ {detail.error} +
+
+ + {detail.record && ( +
+ Record + +
+ )} +
+ ); +} + +export default function DeadLettersPage() { + const [items, setItems] = useState([]); + const [error, setError] = useState(null); + const [loading, setLoading] = useState(false); + const [viewDetail, setViewDetail] = useState(null); + const [actionLoading, setActionLoading] = useState(false); + const [resolvedFilter, setResolvedFilter] = useState("false"); + const [rowSelection, setRowSelection] = useState({}); + + const [cursorStack, setCursorStack] = useState([]); + const [nextCursor, setNextCursor] = useState(null); + + const [columnFilters, setColumnFilters] = useState([]); + const [columnVisibility, setColumnVisibility] = useState({}); + const debounceRef = useRef>(null); + const [debouncedFilters, setDebouncedFilters] = + useState(columnFilters); + + useEffect(() => { + if (debounceRef.current) clearTimeout(debounceRef.current); + debounceRef.current = setTimeout(() => { + setDebouncedFilters(columnFilters); + }, 300); + }, [columnFilters]); + + const collectionFilter = useMemo( + () => + ( + debouncedFilters.find((f) => f.id === "collection")?.value as + | string[] + | undefined + )?.join(",") || undefined, + [debouncedFilters], + ); + + const fetchItems = useCallback( + async (cursor?: string) => { + setLoading(true); + setError(null); + try { + const data = await getDeadLetters({ + collection: collectionFilter, + resolved: resolvedFilter, + cursor, + limit: 50, + }); + setItems(data.dead_letters); + setNextCursor(data.cursor); + setRowSelection({}); + } catch (e: unknown) { + setError(e instanceof Error ? e.message : String(e)); + setItems([]); + setNextCursor(null); + } finally { + setLoading(false); + } + }, + [collectionFilter, resolvedFilter], + ); + + useEffect(() => { + setCursorStack([]); + fetchItems(); + }, [fetchItems]); + + function handleNext() { + if (!nextCursor) return; + setCursorStack((prev) => [...prev, nextCursor]); + fetchItems(nextCursor); + } + + function handlePrevious() { + if (cursorStack.length === 0) return; + const stack = [...cursorStack]; + stack.pop(); + const prevCursor = stack.length > 0 ? stack[stack.length - 1] : undefined; + setCursorStack(stack); + fetchItems(prevCursor); + } + + async function openDetail(row: DeadLetterSummary) { + try { + const detail = await getDeadLetter(row.id); + setViewDetail(detail); + } catch { + setError("Failed to load detail"); + } + } + + async function handleDetailAction(action: "retry" | "reindex" | "dismiss") { + if (!viewDetail) return; + setActionLoading(true); + try { + if (action === "retry") await retryDeadLetter(viewDetail.id); + else if (action === "reindex") await reindexDeadLetter(viewDetail.id); + else await dismissDeadLetter(viewDetail.id); + setViewDetail(null); + fetchItems(); + } catch (e: unknown) { + setError(e instanceof Error ? e.message : String(e)); + } finally { + setActionLoading(false); + } + } + + const selectedIds = useMemo( + () => Object.keys(rowSelection).filter((k) => rowSelection[k]), + [rowSelection], + ); + + async function handleBulkAction( + action: "retry" | "reindex" | "dismiss", + scope: "selected" | "all", + ) { + setLoading(true); + try { + const body = + scope === "all" + ? { all: true, collection: collectionFilter } + : { ids: selectedIds }; + if (action === "retry") await bulkRetryDeadLetters(body); + else if (action === "reindex") await bulkReindexDeadLetters(body); + else await bulkDismissDeadLetters(body); + setRowSelection({}); + fetchItems(); + } catch (e: unknown) { + setError(e instanceof Error ? e.message : String(e)); + } finally { + setLoading(false); + } + } + + const columns = useMemo[]>( + () => [ + { + id: "select", + header: ({ table }) => ( + + table.toggleAllPageRowsSelected(!!value) + } + aria-label="Select all" + /> + ), + cell: ({ row }) => ( +
e.stopPropagation()}> + row.toggleSelected(!!value)} + aria-label="Select row" + /> +
+ ), + enableSorting: false, + enableColumnFilter: false, + size: 32, + }, + { + id: "collection", + accessorKey: "collection", + header: ({ column }) => ( + + ), + cell: ({ row }) => ( + + {row.original.collection} + + ), + enableColumnFilter: true, + enableSorting: false, + meta: { + label: "Collection", + placeholder: "Filter by collection...", + variant: "text", + }, + }, + { + id: "uri", + accessorKey: "uri", + header: ({ column }) => ( + + ), + cell: ({ row }) => ( + + {row.original.uri} + + ), + enableSorting: false, + }, + { + id: "action", + accessorKey: "action", + header: ({ column }) => ( + + ), + cell: ({ row }) => actionBadge(row.original.action), + enableSorting: false, + }, + { + id: "error", + accessorKey: "error", + header: ({ column }) => ( + + ), + cell: ({ row }) => ( + + {row.original.error.split("\n")[0]} + + ), + enableSorting: false, + }, + { + id: "attempts", + accessorKey: "attempts", + header: ({ column }) => ( + + ), + cell: ({ row }) => ( + {row.original.attempts} + ), + enableSorting: false, + }, + { + id: "created_at", + accessorKey: "created_at", + header: ({ column }) => ( + + ), + cell: ({ row }) => ( + + {timeAgo(row.original.created_at)} + + ), + enableSorting: false, + }, + ...(resolvedFilter !== "false" + ? [ + { + id: "status", + accessorKey: "resolved_at" as const, + header: ({ column }) => ( + + ), + cell: ({ row }) => + row.original.resolved_at ? ( + resolved + ) : ( + unresolved + ), + enableSorting: false, + } satisfies ColumnDef, + ] + : []), + ], + [resolvedFilter], + ); + + const table = useReactTable({ + data: items, + columns, + state: { columnFilters, columnVisibility, rowSelection }, + defaultColumn: { enableColumnFilter: false }, + onColumnFiltersChange: setColumnFilters, + onColumnVisibilityChange: setColumnVisibility, + onRowSelectionChange: setRowSelection, + getCoreRowModel: getCoreRowModel(), + getRowId: (row) => row.id, + }); + + return ( + <> + +
+ {error &&

{error}

} + +
+ + +
+ + {selectedIds.length > 0 && ( +
+ + {selectedIds.length} selected + + + + + + + + + + handleBulkAction("retry", "all")} + > + Retry all matching + + handleBulkAction("reindex", "all")} + > + Re-index all matching + + handleBulkAction("dismiss", "all")} + > + Dismiss all matching + + + +
+ )} + + + +
+

+ {items.length} item(s) on this page. +

+
+ + +
+
+ + { + if (!open) setViewDetail(null); + }} + > + + {viewDetail && ( + <> + + + {actionBadge(viewDetail.action)} + + {viewDetail.collection} + + + +
+ +
+ {viewDetail.resolved_at == null && ( + +
+ +
+ + +
+ )} + + )} +
+
+
+ + ); +} diff --git a/web/src/components/app-sidebar.tsx b/web/src/components/app-sidebar.tsx index 6451049..6a11367 100644 --- a/web/src/components/app-sidebar.tsx +++ b/web/src/components/app-sidebar.tsx @@ -1,5 +1,6 @@ "use client"; +import { useEffect, useState } from "react"; import { IconDashboard, IconFileDescription, @@ -17,6 +18,7 @@ import { IconInfoCircle, IconApps, IconArrowUpCircle, + IconSkull, } from "@tabler/icons-react"; import Image from "next/image"; import Link from "next/link"; @@ -52,6 +54,12 @@ const dataItems: NavItem[] = [ { title: "Lexicons", url: "/dashboard/lexicons", icon: IconFileDescription }, { title: "Records", url: "/dashboard/records", icon: IconTable }, { title: "Backfill", url: "/dashboard/backfill", icon: IconDatabase }, + { + title: "Dead Letters", + url: "/dashboard/dead-letters", + icon: IconSkull, + requiredPermissions: ["dead-letters:read"], + }, ]; const accessItems: NavItem[] = [ @@ -123,6 +131,27 @@ export function AppSidebar({ ...props }: React.ComponentProps) { const { hasPermission } = useCurrentUser(); const { hasUpdates } = usePluginUpdates(); + const [deadLetterCount, setDeadLetterCount] = useState(0); + + useEffect(() => { + let cancelled = false; + async function fetchCount() { + try { + const { getDeadLetterCount } = await import("@/lib/api"); + const data = await getDeadLetterCount("false"); + if (!cancelled) setDeadLetterCount(data.count); + } catch { + // sidebar badge is best-effort + } + } + fetchCount(); + const interval = setInterval(fetchCount, 30000); + return () => { + cancelled = true; + clearInterval(interval); + }; + }, []); + function filterByPermission(items: NavItem[]) { return items.filter( (item) => @@ -200,20 +229,29 @@ export function AppSidebar({ ...props }: React.ComponentProps) { Data - {visibleData.map((item) => ( - - - - - {item.title} - - - - ))} + {visibleData.map((item) => { + const showDeadLetterBadge = + item.title === "Dead Letters" && deadLetterCount > 0; + return ( + + + + + {item.title} + {showDeadLetterBadge && ( + + {deadLetterCount} + + )} + + + + ); + })} diff --git a/web/src/lib/api.ts b/web/src/lib/api.ts index 57a052a..ab2949c 100644 --- a/web/src/lib/api.ts +++ b/web/src/lib/api.ts @@ -18,6 +18,12 @@ import type { UnlinkResponse, ConnectResponse, } from "@/types/external-accounts" +import type { + DeadLettersListResponse, + DeadLetterDetail, + DeadLetterCountResponse, + BulkActionResponse, +} from "@/types/dead-letters" export type { ApiKeySummary, CreateApiKeyResponse } from "@/types/api-keys" export type { CollectionStat, StatsResponse } from "@/types/stats" @@ -43,6 +49,13 @@ export type { ConfigSchema, ConfigProperty, } from "@/types/external-accounts" +export type { + DeadLetterSummary, + DeadLetterDetail, + DeadLettersListResponse, + DeadLetterCountResponse, + BulkActionResponse, +} from "@/types/dead-letters" export class ApiError extends Error { status: number @@ -564,3 +577,90 @@ export function previewPlugin(url: string, signal?: AbortSignal) { signal, }) } + +// --------------------------------------------------------------------------- +// Dead Letters +// --------------------------------------------------------------------------- + +export function getDeadLetters(params?: { + collection?: string + resolved?: string + cursor?: string + limit?: number +}) { + const searchParams = new URLSearchParams() + if (params?.collection) searchParams.set("collection", params.collection) + if (params?.resolved) searchParams.set("resolved", params.resolved) + if (params?.cursor) searchParams.set("cursor", params.cursor) + if (params?.limit) searchParams.set("limit", String(params.limit)) + const qs = searchParams.toString() + return apiFetch( + `/admin/dead-letters${qs ? `?${qs}` : ""}`, + ) +} + +export function getDeadLetterCount(resolved?: string) { + const searchParams = new URLSearchParams() + if (resolved) searchParams.set("resolved", resolved) + const qs = searchParams.toString() + return apiFetch( + `/admin/dead-letters/count${qs ? `?${qs}` : ""}`, + ) +} + +export function getDeadLetter(id: string) { + return apiFetch( + `/admin/dead-letters/${encodeURIComponent(id)}`, + ) +} + +export function retryDeadLetter(id: string) { + return apiFetch(`/admin/dead-letters/${encodeURIComponent(id)}/retry`, { + method: "POST", + }) +} + +export function reindexDeadLetter(id: string) { + return apiFetch(`/admin/dead-letters/${encodeURIComponent(id)}/reindex`, { + method: "POST", + }) +} + +export function dismissDeadLetter(id: string) { + return apiFetch(`/admin/dead-letters/${encodeURIComponent(id)}/dismiss`, { + method: "POST", + }) +} + +export function bulkDismissDeadLetters(body: { + ids?: string[] + all?: boolean + collection?: string +}) { + return apiFetch("/admin/dead-letters/bulk/dismiss", { + method: "POST", + body: JSON.stringify(body), + }) +} + +export function bulkRetryDeadLetters(body: { + ids?: string[] + all?: boolean + collection?: string +}) { + return apiFetch("/admin/dead-letters/bulk/retry", { + method: "POST", + body: JSON.stringify(body), + }) +} + +export function bulkReindexDeadLetters(body: { + ids?: string[] + all?: boolean + collection?: string +}) { + return apiFetch("/admin/dead-letters/bulk/reindex", { + method: "POST", + body: JSON.stringify(body), + }) +} diff --git a/web/src/types/dead-letters.ts b/web/src/types/dead-letters.ts new file mode 100644 index 0000000..b32c45c --- /dev/null +++ b/web/src/types/dead-letters.ts @@ -0,0 +1,30 @@ +export interface DeadLetterSummary { + id: string + lexicon_id: string + uri: string + did: string + collection: string + rkey: string + action: string + error: string + attempts: number + created_at: string + resolved_at: string | null +} + +export interface DeadLetterDetail extends DeadLetterSummary { + record: Record | null +} + +export interface DeadLettersListResponse { + dead_letters: DeadLetterSummary[] + cursor: string | null +} + +export interface DeadLetterCountResponse { + count: number +} + +export interface BulkActionResponse { + ok: boolean +}