diff --git a/src/backfill.rs b/src/backfill.rs index e58442f..82bc337 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -13,7 +13,7 @@ pub async fn backfill( dest: mpsc::Sender, source_workers: usize, until: Option
, -) -> anyhow::Result<()> { +) -> anyhow::Result<&'static str> { // queue up the week bundles that should be available let weeks = Arc::new(Mutex::new( until @@ -55,6 +55,10 @@ pub async fn backfill( res.inspect_err(|e| log::error!("problem joining source workers: {e}"))? .inspect_err(|e| log::error!("problem *from* source worker: {e}"))?; } - log::info!("finished fetching backfill in {:?}", t_step.elapsed()); - Ok(()) + log::info!( + "finished fetching backfill in {:?}. senders remaining: {}", + t_step.elapsed(), + dest.strong_count() + ); + Ok("backfill") } diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index 0484905..0311877 100644 --- a/src/bin/backfill.rs +++ b/src/bin/backfill.rs @@ -66,7 +66,7 @@ pub async fn run( catch_up, }: Args, ) -> anyhow::Result<()> { - let mut tasks = JoinSet::new(); + let mut tasks = JoinSet::>::new(); let (bulk_tx, bulk_out) = mpsc::channel(32); // bulk uses big pages @@ -147,10 +147,14 @@ pub async fn run( bulk_out, found_last_tx, )); - tasks.spawn(pages_to_pg(db, full_out)); + if catch_up { + tasks.spawn(pages_to_pg(db, full_out)); + } } else { tasks.spawn(pages_to_stdout(bulk_out, found_last_tx)); - tasks.spawn(pages_to_stdout(full_out, None)); + if catch_up { + tasks.spawn(pages_to_stdout(full_out, None)); + } } } @@ -168,7 +172,9 @@ pub async fn run( log::error!("a joinset task completed with error: {e}"); return Err(e); } - _ => {} + Ok(Ok(name)) => { + log::trace!("a task completed: {name:?}. {} left", tasks.len()); + } } } diff --git a/src/lib.rs b/src/lib.rs index a596b04..ccff5d7 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -83,21 +83,21 @@ impl From<&Op> for OpKey { pub async fn full_pages( mut rx: mpsc::Receiver, tx: mpsc::Sender, -) -> anyhow::Result<()> { +) -> anyhow::Result<&'static str> { while let Some(page) = rx.recv().await { let n = page.ops.len(); if n < 900 { let last_age = page.ops.last().map(|op| chrono::Utc::now() - op.created_at); let Some(age) = last_age else { log::info!("full_pages done, empty final page"); - return Ok(()); + return Ok("full pages (hmm)"); }; if age <= chrono::TimeDelta::hours(6) { log::info!("full_pages done, final page of {n} ops"); } else { log::warn!("full_pages finished with small page of {n} ops, but it's {age} old"); } - return Ok(()); + return Ok("full pages (cool)"); } log::trace!("full_pages: continuing with page of {n} ops"); tx.send(page).await?; @@ -110,7 +110,7 @@ pub async fn full_pages( pub async fn pages_to_stdout( mut rx: mpsc::Receiver, notify_last_at: Option>>, -) -> anyhow::Result<()> { +) -> anyhow::Result<&'static str> { let mut last_at = None; while let Some(page) = rx.recv().await { for op in &page.ops { @@ -128,7 +128,7 @@ pub async fn pages_to_stdout( log::error!("receiver for last_at dropped, can't notify"); }; } - Ok(()) + Ok("pages_to_stdout") } pub fn logo(name: &str) -> String { diff --git a/src/plc_pg.rs b/src/plc_pg.rs index 48a71b5..6b62da5 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -133,7 +133,10 @@ impl Db { } } -pub async fn pages_to_pg(db: Db, mut pages: mpsc::Receiver) -> anyhow::Result<()> { +pub async fn pages_to_pg( + db: Db, + mut pages: mpsc::Receiver, +) -> anyhow::Result<&'static str> { let mut client = db.connect().await?; let ops_stmt = client @@ -176,7 +179,7 @@ pub async fn pages_to_pg(db: Db, mut pages: mpsc::Receiver) -> anyho "no more pages. inserted {ops_inserted} ops and {dids_inserted} dids in {:?}", t0.elapsed() ); - Ok(()) + Ok("pages_to_pg") } /// Dump rows into an empty operations table quickly @@ -197,7 +200,7 @@ pub async fn backfill_to_pg( reset: bool, mut pages: mpsc::Receiver, notify_last_at: Option>>, -) -> anyhow::Result<()> { +) -> anyhow::Result<&'static str> { let mut client = db.connect().await?; let t0 = Instant::now(); @@ -272,6 +275,7 @@ pub async fn backfill_to_pg( last_at = last_at.filter(|&l| l >= s.last_at).or(Some(s.last_at)); } } + log::debug!("finished receiving bulk pages"); if let Some(notify) = notify_last_at { log::trace!("notifying last_at: {last_at:?}"); @@ -314,5 +318,5 @@ pub async fn backfill_to_pg( tx.commit().await?; log::info!("total backfill time: {:?}", t0.elapsed()); - Ok(()) + Ok("backfill_to_pg") } diff --git a/src/poll.rs b/src/poll.rs index 724643a..efcae2f 100644 --- a/src/poll.rs +++ b/src/poll.rs @@ -155,7 +155,7 @@ pub async fn poll_upstream( after: Option
, base: Url, dest: mpsc::Sender, -) -> anyhow::Result<()> { +) -> anyhow::Result<&'static str> { let mut tick = tokio::time::interval(UPSTREAM_REQUEST_INTERVAL); let mut prev_last: Option = after.map(Into::into); let mut boundary_state: Option = None;