From baac66acba4fdacdf82b337ce1615bcc8e08986c Mon Sep 17 00:00:00 2001 From: phil Date: Mon, 15 Sep 2025 18:07:31 -0400 Subject: [PATCH] racing backfills --- src/bin/bundle-weekly.rs | 8 ++++-- src/bin/main.rs | 61 +++++++++++++++++++++++++--------------- src/weekly.rs | 19 ++++++++++--- 3 files changed, 58 insertions(+), 30 deletions(-) diff --git a/src/bin/bundle-weekly.rs b/src/bin/bundle-weekly.rs index 5dd8705..80f0cf1 100644 --- a/src/bin/bundle-weekly.rs +++ b/src/bin/bundle-weekly.rs @@ -1,4 +1,4 @@ -use allegedly::{bin_init, pages_to_weeks, poll_upstream}; +use allegedly::{Week, bin_init, pages_to_weeks, poll_upstream}; use clap::Parser; use std::path::PathBuf; use url::Url; @@ -23,7 +23,7 @@ struct Args { /// /// Must be a week-truncated unix timestamp #[arg(long, env)] - start_at: Option, // TODO!! + start_at: Option, } #[tokio::main] @@ -34,13 +34,15 @@ async fn main() -> anyhow::Result<()> { let mut url = args.upstream; url.set_path("/export"); + let after = args.start_at.map(|n| Week::from_n(n).into()); + log::trace!("ensure weekly output directory exists"); std::fs::create_dir_all(&args.dir)?; let (tx, rx) = flume::bounded(PAGE_QUEUE_SIZE); tokio::task::spawn(async move { - if let Err(e) = poll_upstream(None /*todo*/, url, tx).await { + if let Err(e) = poll_upstream(after, url, tx).await { log::error!("polling failed: {e}"); } else { log::warn!("poller finished ok (weird?)"); diff --git a/src/bin/main.rs b/src/bin/main.rs index 9dc3396..25fc551 100644 --- a/src/bin/main.rs +++ b/src/bin/main.rs @@ -77,28 +77,39 @@ async fn write_pages( rx: flume::Receiver, mut pg_client: tokio_postgres::Client, ) -> Result<(), anyhow::Error> { - let upsert_did = &pg_client - .prepare( - r#" - INSERT INTO dids (did) VALUES ($1) - ON CONFLICT DO NOTHING"#, - ) - .await - .unwrap(); + // TODO: one big upsert at the end from select distinct on the other table + + // let upsert_did = &pg_client + // .prepare( + // r#" + // INSERT INTO dids (did) VALUES ($1) + // ON CONFLICT DO NOTHING"#, + // ) + // .await + // .unwrap(); let insert_op = &pg_client .prepare( r#" INSERT INTO operations (did, operation, cid, nullified, "createdAt") - VALUES ($1, $2, $3, $4, $5)"#, - ) // TODO: check that it hasn't changed + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (did, cid) DO UPDATE + SET nullified = excluded.nullified, + "createdAt" = excluded."createdAt" + WHERE operations.nullified = excluded.nullified + OR operations."createdAt" = excluded."createdAt""#, + ) // idea: op is provable via cid, so leave it out. after did/cid (pk) that leaves nullified and createdAt + // that we want to notice changing. + // normal insert: no conflict, rows changed = 1 + // conflict (exact match): where clause passes, rows changed = 1 + // conflict (mismatch): where clause fails, rows changed = 0 (detect this and warn!) .await .unwrap(); while let Ok(page) = rx.recv_async().await { log::trace!("got a page..."); - let mut tx = pg_client.transaction().await.unwrap(); + let tx = pg_client.transaction().await.unwrap(); // TODO: probably figure out postgres COPY IN // for now just write everything into a transaction @@ -122,12 +133,12 @@ async fn write_pages( log::error!("ayeeeee just ignoring this error for now......"); continue; }; - let client = &tx; + // let client = &tx; - client.execute(upsert_did, &[&op.did]).await.unwrap(); + // client.execute(upsert_did, &[&op.did]).await.unwrap(); - let sp = tx.savepoint("op").await.unwrap(); - if let Err(e) = sp + // let sp = tx.savepoint("op").await.unwrap(); + let inserted = tx .execute( insert_op, &[ @@ -139,15 +150,19 @@ async fn write_pages( ], ) .await - { - if e.code() != Some(&tokio_postgres::error::SqlState::UNIQUE_VIOLATION) { - anyhow::bail!(e); - } - // TODO: assert that the row has not changed - log::warn!("ignoring dup"); - } else { - sp.commit().await.unwrap(); + .unwrap(); + if inserted != 1 { + log::warn!( + "possible log modification: {inserted} rows changed after upserting {op:?}" + ); } + // { + // if e.code() != Some(&tokio_postgres::error::SqlState::UNIQUE_VIOLATION) { + // anyhow::bail!(e); + // } + // // TODO: assert that the row has not changed + // log::warn!("ignoring dup"); + // } } tx.commit().await.unwrap(); diff --git a/src/weekly.rs b/src/weekly.rs index 175809d..03a0152 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -8,6 +8,17 @@ const WEEK_IN_SECONDS: i64 = 7 * 86400; #[derive(Debug, Clone, Copy, PartialEq)] pub struct Week(i64); +impl Week { + pub fn from_n(n: i64) -> Self { + Self(n) + } + pub fn n_ago(&self) -> i64 { + let Self(us) = self; + let Self(cur) = chrono::Utc::now().into(); + (cur - us) / WEEK_IN_SECONDS + } +} + impl From
for Week { fn from(dt: Dt) -> Self { let ts = dt.timestamp(); @@ -35,7 +46,6 @@ pub async fn pages_to_weeks(rx: flume::Receiver, dir: PathBuf) -> an let total_t0 = Instant::now(); let mut week_ops = 0; let mut week_t0 = total_t0; - let mut week = 0; while let Ok(page) = rx.recv_async().await { for mut s in page.ops { @@ -50,7 +60,8 @@ pub async fn pages_to_weeks(rx: flume::Receiver, dir: PathBuf) -> an let now = Instant::now(); log::info!( - "done week {week:3 } ({:10 }): {week_ops:7 } ({:5.0 }/s) ops, {:5 }k total ({:5.0 }/s)", + "done week {:3 } ({:10 }): {week_ops:7 } ({:5.0 }/s) ops, {:5 }k total ({:5.0 }/s)", + current_week.map(|w| -w.n_ago()).unwrap_or(0), current_week.unwrap_or(Week(0)).0, (week_ops as f64) / (now - week_t0).as_secs_f64(), total_ops / 1000, @@ -62,7 +73,6 @@ pub async fn pages_to_weeks(rx: flume::Receiver, dir: PathBuf) -> an current_week = Some(op_week); week_ops = 0; week_t0 = now; - week += 1; } s.push('\n'); // hack log::trace!("writing: {s}"); @@ -76,7 +86,8 @@ pub async fn pages_to_weeks(rx: flume::Receiver, dir: PathBuf) -> an encoder.shutdown().await?; let now = Instant::now(); log::info!( - "done week {week:3 } ({:10 }): {week_ops:7 } ({:5.0 }/s) ops, {:5 }k total ({:5.0 }/s)", + "done week {:3 } ({:10 }): {week_ops:7 } ({:5.0 }/s) ops, {:5 }k total ({:5.0 }/s)", + current_week.map(|w| -w.n_ago()).unwrap_or(0), current_week.unwrap_or(Week(0)).0, (week_ops as f64) / (now - week_t0).as_secs_f64(), total_ops / 1000, -- 2.51.2