diff --git a/Cargo.lock b/Cargo.lock index 63e6c79..06b02d1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -37,7 +37,6 @@ dependencies = [ "governor", "http-body-util", "k256", - "log", "multibase", "native-tls", "opentelemetry", diff --git a/Cargo.toml b/Cargo.toml index 03c5d08..dbdccba 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -22,7 +22,6 @@ fjall = "3.0.2" futures = "0.3.31" governor = "0.10.1" http-body-util = "0.1.3" -log = "0.4.28" native-tls = "0.2.14" opentelemetry = "0.30.0" opentelemetry-otlp = { version = "0.30.0" } diff --git a/src/backfill.rs b/src/backfill.rs index 82bc337..e99218c 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -25,7 +25,7 @@ pub async fn backfill( let mut workers: JoinSet> = JoinSet::new(); let t_step = Instant::now(); - log::info!( + tracing::info!( "fetching backfill for {} weeks with {source_workers} workers...", weeks.lock().await.len() ); @@ -38,12 +38,12 @@ pub async fn backfill( workers.spawn(async move { while let Some(week) = weeks.lock().await.pop() { let when = Into::
::into(week).to_rfc3339(); - log::trace!("worker {w}: fetching week {when} (-{})", week.n_ago()); + tracing::trace!("worker {w}: fetching week {when} (-{})", week.n_ago()); week_to_pages(source.clone(), week, dest.clone()) .await - .inspect_err(|e| log::error!("failing week_to_pages: {e}"))?; + .inspect_err(|e| tracing::error!("failing week_to_pages: {e}"))?; } - log::info!("done with the weeks ig"); + tracing::info!("done with the weeks ig"); Ok(()) }); } @@ -52,10 +52,10 @@ pub async fn backfill( // wait for the big backfill to finish while let Some(res) = workers.join_next().await { - res.inspect_err(|e| log::error!("problem joining source workers: {e}"))? - .inspect_err(|e| log::error!("problem *from* source worker: {e}"))?; + res.inspect_err(|e| tracing::error!("problem joining source workers: {e}"))? + .inspect_err(|e| tracing::error!("problem *from* source worker: {e}"))?; } - log::info!( + tracing::info!( "finished fetching backfill in {:?}. senders remaining: {}", t_step.elapsed(), dest.strong_count() diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index 4a5050a..2f6765c 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -95,7 +95,7 @@ async fn main() -> anyhow::Result<()> { let matches = Cli::command().get_matches(); let name = matches.subcommand().map(|(name, _)| name).unwrap_or("???"); bin_init(args.command.enable_otel()); - log::info!("{}", logo(name)); + tracing::info!("{}", logo(name)); let globals = args.globals.clone(); @@ -116,7 +116,7 @@ async fn main() -> anyhow::Result<()> { .await .expect("to poll upstream") }); - log::trace!("ensuring output directory exists"); + tracing::trace!("ensuring output directory exists"); create_dir_all(&dest) .await .expect("to ensure output dir exists"); @@ -143,6 +143,6 @@ async fn main() -> anyhow::Result<()> { .expect("to write pages to stdout"); } } - log::info!("whew, {:?}. goodbye!", t0.elapsed()); + tracing::info!("whew, {:?}. goodbye!", t0.elapsed()); Ok(()) } diff --git a/src/bin/audit.rs b/src/bin/audit.rs index 3e71670..2234beb 100644 --- a/src/bin/audit.rs +++ b/src/bin/audit.rs @@ -41,19 +41,19 @@ pub async fn run(globals: GlobalArgs, Args { fjall, fix, drop }: Args) -> anyhow while let Some(next) = tasks.join_next().await { match next { Err(e) if e.is_panic() => { - log::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); + tracing::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); return Err(e.into()); } Err(e) => { - log::error!("a joinset task failed to join: {e}"); + tracing::error!("a joinset task failed to join: {e}"); return Err(e.into()); } Ok(Err(e)) => { - log::error!("a joinset task completed with error: {e}"); + tracing::error!("a joinset task completed with error: {e}"); return Err(e); } Ok(Ok(name)) => { - log::trace!("a task completed: {name:?}. {} left", tasks.len()); + tracing::trace!("a task completed: {name:?}. {} left", tasks.len()); } } } @@ -76,7 +76,7 @@ struct CliArgs { async fn main() -> anyhow::Result<()> { let args = CliArgs::parse(); bin_init(args.instrumentation.enable_opentelemetry); - log::info!("{}", logo("audit")); + tracing::info!("{}", logo("audit")); run(args.globals, args.args).await?; Ok(()) } diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index dc76bd6..ec7049c 100644 --- a/src/bin/backfill.rs +++ b/src/bin/backfill.rs @@ -99,13 +99,13 @@ pub async fn run( if no_bulk { // simple mode, just poll upstream from teh beginning if http != DEFAULT_HTTP.parse()? { - log::warn!("ignoring non-default bulk http setting since --no-bulk was set"); + tracing::warn!("ignoring non-default bulk http setting since --no-bulk was set"); } if let Some(d) = dir { - log::warn!("ignoring bulk dir setting ({d:?}) since --no-bulk was set."); + tracing::warn!("ignoring bulk dir setting ({d:?}) since --no-bulk was set."); } if let Some(u) = until { - log::warn!( + tracing::warn!( "ignoring `until` setting ({u:?}) since --no-bulk was set. (feature request?)" ); } @@ -113,9 +113,9 @@ pub async fn run( upstream.set_path("/export"); let throttle = Duration::from_millis(upstream_throttle_ms); if let Some(fjall_path) = to_fjall { - log::trace!("opening fjall db at {fjall_path:?}..."); + tracing::trace!("opening fjall db at {fjall_path:?}..."); let db = FjallDb::open(&fjall_path)?; - log::trace!("opened fjall db"); + tracing::trace!("opened fjall db"); let (poll_tx, poll_out) = mpsc::channel::(128); // normal/small pages let (full_tx, full_out) = mpsc::channel::(1); // don't need to buffer at this filter @@ -166,9 +166,9 @@ pub async fn run( // set up sinks if let Some(pg_url) = to_postgres { - log::trace!("connecting to postgres..."); + tracing::trace!("connecting to postgres..."); let db = Db::new(pg_url.as_str(), postgres_cert).await?; - log::trace!("connected to postgres"); + tracing::trace!("connected to postgres"); tasks.spawn(backfill_to_pg(db.clone(), reset, bulk_out, found_last_tx)); if catch_up { @@ -185,19 +185,19 @@ pub async fn run( while let Some(next) = tasks.join_next().await { match next { Err(e) if e.is_panic() => { - log::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); + tracing::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); return Err(e.into()); } Err(e) => { - log::error!("a joinset task failed to join: {e}"); + tracing::error!("a joinset task failed to join: {e}"); return Err(e.into()); } Ok(Err(e)) => { - log::error!("a joinset task completed with error: {e}"); + tracing::error!("a joinset task completed with error: {e}"); return Err(e); } Ok(Ok(name)) => { - log::trace!("a task completed: {name:?}. {} left", tasks.len()); + tracing::trace!("a task completed: {name:?}. {} left", tasks.len()); } } } @@ -218,7 +218,7 @@ struct CliArgs { async fn main() -> anyhow::Result<()> { let args = CliArgs::parse(); bin_init(false); - log::info!("{}", logo("backfill")); + tracing::info!("{}", logo("backfill")); run(args.globals, args.args).await?; Ok(()) } diff --git a/src/bin/mirror.rs b/src/bin/mirror.rs index c6e3843..d3d68ad 100644 --- a/src/bin/mirror.rs +++ b/src/bin/mirror.rs @@ -105,7 +105,7 @@ pub async fn run( if let Some(ref experimental_domain) = experimental_acme_domain { domains.push(experimental_domain.clone()) } - log::info!("configuring acme for https at {domains:?}..."); + tracing::info!("configuring acme for https at {domains:?}..."); ListenConf::Acme { domains, cache_path, @@ -127,16 +127,16 @@ pub async fn run( if let Some(fjall_path) = wrap_fjall { let db = FjallDb::open(&fjall_path)?; if compact_fjall { - log::info!("compacting fjall..."); + tracing::info!("compacting fjall..."); db.compact()?; } - log::debug!("getting the latest seq from fjall..."); + tracing::debug!("getting the latest seq from fjall..."); let latest_seq = db .get_latest()? .map(|(seq, _)| seq) .expect("there to be at least one op in the db. did you backfill?"); - log::info!("starting seq polling from seq {latest_seq}..."); + tracing::info!("starting seq polling from seq {latest_seq}..."); let (send_page, recv_page) = mpsc::channel::(8); @@ -153,7 +153,7 @@ pub async fn run( tasks.spawn(async move { let mut current_seq = latest_seq; loop { - log::info!("seq polling from seq {current_seq}"); + tracing::info!("seq polling from seq {current_seq}"); let (inner_tx, mut inner_rx) = mpsc::channel::(8); // run poller; it ends only when the channel closes @@ -190,7 +190,7 @@ pub async fn run( current_seq = last_seq_from_poll; // switch to streaming - log::info!("caught up at seq {current_seq}, switching to /export/stream"); + tracing::info!("caught up at seq {current_seq}, switching to /export/stream"); let (stream_inner_tx, mut stream_inner_rx) = mpsc::channel::(8); let stream_task = tokio::spawn(tail_upstream_stream( Some(current_seq), @@ -210,9 +210,9 @@ pub async fn run( // stream ended/errored — loop back to polling to resync match stream_task.await { - Ok(Ok(())) => log::info!("stream closed cleanly, resyncing via poll"), - Ok(Err(e)) => log::warn!("stream error: {e}, resyncing via poll"), - Err(e) => log::warn!("stream task join error: {e}"), + Ok(Ok(())) => tracing::info!("stream closed cleanly, resyncing via poll"), + Ok(Err(e)) => tracing::warn!("stream error: {e}, resyncing via poll"), + Err(e) => tracing::warn!("stream task join error: {e}"), } } }); @@ -230,12 +230,12 @@ pub async fn run( ))?; let db = Db::new(wrap_pg.as_str(), wrap_pg_cert).await?; - log::debug!("getting the latest op from the db..."); + tracing::debug!("getting the latest op from the db..."); let latest = db .get_latest() .await? .expect("there to be at least one op in the db. did you backfill?"); - log::debug!("starting polling from {latest}..."); + tracing::debug!("starting polling from {latest}..."); let (send_page, recv_page) = mpsc::channel(8); @@ -256,19 +256,19 @@ pub async fn run( while let Some(next) = tasks.join_next().await { match next { Err(e) if e.is_panic() => { - log::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); + tracing::error!("a joinset task panicked: {e}. bailing now. (should we panic?)"); return Err(e.into()); } Err(e) => { - log::error!("a joinset task failed to join: {e}"); + tracing::error!("a joinset task failed to join: {e}"); return Err(e.into()); } Ok(Err(e)) => { - log::error!("a joinset task completed with error: {e}"); + tracing::error!("a joinset task completed with error: {e}"); return Err(e); } Ok(Ok(name)) => { - log::trace!("a task completed: {name:?}. {} left", tasks.len()); + tracing::trace!("a task completed: {name:?}. {} left", tasks.len()); } } } @@ -294,7 +294,7 @@ struct CliArgs { async fn main() -> anyhow::Result<()> { let args = CliArgs::parse(); bin_init(args.instrumentation.enable_opentelemetry); - log::info!("{}", logo("mirror")); + tracing::info!("{}", logo("mirror")); run(args.globals, args.args, !args.wrap_mode).await?; Ok(()) } diff --git a/src/cached_value.rs b/src/cached_value.rs index 8cc9e25..b23edbf 100644 --- a/src/cached_value.rs +++ b/src/cached_value.rs @@ -16,10 +16,10 @@ struct ExpiringValue { impl ExpiringValue { fn get(&self, now: Instant) -> Option { if now <= self.expires { - log::trace!("returning val (fresh for {:?})", self.expires - now); + tracing::trace!("returning val (fresh for {:?})", self.expires - now); Some(self.value.clone()) } else { - log::trace!("hiding expired val"); + tracing::trace!("hiding expired val"); None } } @@ -50,7 +50,7 @@ impl> CachedValue { if let Some(v) = val.as_ref().and_then(|v| v.get(now)) { return Ok(v); } - log::debug!( + tracing::debug!( "value {}, fetching...", if val.is_some() { "expired" @@ -62,8 +62,8 @@ impl> CachedValue { .fetcher .fetch() .await - .inspect_err(|e| log::warn!("value fetch failed, next access will retry: {e}"))?; - log::debug!("fetched ok, saving a copy for cache."); + .inspect_err(|e| tracing::warn!("value fetch failed, next access will retry: {e}"))?; + tracing::debug!("fetched ok, saving a copy for cache."); *val = Some(ExpiringValue { value: new.clone(), expires: now + self.validitiy, diff --git a/src/lib.rs b/src/lib.rs index 468f21a..7bfc92d 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -140,17 +140,17 @@ pub async fn full_pages( 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"); + tracing::info!("full_pages done, empty final page"); return Ok("full pages (hmm)"); }; if age <= chrono::TimeDelta::hours(6) { - log::info!("full_pages done, final page of {n} ops"); + tracing::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"); + tracing::warn!("full_pages finished with small page of {n} ops, but it's {age} old"); } return Ok("full pages (cool)"); } - log::trace!("full_pages: continuing with page of {n} ops"); + tracing::trace!("full_pages: continuing with page of {n} ops"); tx.send(page).await?; } Err(anyhow::anyhow!( @@ -167,17 +167,17 @@ pub async fn full_pages_seq( 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"); + tracing::info!("full_pages done, empty final page"); return Ok("full pages (hmm)"); }; if age <= chrono::TimeDelta::hours(6) { - log::info!("full_pages done, final page of {n} ops"); + tracing::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"); + tracing::warn!("full_pages finished with small page of {n} ops, but it's {age} old"); } return Ok("full pages (cool)"); } - log::trace!("full_pages: continuing with page of {n} ops"); + tracing::trace!("full_pages: continuing with page of {n} ops"); tx.send(page).await?; } Err(anyhow::anyhow!( @@ -201,9 +201,9 @@ pub async fn pages_to_stdout( } } if let Some(notify) = notify_last_at { - log::trace!("notifying last_at: {last_at:?}"); + tracing::trace!("notifying last_at: {last_at:?}"); if notify.send(last_at).is_err() { - log::error!("receiver for last_at dropped, can't notify"); + tracing::error!("receiver for last_at dropped, can't notify"); }; } Ok("pages_to_stdout") diff --git a/src/mirror/fjall.rs b/src/mirror/fjall.rs index 68dce3f..5d18ce2 100644 --- a/src/mirror/fjall.rs +++ b/src/mirror/fjall.rs @@ -382,7 +382,7 @@ async fn export_stream( while let Some(op) = op_rx.recv().await { let json = serde_json::to_string(&op.to_sequenced_json()).unwrap(); if let Err(e) = socket.send(Message::Text(json)).await { - log::warn!("closing export stream: {e}"); + tracing::warn!("closing export stream: {e}"); return; } cursor = op.seq + 1; @@ -390,11 +390,11 @@ async fn export_stream( match read.await { Ok(Err(e)) => { - log::error!("stream read failed: {e}"); + tracing::error!("stream read failed: {e}"); return; } Err(e) => { - log::error!("stream read task panicked: {e}"); + tracing::error!("stream read task panicked: {e}"); return; } Ok(Ok(())) => {} @@ -429,7 +429,7 @@ pub async fn serve_fjall( experimental: ExperimentalConf, fjall: FjallDb, ) -> anyhow::Result<&'static str> { - log::info!("starting fjall mirror server..."); + tracing::info!("starting fjall mirror server..."); let client = Client::builder() .user_agent(UA) @@ -461,7 +461,7 @@ pub async fn serve_fjall( .at("/export/stream", get(export_stream)); if experimental.write_upstream { - log::info!("enabling experimental write forwarding to upstream"); + tracing::info!("enabling experimental write forwarding to upstream"); let ip_limiter = IpLimiters::new(Quota::per_hour(10.try_into().unwrap())); let did_limiter = CreatePlcOpLimiter::new(Quota::per_hour(4.try_into().unwrap())); @@ -514,7 +514,7 @@ async fn fjall_forward_create_op_upstream( } let mut headers: reqwest::header::HeaderMap = req.headers().clone(); - log::trace!("original request headers: {headers:?}"); + tracing::trace!("original request headers: {headers:?}"); headers.insert("Host", upstream.host_str().unwrap().parse().unwrap()); let client_ua = headers .get(USER_AGENT) @@ -526,7 +526,7 @@ async fn fjall_forward_create_op_upstream( .parse() .unwrap(), ); - log::trace!("adjusted request headers: {headers:?}"); + tracing::trace!("adjusted request headers: {headers:?}"); let mut target = upstream.clone(); target.set_path(&did); @@ -538,7 +538,7 @@ async fn fjall_forward_create_op_upstream( .send() .await .map_err(|e| { - log::warn!("upstream write fail: {e}"); + tracing::warn!("upstream write fail: {e}"); Error::from_string( failed_to_reach_named("upstream PLC"), StatusCode::BAD_GATEWAY, diff --git a/src/mirror/mod.rs b/src/mirror/mod.rs index 2ec53d4..e153097 100644 --- a/src/mirror/mod.rs +++ b/src/mirror/mod.rs @@ -184,12 +184,12 @@ where } let auto_cert = auto_cert.build().expect("acme config to build"); - log::trace!("auto_cert: {auto_cert:?}"); + tracing::trace!("auto_cert: {auto_cert:?}"); let notice_task = tokio::task::spawn(run_insecure_notice(ipv6)); let listener = TcpListener::bind(if ipv6 { "[::]:443" } else { "0.0.0.0:443" }); let app_res = run(app, listener.acme(auto_cert)).await; - log::warn!("server task ended, aborting insecure server task..."); + tracing::warn!("server task ended, aborting insecure server task..."); notice_task.abort(); app_res?; notice_task.await??; diff --git a/src/mirror/pg.rs b/src/mirror/pg.rs index ce2492b..08e0e8d 100644 --- a/src/mirror/pg.rs +++ b/src/mirror/pg.rs @@ -179,7 +179,7 @@ async fn proxy(req: &Request, Data(state): Data<&State>) -> Result { .send() .await .map_err(|e| { - log::error!("upstream req fail: {e}"); + tracing::error!("upstream req fail: {e}"); Error::from_string( failed_to_reach_named("wrapped reference PLC"), StatusCode::BAD_GATEWAY, @@ -215,7 +215,7 @@ async fn forward_create_op_upstream( } let mut headers: reqwest::header::HeaderMap = req.headers().clone(); - log::trace!("original request headers: {headers:?}"); + tracing::trace!("original request headers: {headers:?}"); headers.insert("Host", upstream.host_str().unwrap().parse().unwrap()); let client_ua = headers .get(USER_AGENT) @@ -227,7 +227,7 @@ async fn forward_create_op_upstream( .parse() .unwrap(), ); - log::trace!("adjusted request headers: {headers:?}"); + tracing::trace!("adjusted request headers: {headers:?}"); let mut target = upstream.clone(); target.set_path(&did); @@ -239,7 +239,7 @@ async fn forward_create_op_upstream( .send() .await .map_err(|e| { - log::warn!("upstream write fail: {e}"); + tracing::warn!("upstream write fail: {e}"); Error::from_string( failed_to_reach_named("upstream PLC"), StatusCode::BAD_GATEWAY, @@ -274,7 +274,7 @@ pub async fn serve( experimental: ExperimentalConf, db: Option, ) -> anyhow::Result<&'static str> { - log::info!("starting server..."); + tracing::info!("starting server..."); let client = Client::builder() .user_agent(UA) @@ -306,7 +306,7 @@ pub async fn serve( .at("/export", get(proxy)); if experimental.write_upstream { - log::info!("enabling experimental write forwarding to upstream"); + tracing::info!("enabling experimental write forwarding to upstream"); let ip_limiter = IpLimiters::new(Quota::per_hour(10.try_into().unwrap())); let did_limiter = CreatePlcOpLimiter::new(Quota::per_hour(4.try_into().unwrap())); diff --git a/src/plc_fjall.rs b/src/plc_fjall.rs index 5c84656..6d7ba0d 100644 --- a/src/plc_fjall.rs +++ b/src/plc_fjall.rs @@ -1031,7 +1031,7 @@ impl FjallDb { }; for err in &errors { - log::warn!("dropping op {} {} (seq {seq}) parse error: {err}", op.did, op.cid); + tracing::warn!("dropping op {} {} (seq {seq}) parse error: {err}", op.did, op.cid); } if !errors.is_empty() { // if parse failed but not fatal, we just dont store it @@ -1069,16 +1069,16 @@ impl FjallDb { .map(|e| e.to_string()) .collect::>() .join("\n"); - log::warn!("dropping op {} {} (seq {seq}) invalid sig:\n{msg}", op.did, op.cid); + tracing::warn!("dropping op {} {} (seq {seq}) invalid sig:\n{msg}", op.did, op.cid); return Ok(0); } } Err(e) => { - log::warn!("dropping op {} {} (seq {seq}): {e}", op.did, op.cid); + tracing::warn!("dropping op {} {} (seq {seq}): {e}", op.did, op.cid); return Ok(0); } } - log::debug!("verified op {} {} (seq {seq})", op.did, op.cid); + tracing::debug!("verified op {} {} (seq {seq})", op.did, op.cid); } let db_op = DbOp { @@ -1104,7 +1104,7 @@ impl FjallDb { self.inner.notify_stream.notify_waiters(); - log::debug!("inserted op {} {} (seq {seq})", op.did, op.cid); + tracing::debug!("inserted op {} {} (seq {seq})", op.did, op.cid); Ok(1) } @@ -1257,7 +1257,7 @@ impl FjallDb { let (seq, by_did_key_bytes) = match (found_seq, found_by_did_key) { (Some(s), Some(k)) => (s, k), _ => { - log::warn!("drop_op: by_did entry not found for {did_str}"); + tracing::warn!("drop_op: by_did entry not found for {did_str}"); return Ok(()); } }; @@ -1310,7 +1310,7 @@ impl FjallDb { }); let prev_cid_ok = op.operation.prev.is_none() || prev_op.is_some(); if !prev_cid_ok { - log::error!("audit: op {did} {cid} prev cid mismatch or missing predecessor, is db corrupted?"); + tracing::error!("audit: op {did} {cid} prev cid mismatch or missing predecessor, is db corrupted?"); failed += 1; send_invalid(); continue; @@ -1325,13 +1325,13 @@ impl FjallDb { .map(|e| e.to_string()) .collect::>() .join("\n "); - log::warn!("audit: invalid op {} {}:\n {msg}", did, cid); + tracing::warn!("audit: invalid op {} {}:\n {msg}", did, cid); failed += 1; send_invalid(); } } Err(e) => { - log::warn!("audit: invalid op {} {}: {e}", did, cid); + tracing::warn!("audit: invalid op {} {}: {e}", did, cid); failed += 1; send_invalid(); } @@ -1437,7 +1437,7 @@ pub async fn backfill_to_fjall( if reset { let db = db.clone(); tokio::task::spawn_blocking(move || db.clear()).await??; - log::warn!("fjall reset: cleared all data"); + tracing::warn!("fjall reset: cleared all data"); } let mut last_at = None; @@ -1492,18 +1492,18 @@ pub async fn backfill_to_fjall( } } } - log::debug!("finished receiving bulk pages"); + tracing::debug!("finished receiving bulk pages"); if let Some(notify) = notify_last_at { - log::trace!("notifying last_at: {last_at:?}"); + tracing::trace!("notifying last_at: {last_at:?}"); if notify.send(last_at).is_err() { - log::error!("receiver for last_at dropped, can't notify"); + tracing::error!("receiver for last_at dropped, can't notify"); }; } tokio::task::spawn_blocking(move || db.persist(PersistMode::SyncAll)).await??; - log::info!( + tracing::info!( "backfill_to_fjall: inserted {ops_inserted} ops in {:?}", t0.elapsed() ); @@ -1515,7 +1515,7 @@ pub async fn seq_pages_to_fjall( db: FjallDb, mut pages: mpsc::Receiver, ) -> anyhow::Result<&'static str> { - log::info!("starting seq_pages_to_fjall writer..."); + tracing::info!("starting seq_pages_to_fjall writer..."); let t0 = Instant::now(); let mut ops_inserted: usize = 0; @@ -1523,7 +1523,7 @@ pub async fn seq_pages_to_fjall( while let Some(page) = pages.recv().await { let first_seq = page.ops.first().map(|op| op.seq); let last_seq = page.ops.last().map(|op| op.seq); - log::debug!( + tracing::debug!( "seq_pages: received page with {} ops, seq {:?}..{:?}", page.ops.len(), first_seq, @@ -1534,7 +1534,7 @@ pub async fn seq_pages_to_fjall( let count = tokio::task::spawn_blocking(move || -> anyhow::Result { let mut count: usize = 0; for seq_op in &page.ops { - log::debug!("seq_pages: processing op {} {} (seq {})", seq_op.did, seq_op.cid, seq_op.seq); + tracing::debug!("seq_pages: processing op {} {} (seq {})", seq_op.did, seq_op.cid, seq_op.seq); let common_op = CommonOp { did: seq_op.did.clone(), cid: seq_op.cid.clone(), @@ -1549,7 +1549,7 @@ pub async fn seq_pages_to_fjall( }) .await??; if count < page_len { - log::warn!( + tracing::warn!( "seq_pages: page seq {:?}..{:?} inserted {count}/{page_len} ops ({} dropped)", first_seq, last_seq, @@ -1559,7 +1559,7 @@ pub async fn seq_pages_to_fjall( ops_inserted += count; } - log::info!( + tracing::info!( "no more seq pages. inserted {ops_inserted} ops in {:?}", t0.elapsed() ); @@ -1570,15 +1570,15 @@ pub async fn audit( db: FjallDb, invalid_ops_tx: mpsc::Sender, ) -> anyhow::Result<&'static str> { - log::info!("starting fjall audit..."); + tracing::info!("starting fjall audit..."); let t0 = std::time::Instant::now(); let (checked, failed) = tokio::task::spawn_blocking(move || db.audit(invalid_ops_tx)).await??; - log::info!( + tracing::info!( "fjall audit complete in {:?}, {checked} ops checked", t0.elapsed() ); if failed > 0 { - log::error!("audit found {failed} invalid operations"); + tracing::error!("audit found {failed} invalid operations"); } Ok("audit_fjall") } @@ -1589,7 +1589,7 @@ pub async fn fix_ops( only_drop: bool, mut invalid_ops_rx: mpsc::Receiver, ) -> anyhow::Result<&'static str> { - log::info!("starting fjall fix ops..."); + tracing::info!("starting fjall fix ops..."); let mut fixed_dids = std::collections::HashSet::new(); let mut count = 0; @@ -1615,7 +1615,7 @@ pub async fn fix_ops( continue; } - log::trace!("fetching upstream ops to fix did: {did}"); + tracing::trace!("fetching upstream ops to fix did: {did}"); let mut url = upstream.clone(); url.set_path(&format!("/{did}/log/audit")); @@ -1626,21 +1626,21 @@ pub async fn fix_ops( StatusCode::OK => match resp.json().await { Ok(ops) => ops, Err(e) => { - log::warn!("failed to parse upstream ops for {did}: {e}"); + tracing::warn!("failed to parse upstream ops for {did}: {e}"); continue; } }, StatusCode::NOT_FOUND => { - log::trace!("did not found upstream: {did}"); + tracing::trace!("did not found upstream: {did}"); Vec::new() // this essentially means drop the whole did } s => { - log::warn!("failed to fetch upstream for {did}: {s}"); + tracing::warn!("failed to fetch upstream for {did}: {s}"); continue; } }; - log::trace!("fetched {} ops for {did}", ops.len()); + tracing::trace!("fetched {} ops for {did}", ops.len()); // we drop all ops first just to be safe let existing = db.ops_for_did(&did)?; @@ -1655,7 +1655,7 @@ pub async fn fix_ops( // if we don't skip these we might miss some ops in between // the latest_at we started with vs the one we ended up with if op.created_at > latest_at { - log::trace!( + tracing::trace!( "skipping op {} for {did} because it is newer than latest_at {latest_at}", op.cid ); @@ -1671,7 +1671,7 @@ pub async fn fix_ops( fixed_dids.insert(did); } - log::info!("fixed {count} ops"); + tracing::info!("fixed {count} ops"); Ok("fix_ops_fjall") } diff --git a/src/plc_pg.rs b/src/plc_pg.rs index ff86176..bd231e5 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -49,7 +49,7 @@ impl Db { pub async fn new(pg_uri: &str, cert: Option) -> 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..."); + tracing::trace!("checking migrations..."); let connector = cert.map(get_tls).transpose()?; @@ -77,7 +77,7 @@ impl Db { drop(client); // make sure the connection worker thing doesn't linger conn_task.await??; - log::info!("db connection succeeded and plc migrations appear as expected"); + tracing::info!("db connection succeeded and plc migrations appear as expected"); Ok(Self { pg_uri: pg_uri.to_string(), @@ -86,7 +86,7 @@ impl Db { } pub async fn connect(&self) -> Result<(Client, JoinHandle>), PgError> { - log::trace!("connecting postgres..."); + tracing::trace!("connecting postgres..."); if let Some(ref connector) = self.cert { get_client_and_task(&self.pg_uri, connector.clone()).await } else { @@ -115,7 +115,7 @@ pub async fn pages_to_pg( db: Db, mut pages: mpsc::Receiver, ) -> anyhow::Result<&'static str> { - log::info!("starting pages_to_pg writer..."); + tracing::info!("starting pages_to_pg writer..."); let (mut client, task) = db.connect().await?; @@ -135,7 +135,7 @@ pub async fn pages_to_pg( let mut dids_inserted = 0; while let Some(page) = pages.recv().await { - log::trace!("writing page with {} ops", page.ops.len()); + tracing::trace!("writing page with {} ops", page.ops.len()); let tx = client.transaction().await?; for op in page.ops { ops_inserted += tx @@ -156,7 +156,7 @@ pub async fn pages_to_pg( } drop(task); - log::info!( + tracing::info!( "no more pages. inserted {ops_inserted} ops and {dids_inserted} dids in {:?}", t0.elapsed() ); @@ -194,7 +194,7 @@ pub async fn backfill_to_pg( if reset { let n = tx.execute(&format!("DELETE FROM {table}"), &[]).await?; if n > 0 { - log::warn!("postgres reset: deleted {n} from {table}"); + tracing::warn!("postgres reset: deleted {n} from {table}"); } } else { let n: i64 = tx @@ -206,23 +206,23 @@ pub async fn backfill_to_pg( } } } - log::trace!("tables clean: {:?}", t_step.elapsed()); + tracing::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()); + tracing::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()); + tracing::trace!("indexes dropped: {:?}", t_step.elapsed()); let t_step = Instant::now(); - log::trace!("starting binary COPY IN..."); + tracing::trace!("starting binary COPY IN..."); let types = &[ Type::TEXT, Type::JSONB, @@ -256,17 +256,17 @@ 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"); + tracing::debug!("finished receiving bulk pages"); if let Some(notify) = notify_last_at { - log::trace!("notifying last_at: {last_at:?}"); + tracing::trace!("notifying last_at: {last_at:?}"); if notify.send(last_at).is_err() { - log::error!("receiver for last_at dropped, can't notify"); + tracing::error!("receiver for last_at dropped, can't notify"); }; } let n = writer.as_mut().finish().await?; - log::trace!("COPY IN wrote {n} ops: {:?}", t_step.elapsed()); + tracing::trace!("COPY IN wrote {n} ops: {:?}", t_step.elapsed()); // CAUTION: these indexes MUST match up exactly with the kysely ones we dropped let t_step = Instant::now(); @@ -280,7 +280,7 @@ pub async fn backfill_to_pg( &[], ) .await?; - log::trace!("indexes recreated: {:?}", t_step.elapsed()); + tracing::trace!("indexes recreated: {:?}", t_step.elapsed()); let t_step = Instant::now(); let n = tx @@ -289,16 +289,16 @@ pub async fn backfill_to_pg( &[], ) .await?; - log::trace!("INSERT wrote {n} dids: {:?}", t_step.elapsed()); + tracing::trace!("INSERT wrote {n} dids: {:?}", t_step.elapsed()); 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()); + tracing::trace!("set tables LOGGED: {:?}", t_step.elapsed()); tx.commit().await?; drop(task); - log::info!("total backfill time: {:?}", t0.elapsed()); + tracing::info!("total backfill time: {:?}", t0.elapsed()); Ok("backfill_to_pg") } diff --git a/src/poll.rs b/src/poll.rs index 546fe3d..5642f4b 100644 --- a/src/poll.rs +++ b/src/poll.rs @@ -143,7 +143,7 @@ pub async fn get_page(url: Url) -> Result<(ExportPage, Option), GetPageE use tokio::io::{AsyncBufReadExt, BufReader}; use tokio_util::compat::FuturesAsyncReadCompatExt; - log::trace!("Getting page: {url}"); + tracing::trace!("Getting page: {url}"); let res = CLIENT.get(url).send().await?.error_for_status()?; let stream = Box::pin( @@ -165,12 +165,12 @@ pub async fn get_page(url: Url) -> Result<(ExportPage, Option), GetPageE } match serde_json::from_str::(line) { Ok(op) => ops.push(op), - Err(e) => log::warn!("failed to parse op: {e} ({line})"), + Err(e) => tracing::warn!("failed to parse op: {e} ({line})"), } } Ok(None) => break, Err(e) => { - log::warn!("transport error mid-page: {}; returning partial page", e); + tracing::warn!("transport error mid-page: {}; returning partial page", e); break; } } @@ -216,7 +216,7 @@ pub async fn poll_upstream( throttle: Duration, dest: mpsc::Sender, ) -> anyhow::Result<&'static str> { - log::info!("starting upstream poller at {base} after {after:?}"); + tracing::info!("starting upstream poller at {base} after {after:?}"); let mut tick = tokio::time::interval(throttle); let mut prev_last: Option = after.map(Into::into); let mut boundary_state: Option = None; @@ -232,7 +232,7 @@ pub async fn poll_upstream( let (mut page, next_last) = match get_page(url).await { Ok(res) => res, Err(e) => { - log::warn!("error polling upstream: {e}"); + tracing::warn!("error polling upstream: {e}"); continue; } }; @@ -246,7 +246,7 @@ pub async fn poll_upstream( match dest.try_send(page) { Ok(()) => {} Err(mpsc::error::TrySendError::Full(page)) => { - log::warn!("export: destination channel full, awaiting..."); + tracing::warn!("export: destination channel full, awaiting..."); dest.send(page).await?; } e => e?, @@ -263,7 +263,7 @@ async fn get_seq_page(url: Url) -> Result { use tokio::io::{AsyncBufReadExt, BufReader}; use tokio_util::compat::FuturesAsyncReadCompatExt; - log::trace!("getting seq page: {url}"); + tracing::trace!("getting seq page: {url}"); let res = CLIENT.get(url).send().await?.error_for_status()?; let stream = Box::pin( @@ -285,12 +285,12 @@ async fn get_seq_page(url: Url) -> Result { } match serde_json::from_str::(line) { Ok(op) => ops.push(op), - Err(e) => log::warn!("failed to parse seq op: {e} ({line})"), + Err(e) => tracing::warn!("failed to parse seq op: {e} ({line})"), } } Ok(None) => break, Err(e) => { - log::warn!( + tracing::warn!( "transport error mid-seq-page: {}; returning partial page", e ); @@ -315,7 +315,7 @@ pub async fn poll_upstream_seq( throttle: Duration, dest: mpsc::Sender, ) -> anyhow::Result<&'static str> { - log::info!("starting seq upstream poller at {base} after {after:?}"); + tracing::info!("starting seq upstream poller at {base} after {after:?}"); let mut tick = tokio::time::interval(throttle); let mut last_seq: u64 = after.unwrap_or(0); @@ -329,7 +329,7 @@ pub async fn poll_upstream_seq( let page = match get_seq_page(url).await { Ok(p) => p, Err(e) => { - log::warn!("error polling upstream (seq): {e}"); + tracing::warn!("error polling upstream (seq): {e}"); continue; } }; @@ -339,7 +339,7 @@ pub async fn poll_upstream_seq( } if !page.is_empty() { - log::debug!( + tracing::debug!( "seq poll: page with {} ops, seq {}..{}", page.ops.len(), page.ops.first().map(|op| op.seq).unwrap_or(0), @@ -348,7 +348,7 @@ pub async fn poll_upstream_seq( match dest.try_send(page) { Ok(()) => {} Err(mpsc::error::TrySendError::Full(page)) => { - log::warn!("seq poll: destination channel full, awaiting..."); + tracing::warn!("seq poll: destination channel full, awaiting..."); dest.send(page).await?; } e => e?, @@ -388,16 +388,16 @@ pub async fn tail_upstream_stream( .append_pair("cursor", &seq.to_string()); } - log::info!("connecting to stream: {url}"); + tracing::info!("connecting to stream: {url}"); let (mut ws, _) = connect_async(url.as_str()).await?; - log::info!("stream connected"); + tracing::info!("stream connected"); while let Some(msg) = ws.next().await { let msg = msg?; let text = match msg { Message::Text(t) => t, Message::Close(_) => { - log::info!("stream closed by server"); + tracing::info!("stream closed by server"); break; } _ => continue, @@ -406,14 +406,14 @@ pub async fn tail_upstream_stream( let op: SeqOp = match serde_json::from_str(&text) { Ok(op) => op, Err(e) => { - log::warn!("failed to parse stream event: {e} ({text})"); + tracing::warn!("failed to parse stream event: {e} ({text})"); continue; } }; let page = SeqPage { ops: vec![op] }; if dest.send(page).await.is_err() { - log::info!("stream dest channel closed, stopping"); + tracing::info!("stream dest channel closed, stopping"); break; } } diff --git a/src/ratelimit.rs b/src/ratelimit.rs index f3eeb44..00a64d6 100644 --- a/src/ratelimit.rs +++ b/src/ratelimit.rs @@ -61,7 +61,7 @@ impl Limiter for CreatePlcOpLimiter { .map_err(|e| e.wait_time_from(CLOCK.now())) } fn housekeep(&self) { - log::debug!( + tracing::debug!( "limiter size before housekeeping: {} dids", self.limiter.len() ); @@ -125,7 +125,7 @@ impl Limiter for IpLimiters { } } fn housekeep(&self) { - log::debug!( + tracing::debug!( "limiter sizes before housekeeping: {}/ip {}/v6_56 {}/v6_48", self.per_ip.len(), self.ip6_56.len(), @@ -205,13 +205,13 @@ where match self.limiters.check_key(&key) { Ok(_) => { - log::debug!("allowing key {key:?}"); + tracing::debug!("allowing key {key:?}"); self.ep.call(req).await } Err(d) => { let wait_time = d.as_secs(); - log::debug!("rate limit exceeded for {key:?}, quota reset in {wait_time}s"); + tracing::debug!("rate limit exceeded for {key:?}, quota reset in {wait_time}s"); let res = Response::builder() .status(StatusCode::TOO_MANY_REQUESTS) diff --git a/src/weekly.rs b/src/weekly.rs index 39dfdbd..d2bc19a 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -97,10 +97,10 @@ impl BundleSource for FolderSource { async fn reader_for(&self, week: Week) -> anyhow::Result { let FolderSource(dir) = self; let path = dir.join(format!("{}.jsonl.gz", week.0)); - log::debug!("opening folder source: {path:?}"); + tracing::debug!("opening folder source: {path:?}"); let file = File::open(path) .await - .inspect_err(|e| log::error!("failed to open file: {e}"))?; + .inspect_err(|e| tracing::error!("failed to open file: {e}"))?; let decoder = GzipDecoder::new(BufReader::new(file)); Ok(decoder) } @@ -151,7 +151,7 @@ pub async fn pages_to_weeks( encoder.shutdown().await?; let now = Instant::now(); - log::info!( + tracing::info!( "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, @@ -170,7 +170,7 @@ pub async fn pages_to_weeks( week_ops = 0; week_t0 = now; } - log::trace!("writing: {op:?}"); + tracing::trace!("writing: {op:?}"); encoder .write_all(serde_json::to_string(&op)?.as_bytes()) .await?; @@ -182,7 +182,7 @@ pub async fn pages_to_weeks( // don't forget the final file encoder.shutdown().await?; let now = Instant::now(); - log::info!( + tracing::info!( "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, @@ -206,7 +206,7 @@ pub async fn week_to_pages( let reader = match source.reader_for(week).await { Ok(r) => r, Err(e) => { - log::warn!( + tracing::warn!( "week_to_pages reader_for failed {e}, retrying in {}s", retry_backoff.as_secs() ); @@ -223,7 +223,7 @@ pub async fn week_to_pages( Ok(Some(c)) => Some(c), Ok(None) => None, Err(e) => { - log::warn!( + tracing::warn!( "failed to get next chunk: {e}, retrying week in {}s", retry_backoff.as_secs() ); @@ -237,7 +237,7 @@ pub async fn week_to_pages( .into_iter() .filter_map(|s| { serde_json::from_str::(&s) - .inspect_err(|e| log::warn!("failed to parse op: {e} ({s})")) + .inspect_err(|e| tracing::warn!("failed to parse op: {e} ({s})")) .ok() }) .collect(); @@ -248,7 +248,7 @@ pub async fn week_to_pages( let page = ExportPage { ops }; if let Err(e) = dest.send(page).await { - log::error!("failed to send page (receiver closed): {e}"); + tracing::error!("failed to send page (receiver closed): {e}"); return Err(e.into()); } }