use crate::api::AppState; #[cfg(feature = "indexer")] use crate::db::keys; use crate::db::registry; #[cfg(feature = "indexer_stream")] use crate::db::types::{DbAction, DbRkey, DbTid, DidKey, TrimmedDid}; #[cfg(feature = "indexer_stream")] use crate::types::{ BroadcastEvent, Commit, RepoMetadata, RepoState, RepoStatus, StoredData, StoredEvent, }; use axum::routing::{get, post}; use axum::{ Json, extract::{Query, State}, http::StatusCode, }; #[cfg(feature = "indexer_stream")] use bytes::Bytes; #[cfg(feature = "indexer_stream")] use jacquard_common::CowStr; #[cfg(feature = "indexer")] use jacquard_common::types::ident::AtIdentifier; #[cfg(feature = "indexer_stream")] use jacquard_common::types::{did::Did, tid::Tid}; #[cfg(feature = "indexer_stream")] use jacquard_repo::{MemoryBlockStore, Mst}; use miette::IntoDiagnostic; use serde::{Deserialize, Serialize}; use serde_json::Value; #[cfg(feature = "indexer_stream")] use std::borrow::Cow; use std::sync::Arc; #[cfg(feature = "indexer")] #[derive(Deserialize)] pub struct DebugCountRequest { pub did: String, pub collection: String, } #[cfg(feature = "indexer")] #[derive(Serialize)] pub struct DebugCountResponse { pub count: usize, } #[cfg(feature = "indexer_stream")] #[derive(Serialize)] pub struct DebugSeedXrpcRepoResponse { pub did: String, pub collection: &'static str, pub rkey: &'static str, pub cid: String, pub commit_cid: String, pub rev: String, pub event_id: u64, pub hostname: &'static str, } pub fn router() -> axum::Router> { let r = axum::Router::new() .route("/debug/get", get(handle_debug_get)) .route("/debug/iter", get(handle_debug_iter)) .route("/debug/compact", post(handle_debug_compact)); #[cfg(any(feature = "indexer_stream", feature = "relay"))] let r = r .route( "/debug/ephemeral_ttl_tick", post(handle_debug_ephemeral_ttl_tick), ) .route("/debug/seed_watermark", post(handle_debug_seed_watermark)) .route("/debug/seed_events", post(handle_debug_seed_events)); #[cfg(all(feature = "relay", feature = "jetstream"))] let r = r.route( "/debug/seed_relay_jetstream_commits", post(handle_debug_seed_relay_jetstream_commits), ); #[cfg(all(feature = "indexer_stream", feature = "backlinks"))] let r = r.route("/debug/seed_backlinks", post(handle_debug_seed_backlinks)); #[cfg(feature = "indexer")] let r = r.route("/debug/count", get(handle_debug_count)); #[cfg(feature = "indexer_stream")] let r = r.route("/debug/seed_xrpc_repo", post(handle_debug_seed_xrpc_repo)); r } #[cfg(feature = "indexer")] pub async fn handle_debug_count( State(state): State>, Query(req): Query, ) -> Result, StatusCode> { let ident = req .did .parse::() .map_err(|_| StatusCode::BAD_REQUEST)?; let did = state .resolver .resolve_did(&ident) .await .map_err(|_| StatusCode::BAD_REQUEST)?; let prefix = keys::record_prefix_collection(&did, &req.collection); let count = state .db .run(move |db| { let start_key = prefix.clone(); // rkeys cannot contain 0xff, so every key in the collection sorts // below `prefix ff` let mut end_key = prefix.clone(); end_key.push(0xff); Ok::<_, miette::Report>(db.indexer.record_range(start_key..end_key).count()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(Json(DebugCountResponse { count })) } #[cfg(feature = "indexer_stream")] pub async fn handle_debug_seed_xrpc_repo( State(state): State>, ) -> Result, StatusCode> { const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz"; const COLLECTION: &str = "app.bsky.feed.post"; const RKEY: &str = "stinkpot"; const REV: &str = "3jzfcijpj2z2a"; const HOSTNAME: &str = "pds.example"; let did: Did = Did::new_static(DID).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let rev = Tid::new(REV).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let body = Bytes::from( serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "text": "eight-byte rkey fixture", })) .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?, ); let cid = jacquard_repo::mst::util::compute_cid(&body) .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let store = Arc::new(MemoryBlockStore::new()); let mut mst = Mst::new(store); mst.add_mut(&format!("{COLLECTION}/{RKEY}"), cid) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; mst.persist() .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let root = mst .get_pointer() .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let mut repo_state = RepoState::backfilling(); repo_state.status = RepoStatus::Synced; repo_state.root = Some(Commit { data: root, prev: None, rev: DbTid::from(&rev), sig: Bytes::from_static(b"fixture-signature"), version: 3, }); repo_state.pds = Some(CowStr::Owned(format!("https://{HOSTNAME}").into())); repo_state.signing_key = Some(DidKey(Cow::Owned( [0xed, 0x01] .into_iter() .chain(std::iter::repeat_n(7, 32)) .collect(), ))); let commit_cid = repo_state .root .clone() .and_then(|commit| commit.into_atp_commit(did.clone())) .ok_or(StatusCode::INTERNAL_SERVER_ERROR)? .to_cid() .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let record_key = keys::record_key(&did, COLLECTION, &crate::db::types::DbRkey::new(RKEY)); let repo_key = keys::repo_key(&did); let metadata_key = keys::repo_metadata_key(&did); let collection_count_key = keys::count_collection_key(&did, COLLECTION); let host_count_key = keys::pds_account_count_key(HOSTNAME); let host_cursor_key = keys::firehose_cursor_key(HOSTNAME); let stored_event = StoredEvent { live: true, did: TrimmedDid::from(&did).into_static(), rev: DbTid::from(&rev), collection: CowStr::Borrowed(COLLECTION), rkey: DbRkey::new(RKEY), action: DbAction::Create, data: StoredData::Ptr(cid), }; let event_id = state .db .run(move |db| { let mut batch = db.inner.batch(); db.indexer .stage_record(&mut batch, record_key, body.as_ref()); batch.insert(&db.repos, repo_key, crate::db::ser_repo_state(&repo_state)?); batch.insert( &db.repo_metadata, metadata_key, crate::db::ser_repo_meta(&RepoMetadata::backfilling(1))?, ); batch.insert(&db.counts, collection_count_key, 1_u64.to_be_bytes()); batch.insert(&db.counts, host_count_key, 1_u64.to_be_bytes()); batch.insert(&db.cursors, host_cursor_key, 42_i64.to_be_bytes()); let event_id = db .stream .next_event_id .fetch_add(1, std::sync::atomic::Ordering::SeqCst); db.stream.stage_event( &mut batch, keys::event_key(event_id), rmp_serde::to_vec(&stored_event).into_diagnostic()?, ); batch.commit().into_diagnostic()?; Ok::<_, miette::Report>(event_id) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let _ = state .db .stream .event_tx .send(BroadcastEvent::Persisted(event_id)); let host_count_key = keys::pds_account_count_key(HOSTNAME); let current_host_count = state.db.get_count_sync(&host_count_key); state .db .update_count(&host_count_key, 1_i64 - current_host_count as i64); crate::pds_meta::PdsMeta::update_host(&state.pds_meta, HOSTNAME, |host| { host.status = crate::pds_meta::HostStatus::Active; }); Ok(Json(DebugSeedXrpcRepoResponse { did: DID.to_string(), collection: COLLECTION, rkey: RKEY, cid: cid.to_string(), commit_cid: commit_cid.to_string(), rev: REV.to_string(), event_id, hostname: HOSTNAME, })) } #[derive(Deserialize)] pub struct DebugGetRequest { pub partition: String, pub key: String, } #[derive(Serialize)] pub struct DebugGetResponse { pub value: Option, } pub async fn handle_debug_get( State(state): State>, Query(req): Query, ) -> Result, StatusCode> { let ks = get_keyspace_by_name(&state.db, &req.partition)?; let key = registry::debug_parse_key(&req.partition, &req.key).ok_or(StatusCode::BAD_REQUEST)?; let partition = req.partition.clone(); let value = state .db .run(move |_| { ks.get(key) .inspect_err(crate::db::check_poisoned) .into_diagnostic() }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)? .map(|v| registry::debug_value(&partition, &v)); Ok(Json(DebugGetResponse { value })) } #[derive(Deserialize)] pub struct DebugIterRequest { pub partition: String, pub start: Option, pub end: Option, pub limit: Option, pub reverse: Option, } #[derive(Serialize)] pub struct DebugIterResponse { pub items: Vec<(String, Value)>, } pub async fn handle_debug_iter( State(state): State>, Query(req): Query, ) -> Result, StatusCode> { let partition = req.partition.clone(); let parse_bound = |s: Option| -> Result>, StatusCode> { s.map(|s| registry::debug_parse_key(&partition, &s).ok_or(StatusCode::BAD_REQUEST)) .transpose() }; let start = parse_bound(req.start)?; let end = parse_bound(req.end)?; let items = state .db .run(move |db| { let ks = get_keyspace_by_name(db, &req.partition) .map_err(|_| miette::miette!("bad request"))?; let limit = req.limit.unwrap_or(50); let collect = |iter: &mut dyn Iterator| { let mut items = Vec::new(); for guard in iter.take(limit) { let (k, v) = guard .into_inner() .map_err(|_| miette::miette!("internal error"))?; let key_str = registry::debug_render_key(&partition, &k); items.push((key_str, registry::debug_value(&partition, &v))); } Ok::<_, miette::Report>(items) }; let start_bound = if let Some(s) = &start { std::ops::Bound::Included(s.as_slice()) } else { std::ops::Bound::Unbounded }; let end_bound = if let Some(e) = &end { std::ops::Bound::Included(e.as_slice()) } else { std::ops::Bound::Unbounded }; if req.reverse == Some(true) { collect( &mut ks .range::<&[u8], (std::ops::Bound<&[u8]>, std::ops::Bound<&[u8]>)>(( start_bound, end_bound, )) .rev(), ) } else { collect( &mut ks.range::<&[u8], (std::ops::Bound<&[u8]>, std::ops::Bound<&[u8]>)>(( start_bound, end_bound, )), ) } }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(Json(DebugIterResponse { items })) } fn get_keyspace_by_name(db: &crate::db::Db, name: &str) -> Result { db.keyspace_by_name(name).ok_or(StatusCode::BAD_REQUEST) } #[derive(Deserialize)] pub struct DebugCompactRequest { pub partition: String, } pub async fn handle_debug_compact( State(state): State>, Query(req): Query, ) -> Result { state .db .run(move |db| { let ks = get_keyspace_by_name(db, &req.partition) .map_err(|_| miette::miette!("bad request"))?; ks.remove(b"dummy_tombstone123").into_diagnostic()?; db.inner .persist(fjall::PersistMode::Buffer) .into_diagnostic()?; ks.rotate_memtable_and_wait().into_diagnostic()?; ks.major_compact().into_diagnostic()?; Ok::<_, miette::Report>(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(StatusCode::OK) } #[cfg(any(feature = "indexer_stream", feature = "relay"))] pub async fn handle_debug_ephemeral_ttl_tick( State(state): State>, ) -> Result { let ephemeral_ttl = state.ephemeral_ttl; state .db .run(move |db| { #[cfg(feature = "indexer_stream")] crate::db::ephemeral::ephemeral_ttl_tick(db, &ephemeral_ttl)?; #[cfg(feature = "relay")] crate::db::ephemeral::relay_events_ttl_tick(db, &ephemeral_ttl)?; #[cfg(feature = "jetstream")] crate::db::ephemeral::jetstream_events_ttl_tick(db, &ephemeral_ttl)?; Ok::<_, miette::Report>(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(StatusCode::OK) } #[derive(Deserialize)] #[cfg(any(feature = "indexer_stream", feature = "relay"))] pub struct DebugSeedWatermarkRequest { /// unix timestamp (seconds) to write the watermark at pub ts: u64, /// event_id the watermark points to, all events before this id will be pruned pub event_id: u64, } /// writes an event watermark entry directly to the cursors keyspace, using identical /// key/value encoding to the real TTL worker. used in tests to plant a past watermark /// so the real `ephemeral_ttl_tick` code path is exercised without waiting 3600 seconds. #[cfg(any(feature = "indexer_stream", feature = "relay"))] pub async fn handle_debug_seed_watermark( State(state): State>, Query(req): Query, ) -> Result { state .db .run(move |db| { #[cfg(feature = "indexer_stream")] db.cursors .insert( crate::db::keys::event_watermark_key(req.ts), req.event_id.to_be_bytes(), ) .into_diagnostic()?; #[cfg(feature = "relay")] db.cursors .insert( crate::db::keys::relay_event_watermark_key(req.ts), req.event_id.to_be_bytes(), ) .into_diagnostic()?; Ok::<_, miette::Report>(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(StatusCode::OK) } #[derive(Deserialize)] #[cfg(any(feature = "indexer_stream", feature = "relay"))] pub struct DebugSeedEventsRequest { pub partition: String, pub count: u64, } #[cfg(any(feature = "indexer_stream", feature = "relay"))] pub async fn handle_debug_seed_events( State(state): State>, Query(req): Query, ) -> Result { if req.partition != "events" && req.partition != "relay_events" { return Err(StatusCode::BAD_REQUEST); } state .db .run(move |db| { let mut batch = db.inner.batch(); if req.partition == "events" { #[cfg(feature = "indexer_stream")] { for _ in 0..req.count { let seq = db .stream .next_event_id .fetch_add(1, std::sync::atomic::Ordering::SeqCst); db.stream.stage_event( &mut batch, crate::db::keys::event_key(seq), b"dummy", ); } } } else if req.partition == "relay_events" { #[cfg(feature = "relay")] { for _ in 0..req.count { let seq = db .relay .next_seq .fetch_add(1, std::sync::atomic::Ordering::SeqCst); batch.insert( &db.relay.events, crate::db::keys::relay_event_key(seq), b"dummy", ); } } } batch.commit().into_diagnostic()?; Ok::<_, miette::Report>(()) }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; Ok(StatusCode::OK) } #[cfg(all(feature = "relay", feature = "jetstream"))] #[derive(Deserialize)] pub struct DebugSeedRelayJetstreamCommitsRequest { pub count: u64, } #[cfg(all(feature = "relay", feature = "jetstream"))] pub async fn handle_debug_seed_relay_jetstream_commits( State(state): State>, Query(req): Query, ) -> Result { let count = usize::try_from(req.count) .ok() .filter(|count| *count <= 10_000) .ok_or(StatusCode::BAD_REQUEST)?; state .db .run(move |db| { use crate::db::types::TrimmedDid; use crate::ingest::stream::{ Account, AccountStatus, Commit, Datetime, Identity, RepoOp, RepoOpAction, SubscribeReposMessage, decode_frame, encode_frame, }; use crate::types::StoredJetstreamEvent; use jacquard_common::CowStr; use jacquard_common::types::cid::CidLink; use jacquard_common::types::string::{Did, Handle, Tid}; use std::sync::atomic::Ordering; let did = Did::new_static("did:plc:ewvi7nxzyoun6zhxrhs64oiz").into_diagnostic()?; let commit_cid = CidLink::from( jacquard_repo::mst::util::compute_cid(b"relay-jetstream-fixture") .into_diagnostic()?, ); let rev: Tid = "3kzbif5moe22m".parse().into_diagnostic()?; let time = Datetime::try_from("2026-08-10T12:34:56Z".to_owned()).into_diagnostic()?; let trimmed_did = TrimmedDid::from(&did).into_static(); let mut batch = db.inner.batch(); let mut broadcasts = Vec::with_capacity(count + 2); for index in 0..count { let relay_seq = db.relay.next_seq.fetch_add(1, Ordering::SeqCst); let rkey = format!("fixture{index}"); let commit = Commit { blobs: Vec::new(), blocks: bytes::Bytes::new(), commit: commit_cid.clone(), ops: vec![RepoOp { action: RepoOpAction::Delete, cid: None, path: CowStr::Owned(format!("app.bsky.feed.post/{rkey}").into()), prev: None, }], prev_data: None, rebase: false, repo: did.clone(), rev: rev.clone(), seq: relay_seq as i64, since: None, time: time.clone(), too_big: false, }; let frame = encode_frame("#commit", &commit)?; let SubscribeReposMessage::Commit(decoded) = decode_frame(frame.as_ref()).map_err(|err| miette::miette!("{err}"))? else { return Err(miette::miette!("relay fixture did not decode as a commit")); }; if decoded.ops.len() != 1 { return Err(miette::miette!("relay fixture lost its operation")); } batch.insert( &db.relay.events, crate::db::keys::relay_event_key(relay_seq), frame.as_ref(), ); broadcasts.push(crate::jetstream::stage_event( &mut batch, db, StoredJetstreamEvent::RelayCommit { did: trimmed_did.clone(), collection: CowStr::Borrowed("app.bsky.feed.post"), relay_seq, op_index: 0, }, Some(crate::jetstream::JetstreamEphemeral { did: did.as_str().to_owned(), rev: rev.as_str().to_owned(), operation: "delete".to_owned(), collection: "app.bsky.feed.post".to_owned(), rkey, record: None, cid: None, live: true, }), )?); } let relay_seq = db.relay.next_seq.fetch_add(1, Ordering::SeqCst); let account = Account { active: false, did: did.clone(), seq: relay_seq as i64, status: Some(AccountStatus::Deactivated), time: time.clone(), }; let frame = encode_frame("#account", &account)?; let SubscribeReposMessage::Account(decoded) = decode_frame(frame.as_ref()).map_err(|err| miette::miette!("{err}"))? else { return Err(miette::miette!( "relay fixture did not decode as an account" )); }; if decoded.active || decoded.status != Some(AccountStatus::Deactivated) { return Err(miette::miette!("relay account fixture lost its state")); } batch.insert( &db.relay.events, crate::db::keys::relay_event_key(relay_seq), frame.as_ref(), ); broadcasts.push(crate::jetstream::stage_event( &mut batch, db, StoredJetstreamEvent::RelayAccount { did: trimmed_did.clone(), relay_seq, }, None, )?); let relay_seq = db.relay.next_seq.fetch_add(1, Ordering::SeqCst); let identity = Identity { did: did.clone(), handle: Some(Handle::new_static("alice.test").into_diagnostic()?), seq: relay_seq as i64, time: time.clone(), }; let frame = encode_frame("#identity", &identity)?; let SubscribeReposMessage::Identity(decoded) = decode_frame(frame.as_ref()).map_err(|err| miette::miette!("{err}"))? else { return Err(miette::miette!( "relay fixture did not decode as an identity" )); }; if decoded.handle.as_ref().map(Handle::as_str) != Some("alice.test") { return Err(miette::miette!("relay identity fixture lost its handle")); } batch.insert( &db.relay.events, crate::db::keys::relay_event_key(relay_seq), frame.as_ref(), ); broadcasts.push(crate::jetstream::stage_event( &mut batch, db, StoredJetstreamEvent::RelayIdentity { did: trimmed_did, relay_seq, }, None, )?); batch.commit().into_diagnostic()?; for broadcast in broadcasts { let _ = db.jetstream.tx.send(broadcast); } Ok::<_, miette::Report>(()) }) .await .map_err(|err| { tracing::error!(%err, "failed to seed relay Jetstream fixture"); StatusCode::INTERNAL_SERVER_ERROR })?; Ok(StatusCode::OK) } #[cfg(all(feature = "indexer_stream", feature = "backlinks"))] #[derive(Deserialize)] pub struct DebugSeedBacklinksRequest { pub like_count: u64, pub post_count: u64, } #[cfg(all(feature = "indexer_stream", feature = "backlinks"))] #[derive(Serialize)] pub struct DebugSeedBacklinksResponse { pub did: String, pub subject: String, pub like_count: u64, pub post_count: u64, } #[cfg(all(feature = "indexer_stream", feature = "backlinks"))] pub async fn handle_debug_seed_backlinks( State(state): State>, Query(req): Query, ) -> Result, StatusCode> { req.like_count .checked_add(req.post_count) .filter(|count| *count <= 10_000) .ok_or(StatusCode::BAD_REQUEST)?; state .db .run(move |db| { use crate::db::types::DbRkey; use jacquard_common::types::string::Did; const DID: &str = "did:plc:wqstj3k5tslmm246baaf3tpa"; const SUBJECT: &str = "at://did:plc:ewvi7nxzyoun6zhxrhs64oiz/app.bsky.feed.post/root"; const SUBJECT_CID: &str = "bafyreicldnax5mfobeba2pnv7fgsdfz2563gc72gyi5cnhyuqfnk3jz7aa"; let did = Did::new_static(DID).into_diagnostic()?; let mut batch = db.inner.batch(); for (collection, count) in [ ("app.bsky.feed.like", req.like_count), ("app.bsky.feed.post", req.post_count), ] { for index in 0..count { let rkey = if collection == "app.bsky.feed.like" && index == 0 { "stinkpot".to_owned() } else { let prefix = if collection == "app.bsky.feed.like" { "like" } else { "post" }; format!("{prefix}{index:04}") }; let body = serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": collection, "subject": { "uri": SUBJECT, "cid": SUBJECT_CID, }, "createdAt": "2026-08-10T12:34:56Z", })) .into_diagnostic()?; let value = serde_ipld_dagcbor::from_slice::(&body) .into_diagnostic()?; let record_key = crate::db::keys::record_key(&did, collection, &DbRkey::new(&rkey)); db.indexer .stage_record(&mut batch, record_key, body.as_slice()); crate::backlinks::store::index_record( &mut batch, &db.backlinks, did.as_str(), collection, &rkey, &value, )?; } } batch.commit().into_diagnostic()?; Ok::<_, miette::Report>(DebugSeedBacklinksResponse { did: DID.to_owned(), subject: SUBJECT.to_owned(), like_count: req.like_count, post_count: req.post_count, }) }) .await .map(Json) .map_err(|err| { tracing::error!(%err, "failed to seed backlinks fixture"); StatusCode::INTERNAL_SERVER_ERROR }) }