diff --git a/src/backfill.rs b/src/backfill.rs
index a53159c..6d17600 100644
--- a/src/backfill.rs
+++ b/src/backfill.rs
@@ -26,16 +26,16 @@ pub async fn backfill(
let dest = dest.clone();
let source = source.clone();
workers.spawn(async move {
- log::info!("about to get weeks...");
+ log::trace!("about to get weeks...");
while let Some(week) = weeks.lock().await.pop() {
- log::info!(
+ log::trace!(
"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::trace!("week {}", Into::::into(week).to_rfc3339());
}
log::info!("done with the weeks ig");
Ok(())
diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs
index 034af95..5e475a0 100644
--- a/src/bin/allegedly.rs
+++ b/src/bin/allegedly.rs
@@ -37,6 +37,11 @@ enum Commands {
/// Pass a postgres connection url like "postgresql://localhost:5432"
#[arg(long)]
to_postgres: Option,
+ /// Delete all operations from the postgres db before starting
+ ///
+ /// only used if `--to-postgres` is present
+ #[arg(long, action)]
+ postgres_reset: bool,
/// Stop at the week ending before this date
#[arg(long)]
until: Option,
@@ -91,6 +96,7 @@ async fn main() {
dir,
source_workers,
to_postgres,
+ postgres_reset,
until,
} => {
let (tx, rx) = flume::bounded(32); // big pages
@@ -109,7 +115,7 @@ async fn main() {
});
if let Some(url) = to_postgres {
let db = Db::new(url.as_str()).await.unwrap();
- pages_to_pg(db, rx).await.unwrap();
+ pages_to_pg(db, rx, postgres_reset).await.unwrap();
} else {
pages_to_stdout(rx).await.unwrap();
}
diff --git a/src/plc_pg.rs b/src/plc_pg.rs
index be42650..a139ee6 100644
--- a/src/plc_pg.rs
+++ b/src/plc_pg.rs
@@ -71,47 +71,76 @@ impl Db {
}
}
-pub async fn write_bulk(db: Db, pages: flume::Receiver) -> Result<(), PgError> {
+/// Dump rows into an empty operations table quickly
+///
+/// you must run this after initializing the db with kysely migrations from the
+/// typescript app, but before inserting any content.
+///
+/// it's an invasive process: it will drop the indexes that kysely created (and
+/// restore them after) in order to get the backfill in as quickly as possible.
+///
+/// fails: if the backfill data violates the primary key constraint (unique did*cid)
+///
+/// panics: if the operations or dids tables are not empty, unless reset is true
+///
+/// recommended postgres setting: `max_wal_size=4GB` (or more)
+pub async fn write_bulk(
+ db: Db,
+ pages: flume::Receiver,
+ reset: bool,
+) -> 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 t0 = Instant::now();
let tx = client.transaction().await?;
- 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?;
+ let t_step = Instant::now();
+ for table in ["operations", "dids"] {
+ if reset {
+ let n = tx.execute(&format!("DELETE FROM {table}"), &[]).await?;
+ if n > 0 {
+ log::warn!("postgres reset: deleted {n} from {table}");
+ }
+ } else {
+ let n: i64 = tx
+ .query_one(&format!("SELECT count(*) FROM {table}"), &[])
+ .await?
+ .get(0);
+ if n > 0 {
+ panic!("postgres: {table} table was not empty and `reset` not requested");
+ }
+ }
+ }
+ log::trace!("tables clean: {:?}", t_step.elapsed());
+
+ let t_step = Instant::now();
+ tx.execute("ALTER TABLE operations SET UNLOGGED", &[])
+ .await?;
+ tx.execute("ALTER TABLE dids SET UNLOGGED", &[]).await?;
+ log::trace!("set tables unlogged: {:?}", t_step.elapsed());
+
+ let t_step = Instant::now();
+ tx.execute(r#"DROP INDEX "operations_createdAt_index""#, &[])
+ .await?;
+ tx.execute("DROP INDEX operations_did_createdat_idx", &[])
+ .await?;
+ log::trace!("indexes dropped: {:?}", t_step.elapsed());
+ let t_step = Instant::now();
+ log::trace!("starting binary COPY IN...");
let types = &[
- Type::TEXT,
Type::TEXT,
Type::JSONB,
+ Type::TEXT,
Type::BOOL,
Type::TIMESTAMPTZ,
];
- let t0 = Instant::now();
-
let sync = tx
- .copy_in("COPY allegedly_backfill FROM STDIN BINARY")
+ .copy_in(
+ r#"COPY operations (did, operation, cid, nullified, "createdAt") FROM STDIN BINARY"#,
+ )
.await?;
let mut writer = pin!(BinaryCopyInWriter::new(sync, types));
-
while let Ok(page) = pages.recv_async().await {
for s in page.ops {
let Ok(op) = serde_json::from_str::(&s) else {
@@ -122,50 +151,46 @@ pub async fn write_bulk(db: Db, pages: flume::Receiver) -> Result<()
.as_mut()
.write(&[
&op.did,
- &op.cid,
&Json(op.operation),
+ &op.cid,
&op.nullified,
&op.created_at,
])
.await?;
}
}
-
let n = writer.as_mut().finish().await?;
- log::info!("copied in {n} rows");
-
- tx.commit().await?;
- log::info!("copy in time: {:?}", t0.elapsed());
+ log::trace!("COPY IN wrote {n} ops: {:?}", t_step.elapsed());
- 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());
+ // CAUTION: these indexes MUST match up exactly with the kysely ones we dropped
+ let t_step = Instant::now();
+ tx.execute(
+ r#"CREATE INDEX operations_did_createdat_idx ON operations (did, "createdAt")"#,
+ &[],
+ )
+ .await?;
+ tx.execute(
+ r#"CREATE INDEX "operations_createdAt_index" ON operations ("createdAt")"#,
+ &[],
+ )
+ .await?;
+ log::trace!("indexes recreated: {:?}", t_step.elapsed());
- log::info!("copying ops into plc table...");
- let n = client
+ let t_step = Instant::now();
+ let n = tx
.execute(
- r#"
- INSERT INTO operations (did, cid, operation, nullified, "createdAt")
- SELECT did, cid, operation, nullified, "createdAt" FROM allegedly_backfill
- ON CONFLICT do nothing"#,
+ r#"INSERT INTO dids SELECT distinct did FROM operations"#,
&[],
)
.await?;
- log::info!("{n} inserted; elapsed: {:?}", t0.elapsed());
+ log::trace!("INSERT wrote {n} dids: {:?}", t_step.elapsed());
- log::info!("clean up backfill table...");
- client
- .execute(r#"DROP TABLE allegedly_backfill"#, &[])
- .await?;
+ let t_step = Instant::now();
+ tx.execute("ALTER TABLE dids SET LOGGED", &[]).await?;
+ tx.execute("ALTER TABLE operations SET LOGGED", &[]).await?;
+ log::trace!("set tables LOGGED: {:?}", t_step.elapsed());
+ tx.commit().await?;
log::info!("total backfill time: {:?}", t0.elapsed());
Ok(())