From d8c32d375421595b89168e3fbf9477eb30dcf981 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 9 Jul 2026 12:16:43 +0800 Subject: [PATCH] fix(ecstore): retry idempotent reads to survive read-after-write races (#4597) --- .../docker-compose.cluster.local-build.yml | 8 +- crates/config/src/constants/internode.rs | 19 +- crates/ecstore/src/cluster/rpc/remote_disk.rs | 197 ++++++++++++++++-- .../src/cluster/rpc/runtime_sources.rs | 20 +- 4 files changed, 211 insertions(+), 33 deletions(-) diff --git a/.docker/compose/docker-compose.cluster.local-build.yml b/.docker/compose/docker-compose.cluster.local-build.yml index 615219162..262e1b9ec 100644 --- a/.docker/compose/docker-compose.cluster.local-build.yml +++ b/.docker/compose/docker-compose.cluster.local-build.yml @@ -46,7 +46,7 @@ services: - RUSTFS_INTERNODE_OFFLINE_BYPASS=${RUSTFS_INTERNODE_OFFLINE_BYPASS:-false} - RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS=${RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS:-5} - RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD=${RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD:-3} - - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-0} + - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-1} extra_hosts: - "host.docker.internal:host-gateway" volumes: @@ -92,7 +92,7 @@ services: - RUSTFS_INTERNODE_OFFLINE_BYPASS=${RUSTFS_INTERNODE_OFFLINE_BYPASS:-false} - RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS=${RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS:-5} - RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD=${RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD:-3} - - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-0} + - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-1} extra_hosts: - "host.docker.internal:host-gateway" volumes: @@ -138,7 +138,7 @@ services: - RUSTFS_INTERNODE_OFFLINE_BYPASS=${RUSTFS_INTERNODE_OFFLINE_BYPASS:-false} - RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS=${RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS:-5} - RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD=${RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD:-3} - - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-0} + - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-1} extra_hosts: - "host.docker.internal:host-gateway" volumes: @@ -184,7 +184,7 @@ services: - RUSTFS_INTERNODE_OFFLINE_BYPASS=${RUSTFS_INTERNODE_OFFLINE_BYPASS:-false} - RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS=${RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS:-5} - RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD=${RUSTFS_INTERNODE_OFFLINE_FAILURE_THRESHOLD:-3} - - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-0} + - RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES=${RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES:-1} extra_hosts: - "host.docker.internal:host-gateway" volumes: diff --git a/crates/config/src/constants/internode.rs b/crates/config/src/constants/internode.rs index 40d39d64e..bddb8c286 100644 --- a/crates/config/src/constants/internode.rs +++ b/crates/config/src/constants/internode.rs @@ -148,14 +148,18 @@ pub const DEFAULT_INTERNODE_OFFLINE_REPROBE_SECS: u64 = 5; const _: () = assert!(!DEFAULT_INTERNODE_PREWARM); const _: () = assert!(!DEFAULT_INTERNODE_OFFLINE_BYPASS); -/// Extra attempts for idempotent, read-only/reentrant control-plane RPCs (e.g. `DiskInfo`) on -/// transient network failures, with exponential backoff (grpc-optimization P3-3). +/// Extra attempts for idempotent, read-only/reentrant control-plane and object-read RPCs +/// (`DiskInfo`, `ReadAll`, `ReadMetadata`, `ReadVersion`) on transient network failures, with +/// exponential backoff (grpc-optimization P3-3). /// -/// Defaults to `0` (disabled) — retries change latency-on-failure behavior and are opt-in. **Only** -/// idempotent reads use this; write/lock RPCs (`WriteAll`/`RenameData`/`Delete*`/`Lock`/`UnLock`) -/// must never auto-retry, to preserve quorum and idempotency semantics (see `CLAUDE.md`). +/// Defaults to `1`: a freshly-committed object lives on only `write_quorum` disks until heal +/// reconstructs the rest, so a single reset-by-peer on a metadata read during that +/// read-after-write window can erode the read quorum and surface as a spurious +/// `InsufficientReadQuorum` (issue #2761). One bounded retry absorbs that transient failure. +/// **Only** idempotent reads use this; write/lock RPCs (`WriteAll`/`RenameData`/`Delete*`/`Lock` +/// /`UnLock`) must never auto-retry, to preserve quorum and idempotency semantics (see `CLAUDE.md`). pub const ENV_INTERNODE_IDEMPOTENT_READ_RETRIES: &str = "RUSTFS_INTERNODE_IDEMPOTENT_READ_RETRIES"; -pub const DEFAULT_INTERNODE_IDEMPOTENT_READ_RETRIES: usize = 0; +pub const DEFAULT_INTERNODE_IDEMPOTENT_READ_RETRIES: usize = 1; // ── Control/bulk channel isolation (P1) ── // Large `bytes`-carrying unary RPCs (ReadAll/WriteAll/ReadMultiple/BatchReadVersion) otherwise @@ -283,7 +287,8 @@ mod tests { fn internode_p3_lifecycle_defaults_are_opt_in() { // The opt-in (default-off) invariants are asserted at compile time next to the definitions. assert_eq!(DEFAULT_INTERNODE_OFFLINE_REPROBE_SECS, 5); - assert_eq!(DEFAULT_INTERNODE_IDEMPOTENT_READ_RETRIES, 0); + // One bounded retry for idempotent reads to absorb read-after-write transient failures (#2761). + assert_eq!(DEFAULT_INTERNODE_IDEMPOTENT_READ_RETRIES, 1); assert_eq!(ENV_INTERNODE_PREWARM, "RUSTFS_INTERNODE_PREWARM"); assert_eq!(ENV_INTERNODE_OFFLINE_BYPASS, "RUSTFS_INTERNODE_OFFLINE_BYPASS"); assert_eq!(ENV_INTERNODE_OFFLINE_REPROBE_SECS, "RUSTFS_INTERNODE_OFFLINE_REPROBE_SECS"); diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 0b518ab45..b81e6c159 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -79,6 +79,8 @@ enum FailureHealthAction { const REMOTE_DISK_OPEN_WRITE_MAX_ATTEMPTS: usize = 2; const REMOTE_DISK_OPEN_WRITE_RETRY_BACKOFF: Duration = Duration::from_millis(20); +const REMOTE_DISK_OPEN_READ_MAX_ATTEMPTS: usize = 2; +const REMOTE_DISK_OPEN_READ_RETRY_BACKOFF: Duration = Duration::from_millis(20); /// Base backoff for idempotent read-only RPC retries (grpc-optimization P3-3); doubles per attempt. const REMOTE_DISK_READ_RETRY_BASE_BACKOFF: Duration = Duration::from_millis(50); const ENV_RUSTFS_METADATA_BATCH_READ: &str = "RUSTFS_METADATA_BATCH_READ"; @@ -227,8 +229,9 @@ pub(crate) fn internode_offline_bypass_reason(addr: &str) -> Option { None } -/// Number of extra attempts for idempotent read-only control-plane RPCs on transient network -/// failures (grpc-optimization P3-3). `0` (default) disables retries. Write/lock RPCs never retry. +/// Number of extra attempts for idempotent read-only control-plane and object-read RPCs on +/// transient network failures (grpc-optimization P3-3). Defaults to `1` so a single reset-by-peer +/// during the read-after-write window does not erode read quorum (#2761). Write/lock RPCs never retry. fn internode_idempotent_read_retries() -> usize { rustfs_utils::get_env_usize( rustfs_config::ENV_INTERNODE_IDEMPOTENT_READ_RETRIES, @@ -285,6 +288,14 @@ impl RemoteDisk { err.is_retryable_internode_write_failure() } + fn is_retryable_open_read_error(err: &DiskError) -> bool { + // Opening a read stream is idempotent (no data has been consumed yet), so any transient + // internode transport failure classified as retryable can be safely re-dialed. The + // classifier is direction-agnostic — it inspects the InternodeHttpError kind — so it is + // reused here from the write path. + err.is_retryable_internode_write_failure() + } + pub(crate) async fn new(ep: &Endpoint, opt: &DiskOption, data_transport: Arc) -> Result { let addr = if let Some(port) = ep.url.port() { format!("{}://{}:{}", ep.url.scheme(), ep.url.host_str().expect("operation should succeed"), port) @@ -364,6 +375,49 @@ impl RemoteDisk { } } + /// Open a remote read stream with a bounded retry on transient transport failures, mirroring + /// [`Self::open_write_with_retry`]. A freshly-committed object exists on only `write_quorum` + /// disks until background heal reconstructs the rest, so a single transient network error + /// (BrokenPipe/ConnectionReset/reset-by-peer) on a shard read during that read-after-write + /// window can drop the readable-shard count below `data_shards` and surface as a spurious + /// `InsufficientReadQuorum`. Opening the stream is idempotent, so re-dialing once absorbs the + /// transient failure instead of eroding quorum. See issue #2761. + async fn open_read_with_retry(&self, request: ReadStreamRequest) -> Result { + let mut attempt = 1; + let mut last_retry_classification = None; + loop { + match self.data_transport.open_read(request.clone()).await { + Ok(reader) => { + if attempt > 1 + && let Some(classification) = last_retry_classification + { + crate::cluster::rpc::runtime_sources::record_remote_disk_open_read_retry_success(classification); + } + return Ok(reader); + } + Err(err) if attempt < REMOTE_DISK_OPEN_READ_MAX_ATTEMPTS && Self::is_retryable_open_read_error(&err) => { + if let Some(classification) = err.internode_http_error_kind() { + let classification = classification.metric_label(); + crate::cluster::rpc::runtime_sources::record_remote_disk_open_read_retry(classification); + last_retry_classification = Some(classification); + } + debug!( + endpoint = %request.endpoint, + volume = %request.volume, + path = %request.path, + offset = request.offset, + length = request.length, + attempt, + "retrying remote open_read after retryable transport error" + ); + tokio::time::sleep(REMOTE_DISK_OPEN_READ_RETRY_BACKOFF).await; + attempt += 1; + } + Err(err) => return Err(err), + } + } + } + pub fn record_capacity_probe(&self, total: u64, used: u64, free: u64) { self.health.record_capacity_probe(total, used, free); } @@ -685,9 +739,10 @@ impl RemoteDisk { /// Execute an **idempotent, read-only/reentrant** RPC with a bounded number of retries on /// transient network failures, with exponential backoff (grpc-optimization P3-3). Retries - /// default to 0 (disabled). MUST NOT be used for write/lock RPCs — those must never auto-retry - /// (quorum/idempotency safety). The `operation` closure is re-invoked per attempt, so it must be - /// `Fn` (rebuild the request from borrowed inputs, do not move captured state out). + /// default to 1 (see [`internode_idempotent_read_retries`]). MUST NOT be used for write/lock + /// RPCs — those must never auto-retry (quorum/idempotency safety). The `operation` closure is + /// re-invoked per attempt, so it must be `Fn` (rebuild the request from borrowed inputs, do not + /// move captured state out). async fn execute_read_with_retry(&self, op: &'static str, operation: F, timeout_duration: Duration) -> Result where F: Fn() -> Fut, @@ -1640,7 +1695,10 @@ impl DiskAPI for RemoteDisk { } async fn read_metadata(&self, volume: &str, path: &str) -> Result { - self.execute_with_timeout_for_op( + // Idempotent metadata read: eligible for the bounded transient-network retry so a single + // reset-by-peer during the read-after-write window does not erode the metadata read + // quorum (see issue #2761). + self.execute_read_with_retry( "read_metadata", || async { let disk = self.disk_ref().await; @@ -1740,8 +1798,15 @@ impl DiskAPI for RemoteDisk { let opts_str = compat_json(opts)?; let opts_bin = encode_msgpack(opts)?; - self.execute_with_timeout( - move || async { + // Idempotent version read: eligible for the bounded transient-network retry so a single + // reset-by-peer during the read-after-write window does not erode the metadata read + // quorum (see issue #2761). The request payload is rebuilt (cloned) per attempt so the + // operation closure stays re-invocable (`Fn`). + self.execute_read_with_retry( + "read_version", + || async { + let opts_str = opts_str.clone(); + let opts_bin = opts_bin.clone(); let disk = self.disk_ref().await; let mut client = self .get_client() @@ -2101,17 +2166,16 @@ impl DiskAPI for RemoteDisk { } let disk = self.disk_ref().await; let stall_timeout = get_object_disk_read_timeout(); - self.data_transport - .open_read(ReadStreamRequest { - endpoint: self.endpoint.grid_host(), - disk, - volume: volume.to_string(), - path: path.to_string(), - offset, - length, - stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout), - }) - .await + self.open_read_with_retry(ReadStreamRequest { + endpoint: self.endpoint.grid_host(), + disk, + volume: volume.to_string(), + path: path.to_string(), + offset, + length, + stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout), + }) + .await } /// Buffered read for remote disks. @@ -2555,7 +2619,11 @@ impl DiskAPI for RemoteDisk { "Remote disk RPC started" ); - self.execute_with_timeout( + // Idempotent full read: eligible for the bounded transient-network retry so a single + // reset-by-peer during the read-after-write window does not erode the read quorum + // (see issue #2761). + self.execute_read_with_retry( + "read_all", || async { let disk = self.disk_ref().await; let mut client = self.get_bulk_client().await.map_err(|err| { @@ -2776,6 +2844,30 @@ mod tests { } } + #[derive(Debug, Clone, Default)] + struct RetryingOpenReadInternodeDataTransport { + calls: Arc>>, + // Reuses OpenWriteTestStep as a transport-agnostic open outcome (Error | Success). + steps: Arc>>, + } + + impl RetryingOpenReadInternodeDataTransport { + fn with_steps(steps: Vec) -> Self { + Self { + calls: Arc::new(StdMutex::new(Vec::new())), + steps: Arc::new(StdMutex::new(steps)), + } + } + + fn calls(&self) -> Vec { + self.calls.lock().expect("recorded transport calls lock poisoned").clone() + } + + fn record(&self, call: RecordedTransportCall) { + self.calls.lock().expect("recorded transport calls lock poisoned").push(call); + } + } + #[derive(Debug, Default)] struct EmptyTestReader; @@ -3176,6 +3268,34 @@ mod tests { } } + #[async_trait::async_trait] + impl InternodeDataTransport for RetryingOpenReadInternodeDataTransport { + async fn open_read(&self, request: ReadStreamRequest) -> Result { + self.record(RecordedTransportCall::Read(request)); + let step = self.steps.lock().expect("open_read retry steps lock poisoned").remove(0); + match step { + OpenWriteTestStep::Error(err) => Err(err), + OpenWriteTestStep::Success => Ok(Box::new(EmptyTestReader)), + } + } + + async fn open_write(&self, _request: WriteStreamRequest) -> Result { + panic!("open_write should not be used in open_read retry test"); + } + + async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result { + panic!("open_walk_dir should not be used in open_read retry test"); + } + + fn name(&self) -> &'static str { + "retrying-open-read" + } + + fn capabilities(&self) -> InternodeDataTransportCapabilities { + InternodeDataTransportCapabilities::tcp_http() + } + } + async fn new_remote_disk_with_transport(data_transport: Arc) -> RemoteDisk { let endpoint = Endpoint { url: url::Url::parse("http://remote-node:9000/data/rustfs0").expect("operation should succeed"), @@ -3701,6 +3821,43 @@ mod tests { assert_eq!(transport.calls().len(), 1, "append_file should not retry non-retryable errors"); } + #[tokio::test] + async fn test_remote_disk_read_file_stream_retries_once_on_retryable_open_read_error() { + // A transient reset-by-peer on a shard read during the read-after-write window must be + // absorbed by one re-dial rather than eroding read quorum (issue #2761). + let transport = RetryingOpenReadInternodeDataTransport::with_steps(vec![ + OpenWriteTestStep::Error(DiskError::from(rustfs_rio::new_test_internode_http_io_error( + rustfs_rio::InternodeHttpErrorKind::ConnectionReset, + ))), + OpenWriteTestStep::Success, + ]); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + crate::cluster::rpc::runtime_sources::reset_internode_metrics_for_test(); + + let _reader = remote_disk + .read_file_stream("bucket", "object/part.1", 0, 4096) + .await + .expect("retryable open_read error should recover"); + + assert_eq!(transport.calls().len(), 2, "read_file_stream should retry exactly once"); + } + + #[tokio::test] + async fn test_remote_disk_read_file_stream_does_not_retry_non_retryable_open_read_error() { + let transport = RetryingOpenReadInternodeDataTransport::with_steps(vec![OpenWriteTestStep::Error(DiskError::from( + rustfs_rio::new_test_internode_http_io_error(rustfs_rio::InternodeHttpErrorKind::Unknown), + ))]); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + + let err = match remote_disk.read_file_stream("bucket", "object/part.2", 0, 4096).await { + Ok(_) => panic!("non-retryable open_read error should be returned directly"), + Err(err) => err, + }; + + assert_eq!(err.internode_http_error_kind(), Some(rustfs_rio::InternodeHttpErrorKind::Unknown)); + assert_eq!(transport.calls().len(), 1, "read_file_stream should not retry non-retryable errors"); + } + #[tokio::test] async fn test_remote_disk_walk_dir_uses_configured_data_transport() { let transport = RecordingInternodeDataTransport::default(); diff --git a/crates/ecstore/src/cluster/rpc/runtime_sources.rs b/crates/ecstore/src/cluster/rpc/runtime_sources.rs index 0b21769c2..026e7fa59 100644 --- a/crates/ecstore/src/cluster/rpc/runtime_sources.rs +++ b/crates/ecstore/src/cluster/rpc/runtime_sources.rs @@ -14,8 +14,8 @@ use rustfs_io_metrics::internode_metrics::{ INTERNODE_MSGPACK_DIRECTION_RESPONSE, INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_OPERATION_GRPC_READ_MULTIPLE, - INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_OPERATION_PUT_FILE_STREAM, INTERNODE_TRANSPORT_BACKEND_GRPC, - INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, global_internode_metrics, + INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_OPERATION_PUT_FILE_STREAM, INTERNODE_OPERATION_READ_FILE_STREAM, + INTERNODE_TRANSPORT_BACKEND_GRPC, INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, global_internode_metrics, }; #[cfg(test)] @@ -37,6 +37,22 @@ pub(crate) fn record_remote_disk_open_write_retry_success(classification: &'stat ); } +pub(crate) fn record_remote_disk_open_read_retry(classification: &'static str) { + global_internode_metrics().record_retry_for_operation_and_backend( + INTERNODE_OPERATION_READ_FILE_STREAM, + INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, + classification, + ); +} + +pub(crate) fn record_remote_disk_open_read_retry_success(classification: &'static str) { + global_internode_metrics().record_retry_success_for_operation_and_backend( + INTERNODE_OPERATION_READ_FILE_STREAM, + INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, + classification, + ); +} + pub(crate) fn record_remote_disk_grpc_write_all_error() { global_internode_metrics() .record_error_for_operation_and_backend(INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC);