use crate::config::BackfillConfig; use crate::db; use crate::indexer::types::{AggregateDeltaStore, BackfillItem, BackfillItemInner}; use crate::indexer::{self, records}; use chrono::prelude::*; use deadpool_postgres::{Object, Pool, Transaction}; use did_resolver::Resolver; use ipld_core::cid::Cid; use lexica::StrongRef; use metrics::counter; use parakeet_db::types::{ActorStatus, ActorSyncState}; use redis::aio::MultiplexedConnection; use redis::AsyncTypedCommands; use reqwest::{Client, StatusCode}; use std::path::PathBuf; use std::str::FromStr; use std::sync::Arc; use tokio::sync::watch::Receiver as WatchReceiver; use tokio::sync::Semaphore; use tokio_util::task::TaskTracker; use tracing::instrument; mod downloader; mod repo; mod types; const DL_DONE_KEY: &str = "bf_download_complete"; const PDS_SERVICE_ID: &str = "#atproto_pds"; // There's a 4MiB limit on parakeet-index, so break delta batches up if there's loads. // this should be plenty low enough to not trigger the size limit. (59k did slightly) const DELTA_BATCH_SIZE: usize = 32 * 1024; #[derive(Clone)] pub struct BackfillManagerInner { index_client: Option, tmp_dir: PathBuf, } pub struct BackfillManager { pool: Pool, redis: MultiplexedConnection, resolver: Arc, semaphore: Arc, opts: BackfillConfig, inner: BackfillManagerInner, } impl BackfillManager { pub async fn new( pool: Pool, redis: MultiplexedConnection, resolver: Arc, index_client: Option, opts: BackfillConfig, ) -> eyre::Result { let semaphore = Arc::new(Semaphore::new(opts.workers as usize)); Ok(BackfillManager { pool, redis, resolver, semaphore, inner: BackfillManagerInner { index_client, tmp_dir: PathBuf::from_str(&opts.download_tmp_dir)?, }, opts, }) } pub async fn run(mut self, stop: WatchReceiver) -> eyre::Result<()> { let tracker = TaskTracker::new(); tracker.spawn(downloader::downloader( self.redis.clone(), self.pool.clone(), self.resolver, self.inner.tmp_dir.clone(), self.opts.download_workers, self.opts.download_buffer, tracker.clone(), stop.clone(), )); loop { if stop.has_changed().unwrap_or(true) { tracker.close(); tracing::info!("stopping backfiller"); break; } let Some(job): Option = self.redis.lpop(DL_DONE_KEY, None).await? else { tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; continue; }; let p = self.semaphore.clone().acquire_owned().await?; let mut inner = self.inner.clone(); let mut conn = self.pool.get().await?; let mut rc = self.redis.clone(); tracker.spawn(async move { let _p = p; tracing::trace!("backfilling {job}"); if let Err(e) = backfill_actor(&mut conn, &mut rc, &mut inner, &job).await { tracing::error!(did = &job, "backfill failed: {e}"); counter!("backfill_failure").increment(1); db::backfill_job_write(&mut conn, &job, "failed.write") .await .unwrap(); } else { counter!("backfill_success").increment(1); db::backfill_job_write(&mut conn, &job, "successful") .await .unwrap(); } if let Err(e) = tokio::fs::remove_file(inner.tmp_dir.join(&job)).await { tracing::error!(did = &job, "failed to remove file: {e}"); } }); } tracker.wait().await; Ok(()) } } #[instrument(skip(conn, inner))] async fn backfill_actor( conn: &mut Object, rc: &mut MultiplexedConnection, inner: &mut BackfillManagerInner, did: &str, ) -> eyre::Result<()> { let mut t = conn.transaction().await?; t.execute("SET CONSTRAINTS ALL DEFERRED", &[]).await?; tracing::trace!("loading repo"); let (commit, mut deltas, copies) = repo::insert_repo(&mut t, rc, &inner.tmp_dir, did).await?; db::actor_set_repo_state(&mut t, did, &commit.rev, commit.data).await?; copies.submit(&mut t, did).await?; t.execute( "UPDATE actors SET sync_state=$2, last_indexed=$3 WHERE did=$1", &[&did, &ActorSyncState::Synced, &Utc::now().naive_utc()], ) .await?; handle_backfill_rows(&mut t, rc, &mut deltas, did, &commit.rev).await?; tracing::trace!("insertion finished"); if let Some(index_client) = &mut inner.index_client { // submit the deltas let delta_store = deltas .into_iter() .map(|((uri, typ), delta)| parakeet_index::AggregateDeltaReq { typ, uri: uri.to_string(), delta, }) .collect::>(); let mut read = 0; while read < delta_store.len() { let rem = delta_store.len() - read; let take = DELTA_BATCH_SIZE.min(rem); tracing::debug!("reading & submitting {take} deltas"); let deltas = delta_store[read..read + take].to_vec(); index_client .submit_aggregate_delta_batch(parakeet_index::AggregateDeltaBatchReq { deltas }) .await?; read += take; tracing::debug!("read {read} of {} deltas", delta_store.len()); } } t.commit().await?; Ok(()) } async fn handle_backfill_rows( conn: &mut Transaction<'_>, rc: &mut MultiplexedConnection, deltas: &mut impl AggregateDeltaStore, repo: &str, rev: &str, ) -> eyre::Result<()> { // `pull_backfill_rows` filters out anything before the last commit we pulled let backfill_rows = db::backfill_rows_get(conn, repo, rev).await?; for row in backfill_rows { // blindly unwrap-ing this CID as we've already parsed it and re-serialized it let repo_cid = Cid::from_str(&row.cid)?; db::actor_set_repo_state(conn, repo, &row.repo_ver, repo_cid).await?; // again, we've serialized this. let items: Vec = serde_json::from_value(row.data)?; for item in items { let Some((_, rkey)) = item.at_uri.rsplit_once("/") else { return Ok(()); }; match item.inner { BackfillItemInner::Create(record) | BackfillItemInner::Update(record) => { let Some(cid) = item.cid else { continue; }; indexer::index_op(conn, rc, deltas, repo, cid, record, &item.at_uri, rkey) .await? } BackfillItemInner::Delete => { indexer::index_op_delete( conn, rc, deltas, repo, item.collection, &item.at_uri, rkey, ) .await? } } } } // finally, clear the backfill table entries for this actor db::backfill_delete_rows(conn, repo).await?; Ok(()) } async fn check_pds_repo_status( client: &Client, pds: &str, repo: &str, ) -> eyre::Result> { let res = client .get(format!( "{pds}/xrpc/com.atproto.sync.getRepoStatus?did={repo}" )) .send() .await?; if [StatusCode::NOT_FOUND, StatusCode::BAD_REQUEST].contains(&res.status()) { return Ok(None); } Ok(res.json().await?) } #[derive(Debug, Default)] struct CopyStore { likes: Vec<(String, StrongRef, Option, DateTime)>, posts: Vec<(String, Cid, records::AppBskyFeedPost)>, reposts: Vec<(String, StrongRef, Option, DateTime)>, blocks: Vec<(String, String, DateTime)>, follows: Vec<(String, String, DateTime)>, list_items: Vec<(String, records::AppBskyGraphListItem)>, verifications: Vec<(String, Cid, records::AppBskyGraphVerification)>, records: Vec<(String, Cid)>, } impl CopyStore { async fn submit(self, t: &mut Transaction<'_>, did: &str) -> Result<(), tokio_postgres::Error> { db::copy::copy_likes(t, did, self.likes).await?; db::copy::copy_posts(t, did, self.posts).await?; db::copy::copy_reposts(t, did, self.reposts).await?; db::copy::copy_blocks(t, did, self.blocks).await?; db::copy::copy_follows(t, did, self.follows).await?; db::copy::copy_list_items(t, self.list_items).await?; db::copy::copy_verification(t, did, self.verifications).await?; db::copy::copy_records(t, did, self.records).await?; Ok(()) } fn push_record(&mut self, at_uri: &str, cid: Cid) { self.records.push((at_uri.to_string(), cid)) } }