From 672bd26ead90723a27f1fbee65d37b119aa80609 Mon Sep 17 00:00:00 2001 From: afterlifepro Date: Thu, 11 Dec 2025 23:29:20 +0000 Subject: [PATCH] migrate (getit) to using sqlx migrations also remove stupid car rev nonsense bc lowk was stupid lol its only on startup the backfill overhead is NAWT that high (like one second to parse and save on my pc + only on startup its _nawt_ that bad. bottleneck is network atp) optomizations to reduce downloaded car size (mst diffs) may be used later but i cba rn lol i also need to figure out some rev logic to dedupe repo ops :3 --- build.rs | 5 ++++ migrations/0001_init.sql | 14 ++++++++++ src/backfill/load_car.rs | 41 +++++++--------------------- src/backfill/mod.rs | 58 ++-------------------------------------- src/db.rs | 54 +++++-------------------------------- src/main.rs | 2 +- 6 files changed, 37 insertions(+), 137 deletions(-) create mode 100644 build.rs create mode 100644 migrations/0001_init.sql diff --git a/build.rs b/build.rs new file mode 100644 index 0000000..d506869 --- /dev/null +++ b/build.rs @@ -0,0 +1,5 @@ +// generated by `sqlx migrate build-script` +fn main() { + // trigger recompilation when a new migration is added + println!("cargo:rerun-if-changed=migrations"); +} diff --git a/migrations/0001_init.sql b/migrations/0001_init.sql new file mode 100644 index 0000000..b15fb32 --- /dev/null +++ b/migrations/0001_init.sql @@ -0,0 +1,14 @@ +-- Add migration script here +CREATE TABLE IF NOT EXISTS records ( + collection TEXT, + rkey TEXT, + record JSON NOT NULL, + PRIMARY KEY (collection, rkey) +); + +CREATE TABLE IF NOT EXISTS blobs ( + did TEXT, + cid TEXT, + blob bytea NOT NULL, + PRIMARY KEY (did, cid) +); \ No newline at end of file diff --git a/src/backfill/load_car.rs b/src/backfill/load_car.rs index 021dc6f..ecb683a 100644 --- a/src/backfill/load_car.rs +++ b/src/backfill/load_car.rs @@ -1,10 +1,9 @@ use jacquard::api::com_atproto; use jacquard::client::Agent; -use jacquard::types::{did::Did, tid::Tid}; +use jacquard::types::did::Did; use jacquard::url::Url; use jacquard::xrpc::XrpcExt; -use jacquard_repo::commit::Commit; -use jacquard_repo::{BlockStore, MemoryBlockStore, Mst}; +use jacquard_repo::{BlockStore, MemoryBlockStore, Mst, commit::Commit}; use std::sync::Arc; use thiserror::Error; @@ -14,32 +13,13 @@ pub enum Error { Client(#[from] jacquard::error::ClientError), #[error("Error loading car file: {}", .0)] Repo(#[from] jacquard_repo::RepoError), - #[error("Missing root block from car file (malformed car file)")] + #[error("Missing root from car file (malformed)")] MissingRoot, } pub struct Car { pub storage: MemoryBlockStore, pub mst: Mst, - pub rev: Tid, -} - -impl PartialEq for Car { - fn eq(&self, other: &Tid) -> bool { - &self.rev == other - } -} - -use std::cmp::Ordering; -impl PartialOrd for Car { - fn partial_cmp(&self, other: &Tid) -> Option { - match self.rev.compare_to(other) { - 1 => Some(Ordering::Greater), - 0 => Some(Ordering::Equal), - -1 => Some(Ordering::Less), - _ => None, - } - } } pub async fn load_car(user: Did<'_>, pds: Url) -> Result { @@ -54,15 +34,12 @@ pub async fn load_car(user: Did<'_>, pds: Url) -> Result { let storage = jacquard_repo::storage::MemoryBlockStore::new_from_blocks(car.blocks); - let root = storage - .get(&car.root) - .await? - .ok_or_else(|| Error::MissingRoot)?; - let root = Commit::from_cbor(&root)?; - let rev = root.rev().to_owned(); - let root = root.data(); + let Some(root) = storage.get(&car.root).await? else { + return Err(Error::MissingRoot); + }; + let root = Commit::from_cbor(&root)?.data; - let mst = Mst::load(Arc::new(storage.clone()), *root, None); + let mst = Mst::load(Arc::new(storage.clone()), root, None); - Ok(Car { storage, mst, rev }) + Ok(Car { storage, mst }) } diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 37eb04a..7247f1e 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -6,9 +6,9 @@ //! 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 std::str::FromStr; -use jacquard::{types::tid::Tid, url::Url}; +use jacquard::url::Url; use sqlx::{Pool, Postgres, query}; use thiserror::Error; @@ -28,12 +28,6 @@ pub enum Error { TidParse(#[from] jacquard::types::string::AtStrError), #[error("{}", .0)] GetCar(#[from] crate::backfill::load_car::Error), - #[error( - "The database claims to be more up to date than the PDS. -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." - )] - DbTidTooLow, #[error("Database error: {}", .0)] Db(#[from] sqlx::Error), #[error("{}", .0)] @@ -45,20 +39,6 @@ pub async fn backfill( 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() - ) - .fetch_one(conn) - .await - .ok() - .and_then(|x| x.rev) - { - Tid::from_str(&rev)? - } else { - Tid::from_time(0, 0) - }; - let pds = Url::from_str(&format!("https://{pds}/")).unwrap(); let car = load_car(config::USER.clone(), pds).await?; @@ -67,25 +47,6 @@ pub async fn backfill( } 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 => {} - // revisions are the same so we can skip backfill - Ordering::Equal => return Ok(()), - // db rev newer than car rev - // this means the db or car file is borked - // panic out and let the user deal with things - Ordering::Less => return Err(Error::DbTidTooLow), - // panic!( - // r"The database claims to be more up to date than the PDS. - // 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." - // ), - }; - }; - // erase all old records and return if it fails // we dont use diffs bc theyre complex and the overhead is minimal rn // only real overhead is network latency which would be ~= anyway @@ -123,20 +84,5 @@ pub async fn backfill( println!("Saved to database ({:?})", time.elapsed()); } - if let Err(err) = query!( - "UPDATE meta SET rev = $1 WHERE did = $2", - car.rev.to_string(), - config::USER.to_string() - ) - .execute(conn) - .await - { - // couldnt save tid so go nuclear - // this is program startup so its prolly safe lol - println!("Got error \"{}\"\nDeleting records and exiting...", err); - let _ = query!("DELETE FROM records").execute(conn).await?; - panic!() - }; - Ok(()) } diff --git a/src/db.rs b/src/db.rs index 7321c2d..ec79dba 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1,11 +1,10 @@ //! create a connection pool and setup tables before making avaliable use crate::config; -use sqlx::{Pool, Postgres, postgres::PgPool, query}; +use sqlx::{Pool, Postgres, migrate, postgres::PgPool}; -pub async fn init() -> Pool { - let conn = PgPool::connect(&config::POSTGRES_URL).await; - let conn = match conn { +pub async fn conn() -> Pool { + let conn = match PgPool::connect(&config::POSTGRES_URL).await { Ok(val) => val, Err(err) => { println!("Could not connect to the database. Got error {err}"); @@ -13,50 +12,9 @@ pub async fn init() -> Pool { } }; - // initialise db tables - if let Err(err) = query!( - "CREATE TABLE IF NOT EXISTS records ( - collection TEXT, - rkey TEXT, - record JSON NOT NULL, - PRIMARY KEY (collection, rkey) - );" - ) - .execute(&conn) - .await - { - println!("Creating table `records`: \n{err}"); - panic!("Could not instantiate db"); - }; - - if let Err(err) = query!( - "CREATE TABLE IF NOT EXISTS blobs ( - did TEXT, - cid TEXT, - blob bytea NOT NULL, - PRIMARY KEY (did, cid) - )" - ) - .execute(&conn) - .await - { - println!("Creating table `blobs`: \n{err}"); - panic!(); - }; - - if let Err(err) = query!( - "CREATE TABLE IF NOT EXISTS meta ( - did TEXT, - rev TEXT, - PRIMARY KEY (did) - );" - ) - .execute(&conn) - .await - { - println!("Creating table `meta`: \n{err}"); - panic!(); - }; + migrate!().run(&conn).await.unwrap_or_else(|err| { + panic!("Error running migrations: {}", err); + }); conn } diff --git a/src/main.rs b/src/main.rs index 7afe7da..a43aadd 100644 --- a/src/main.rs +++ b/src/main.rs @@ -14,7 +14,7 @@ struct Error; async fn main() -> Result<(), Error> { env_logger::init(); println!("User: {}", *config::USER); - let conn: Pool = db::init().await; + let conn: Pool = db::conn().await; println!("Database connected and initialized"); let pds = match utils::resolver::resolve(&config::USER).await { -- 2.51.2