Merge branch 'main' into fix/list-parts-pagination-infinite-loop

This commit is contained in:
Hauser
2026-10-02 13:38:02 +08:00
committed by GitHub
10 changed files with 19 additions and 30 deletions
@@ -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)?;
+1 -1
View File
@@ -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);
}
@@ -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);
-2
View File
@@ -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<W>(
&self,
@@ -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),
+5 -9
View File
@@ -908,7 +908,7 @@ fn explicit_url_port(raw: &str) -> Result<Option<u16>> {
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<Option<u16>> {
.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<Host<String>> {
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('*'))
+2 -2
View File
@@ -396,7 +396,7 @@ impl InstanceContext {
pub(crate) fn begin_namespace_commit(self: &Arc<Self>) -> Arc<NamespaceCommitGuard> {
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 {
+1 -1
View File
@@ -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))
})
-2
View File
@@ -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<Plan> {
match statement {
ExtStatement::SqlStatement(stmt) => self.df_sql_to_plan(*stmt, session).await,
@@ -33,7 +33,6 @@ use bytes::Bytes;
use futures::StreamExt as _;
use p256::ecdsa::{Signature, SigningKey, signature::Signer as _};
use p256::pkcs8::DecodePrivateKey as _;
use percent_encoding::{NON_ALPHANUMERIC, utf8_percent_encode};
use reqwest::{Client, Method, Response, StatusCode, Url};
use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256};
@@ -1317,7 +1316,7 @@ async fn drain_response(
fn object_url(endpoint: &Url, bucket: &str, key: &str, version_id: Option<&str>) -> Result<Url, SiteReplicationProbeError> {
let mut url = endpoint
.join(&format!("{bucket}/{}", encode_path(key)))
.join(&format!("{bucket}/{key}"))
.map_err(|_| SiteReplicationProbeError::ProtocolFailure)?;
if let Some(version_id) = version_id {
url.query_pairs_mut().append_pair("versionId", version_id);
@@ -1333,14 +1332,6 @@ fn list_versions_url(endpoint: &Url, bucket: &str, key: &str) -> Result<Url, Sit
Ok(url)
}
fn encode_path(value: &str) -> String {
value
.split('/')
.map(|segment| utf8_percent_encode(segment, NON_ALPHANUMERIC).to_string())
.collect::<Vec<_>>()
.join("/")
}
fn deployment_endpoint(value: &str) -> Result<Url, SiteReplicationPerformanceError> {
let mut url = Url::parse(value).map_err(|_| SiteReplicationPerformanceError::InvalidEndpoint)?;
let local_http = url.scheme() == "http"
@@ -1655,6 +1646,12 @@ mod tests {
#[test]
fn cleanup_queries_are_version_specific_and_task_scoped() {
let endpoint = Url::parse("https://source.example/").expect("endpoint");
let scoped = object_url(&endpoint, "scratch-bucket", "rustfs-connect/site-replication/019c-1234", None)
.expect("scoped object URL");
assert_eq!(
scoped.as_str(),
"https://source.example/scratch-bucket/rustfs-connect/site-replication/019c-1234"
);
let object = object_url(&endpoint, "scratch-bucket", "path/a b", Some("version+1")).expect("object URL");
assert_eq!(object.as_str(), "https://source.example/scratch-bucket/path/a%20b?versionId=version%2B1");
let list = list_versions_url(&endpoint, "scratch-bucket", "path/a b").expect("list URL");