diff --git a/src/backfill.rs b/src/backfill.rs index 744a0f3..526f86c 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -1,5 +1,6 @@ use crate::{BundleSource, Dt, ExportPage, Week, week_to_pages}; -use tokio::task::JoinSet; +use std::sync::Arc; +use tokio::{sync::Mutex, task::JoinSet}; const FIRST_WEEK: Week = Week::from_n(1668643200); @@ -9,33 +10,35 @@ pub async fn backfill( source_workers: usize, ) -> anyhow::Result<()> { // queue up the week bundles that should be available - let (week_tx, week_rx) = flume::bounded(1024); // work queue - let mut week = FIRST_WEEK; - while week.is_immutable() { - week_tx.try_send(week)?; // if this fails, something has gone really wrong or we're farrrr in the future - week = week.next(); - } + let weeks = Arc::new(Mutex::new(Week::range(FIRST_WEEK..))); + weeks.lock().await.reverse(); let mut workers: JoinSet> = JoinSet::new(); // spin up the fetchers to work in parallel for w in 0..source_workers { - let weeks = week_rx.clone(); + let weeks = weeks.clone(); let dest = dest.clone(); let source = source.clone(); workers.spawn(async move { - while let Ok(week) = weeks.recv_async().await { + log::info!("about to get weeks..."); + + while let Some(week) = weeks.lock().await.pop() { log::info!( "worker {w}: fetching week {} (-{})", Into::
::into(week).to_rfc3339(), week.n_ago(), ); week_to_pages(source.clone(), week, dest.clone()).await?; + log::info!("done a week"); } + log::info!("done with the weeks ig"); Ok(()) }); } + // TODO: handle missing/failed weeks + // wait for them to finish while let Some(res) = workers.join_next().await { res??; diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index 82e71e9..5ac5745 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -1,4 +1,7 @@ -use allegedly::{Db, Dt, ExportPage, FolderSource, HttpSource, backfill, bin_init, pages_to_weeks, poll_upstream, write_bulk as pages_to_pg}; +use allegedly::{ + Db, Dt, ExportPage, FolderSource, HttpSource, backfill, bin_init, pages_to_weeks, + poll_upstream, write_bulk as pages_to_pg, +}; use clap::{Parser, Subcommand}; use std::path::PathBuf; use url::Url; @@ -25,9 +28,10 @@ enum Commands { #[arg(long)] dir: Option, /// Parallel bundle fetchers + /// + /// Default: 4 for http fetches, 1 for local folder #[arg(long)] - #[clap(default_value = "4")] - source_workers: usize, + source_workers: Option, /// Bulk load into did-method-plc-compatible postgres instead of stdout /// /// Pass a postgres connection url like "postgresql://localhost:5432" @@ -64,11 +68,12 @@ enum Commands { } async fn pages_to_stdout(rx: flume::Receiver) -> Result<(), flume::RecvError> { - loop { - for op in rx.recv_async().await?.ops { + while let Ok(page) = rx.recv_async().await { + for op in page.ops { println!("{op}") } } + Ok(()) } #[tokio::main] @@ -84,22 +89,22 @@ async fn main() { source_workers, to_postgres, } => { - let (tx, rx) = flume::bounded(1024); // big pages + 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) + backfill(FolderSource(dir), tx, source_workers.unwrap_or(1)) .await .unwrap(); } else { log::info!("Fetching weekly bundles from from {http}"); - backfill(HttpSource(http), tx, source_workers) + backfill(HttpSource(http), tx, source_workers.unwrap_or(4)) .await .unwrap(); } }); if let Some(url) = to_postgres { - let db = Db::new(url.as_str()); + let db = Db::new(url.as_str()).await.unwrap(); pages_to_pg(db, rx).await.unwrap(); } else { pages_to_stdout(rx).await.unwrap(); @@ -122,7 +127,7 @@ async fn main() { let mut url = args.upstream; url.set_path("/export"); let start_at = after.or_else(|| Some(chrono::Utc::now())); - let (tx, rx) = flume::bounded(0); // rendezvous, don't read ahead + let (tx, rx) = flume::bounded(1); tokio::task::spawn(async move { poll_upstream(start_at, url, tx).await.unwrap() }); pages_to_stdout(rx).await.unwrap(); } diff --git a/src/plc_pg.rs b/src/plc_pg.rs index 70867d1..311de94 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -1,7 +1,12 @@ use crate::{ExportPage, Op}; -use tokio_postgres::{Client, types::{Type, Json}, Error as PgError, NoTls, connect, binary_copy::BinaryCopyInWriter}; use std::pin::pin; - +use std::time::Instant; +use tokio_postgres::{ + Client, Error as PgError, NoTls, + binary_copy::BinaryCopyInWriter, + connect, + types::{Json, Type}, +}; /// a little tokio-postgres helper /// @@ -13,10 +18,40 @@ pub struct Db { } impl Db { - pub fn new(pg_uri: &str) -> Self { - Self { + pub async fn new(pg_uri: &str) -> Result { + // we're going to interact with did-method-plc's database, so make sure + // it's what we expect: check for db migrations. + log::trace!("checking migrations..."); + let (client, connection) = connect(pg_uri, NoTls).await?; + let connection_task = tokio::task::spawn(async move { + connection + .await + .inspect_err(|e| log::error!("connection ended with error: {e}")) + .unwrap(); + }); + let migrations: Vec = client + .query("SELECT name FROM kysely_migration ORDER BY name", &[]) + .await? + .iter() + .map(|row| row.get(0)) + .collect(); + assert_eq!( + &migrations, + &[ + "_20221020T204908820Z", + "_20230223T215019669Z", + "_20230406T174552885Z", + "_20231128T203323431Z", + ] + ); + drop(client); + // make sure the connection worker thing doesn't linger + connection_task.await?; + log::info!("db connection succeeded and plc migrations appear as expected"); + + Ok(Self { pg_uri: pg_uri.to_string(), - } + }) } pub async fn connect(&self) -> Result { @@ -36,24 +71,32 @@ impl Db { } } -pub async fn write_bulk( - db: Db, - pages: flume::Receiver, -) -> Result<(), PgError> { +pub async fn write_bulk(db: Db, pages: flume::Receiver) -> Result<(), PgError> { let mut client = db.connect().await?; + + // TODO: maybe we want to be more cautious + client + .execute( + r#" + DROP TABLE IF EXISTS allegedly_backfill"#, + &[], + ) + .await?; + let tx = client.transaction().await?; - tx - .execute(r#" - CREATE TABLE backfill ( + tx.execute( + r#" + CREATE UNLOGGED TABLE allegedly_backfill ( did text not null, cid text not null, operation jsonb not null, nullified boolean not null, createdAt timestamptz not null - )"#, &[]) - .await?; - + )"#, + &[], + ) + .await?; let types = &[ Type::TEXT, @@ -62,8 +105,11 @@ pub async fn write_bulk( Type::BOOL, Type::TIMESTAMPTZ, ]; + let t0 = Instant::now(); - let sync = tx.copy_in("COPY backfill FROM STDIN BINARY").await?; + let sync = tx + .copy_in("COPY allegedly_backfill FROM STDIN BINARY") + .await?; let mut writer = pin!(BinaryCopyInWriter::new(sync, types)); while let Ok(page) = pages.recv_async().await { @@ -72,13 +118,16 @@ pub async fn write_bulk( log::warn!("ignoring unparseable op: {s:?}"); continue; }; - writer.as_mut().write(&[ - &op.did, - &op.cid, - &Json(op.operation), - &op.nullified, - &op.created_at, - ]).await?; + writer + .as_mut() + .write(&[ + &op.did, + &op.cid, + &Json(op.operation), + &op.nullified, + &op.created_at, + ]) + .await?; } } @@ -87,5 +136,8 @@ pub async fn write_bulk( tx.commit().await?; + let dt = t0.elapsed(); + log::info!("backfill time: {dt:?}"); + Ok(()) } diff --git a/src/weekly.rs b/src/weekly.rs index 64d80f4..22abf61 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -3,6 +3,7 @@ use async_compression::tokio::bufread::GzipDecoder; use async_compression::tokio::write::GzipEncoder; use core::pin::pin; use std::future::Future; +use std::ops::{Bound, RangeBounds}; use std::path::PathBuf; use tokio::{ fs::File, @@ -14,29 +15,56 @@ use url::Url; const WEEK_IN_SECONDS: i64 = 7 * 86_400; -#[derive(Debug, Clone, Copy, PartialEq)] +#[derive(Debug, Clone, Copy, PartialEq, PartialOrd)] pub struct Week(i64); impl Week { pub const fn from_n(n: i64) -> Self { Self(n) } + pub fn range(r: impl RangeBounds) -> Vec { + let first = match r.start_bound() { + Bound::Included(week) => *week, + Bound::Excluded(week) => week.next(), + Bound::Unbounded => panic!("week range must have a defined start bound"), + }; + let last = match r.end_bound() { + Bound::Included(week) => *week, + Bound::Excluded(week) => week.prev(), + Bound::Unbounded => Self(Self::nullification_cutoff()).prev(), + }; + let mut out = Vec::new(); + let mut current = first; + while current <= last { + out.push(current); + current = current.next(); + } + out + } pub fn n_ago(&self) -> i64 { - let Self(us) = self; - let Self(cur) = chrono::Utc::now().into(); - (cur - us) / WEEK_IN_SECONDS + let now = chrono::Utc::now().timestamp(); + (now - self.0) / WEEK_IN_SECONDS + } + pub fn n_until(&self, other: Week) -> i64 { + let Self(until) = other; + (until - self.0) / WEEK_IN_SECONDS } pub fn next(&self) -> Week { Self(self.0 + WEEK_IN_SECONDS) } - /// is the plc log for this week entirely outside the 72h nullification window + pub fn prev(&self) -> Week { + Self(self.0 - WEEK_IN_SECONDS) + } + /// whether the plc log for this week outside the 72h nullification window /// /// plus one hour for safety (week must have ended > 73 hours ago) pub fn is_immutable(&self) -> bool { + self.next().0 <= Self::nullification_cutoff() + } + fn nullification_cutoff() -> i64 { const HOUR_IN_SECONDS: i64 = 3600; let now = chrono::Utc::now().timestamp(); - let nullification_cutoff = now - (73 * HOUR_IN_SECONDS); - self.next().0 <= nullification_cutoff + now - (73 * HOUR_IN_SECONDS) } }