diff --git a/src/backfill.rs b/src/backfill.rs index 526f86c..a53159c 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -8,9 +8,14 @@ pub async fn backfill( source: impl BundleSource + Send + 'static, dest: flume::Sender, source_workers: usize, + until: Option
, ) -> anyhow::Result<()> { // queue up the week bundles that should be available - let weeks = Arc::new(Mutex::new(Week::range(FIRST_WEEK..))); + let weeks = Arc::new(Mutex::new( + until + .map(|u| Week::range(FIRST_WEEK..u.into())) + .unwrap_or(Week::range(FIRST_WEEK..)), + )); weeks.lock().await.reverse(); let mut workers: JoinSet> = JoinSet::new(); diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index 5ac5745..034af95 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -37,6 +37,9 @@ enum Commands { /// Pass a postgres connection url like "postgresql://localhost:5432" #[arg(long)] to_postgres: Option, + /// Stop at the week ending before this date + #[arg(long)] + until: Option
, }, /// Scrape a PLC server, collecting ops into weekly bundles /// @@ -88,17 +91,18 @@ async fn main() { dir, source_workers, to_postgres, + until, } => { let (tx, rx) = flume::bounded(32); // big pages tokio::task::spawn(async move { if let Some(dir) = dir { log::info!("Reading weekly bundles from local folder {dir:?}"); - backfill(FolderSource(dir), tx, source_workers.unwrap_or(1)) + backfill(FolderSource(dir), tx, source_workers.unwrap_or(1), until) .await .unwrap(); } else { log::info!("Fetching weekly bundles from from {http}"); - backfill(HttpSource(http), tx, source_workers.unwrap_or(4)) + backfill(HttpSource(http), tx, source_workers.unwrap_or(4), until) .await .unwrap(); } diff --git a/src/plc_pg.rs b/src/plc_pg.rs index 311de94..be42650 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -92,7 +92,7 @@ pub async fn write_bulk(db: Db, pages: flume::Receiver) -> Result<() cid text not null, operation jsonb not null, nullified boolean not null, - createdAt timestamptz not null + "createdAt" timestamptz not null )"#, &[], ) @@ -135,9 +135,38 @@ pub async fn write_bulk(db: Db, pages: flume::Receiver) -> Result<() log::info!("copied in {n} rows"); tx.commit().await?; + log::info!("copy in time: {:?}", t0.elapsed()); - let dt = t0.elapsed(); - log::info!("backfill time: {dt:?}"); + log::info!("copying dids into plc table..."); + let n = client + .execute( + r#" + INSERT INTO dids + SELECT distinct did FROM allegedly_backfill + ON CONFLICT do nothing"#, + &[], + ) + .await?; + log::info!("{n} inserted; elapsed: {:?}", t0.elapsed()); + + log::info!("copying ops into plc table..."); + let n = client + .execute( + r#" + INSERT INTO operations (did, cid, operation, nullified, "createdAt") + SELECT did, cid, operation, nullified, "createdAt" FROM allegedly_backfill + ON CONFLICT do nothing"#, + &[], + ) + .await?; + log::info!("{n} inserted; elapsed: {:?}", t0.elapsed()); + + log::info!("clean up backfill table..."); + client + .execute(r#"DROP TABLE allegedly_backfill"#, &[]) + .await?; + + log::info!("total backfill time: {:?}", t0.elapsed()); Ok(()) }