Something went wrong. Try again.
Rust AppView - highly experimental! forked from parakeet.at/parakeet
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299use 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<parakeet_index::Client>, tmp_dir: PathBuf,}
pub struct BackfillManager { pool: Pool, redis: MultiplexedConnection, resolver: Arc<Resolver>, semaphore: Arc<Semaphore>, opts: BackfillConfig, inner: BackfillManagerInner,}
impl BackfillManager { pub async fn new( pool: Pool, redis: MultiplexedConnection, resolver: Arc<Resolver>, index_client: Option<parakeet_index::Client>, opts: BackfillConfig, ) -> eyre::Result<Self> { 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<bool>) -> 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<String> = 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, rc, 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::<Vec<_>>();
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<BackfillItem> = 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<Option<types::GetRepoStatusRes>> { 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<StrongRef>, DateTime<Utc>)>, posts: Vec<(String, Cid, records::AppBskyFeedPost)>, reposts: Vec<(String, StrongRef, Option<StrongRef>, DateTime<Utc>)>, blocks: Vec<(String, String, DateTime<Utc>)>, follows: Vec<(String, String, DateTime<Utc>)>, 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)) }}