diff --git a/crates/ecstore/src/erasure/coding/decode.rs b/crates/ecstore/src/erasure/coding/decode.rs index 20eb535c4..a517e1e83 100644 --- a/crates/ecstore/src/erasure/coding/decode.rs +++ b/crates/ecstore/src/erasure/coding/decode.rs @@ -1069,13 +1069,11 @@ where } // Pre-claim per-slot buffers so the `self.readers` borrow below stays - // disjoint from `self.buffers`. - let participating: Vec = (0..num_readers) - .map(|i| self.engaged[i] && self.readers[i].is_some()) - .collect(); + // disjoint from `self.buffers`; `Some(buffer)` also records which slots + // participate, avoiding a per-stripe sidecar allocation. let mut bufs: Vec>> = Vec::with_capacity(num_readers); - for (i, participates) in participating.iter().enumerate() { - bufs.push(if *participates { + for i in 0..num_readers { + bufs.push(if self.engaged[i] && self.readers[i].is_some() { Some(self.buffers.take(i, shard_size)) } else { None @@ -1085,7 +1083,6 @@ where let data_shards = self.data_shards; let read_timeout = self.read_timeout; let metrics_path = self.metrics_path; - let read_costs = self.read_costs.clone(); let locality_preference_enabled = self.locality_preference_enabled; let stripe_read_start = metrics_path.map(|_| Instant::now()); @@ -1100,19 +1097,21 @@ where // before the retirement pass mutates `self.readers` below. { let mut sets = FuturesUnordered::new(); - let reader_iter = ReaderLaunchIter::new(&mut self.readers, &read_costs, locality_preference_enabled); + let reader_iter = ReaderLaunchIter::new(&mut self.readers, self.read_costs.as_slice(), locality_preference_enabled); for (i, reader) in reader_iter { - if reader.is_none() || !participating[i] { + if reader.is_none() { continue; } - let read_cost = read_costs.get(i).copied().unwrap_or(ShardReadCost::Unknown); - let recycled_buf = bufs[i].take(); + let Some(recycled_buf) = bufs[i].take() else { + continue; + }; + let read_cost = self.read_costs.get(i).copied().unwrap_or(ShardReadCost::Unknown); scheduled += 1; sets.push(read_shard( i, read_cost, reader, - recycled_buf, + Some(recycled_buf), shard_size, data_shards, read_timeout, @@ -1208,7 +1207,7 @@ where // covered by the stripe-aligned parity substitution below. if hedged { for i in 0..num_readers { - if participating[i] && shards[i].is_none() && errs[i].is_none() { + if self.engaged[i] && self.readers[i].is_some() && shards[i].is_none() && errs[i].is_none() { errs[i] = Some(Error::from(io::Error::new(ErrorKind::TimedOut, "shard read hedged after a slow shard"))); retire_readers.push(i); } @@ -1237,7 +1236,7 @@ where if !self.try_engage_parity(idx, stripe_index) { continue; } - let read_cost = read_costs.get(idx).copied().unwrap_or(ShardReadCost::Unknown); + let read_cost = self.read_costs.get(idx).copied().unwrap_or(ShardReadCost::Unknown); let recycled_buf = Some(self.buffers.take(idx, shard_size)); scheduled += 1; let (i, _read_cost, result, _should_retire) = read_shard( diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 1d31c0e27..548f6c655 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -48,6 +48,7 @@ use metrics::counter; use std::{ collections::{HashMap, VecDeque}, future::Future, + io::IoSlice, pin::Pin, sync::OnceLock, task::{Context, Poll}, @@ -76,6 +77,16 @@ impl AsyncWrite for GetObjectDownstreamWriter { .map(|result| result.map_err(mark_get_object_downstream_closed)) } + fn poll_write_vectored(mut self: Pin<&mut Self>, cx: &mut Context<'_>, bufs: &[IoSlice<'_>]) -> Poll> { + Pin::new(&mut self.inner) + .poll_write_vectored(cx, bufs) + .map(|result| result.map_err(mark_get_object_downstream_closed)) + } + + fn is_write_vectored(&self) -> bool { + self.inner.is_write_vectored() + } + fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { Pin::new(&mut self.inner) .poll_flush(cx) @@ -3101,7 +3112,7 @@ mod metadata_cache_tests { mod tests { use super::*; use crate::erasure::coding::BitrotWriter; - use std::io::{Cursor, ErrorKind}; + use std::io::{Cursor, ErrorKind, IoSlice}; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, @@ -3128,6 +3139,63 @@ mod tests { ); } + #[tokio::test] + async fn downstream_writer_preserves_vectored_write_support() { + #[derive(Default)] + struct VectoredSink { + writes: usize, + vectored_writes: usize, + bytes: Vec, + } + + impl AsyncWrite for VectoredSink { + fn poll_write(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll> { + self.writes += 1; + self.bytes.extend_from_slice(buf); + Poll::Ready(Ok(buf.len())) + } + + fn poll_write_vectored( + mut self: Pin<&mut Self>, + _cx: &mut Context<'_>, + bufs: &[IoSlice<'_>], + ) -> Poll> { + self.vectored_writes += 1; + let mut written = 0; + for buf in bufs { + written += buf.len(); + self.bytes.extend_from_slice(buf); + } + Poll::Ready(Ok(written)) + } + + fn is_write_vectored(&self) -> bool { + true + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + } + + let mut writer = GetObjectDownstreamWriter::new(VectoredSink::default()); + assert!(writer.is_write_vectored(), "downstream writer must preserve vectored-write capability"); + + let written = writer + .write_vectored(&[IoSlice::new(b"hello "), IoSlice::new(b"world")]) + .await + .expect("vectored write through downstream adapter must succeed"); + + assert_eq!(written, 11); + assert_eq!(writer.inner.vectored_writes, 1); + assert_eq!(writer.inner.writes, 0); + assert_eq!(writer.inner.bytes, b"hello world"); + } + async fn local_test_disks(count: usize, bucket: &str) -> (Vec, Vec>) { let mut dirs = Vec::with_capacity(count); let mut disks = Vec::with_capacity(count); diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index 4798a22df..e813b4001 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -5109,7 +5109,7 @@ impl DefaultObjectUsecase { part_number: Option, has_range: bool, encryption_applied: bool, - buffered_body: Option, + mut buffered_body: Option, cache_hook_served: bool, cache_hook_probed: bool, cache_fill_allowed: bool, @@ -5125,7 +5125,7 @@ impl DefaultObjectUsecase { // ODC-16 (backlog#1121): when the ecstore hook or shared cold fill // already supplied this body, the request-level plan was built before // the authoritative lookup. Serve it without planning a second time. - if cache_hook_served && let Some(bytes) = buffered_body.clone() { + if cache_hook_served && let Some(bytes) = buffered_body.take() { return Ok(Self::build_memory_bytes_blob( bytes, response_content_length,