From 1c47ca22ebfac57c9cf1f20ee212b1b60d502d9c Mon Sep 17 00:00:00 2001 From: overtrue Date: Mon, 14 Sep 2026 22:11:05 +0800 Subject: [PATCH] fix(ecstore): preserve ready read errors during encoding (#7798) Backport the erasure fix from c95b65a7fc02a493d90d42cae2c4b1a72be9d55b while retaining release shard integrity support. Dependency upgrades and the main-only Connect proxy and diagnostic fixes are excluded. --- crates/ecstore/src/erasure/coding/encode.rs | 151 +++++++++++++++++++- 1 file changed, 147 insertions(+), 4 deletions(-) diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index a8a851682..808105cfe 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -124,6 +124,10 @@ impl AbortOnDropTask { async fn join(&mut self) -> Result { (&mut self.0).await } + + fn is_finished(&self) -> bool { + self.0.is_finished() + } } impl Drop for AbortOnDropTask { @@ -132,6 +136,24 @@ impl Drop for AbortOnDropTask { } } +enum ReadyProducerState { + ReadError(std::io::Error), + Finished, +} + +async fn ready_producer_state(task: &mut AbortOnDropTask>) -> Option { + if !task.is_finished() { + tokio::task::yield_now().await; + } + if !task.is_finished() { + return None; + } + match task.join().await { + Ok(Err(err)) => Some(ReadyProducerState::ReadError(err)), + Ok(Ok(_)) | Err(_) => Some(ReadyProducerState::Finished), + } +} + /// Read up to `limit` bytes into `buf`'s uninitialized spare capacity, appending after its /// current length, and distinguish a clean EOF from a short read. /// @@ -879,8 +901,12 @@ impl Erasure { record_internal_stage_if_enabled("erasure_encode_write", write_stage_start); } - if let Some(err) = write_err { - task.abort_and_wait().await; + if let Some(mut err) = write_err { + match ready_producer_state(&mut task).await { + Some(ReadyProducerState::ReadError(read_err)) => err = read_err, + Some(ReadyProducerState::Finished) => {} + None => task.abort_and_wait().await, + } drop(rx); let shutdown_stage_start = stage_timer_if_enabled(); if let Err(shutdown_err) = writers.shutdown().await { @@ -1020,8 +1046,12 @@ impl Erasure { } } - if let Some(err) = write_err { - task.abort_and_wait().await; + if let Some(mut err) = write_err { + match ready_producer_state(&mut task).await { + Some(ReadyProducerState::ReadError(read_err)) => err = read_err, + Some(ReadyProducerState::Finished) => {} + None => task.abort_and_wait().await, + } drop(rx); let shutdown_stage_start = stage_timer_if_enabled(); if let Err(shutdown_err) = writers.shutdown().await { @@ -1217,6 +1247,52 @@ mod tests { } } + #[derive(Debug)] + struct ReadyReaderFailure; + + impl std::fmt::Display for ReadyReaderFailure { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("injected reader failure") + } + } + + impl std::error::Error for ReadyReaderFailure {} + + #[derive(Debug)] + struct BlocksThenErrorReader { + blocks_remaining: usize, + block: Vec, + failed: Option>, + } + + impl BlocksThenErrorReader { + fn new(blocks_remaining: usize, block_size: usize) -> (Self, oneshot::Receiver<()>) { + let (failed_tx, failed_rx) = oneshot::channel(); + ( + Self { + blocks_remaining, + block: vec![0x6b; block_size], + failed: Some(failed_tx), + }, + failed_rx, + ) + } + } + + impl AsyncRead for BlocksThenErrorReader { + fn poll_read(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + if self.blocks_remaining > 0 { + self.blocks_remaining -= 1; + buf.put_slice(&self.block); + return Poll::Ready(Ok(())); + } + if let Some(failed) = self.failed.take() { + let _ = failed.send(()); + } + Poll::Ready(Err(std::io::Error::other(ReadyReaderFailure))) + } + } + fn erasure_with_zero_block_size() -> Erasure { let mut erasure = Erasure::default(); erasure.data_shards = 1; @@ -1298,6 +1374,27 @@ mod tests { } } + struct FailAfterReaderErrorWriter { + reader_failed: oneshot::Receiver<()>, + } + + impl AsyncWrite for FailAfterReaderErrorWriter { + fn poll_write(mut self: Pin<&mut Self>, cx: &mut Context<'_>, _buf: &[u8]) -> Poll> { + match Pin::new(&mut self.reader_failed).poll(cx) { + Poll::Pending => Poll::Pending, + Poll::Ready(_) => Poll::Ready(Err(std::io::Error::other("injected write failure after reader error"))), + } + } + + 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(())) + } + } + struct StallOnWriteWithSignal { entered: Option>, writes: Arc, @@ -1608,6 +1705,34 @@ mod tests { rustfs_io_metrics::set_put_stage_metrics_enabled(false); } + async fn ready_reader_error_wins_over_writer_error(pipeline: EncodePipeline) { + const BLOCK_SIZE: usize = 16; + + let erasure = Arc::new(Erasure::new(1, 0, BLOCK_SIZE)); + let batch_blocks = encode_batch_block_count().min(encode_channel_capacity( + erasure.shard_size().saturating_mul(erasure.total_shard_count()), + erasure_encode_max_inflight_bytes(), + )); + let blocks_before_error = match pipeline { + EncodePipeline::Batched => batch_blocks, + EncodePipeline::Vec | EncodePipeline::BytesMut => 1, + }; + let (reader, reader_failed) = BlocksThenErrorReader::new(blocks_before_error, BLOCK_SIZE); + let mut writers = vec![Some(bitrot_writer(FailAfterReaderErrorWriter { reader_failed }, BLOCK_SIZE))]; + + let result = match pipeline { + EncodePipeline::Vec => erasure.encode_with_ingest_mode(reader, &mut writers, 1, false, None).await, + EncodePipeline::BytesMut => erasure.encode_with_ingest_mode(reader, &mut writers, 1, true, None).await, + EncodePipeline::Batched => erasure.encode_batched(reader, &mut writers, 1).await, + }; + + let err = result.expect_err("reader and writer failures must fail the encode pipeline"); + assert!( + err.get_ref().is_some_and(|source| source.is::()), + "ready reader failures must keep precedence over writer failures: {err:?}" + ); + } + async fn aborting_full_queue_settles_pending_send() { const BLOCK_SIZE: usize = 16; @@ -1731,6 +1856,24 @@ mod tests { writer_error_aborts_blocked_producer(EncodePipeline::Batched).await; } + #[tokio::test] + #[serial_test::serial] + async fn vec_ready_reader_error_wins_over_writer_error() { + ready_reader_error_wins_over_writer_error(EncodePipeline::Vec).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn bytesmut_ready_reader_error_wins_over_writer_error() { + ready_reader_error_wins_over_writer_error(EncodePipeline::BytesMut).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn batched_ready_reader_error_wins_over_writer_error() { + ready_reader_error_wins_over_writer_error(EncodePipeline::Batched).await; + } + #[tokio::test] #[serial_test::serial] async fn cancelling_full_queue_settles_pending_send() {