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 {