diff --git a/migrations/0002_cid-for-records.sql b/migrations/0002_cid-for-records.sql new file mode 100644 index 0000000..5974744 --- /dev/null +++ b/migrations/0002_cid-for-records.sql @@ -0,0 +1,7 @@ +-- Add migration script here + +ALTER TABLE records ADD cid TEXT; + +UPDATE records SET cid = '' WHERE cid IS NULL; + +ALTER TABLE records ALTER COLUMN cid SET NOT NULL; \ No newline at end of file diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 7247f1e..4bf23f4 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -8,6 +8,7 @@ use std::str::FromStr; +use ipld_core::cid::multibase::Base; use jacquard::url::Url; use sqlx::{Pool, Postgres, query}; use thiserror::Error; @@ -20,8 +21,6 @@ use crate::{ pub mod load_car; pub mod parse_car; -const DB_MAX_REQ: usize = 65535; - #[derive(Error, Debug)] pub enum Error { #[error("Error parsing TID: {}", .0)] @@ -32,6 +31,8 @@ pub enum Error { Db(#[from] sqlx::Error), #[error("{}", .0)] ParseCar(#[from] crate::backfill::parse_car::Error), + #[error("Error processing cid: {}", .0)] + Cid(#[from] ipld_core::cid::Error), } pub async fn backfill( @@ -52,8 +53,19 @@ pub async fn backfill( // only real overhead is network latency which would be ~= anyway let _ = query!("DELETE FROM records").execute(conn).await?; - let data = parse_car(&car).await?; - let data = data.chunks(DB_MAX_REQ / 4); + let data = parse_car(&car) + .await? + .into_iter() + .map(|(collection, rkey, cid, value)| { + Ok::<_, Error>(( + collection, + rkey, + cid.to_string_of_base(Base::Base32Lower)?, + value, + )) + }) + .collect::, _>>()?; + let data = data.chunks(config::DB_MAX_REQ / 4); if let Some(time) = time { println!("Parsed car file ({:?})", time.elapsed()); @@ -61,13 +73,15 @@ pub async fn backfill( let time = time.map(|_| std::time::Instant::now()); for data in data { - let mut query = sqlx::QueryBuilder::new("INSERT INTO records(collection, rkey, record) "); + let mut query = + sqlx::QueryBuilder::new("INSERT INTO records(collection, rkey, cid, record) "); query.push_values( - data, + data.to_owned(), |mut b: sqlx::query_builder::Separated<'_, '_, Postgres, &'static str>, data| { - b.push_bind(data.0.0.clone()) - .push_bind(data.0.1.clone()) - .push_bind(data.1.clone()); + b.push_bind(data.0) + .push_bind(data.1) + .push_bind(data.2) + .push_bind(data.3); }, ); diff --git a/src/backfill/parse_car.rs b/src/backfill/parse_car.rs index ee53e8d..9f66bce 100644 --- a/src/backfill/parse_car.rs +++ b/src/backfill/parse_car.rs @@ -20,12 +20,14 @@ pub enum Error { IpldToJson(#[from] crate::utils::ipld_json::Error), #[error("Could not break {} into a collection and rkey", .0)] MalformedRecordKey(SmolStr), + #[error("Could not generate cid for commit: {}", .0)] + Cid(#[from] ipld_core::cid::Error), } -pub type AccountData = Vec<((String, String), Value)>; +pub type AccountData = Vec<(String, String, CidGeneric<64>, Value)>; pub async fn parse_car(car: &Car) -> Result { - let (keys, records): (Vec, Vec>) = + let (keys, record_cids): (Vec, Vec>) = car.mst.leaves().await?.into_iter().unzip(); // convert keys into (collection, rkey) @@ -47,23 +49,26 @@ pub async fn parse_car(car: &Car) -> Result { .collect::, _>>()?; // convert records into Value - let records = &records[..]; + let record_cids = &record_cids[..]; let records = car .storage - .get_many(records) + .get_many(record_cids) .await? .into_iter() .collect::>>() - .ok_or_else(|| Error::MissingCid)? - .into_iter() - .map(|x| { - let data = serde_ipld_dagcbor::from_slice::(&x)?; + .ok_or_else(|| Error::MissingCid)?; + + let records = zip(records, record_cids) + .map(|(bytes, cid)| { + let data = serde_ipld_dagcbor::from_slice::(&bytes)?; let value = ipld_to_json_value(&data)?; - Ok::<_, Error>(value) + Ok::<(CidGeneric<64>, Value), Error>((cid.to_owned(), value)) }) .collect::, _>>()?; - let data = zip(keys, records).collect::>(); + let data = zip(keys, records) + .map(|((collection, rkey), (cid, record))| (collection, rkey, cid, record)) + .collect(); Ok(data) } diff --git a/src/config.rs b/src/config.rs index 95e831d..784f91b 100644 --- a/src/config.rs +++ b/src/config.rs @@ -7,6 +7,8 @@ use jacquard::types::string::Did; use std::env; use std::sync::LazyLock; +pub const DB_MAX_REQ: usize = 65535; + // this should be loaded before the program starts any threads // if this panics threads that access it will be poisoned pub static USER: LazyLock> = LazyLock::new(|| {