From 0817cd9e5b389f109754e18c6343fe18e2094594 Mon Sep 17 00:00:00 2001 From: "oyster.cafe" Date: Thu, 7 May 2026 17:30:27 +0000 Subject: [PATCH] [hydrant] clippy fixes clipy was complaining!! --- Cargo.toml | 7 ++ examples/statusphere.rs | 2 +- src/api/repos.rs | 2 +- src/api/ws.rs | 1 + src/api/xrpc/com_atproto_describe_repo.rs | 2 +- src/api/xrpc/describe_repo.rs | 2 +- src/api/xrpc/mod.rs | 4 +- src/backfill/manager.rs | 80 +++++++++++------------ src/backfill/mod.rs | 12 ++-- src/backlinks/mod.rs | 6 +- src/control/crawler.rs | 1 + src/control/mod.rs | 12 ++-- src/control/pds.rs | 9 ++- src/control/repos/indexer.rs | 46 +++++++------ src/control/repos/mod.rs | 11 ++-- src/control/seed.rs | 2 +- src/control/stream.rs | 2 +- src/crawler/list_repos.rs | 10 +-- src/crawler/worker.rs | 6 +- src/db/keys/indexer.rs | 2 +- src/db/keys/mod.rs | 4 +- src/db/mod.rs | 17 ++--- src/db/types.rs | 43 +++++------- src/ingest/firehose.rs | 18 ++--- src/ingest/indexer.rs | 11 +++- src/ingest/relay.rs | 44 ++++++------- src/ingest/validation.rs | 9 +-- src/ops.rs | 17 ++--- src/pds_meta.rs | 9 +-- src/resolver.rs | 15 ++--- src/state.rs | 2 +- src/types.rs | 10 +-- src/util/mod.rs | 24 +++---- src/util/throttle.rs | 2 +- 34 files changed, 219 insertions(+), 225 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 7262ee0..99089ae 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -69,6 +69,10 @@ ahash = "0.8.12" [dev-dependencies] tempfile = "3.26.0" +[[example]] +name = "statusphere" +required-features = ["indexer_stream"] + [profile.dev] opt-level = 1 @@ -79,3 +83,6 @@ lto = "thin" [profile.release] lto = "thin" codegen-units = 1 + +[lints.clippy] +obfuscated_if_else = "allow" diff --git a/examples/statusphere.rs b/examples/statusphere.rs index 505e140..dae7036 100644 --- a/examples/statusphere.rs +++ b/examples/statusphere.rs @@ -75,7 +75,7 @@ impl StatusIndex { true }); let mut ranked: Vec<_> = counts.into_iter().collect(); - ranked.sort_by(|a, b| b.1.cmp(&a.1)); + ranked.sort_by_key(|b| std::cmp::Reverse(b.1)); ranked.truncate(n); ranked } diff --git a/src/api/repos.rs b/src/api/repos.rs index bd7226a..e9a670b 100644 --- a/src/api/repos.rs +++ b/src/api/repos.rs @@ -76,7 +76,7 @@ pub async fn handle_get_repos( let json = serde_json::to_string(&item).ok()?; Some(Ok::<_, std::io::Error>(format!("{json}\n"))) })) - .filter_map(|x| futures::future::ready(x)); + .filter_map(futures::future::ready); let body = Body::from_stream(stream); diff --git a/src/api/ws.rs b/src/api/ws.rs index 8cc6f84..e53f25b 100644 --- a/src/api/ws.rs +++ b/src/api/ws.rs @@ -11,6 +11,7 @@ const CLOSE_TIMEOUT: Duration = Duration::from_secs(1); pub(super) enum WsAction { Send(Message), + #[cfg_attr(not(feature = "indexer_stream"), allow(dead_code))] Skip, Close(Option), } diff --git a/src/api/xrpc/com_atproto_describe_repo.rs b/src/api/xrpc/com_atproto_describe_repo.rs index b5d2f06..2c6468a 100644 --- a/src/api/xrpc/com_atproto_describe_repo.rs +++ b/src/api/xrpc/com_atproto_describe_repo.rs @@ -45,7 +45,7 @@ pub async fn handle( handle, handle_is_correct, did_doc, - collections: collections.into_iter().map(|(k, _)| k).collect(), + collections: collections.into_keys().collect(), extra_data: Default::default(), })) } diff --git a/src/api/xrpc/describe_repo.rs b/src/api/xrpc/describe_repo.rs index 021f9ef..85734ff 100644 --- a/src/api/xrpc/describe_repo.rs +++ b/src/api/xrpc/describe_repo.rs @@ -80,6 +80,6 @@ pub async fn handle( handle: doc.handle, pds: doc.pds.to_smolstr(), signing_key: doc.signing_key, - collections: collections.into_iter().map(|(k, _)| k).collect(), + collections: collections.into_keys().collect(), })) } diff --git a/src/api/xrpc/mod.rs b/src/api/xrpc/mod.rs index 84842d8..94d2dd1 100644 --- a/src/api/xrpc/mod.rs +++ b/src/api/xrpc/mod.rs @@ -64,7 +64,9 @@ mod request_crawl; #[cfg(feature = "relay")] mod subscribe_repos; -pub fn router(blocks_available: bool) -> Router { +pub fn router( + #[cfg_attr(not(feature = "indexer"), allow(unused_variables))] blocks_available: bool, +) -> Router { let r = Router::new() .route(GetHostStatusRequest::PATH, get(get_host_status::handle)) .route(ListHostsRequest::PATH, get(list_hosts::handle)) diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index 3483251..db9e5ca 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -25,46 +25,46 @@ pub fn queue_gone_backfills(state: &Arc) -> Result<()> { } }; - if let Ok(resync_state) = rmp_serde::from_slice::(&val) { - if matches!(resync_state, ResyncState::Gone { .. }) { - debug!(did = %did, "queuing retry for gone repo"); - - let metadata_key = keys::repo_metadata_key(&did); - let metadata_bytes = match state - .db - .repo_metadata - .get(&metadata_key) - .map(|b| b.ok_or_else(|| miette::miette!("repo metadata not found"))) - .into_diagnostic() - .flatten() - { - Ok(b) => b, - Err(e) => { - error!(did = %did, err = %e, "failed to get repo metadata"); - continue; - } - }; - let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; - - // move from resync back into pending - batch.remove(&state.db.resync, key.clone()); - let old_pending = keys::pending_key(metadata.index_id); - batch.remove(&state.db.pending, old_pending); - metadata.index_id = rand::random::(); - batch.insert( - &state.db.pending, - keys::pending_key(metadata.index_id), - key.clone(), - ); - batch.insert( - &state.db.repo_metadata, - &metadata_key, - crate::db::ser_repo_meta(&metadata)?, - ); - - count_deltas.add_gauge_diff(&GaugeState::Resync(None), &GaugeState::Pending); - transitions += 1; - } + if let Ok(resync_state) = rmp_serde::from_slice::(&val) + && matches!(resync_state, ResyncState::Gone { .. }) + { + debug!(did = %did, "queuing retry for gone repo"); + + let metadata_key = keys::repo_metadata_key(&did); + let metadata_bytes = match state + .db + .repo_metadata + .get(&metadata_key) + .map(|b| b.ok_or_else(|| miette::miette!("repo metadata not found"))) + .into_diagnostic() + .flatten() + { + Ok(b) => b, + Err(e) => { + error!(did = %did, err = %e, "failed to get repo metadata"); + continue; + } + }; + let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; + + // move from resync back into pending + batch.remove(&state.db.resync, key.clone()); + let old_pending = keys::pending_key(metadata.index_id); + batch.remove(&state.db.pending, old_pending); + metadata.index_id = rand::random::(); + batch.insert( + &state.db.pending, + keys::pending_key(metadata.index_id), + key.clone(), + ); + batch.insert( + &state.db.repo_metadata, + &metadata_key, + crate::db::ser_repo_meta(&metadata)?, + ); + + count_deltas.add_gauge_diff(&GaugeState::Resync(None), &GaugeState::Pending); + transitions += 1; } } diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 3eae5fb..750ebea 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -187,9 +187,9 @@ async fn did_task( ) -> Result<(), BackfillError> { let db = &state.db; - match process_did(&state, &http, &did, verify_signatures).await { + match process_did(state, &http, did, verify_signatures).await { Ok(Some(_repo_state)) => { - let did_key = keys::repo_key(&did); + let did_key = keys::repo_key(did); // determine old gauge state // if it was error/suspended etc, we need to know which error kind it was to decrement correctly. @@ -263,7 +263,7 @@ async fn did_task( BackfillError::Deleted => unreachable!("already handled"), }; - let did_key = keys::repo_key(&did); + let did_key = keys::repo_key(did); // 1. get current retry count let existing_state = Db::get(db.resync.clone(), &did_key).await.and_then(|b| { @@ -284,14 +284,14 @@ async fn did_task( let next_retry = ResyncState::next_backoff(retry_count); let resync_state = ResyncState::Error { - kind: error_kind.clone(), + kind: error_kind, retry_count, next_retry, }; let old_gauge = prev_kind .map(|k| GaugeState::Resync(Some(k))) .unwrap_or(GaugeState::Pending); - let new_gauge = GaugeState::Resync(Some(error_kind.clone())); + let new_gauge = GaugeState::Resync(Some(error_kind)); let mut count_deltas = CountDeltas::default(); count_deltas.add_gauge_diff(&old_gauge, &new_gauge); @@ -384,7 +384,7 @@ impl From for BackfillError { } } -async fn process_did<'i>( +async fn process_did( app_state: &Arc, http: &reqwest::Client, did: &Did<'static>, diff --git a/src/backlinks/mod.rs b/src/backlinks/mod.rs index f9cf069..1753cff 100644 --- a/src/backlinks/mod.rs +++ b/src/backlinks/mod.rs @@ -171,10 +171,8 @@ impl BacklinksFetch { continue; }; - if !self.dids.is_empty() { - if !self.dids.iter().any(|d| d == did) { - continue; - } + if !self.dids.is_empty() && !self.dids.iter().any(|d| d == did) { + continue; } let Some(cid) = store::lookup_cid_from_ks(&db.records, did, col, rkey) else { diff --git a/src/control/crawler.rs b/src/control/crawler.rs index 2d64852..170b7a2 100644 --- a/src/control/crawler.rs +++ b/src/control/crawler.rs @@ -34,6 +34,7 @@ pub struct CrawlerSourceInfo { pub mode: crate::config::CrawlerMode, } +#[allow(clippy::too_many_arguments)] pub(super) fn spawn_crawler_producer( source: &crate::config::CrawlerSource, http: &reqwest::Client, diff --git a/src/control/mod.rs b/src/control/mod.rs index d43887e..bd2aa33 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -358,11 +358,11 @@ impl Hydrant { state.firehose_cursors.iter_sync(|relay, cursor| { let seq = cursor.load(Ordering::SeqCst); - if seq > 0 { - if let Err(e) = db::set_firehose_cursor(&state.db, relay, seq) { - error!(relay = %relay, err = %e, "failed to save cursor"); - db::check_poisoned_report(&e); - } + if seq > 0 + && let Err(e) = db::set_firehose_cursor(&state.db, relay, seq) + { + error!(relay = %relay, err = %e, "failed to save cursor"); + db::check_poisoned_report(&e); } true }); @@ -898,7 +898,7 @@ impl Hydrant { let status = state.pds_meta.load().status(&hostname); Ok(Some(Host { - name: hostname.into(), + name: hostname, seq, account_count, status, diff --git a/src/control/pds.rs b/src/control/pds.rs index 6813b9a..d9802e4 100644 --- a/src/control/pds.rs +++ b/src/control/pds.rs @@ -104,9 +104,8 @@ impl PdsControl { snapshot .hosts .iter() - .filter_map(|(host, desc): (&String, &crate::pds_meta::HostDesc)| { - matches!(desc.status, HostStatus::Banned).then(|| host.clone()) - }) + .filter(|(_, desc)| matches!(desc.status, HostStatus::Banned)) + .map(|(host, _)| host.clone()) .collect() } @@ -151,7 +150,7 @@ impl PdsControl { self.update( move |batch, ks| { - let _ = db_pds::set_tier(batch, ks, &host_clone, &tier_clone); + db_pds::set_tier(batch, ks, &host_clone, &tier_clone); if let Some(status) = maybe_status { let _ = db_pds::set_status(batch, ks, &host_clone, status); } @@ -190,7 +189,7 @@ impl PdsControl { self.update( move |batch, ks| { - let _ = db_pds::remove_tier(batch, ks, &host_clone); + db_pds::remove_tier(batch, ks, &host_clone); if let Some(status) = maybe_status { let _ = db_pds::set_status(batch, ks, &host_clone, status); } diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs index 6032895..cfa976f 100644 --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -52,8 +52,7 @@ impl ReposControl { repo_state_to_info(did, repo_state.into_static(), metadata.tracked), ))) }) - .map(|b| b.transpose()) - .flatten() + .filter_map(|b| b.transpose()) } #[allow(dead_code)] @@ -97,8 +96,7 @@ impl ReposControl { metadata.tracked, ))) }) - .map(|b| b.transpose()) - .flatten() + .filter_map(|b| b.transpose()) } pub(crate) fn _resync( @@ -136,7 +134,7 @@ impl ReposControl { metadata.tracked = true; // insert into pending with new index_id let old_pending = keys::pending_key(metadata.index_id); - batch.remove(&db.pending, &old_pending); + batch.remove(&db.pending, old_pending); metadata.index_id = rand::Rng::next_u64(&mut rand::rng()); batch.insert(&db.pending, keys::pending_key(metadata.index_id), &did_key); batch.remove(&db.resync, &did_key); @@ -289,23 +287,23 @@ impl ReposControl { .map(|b| crate::db::deser_repo_meta(&b)) .transpose()?; - if let Some(mut metadata) = existing_metadata { - if metadata.tracked { - let resync = db.resync.get(&did_key).into_diagnostic()?; - let old = db::Db::repo_gauge_state(&repo_state, resync.as_deref()); - metadata.tracked = false; - batch.insert( - &db.repo_metadata, - &metadata_key, - crate::db::ser_repo_meta(&metadata)?, - ); - batch.remove(&db.pending, keys::pending_key(metadata.index_id)); - batch.remove(&db.resync, &did_key); - if old != GaugeState::Synced { - count_deltas.add_gauge_diff(&old, &GaugeState::Synced); - } - untracked.push(did); + if let Some(mut metadata) = existing_metadata + && metadata.tracked + { + let resync = db.resync.get(&did_key).into_diagnostic()?; + let old = db::Db::repo_gauge_state(&repo_state, resync.as_deref()); + metadata.tracked = false; + batch.insert( + &db.repo_metadata, + &metadata_key, + crate::db::ser_repo_meta(&metadata)?, + ); + batch.remove(&db.pending, keys::pending_key(metadata.index_id)); + batch.remove(&db.resync, &did_key); + if old != GaugeState::Synced { + count_deltas.add_gauge_diff(&old, &GaugeState::Synced); } + untracked.push(did); } } } @@ -437,7 +435,7 @@ impl<'i> RepoHandle<'i> { if let Ok(Some(block_bytes)) = state .db .blocks - .get(&keys::block_key(collection.as_str(), &cid_bytes)) + .get(keys::block_key(collection.as_str(), &cid_bytes)) { let value: Data = serde_ipld_dagcbor::from_slice(&block_bytes).unwrap_or(Data::Null); @@ -471,9 +469,9 @@ impl<'i> RepoHandle<'i> { /// /// ## notes /// - calling this if you are using collection allowlist will always result - /// in an error since the commit root won't match the reconstructed CID. + /// in an error since the commit root won't match the reconstructed CID. /// - calling this for big repositories will incur more resource cost due to - /// hydrant's structure, the whole MST is always reconstructed. + /// hydrant's structure, the whole MST is always reconstructed. pub async fn generate_car( &self, ) -> Result> + Send + 'static>> diff --git a/src/control/repos/mod.rs b/src/control/repos/mod.rs index 40bbb6f..e279efa 100644 --- a/src/control/repos/mod.rs +++ b/src/control/repos/mod.rs @@ -112,16 +112,17 @@ impl ReposControl { #[cfg(feature = "indexer")] let state = self.0.clone(); self.iter_states(cursor).map(move |r| { + #[allow(clippy::bind_instead_of_map)] r.and_then(|(did, s)| { #[cfg(feature = "indexer")] let tracked = state .db .repo_metadata - .get(&keys::repo_metadata_key(&did)) + .get(keys::repo_metadata_key(&did)) .into_diagnostic()? .map(|b| crate::db::deser_repo_meta(&b)) .transpose()? - .map_or(true, |m| m.tracked); + .is_none_or(|m| m.tracked); #[cfg(not(feature = "indexer"))] let tracked = true; Ok(repo_state_to_info(did, s, tracked)) @@ -263,7 +264,7 @@ impl<'i> RepoHandle<'i> { .into_diagnostic()? .map(|b| crate::db::deser_repo_meta(&b)) .transpose()? - .map_or(true, |m| m.tracked); + .is_none_or(|m| m.tracked); #[cfg(not(feature = "indexer"))] let tracked = true; @@ -325,12 +326,12 @@ impl<'i> RepoHandle<'i> { }; let metadata = crate::db::deser_repo_meta(metadata_bytes.as_ref())?; - return Ok(app_state + Ok(app_state .db .pending .get(crate::db::keys::pending_key(metadata.index_id)) .into_diagnostic()? - .is_some()); + .is_some()) }) .await .map_err(|e| MiniDocError::Other(miette::miette!(e)))? diff --git a/src/control/seed.rs b/src/control/seed.rs index 24f4ee1..2fb73cf 100644 --- a/src/control/seed.rs +++ b/src/control/seed.rs @@ -42,7 +42,7 @@ pub(crate) async fn seed_from_list_hosts( }) .buffer_unordered(MAX_CONCURRENT_SEEDS); - while let Some(_) = futs.next().await {} + while futs.next().await.is_some() {} } #[tracing::instrument(skip_all, fields(seed_url = %seed_url))] diff --git a/src/control/stream.rs b/src/control/stream.rs index efbe1d9..e188daf 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -787,7 +787,7 @@ fn stored_to_event( let block = state .db .blocks - .get(&keys::block_key(collection.as_str(), &cid.to_bytes())); + .get(keys::block_key(collection.as_str(), &cid.to_bytes())); match block { Ok(Some(bytes)) => { match serde_ipld_dagcbor::from_slice::(bytes.as_ref()) { diff --git a/src/crawler/list_repos.rs b/src/crawler/list_repos.rs index 8c3fca8..933eb16 100644 --- a/src/crawler/list_repos.rs +++ b/src/crawler/list_repos.rs @@ -113,15 +113,15 @@ fn is_throttle_worthy(e: &reqwest::Error) -> bool { let mut src = e.source(); while let Some(s) = src { - if let Some(e) = s.downcast_ref::() { - if is_io_error_their_fault(&e) || is_tls_cert_error(e) { - return true; - } + if let Some(e) = s.downcast_ref::() + && (is_io_error_their_fault(e) || is_tls_cert_error(e)) + { + return true; } src = s.source(); } - e.status().map_or(false, |s| { + e.status().is_some_and(|s| { // lets not consider internal server errors for throttling // a server might be having a hard time on one request, but not on the rest s.as_u16() != 500 && is_status_their_fault(s.as_u16()) diff --git a/src/crawler/worker.rs b/src/crawler/worker.rs index 1b61434..8540355 100644 --- a/src/crawler/worker.rs +++ b/src/crawler/worker.rs @@ -146,8 +146,8 @@ impl CrawlerWorker { let mut surviving = Vec::new(); let mut count_deltas = CountDeltas::default(); for guard in guards { - let did_key = keys::repo_key(&*guard); - let metadata_key = keys::repo_metadata_key(&*guard); + let did_key = keys::repo_key(&guard); + let metadata_key = keys::repo_metadata_key(&guard); if app_state .db .repos @@ -171,7 +171,7 @@ impl CrawlerWorker { &did_key, ); // clear any stale retry entry, this DID is confirmed and being enqueued - batch.remove(&app_state.db.crawler, keys::crawler_retry_key(&*guard)); + batch.remove(&app_state.db.crawler, keys::crawler_retry_key(&guard)); trace!(did = %*guard, "enqueuing repo"); count_deltas.add("repos", 1); #[cfg(feature = "indexer")] diff --git a/src/db/keys/indexer.rs b/src/db/keys/indexer.rs index 506fd42..7caa207 100644 --- a/src/db/keys/indexer.rs +++ b/src/db/keys/indexer.rs @@ -96,7 +96,7 @@ pub fn resync_buffer_key(did: &Did, rev: DbTid) -> Vec { let mut key = Vec::with_capacity(repo.len() + 1 + 8); repo.write_to_vec(&mut key); key.push(SEP); - key.extend_from_slice(&rev.as_bytes()); + key.extend_from_slice(rev.as_bytes()); key } diff --git a/src/db/keys/mod.rs b/src/db/keys/mod.rs index 9b63d8a..c8e00b5 100644 --- a/src/db/keys/mod.rs +++ b/src/db/keys/mod.rs @@ -21,7 +21,7 @@ pub const RELAY_EVENT_WATERMARK_PREFIX: &[u8] = b"rwm|"; pub const VERSIONING_KEY: &[u8] = b"db_version"; // key format: {DID} -pub fn repo_key<'a>(did: &'a Did) -> Vec { +pub fn repo_key(did: &Did) -> Vec { let mut vec = Vec::with_capacity(32); TrimmedDid::from(did).write_to_vec(&mut vec); vec @@ -29,7 +29,7 @@ pub fn repo_key<'a>(did: &'a Did) -> Vec { pub const REPO_METADATA_PREFIX: &[u8] = b"rm|"; -pub fn repo_metadata_key<'a>(did: &'a Did) -> Vec { +pub fn repo_metadata_key(did: &Did) -> Vec { let mut vec = Vec::with_capacity(REPO_METADATA_PREFIX.len() + 32); vec.extend_from_slice(REPO_METADATA_PREFIX); TrimmedDid::from(did).write_to_vec(&mut vec); diff --git a/src/db/mod.rs b/src/db/mod.rs index 5ed6d85..3b3c9ce 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -235,14 +235,14 @@ impl Db { let load_dict = |name: &str| -> Option> { let path = cfg.database_path.join(format!("dict_{name}.bin")); - if path.exists() { - if let Ok(bytes) = std::fs::read(&path) { - tracing::debug!( - "loaded zstd dictionary for keyspace {name} ({} bytes)", - bytes.len() - ); - return Some(bytes.into()); - } + if path.exists() + && let Ok(bytes) = std::fs::read(&path) + { + tracing::debug!( + "loaded zstd dictionary for keyspace {name} ({} bytes)", + bytes.len() + ); + return Some(bytes.into()); } None }; @@ -961,6 +961,7 @@ pub async fn get_firehose_cursor(db: &Db, relay: &Url) -> Result> { .transpose() } +#[cfg(feature = "indexer")] pub fn ser_repo_meta(state: &RepoMetadata) -> Result> { rmp_serde::to_vec(&state).into_diagnostic() } diff --git a/src/db/types.rs b/src/db/types.rs index c0f8420..03870c8 100644 --- a/src/db/types.rs +++ b/src/db/types.rs @@ -67,21 +67,6 @@ impl<'s> TrimmedDid<'s> { TrimmedDid::Other(s) => buf.extend_from_slice(s.as_bytes()), } } - - pub fn to_string(&self) -> String { - match self { - TrimmedDid::Plc(bytes) => { - let mut s = String::with_capacity(28); - s.push_str("plc:"); - s.push_str(&BASE32_NOPAD.encode(bytes).to_ascii_lowercase()); - s - } - TrimmedDid::Web(s) => { - format!("web:{}", s) - } - TrimmedDid::Other(s) => s.to_string(), - } - } } impl Display for TrimmedDid<'_> { @@ -106,10 +91,10 @@ impl<'a> From<&'a Did<'a>> for TrimmedDid<'a> { if let Some(rest) = s.strip_prefix("did:plc:") { if rest.len() == 24 { // decode - if let Ok(bytes) = BASE32_NOPAD.decode(rest.to_ascii_uppercase().as_bytes()) { - if bytes.len() == 15 { - return TrimmedDid::Plc(bytes.try_into().unwrap()); - } + if let Ok(bytes) = BASE32_NOPAD.decode(rest.to_ascii_uppercase().as_bytes()) + && bytes.len() == 15 + { + return TrimmedDid::Plc(bytes.try_into().unwrap()); } } } else if let Some(rest) = s.strip_prefix("did:web:") { @@ -151,18 +136,18 @@ impl<'a> TryFrom<&'a [u8]> for TrimmedDid<'a> { } } -impl<'a> Into for TrimmedDid<'a> { - fn into(self) -> UserKey { +impl<'a> From> for UserKey { + fn from(val: TrimmedDid<'a>) -> Self { let mut vec = Vec::with_capacity(32); - self.write_to_vec(&mut vec); + val.write_to_vec(&mut vec); UserKey::new(&vec) } } -impl<'a> Into for &TrimmedDid<'a> { - fn into(self) -> UserKey { +impl<'a> From<&TrimmedDid<'a>> for UserKey { + fn from(val: &TrimmedDid<'a>) -> Self { let mut vec = Vec::with_capacity(32); - self.write_to_vec(&mut vec); + val.write_to_vec(&mut vec); UserKey::new(&vec) } } @@ -220,11 +205,11 @@ impl DbTid { &self.0 } - pub fn to_tid(&self) -> Tid { + pub fn to_tid(self) -> Tid { Tid::raw(self.to_smolstr()) } - fn to_smolstr(&self) -> SmolStr { + fn to_smolstr(self) -> SmolStr { let mut i = u64::from_be_bytes(self.0); let mut s = SmolStrBuilder::new(); for _ in 0..13 { @@ -269,12 +254,14 @@ impl Display for DbTid { #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[repr(u8)] +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub enum DbAction { Create = 0, Update = 1, Delete = 2, } +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] impl DbAction { pub fn as_str(&self) -> &'static str { match self { @@ -306,11 +293,13 @@ impl Display for DbAction { #[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(untagged)] +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub enum DbRkey { Tid(DbTid), Str(SmolStr), } +#[cfg_attr(not(feature = "indexer"), allow(dead_code))] impl DbRkey { pub fn new(s: &str) -> Self { if let Ok(tid) = Tid::new(s) { diff --git a/src/ingest/firehose.rs b/src/ingest/firehose.rs index 36a4304..ea611f3 100644 --- a/src/ingest/firehose.rs +++ b/src/ingest/firehose.rs @@ -68,7 +68,7 @@ impl FirehoseIngestor { #[tracing::instrument(skip(self), fields(host = %self.relay_host))] pub async fn run(mut self) -> Result<()> { let host = self.relay_host.host_str().unwrap_or(""); - let count_key = crate::db::keys::pds_account_count_key(&host); + let count_key = crate::db::keys::pds_account_count_key(host); let mut rng: SmallRng = rand::make_rng(); @@ -188,10 +188,10 @@ impl FirehoseIngestor { "active_sleep: computed status transition" ); - if current_status != new_status { - if let Err(e) = self.set_host_status(new_status) { - error!(err = %e, "failed to update host status"); - } + if current_status != new_status + && let Err(e) = self.set_host_status(new_status) + { + error!(err = %e, "failed to update host status"); } } } @@ -255,10 +255,10 @@ impl FirehoseIngestor { let failures = self.throttle.consecutive_failures(); if failures >= MAX_FAILURES { warn!(failures, "too many consecutive failures, giving up on host"); - if self.is_pds { - if let Err(e) = self.set_host_status(HostStatus::Offline) { - error!(err = %e, "failed to update host status to offline"); - } + if self.is_pds + && let Err(e) = self.set_host_status(HostStatus::Offline) + { + error!(err = %e, "failed to update host status to offline"); } return None; } diff --git a/src/ingest/indexer.rs b/src/ingest/indexer.rs index 12d4d3c..822218e 100644 --- a/src/ingest/indexer.rs +++ b/src/ingest/indexer.rs @@ -53,6 +53,7 @@ pub struct IndexerAccountData { } #[derive(Debug)] +#[allow(clippy::large_enum_variant)] pub enum IndexerEventData { Commit(IndexerCommitData), Identity(IndexerIdentityData), @@ -107,6 +108,7 @@ impl From for IngestError { } #[derive(Debug)] +#[allow(clippy::large_enum_variant)] enum RepoProcessResult<'s, 'c> { // message processed successfully, here is the (possibly updated) state Ok(RepoState<'s>), @@ -673,13 +675,18 @@ impl FirehoseWorker { .map(|b| crate::db::deser_repo_meta(&b)) .transpose()?; let had_metadata = existing_metadata.is_some(); - let mut metadata = existing_metadata.unwrap_or_else(|| RepoMetadata { + let mut metadata = existing_metadata.unwrap_or(RepoMetadata { index_id: 0, // this is set later tracked: true, }); let old_pkey = keys::pending_key(metadata.index_id); - let was_pending = had_metadata && db.pending.get(&old_pkey).into_diagnostic()?.is_some(); + let was_pending = had_metadata + && db + .pending + .get(old_pkey.as_slice()) + .into_diagnostic()? + .is_some(); // remove old pending entry and insert new one with fresh index_id if had_metadata { // only remove if we had one so we dont delete a random entry diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index 9708f8e..2818d7b 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -147,7 +147,7 @@ impl RelayWorker { InfoName::Other(name) => { let message = inf .message - .unwrap_or_else(|| CowStr::Borrowed("")); + .unwrap_or(CowStr::Borrowed("")); info!(name = %name, "relay sent info: {message}"); } } @@ -190,6 +190,7 @@ impl RelayWorker { Err(miette::miette!("relay worker dispatcher shutting down")) } + #[allow(clippy::too_many_arguments)] fn shard( id: usize, mut rx: mpsc::Receiver, @@ -553,23 +554,21 @@ impl RelayWorker { } // update per-PDS active account count on transitions - if is_pds { - if let Some(host) = firehose.host_str() { - let count_key = pds_account_count_key(host); - let delta = if !was_active && repo_state.active { - 1 - } else if was_active && !repo_state.active { - -1 - } else { - 0 - }; + if is_pds && let Some(host) = firehose.host_str() { + let count_key = pds_account_count_key(host); + let delta = if !was_active && repo_state.active { + 1 + } else if was_active && !repo_state.active { + -1 + } else { + 0 + }; - if delta != 0 { - ctx.count_deltas.add(&count_key, delta); - let count = ctx.count_deltas.projected_count(&ctx.state.db, &count_key); - ctx.state - .apply_host_limit_status(&mut ctx.batch, host, count); - } + if delta != 0 { + ctx.count_deltas.add(&count_key, delta); + let count = ctx.count_deltas.projected_count(&ctx.state.db, &count_key); + ctx.state + .apply_host_limit_status(&mut ctx.batch, host, count); } } @@ -829,7 +828,7 @@ impl WorkerContext<'_> { .map(|bytes| db::deser_repo_meta(&bytes)) .transpose()?; - if metadata.map_or(false, |m| !m.tracked) { + if metadata.is_some_and(|m| !m.tracked) { trace!(did = %did, "ignoring message, repo is explicitly untracked"); return Ok(None); } @@ -925,10 +924,11 @@ impl WorkerContext<'_> { self.count_deltas.add("repos", 1); // track initial active state for per-PDS rate limiting - if msg.is_pds && repo_state.active { - if let Some(host) = msg.firehose.host_str() { - self.count_deltas.add(&pds_account_count_key(host), 1); - } + if msg.is_pds + && repo_state.active + && let Some(host) = msg.firehose.host_str() + { + self.count_deltas.add(&pds_account_count_key(host), 1); } Ok(Some(repo_state)) diff --git a/src/ingest/validation.rs b/src/ingest/validation.rs index d2fd431..95832b2 100644 --- a/src/ingest/validation.rs +++ b/src/ingest/validation.rs @@ -98,6 +98,7 @@ impl ChainBreak { /// a successfully validated `#commit` message, carrying pre-parsed data for apply_commit pub struct ValidatedCommit<'c> { + #[cfg_attr(not(feature = "indexer"), allow(dead_code))] pub commit: &'c Commit<'c>, /// result of parse_car_bytes, already done so apply_commit does not re-parse pub parsed_blocks: ParsedCar, @@ -185,10 +186,10 @@ pub fn validate_commit<'c>( } // 2. stale rev, skip if msg.rev <= last known rev (lexicographic order) - if let Some(root) = &repo_state.root { - if msg.rev.as_str() <= root.rev.to_tid().as_str() { - return Err(CommitValidationError::StaleRev); - } + if let Some(root) = &repo_state.root + && msg.rev.as_str() <= root.rev.to_tid().as_str() + { + return Err(CommitValidationError::StaleRev); } // 3. future rev, reject if rev timestamp is more than clock_skew_secs ahead of now diff --git a/src/ops.rs b/src/ops.rs index d3c7f80..c2b3dfb 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -124,8 +124,8 @@ pub fn delete_repo( Ok(()) } -pub fn transition_repo<'batch, 's>( - batch: &'batch mut OwnedWriteBatch, +pub fn transition_repo<'s>( + batch: &mut OwnedWriteBatch, db: &Db, did: &Did, mut repo_state: RepoState<'s>, @@ -144,12 +144,12 @@ pub fn transition_repo<'batch, 's>( // manage queues match &new_status { RepoStatus::Synced => { - batch.remove(&db.pending, &pending_key); + batch.remove(&db.pending, pending_key.as_slice()); // we dont have to remove from resync here because it has to transition resync -> pending first } RepoStatus::Error(msg) => { tracing::warn!("transitioning to error: {msg}"); - batch.remove(&db.pending, &pending_key); + batch.remove(&db.pending, pending_key.as_slice()); // TODO: we need to make errors have kind instead of "message" in repo status // and then pass it to resync error kind let resync_state = crate::types::ResyncState::Error { @@ -177,12 +177,12 @@ pub fn transition_repo<'batch, 's>( } RepoStatus::Deleted => { // terminal state: remove from queues, no resync entry needed - batch.remove(&db.pending, &pending_key); + batch.remove(&db.pending, pending_key.as_slice()); batch.remove(&db.resync, &repo_key); } RepoStatus::Desynchronized | RepoStatus::Throttled => { // like an error: remove from pending and schedule a resync attempt - batch.remove(&db.pending, &pending_key); + batch.remove(&db.pending, pending_key.as_slice()); let resync_state = crate::types::ResyncState::Error { kind: crate::types::ResyncErrorKind::Generic, retry_count: 0, @@ -236,6 +236,7 @@ pub fn apply_commit<'s>( let mut records_delta = 0; let mut blocks_count = 0; let mut collection_deltas: HashMap<&str, i64> = HashMap::new(); + #[cfg(feature = "indexer_stream")] let rev = DbTid::from(&commit.rev); #[cfg(feature = "indexer_stream")] @@ -275,7 +276,7 @@ pub fn apply_commit<'s>( .wrap_err("expected valid cid from relay")?; #[cfg(feature = "indexer_stream")] { - cid_for_event = Some(cid_ipld.clone()); + cid_for_event = Some(cid_ipld); } let Some(bytes) = parsed.blocks.get(&cid_ipld) else { @@ -348,7 +349,7 @@ pub fn apply_commit<'s>( .map(StoredData::Block) .or_else(|| { (!only_index_links) - .then(|| cid_for_event.clone().map(StoredData::Ptr)) + .then(|| cid_for_event.map(StoredData::Ptr)) .flatten() }) .unwrap_or(StoredData::Nothing); diff --git a/src/pds_meta.rs b/src/pds_meta.rs index 1e99ad7..5039d5b 100644 --- a/src/pds_meta.rs +++ b/src/pds_meta.rs @@ -6,8 +6,9 @@ use std::collections::HashMap; use std::sync::Arc; use tracing::debug; -#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)] pub enum HostStatus { + #[default] Active, Idle, Offline, @@ -48,12 +49,6 @@ impl HostStatus { } } -impl Default for HostStatus { - fn default() -> Self { - Self::Active - } -} - #[derive(Debug, Default, Clone)] pub struct HostDesc { pub tier: Option, diff --git a/src/resolver.rs b/src/resolver.rs index 544c7e7..9982c26 100644 --- a/src/resolver.rs +++ b/src/resolver.rs @@ -75,11 +75,13 @@ impl Resolver { let mut jacquards = Vec::with_capacity(plc_urls.len()); for url in plc_urls { - let mut opts = ResolverOptions::default(); - opts.plc_source = PlcSource::PlcDirectory { - base: url_to_fluent_uri(&url), + let opts = ResolverOptions { + plc_source: PlcSource::PlcDirectory { + base: url_to_fluent_uri(&url), + }, + request_timeout: Some(Duration::from_secs(3)), + ..Default::default() }; - opts.request_timeout = Some(Duration::from_secs(3)); jacquards.push(JacquardResolver::new(http.clone(), opts).with_system_dns()); } @@ -200,10 +202,7 @@ impl Resolver { ) -> Result, ResolverError> { let did = did.clone().into_static(); let mini = self.resolve_doc(&did).await?; - Ok(mini - .key - .ok_or_else(|| NoSigningKeyError(did)) - .into_diagnostic()?) + Ok(mini.key.ok_or(NoSigningKeyError(did)).into_diagnostic()?) } /// resolves the full DID document as raw [`Data`] without caching, and diff --git a/src/state.rs b/src/state.rs index a8bbd46..165b558 100644 --- a/src/state.rs +++ b/src/state.rs @@ -116,7 +116,7 @@ impl AppState { #[cfg(feature = "indexer")] backfill_enabled, ephemeral: config.ephemeral, - ephemeral_ttl: config.ephemeral_ttl.clone(), + ephemeral_ttl: config.ephemeral_ttl, only_index_links: config.only_index_links, throttler: Throttler::new(), #[cfg(feature = "relay")] diff --git a/src/types.rs b/src/types.rs index 3baae5c..a6a7546 100644 --- a/src/types.rs +++ b/src/types.rs @@ -375,8 +375,9 @@ pub struct AccountEvt<'i> { } #[cfg(feature = "indexer_stream")] -#[derive(Serialize, Deserialize, Clone)] +#[derive(Serialize, Deserialize, Clone, Default)] pub(crate) enum StoredData { + #[default] Nothing, Ptr(IpldCid), #[serde(with = "jacquard_common::serde_bytes_helper")] @@ -390,13 +391,6 @@ impl StoredData { } } -#[cfg(feature = "indexer_stream")] -impl Default for StoredData { - fn default() -> Self { - Self::Nothing - } -} - #[cfg(feature = "indexer_stream")] impl Debug for StoredData { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { diff --git a/src/util/mod.rs b/src/util/mod.rs index 6f3419b..a187884 100644 --- a/src/util/mod.rs +++ b/src/util/mod.rs @@ -20,15 +20,15 @@ pub fn is_timeout(err: &dyn std::error::Error) -> bool { let mut source = err.source(); while let Some(err) = source { - if let Some(hyper_err) = err.downcast_ref::() { - if hyper_err.is_timeout() { - return true; - } + if let Some(hyper_err) = err.downcast_ref::() + && hyper_err.is_timeout() + { + return true; } - if let Some(io) = err.downcast_ref::() { - if io.kind() == std::io::ErrorKind::TimedOut { - return true; - } + if let Some(io) = err.downcast_ref::() + && io.kind() == std::io::ErrorKind::TimedOut + { + return true; } source = err.source(); } @@ -91,9 +91,9 @@ pub fn is_tls_error_their_fault(e: &rustls::Error) -> bool { // use this for public (unauth) xrpc errors pub fn is_status_their_fault(status: u16) -> bool { - return (status >= 100 && status < 200) // informational, why are we here? - || (status >= 500 && status < 600) // server error :> - || (status >= 300 && status < 400) // any 3xx error doesnt make sense in the context of a pds / relay + (100..200).contains(&status) // informational, why are we here? + || (500..600).contains(&status) // server error :> + || (300..400).contains(&status) // any 3xx error doesnt make sense in the context of a pds / relay || matches!( status, 404 // NOT FOUND: we know its not our fault because we use known xrpcs.. @@ -101,7 +101,7 @@ pub fn is_status_their_fault(status: u16) -> bool { | 403 // FORBIDDEN: sob | 401 // UNAUTHORIZED: sob | 410 // GONE: sob - ); + ) } /// outcome of [`RetryWithBackoff::retry`] when the operation does not succeed. diff --git a/src/util/throttle.rs b/src/util/throttle.rs index b13455d..62c3c03 100644 --- a/src/util/throttle.rs +++ b/src/util/throttle.rs @@ -146,7 +146,7 @@ impl ThrottleHandle { } pub fn consecutive_failures(&self) -> usize { - self.state.consecutive_failures.load(Ordering::Acquire) as usize + self.state.consecutive_failures.load(Ordering::Acquire) } /// returns whether the timeout attempts are exhausted -- 2.51.2