diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 7b35fa4..94ed9a1 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -1,3 +1,11 @@ +//! backfill works as follows (https://docs.bsky.app/docs/advanced-guides/backfill) +//! +//! 1. resolve did -> pds +//! 2. get a car file from com.atproto.sync.getRepo +//! 3. extract collection, rkey, and cbor data from each leaf +//! 4. convert cbor data to json +//! 5. store in db (limit to DB_MAX_REQ / 4 to avoid err) + use std::{cmp::Ordering, str::FromStr}; use jacquard::{types::tid::Tid, url::Url}; @@ -32,21 +40,11 @@ Check your PDS repo is working and/or drop the database." ParseCarError(#[from] crate::backfill::parse_car::Error), } -/// backfill works as follows (https://docs.bsky.app/docs/advanced-guides/backfill) -/// -/// 1. resolve did -> pds -/// 2. stream com.atproto.sync.subscribeRepos to a buffer -/// 3. get a car file from com.atproto.sync.getRepo (diff if a rev is stored in database) -/// 4. apply car file diff to database (incl rev) -/// 5. start playing events from buffer -/// 1. drop all events from other users -/// 2. drop all events with a lower rev than current rev -/// 3. apply event & update rev -/// 4. (non blocking) get blobs if missing -/// 5. (non blocking) parse for strongref and store strongrefs -/// 6. (non blocking) trigger garbage collection of blobs and strongref -/// 6. once buffer is empty, parse events live -pub async fn backfill(pds: &str, conn: &Pool) -> Result<(), Error> { +pub async fn backfill( + pds: &str, + conn: &Pool, + time: Option, +) -> Result<(), Error> { let db_rev = if let Some(rev) = query!( "SELECT (rev) FROM meta WHERE did = $1", config::USER.to_string() @@ -64,8 +62,13 @@ pub async fn backfill(pds: &str, conn: &Pool) -> Result<(), Error> { let pds = Url::from_str(&format!("https://{pds}/")).unwrap(); let car = load_car(config::USER.clone(), pds).await?; - match car.partial_cmp(&db_rev) { - Some(val) => match val { + if let Some(time) = time { + println!("Downloaded car file ({:?})", time.elapsed()); + } + let time = time.map(|_| std::time::Instant::now()); + + if let Some(val) = car.partial_cmp(&db_rev) { + match val { // car rev newer than db rev // continue on; every other branch diverges Ordering::Greater => {} @@ -80,9 +83,7 @@ pub async fn backfill(pds: &str, conn: &Pool) -> Result<(), Error> { // Most likely either the PDS or repo is broken, or the database has been corrupted. // Check your PDS repo is working and/or drop the database." // ), - }, - // cant compare rev so assume all is ok and continue - None => {} + }; }; // erase all old records and return if it fails @@ -93,6 +94,11 @@ pub async fn backfill(pds: &str, conn: &Pool) -> Result<(), Error> { let data = parse_car(&car).await?; let mut data = data.chunks(DB_MAX_REQ / 4); + if let Some(time) = time { + println!("Parsed car file ({:?})", time.elapsed()); + } + let time = time.map(|_| std::time::Instant::now()); + while let Some(data) = data.next() { let mut query = sqlx::QueryBuilder::new("INSERT INTO records(collection, rkey, record) "); query.push_values( @@ -116,6 +122,10 @@ pub async fn backfill(pds: &str, conn: &Pool) -> Result<(), Error> { }; } + if let Some(time) = time { + println!("Saved to database ({:?})", time.elapsed()); + } + match query!( "UPDATE meta SET rev = $1 WHERE did = $2", car.rev.to_string(), diff --git a/src/config.rs b/src/config.rs index 4d69b72..95e831d 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,3 +1,8 @@ +//! get static and parsed environment variables +//! +//! USER is from env variable USER and parsed into a jacquard Did +//! POSTGRES_URL is from POSTGRES_USER, POSTGRES_PASSWORD, and POSTGRES_HOST + use jacquard::types::string::Did; use std::env; use std::sync::LazyLock; diff --git a/src/db.rs b/src/db.rs index 260ccfb..7321c2d 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1,3 +1,5 @@ +//! create a connection pool and setup tables before making avaliable + use crate::config; use sqlx::{Pool, Postgres, postgres::PgPool, query}; @@ -27,22 +29,6 @@ pub async fn init() -> Pool { panic!("Could not instantiate db"); }; - if let Err(err) = query!( - "CREATE TABLE IF NOT EXISTS foreign_records ( - did TEXT, - collection TEXT, - rkey TEXT, - record JSON NOT NULL, - PRIMARY KEY (did, collection, rkey) - );" - ) - .execute(&conn) - .await - { - println!("Creating table `foreign_records`: \n{err}"); - panic!(); - }; - if let Err(err) = query!( "CREATE TABLE IF NOT EXISTS blobs ( did TEXT, diff --git a/src/main.rs b/src/main.rs index 6c1f097..7afe7da 100644 --- a/src/main.rs +++ b/src/main.rs @@ -7,8 +7,11 @@ mod config; mod db; mod utils; +#[derive(Debug)] +struct Error; + #[tokio::main] -async fn main() -> Result<(), ()> { +async fn main() -> Result<(), Error> { env_logger::init(); println!("User: {}", *config::USER); let conn: Pool = db::init().await; @@ -19,12 +22,16 @@ async fn main() -> Result<(), ()> { Err(err) => panic!("{}", err), }; - let backfilled = backfill(&pds, &conn).await; - if let Err(err) = backfilled { + println!("Starting backfill"); + let timer = std::time::Instant::now(); + + if let Err(err) = backfill(&pds, &conn, Some(timer)).await { println!("{}", err); - return Err(()); + return Err(Error); }; + println!("Backfill complete. Took {:?}", timer.elapsed()); + println!("Completed sucessfully!"); Ok(()) } diff --git a/src/utils/ipld_json.rs b/src/utils/ipld_json.rs index f386ce8..e2eff23 100644 --- a/src/utils/ipld_json.rs +++ b/src/utils/ipld_json.rs @@ -1,3 +1,13 @@ +//! convert an ipld_core::ipld::Ipld enum into a serde_json::value::Value in the atproto data model +//! +//! a specific helper is required for this as Bytes and Link have differing representations to how serde_json handles them by default +//! +//! in general. types are naievely converted. the following types have special cases: +//! - `integer`: this could throw an error if the number is `x` in `i64::MIN < x < u64::MAX` +//! - `float`: always issues a warning since this is technically illegal. If its NaN or infinity, this errors as they cant be represented in json +//! - `bytes`: atproto JSON represents them as `{"$bytes": "BASE 64 NO PADDING"}`, but serde_json defaults to `[u8]` +//! - `link`: atproto JSON represents them as `{"$link": "BASE 32 NO PADDING"}`, but serde_json defaults to `[u8]` + use base64::{Engine, prelude::BASE64_STANDARD_NO_PAD}; use ipld_core::{cid::multibase::Base, ipld::Ipld}; use log::warn; @@ -48,7 +58,8 @@ pub fn ipld_to_json_value(data: &Ipld) -> Result { .map(|(k, v)| Ok::<_, Error>((k.clone(), ipld_to_json_value(v)?))) .collect::, _>>()?, ), - Ipld::Link(cid) => json!({"$link": - cid.to_string_of_base(Base::Base32Lower)? }), + Ipld::Link(cid) => json!({ + "$link": cid.to_string_of_base(Base::Base32Lower)? + }), }) } diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 6b6c0fc..5795a1c 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -1,2 +1,6 @@ +//! contains utility functions +//! +//! see sub modules for more details + pub mod ipld_json; pub mod resolver; diff --git a/src/utils/resolver.rs b/src/utils/resolver.rs index d68adeb..26b7dc6 100644 --- a/src/utils/resolver.rs +++ b/src/utils/resolver.rs @@ -1,3 +1,5 @@ +//! resolve a Did to a pds domain + use jacquard::prelude::IdentityResolver; use jacquard::types::did::Did; use thiserror::Error;