very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804use 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<Arc<AppState>> { 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<Arc<AppState>>, Query(req): Query<DebugCountRequest>,) -> Result<Json<DebugCountResponse>, StatusCode> { let ident = req .did .parse::<AtIdentifier>() .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<Arc<AppState>>,) -> Result<Json<DebugSeedXrpcRepoResponse>, 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<Value>,}
pub async fn handle_debug_get( State(state): State<Arc<AppState>>, Query(req): Query<DebugGetRequest>,) -> Result<Json<DebugGetResponse>, 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<String>, pub end: Option<String>, pub limit: Option<usize>, pub reverse: Option<bool>,}
#[derive(Serialize)]pub struct DebugIterResponse { pub items: Vec<(String, Value)>,}
pub async fn handle_debug_iter( State(state): State<Arc<AppState>>, Query(req): Query<DebugIterRequest>,) -> Result<Json<DebugIterResponse>, StatusCode> { let partition = req.partition.clone();
let parse_bound = |s: Option<String>| -> Result<Option<Vec<u8>>, 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<Item = fjall::Guard>| { 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<fjall::Keyspace, StatusCode> { 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<Arc<AppState>>, Query(req): Query<DebugCompactRequest>,) -> Result<StatusCode, StatusCode> { 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<Arc<AppState>>,) -> Result<StatusCode, StatusCode> { 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<Arc<AppState>>, Query(req): Query<DebugSeedWatermarkRequest>,) -> Result<StatusCode, StatusCode> { 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<Arc<AppState>>, Query(req): Query<DebugSeedEventsRequest>,) -> Result<StatusCode, StatusCode> { 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<Arc<AppState>>, Query(req): Query<DebugSeedRelayJetstreamCommitsRequest>,) -> Result<StatusCode, StatusCode> { 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<Arc<AppState>>, Query(req): Query<DebugSeedBacklinksRequest>,) -> Result<Json<DebugSeedBacklinksResponse>, 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::<jacquard_common::Data>(&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 })}