diff --git a/crates/ecstore/src/cluster/rpc/network_probe.rs b/crates/ecstore/src/cluster/rpc/network_probe.rs index b4e58baad..75d9049ea 100644 --- a/crates/ecstore/src/cluster/rpc/network_probe.rs +++ b/crates/ecstore/src/cluster/rpc/network_probe.rs @@ -56,7 +56,7 @@ impl DiagnosticPacing { // Reservation is never refunded: a failed or dropped RPC may have sent some payload. let previous = self .charged_bytes - .fetch_update(Ordering::AcqRel, Ordering::Acquire, |charged| { + .try_update(Ordering::AcqRel, Ordering::Acquire, |charged| { charged.checked_add(bytes).filter(|total| *total <= MAX_NETWORK_PROBE_BYTES) }) .map_err(|_| NetworkPeerProbeError::LimitExceeded)?; diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index b40b520ff..66d21ff25 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -12976,7 +12976,7 @@ impl ECStore { .targets .iter() .find(|target| target.pool_index == permitted_target_pool_index) - .and_then(&candidate) + .and_then(candidate) { return Ok(pool_index); } diff --git a/crates/ecstore/src/data_movement/scanner_backlog.rs b/crates/ecstore/src/data_movement/scanner_backlog.rs index c00dec987..22c4a925c 100644 --- a/crates/ecstore/src/data_movement/scanner_backlog.rs +++ b/crates/ecstore/src/data_movement/scanner_backlog.rs @@ -306,7 +306,7 @@ pub mod test_util { let Some(fault) = fault else { return Ok(None) }; if fault .remaining - .fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1)) + .try_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1)) != Ok(1) { return Ok(None); diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 1c62872f3..2bf70a791 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -6798,7 +6798,6 @@ impl LocalDisk { } #[tracing::instrument(name = "delete_file", level = "trace", skip_all)] - #[async_recursion::async_recursion] async fn delete_file_with_namespace_owner( &self, base_path: &PathBuf, @@ -7952,7 +7951,6 @@ impl LocalDisk { Err(DiskError::FileCorrupt) } - #[async_recursion::async_recursion] #[allow(clippy::too_many_arguments)] async fn scan_dir( &self, diff --git a/crates/ecstore/src/disk/uring_result_budget.rs b/crates/ecstore/src/disk/uring_result_budget.rs index fa1d844c5..b00959325 100644 --- a/crates/ecstore/src/disk/uring_result_budget.rs +++ b/crates/ecstore/src/disk/uring_result_budget.rs @@ -73,7 +73,7 @@ impl ProcessResultBudget { } inner .available - .fetch_update(Ordering::AcqRel, Ordering::Acquire, |available| available.checked_sub(bytes)) + .try_update(Ordering::AcqRel, Ordering::Acquire, |available| available.checked_sub(bytes)) .map(|_| { Some(ResultBudgetReservation { inner: Arc::clone(inner), diff --git a/crates/ecstore/src/layout/endpoints.rs b/crates/ecstore/src/layout/endpoints.rs index 1d9ad5faa..d6f17d9b3 100644 --- a/crates/ecstore/src/layout/endpoints.rs +++ b/crates/ecstore/src/layout/endpoints.rs @@ -908,7 +908,7 @@ fn explicit_url_port(raw: &str) -> Result> { return Err(invalid()); } - let (scheme, remainder) = raw.split_once("://").ok_or_else(&invalid)?; + let (scheme, remainder) = raw.split_once("://").ok_or_else(invalid)?; if !scheme.eq_ignore_ascii_case("http") && !scheme.eq_ignore_ascii_case("https") { return Err(invalid()); } @@ -916,17 +916,17 @@ fn explicit_url_port(raw: &str) -> Result> { .split('/') .next() .filter(|authority| !authority.is_empty()) - .ok_or_else(&invalid)?; + .ok_or_else(invalid)?; if authority.contains('@') { return Err(invalid()); } let port = if let Some(bracketed) = authority.strip_prefix('[') { - let (_, suffix) = bracketed.split_once(']').ok_or_else(&invalid)?; + let (_, suffix) = bracketed.split_once(']').ok_or_else(invalid)?; if suffix.is_empty() { return Ok(None); } - suffix.strip_prefix(':').ok_or_else(&invalid)? + suffix.strip_prefix(':').ok_or_else(invalid)? } else if let Some((_, port)) = authority.rsplit_once(':') { port } else { @@ -1001,11 +1001,7 @@ fn parse_explicit_local_endpoint_host(raw: &str) -> Result> { let host = Host::parse(raw).map_err(|_| invalid())?; let host = match host { - Host::Domain(domain) => Host::Domain( - domain_without_optional_trailing_dot(&domain) - .ok_or_else(&invalid)? - .to_string(), - ), + Host::Domain(domain) => Host::Domain(domain_without_optional_trailing_dot(&domain).ok_or_else(invalid)?.to_string()), host => host, }; if matches!(&host, Host::Domain(domain) if domain.contains('*')) diff --git a/crates/ecstore/src/runtime/instance.rs b/crates/ecstore/src/runtime/instance.rs index 16a974af7..7e187c3a0 100644 --- a/crates/ecstore/src/runtime/instance.rs +++ b/crates/ecstore/src/runtime/instance.rs @@ -396,7 +396,7 @@ impl InstanceContext { pub(crate) fn begin_namespace_commit(self: &Arc) -> Arc { let counted = self .namespace_commits - .fetch_update(Ordering::AcqRel, Ordering::Acquire, |count| count.checked_add(1)) + .try_update(Ordering::AcqRel, Ordering::Acquire, |count| count.checked_add(1)) .is_ok(); if counted { self.advance_namespace_commit_generation(); @@ -412,7 +412,7 @@ impl InstanceContext { fn advance_namespace_commit_generation(&self) { let _ = self .namespace_commit_generation - .fetch_update(Ordering::AcqRel, Ordering::Acquire, |generation| Some(generation.saturating_add(1))); + .try_update(Ordering::AcqRel, Ordering::Acquire, |generation| Some(generation.saturating_add(1))); } pub(crate) fn namespace_commit_generation(&self) -> u64 { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index a18bffe04..c7654c1ac 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -4362,7 +4362,7 @@ mod tests { // successful attempts. Only injected faults spend this global budget. // A real failure may consume an attempt, so preserve the final chance. faults - .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |faults| { + .try_update(Ordering::SeqCst, Ordering::SeqCst, |faults| { (faults < crate::core::pools::DECOMMISSION_VERSION_COPY_ATTEMPTS.saturating_sub(1)) .then_some(faults.saturating_add(1)) }) diff --git a/crates/s3select-query/src/sql/planner.rs b/crates/s3select-query/src/sql/planner.rs index b847c07d9..57feec67c 100644 --- a/crates/s3select-query/src/sql/planner.rs +++ b/crates/s3select-query/src/sql/planner.rs @@ -14,7 +14,6 @@ use std::{convert::Infallible, ops::ControlFlow}; -use async_recursion::async_recursion; use async_trait::async_trait; use datafusion::sql::{ planner::{IdentNormalizer, SqlToRel}, @@ -58,7 +57,6 @@ impl<'a, S: ContextProviderExtension + Send + Sync + 'a> SqlPlanner<'a, S> { } /// Generate a logical plan from an Extent SQL statement - #[async_recursion] pub(crate) async fn statement_to_plan(&self, statement: ExtStatement, session: &SessionCtx) -> QueryResult { match statement { ExtStatement::SqlStatement(stmt) => self.df_sql_to_plan(*stmt, session).await,