diff --git a/crates/ecstore/src/erasure/coding/decode.rs b/crates/ecstore/src/erasure/coding/decode.rs index d073eea5c..daf27c93a 100644 --- a/crates/ecstore/src/erasure/coding/decode.rs +++ b/crates/ecstore/src/erasure/coding/decode.rs @@ -1645,15 +1645,28 @@ impl Erasure { } Err(e) => { record_get_stage_duration_if_enabled(GET_OBJECT_PATH_LEGACY_DUPLEX, GET_STAGE_EMIT, emit_stage_start); - error!( - block_offset, - block_length, - bytes_written = *written, - stage = GET_STAGE_EMIT, - reason = classify_io_error(&e).as_str(), - error = ?e, - "Erasure decode failed to emit reconstructed data" - ); + let reason = classify_io_error(&e); + if reason == GetObjectFailureReason::DownstreamClosed { + debug!( + block_offset, + block_length, + bytes_written = *written, + stage = GET_STAGE_EMIT, + reason = reason.as_str(), + error = ?e, + "Erasure decode stopped after downstream closed" + ); + } else { + error!( + block_offset, + block_length, + bytes_written = *written, + stage = GET_STAGE_EMIT, + reason = reason.as_str(), + error = ?e, + "Erasure decode failed to emit reconstructed data" + ); + } *ret_err = Some(e); return StripeFlow::Stop; } @@ -1945,7 +1958,7 @@ mod tests { use std::io::Cursor; use std::pin::Pin; use std::sync::{ - Arc, + Arc, Mutex, atomic::{AtomicUsize, Ordering}, }; use std::task::{Context, Poll}; @@ -2120,6 +2133,59 @@ mod tests { } } + struct DownstreamClosedWriter; + + impl AsyncWrite for DownstreamClosedWriter { + fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &[u8]) -> Poll> { + Poll::Ready(Err(crate::diagnostics::get::mark_get_object_downstream_closed(io::Error::new( + ErrorKind::BrokenPipe, + "injected downstream close", + )))) + } + + 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(())) + } + } + + #[derive(Clone, Default)] + struct CapturedLogs(Arc>>); + + struct CapturedLogWriter(Arc>>); + + impl CapturedLogs { + fn contents(&self) -> String { + String::from_utf8(self.0.lock().expect("captured logs mutex should not be poisoned").clone()) + .expect("captured logs should be valid UTF-8") + } + } + + impl std::io::Write for CapturedLogWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0 + .lock() + .expect("captured logs mutex should not be poisoned") + .extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogWriter; + + fn make_writer(&'a self) -> Self::Writer { + CapturedLogWriter(Arc::clone(&self.0)) + } + } + #[test] fn parallel_reader_constructor_variants_preserve_read_cost_and_verification_flags() { let erasure = Erasure::new(2, 1, 64); @@ -2215,6 +2281,47 @@ mod tests { assert_eq!(err.to_string(), "injected emit failure"); } + #[tokio::test(flavor = "current_thread")] + async fn erasure_decode_logs_reconstructed_downstream_close_at_debug() { + let logs = CapturedLogs::default(); + let subscriber = tracing_subscriber::fmt() + .with_max_level(tracing::Level::DEBUG) + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + + let erasure = Erasure::new(2, 1, 64); + let data: Vec = (0..64).collect(); + let shard_size = erasure.shard_size(); + let encoded = erasure.encode_data(&data).expect("test data should encode"); + let readers = vec![ + None, + Some(BitrotReader::new( + Cursor::new(encoded[1].to_vec()), + shard_size, + HashAlgorithm::None, + false, + )), + Some(BitrotReader::new( + Cursor::new(encoded[2].to_vec()), + shard_size, + HashAlgorithm::None, + false, + )), + ]; + + let mut writer = DownstreamClosedWriter; + let (written, err) = erasure.decode(&mut writer, readers, 0, data.len(), data.len()).await; + + assert_eq!(written, 0); + assert_eq!(err.expect("downstream close must still terminate the GET").kind(), ErrorKind::BrokenPipe); + let captured = logs.contents(); + assert!(captured.contains("Erasure decode stopped after downstream closed")); + assert!(!captured.contains("Erasure decode failed to emit reconstructed data")); + } + #[tokio::test] async fn test_erasure_decode_rejects_reader_count_and_range_overflow() { let erasure = Erasure::new(2, 1, 64); diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index ab5fb52b9..2c41676ac 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -4861,10 +4861,7 @@ async fn poll_merge_head(rx: &CancellationToken, in_channels: &mut [Receiver, entry: MetaCacheEntry) -> Result { tokio::select! { - result = out_channel.send(entry) => { - result.map_err(Error::other)?; - Ok(true) - } + result = out_channel.send(entry) => Ok(result.is_ok()), _ = rx.cancelled() => Ok(false), } } @@ -9858,6 +9855,28 @@ mod test { assert_eq!(results, vec!["obj-a", "obj-b"]); } + #[tokio::test] + async fn merge_entry_channels_treats_dropped_output_receiver_as_completion() { + let (tx_a, rx_a) = mpsc::channel(4); + let (tx_b, rx_b) = mpsc::channel(4); + let (out_tx, out_rx) = mpsc::channel(1); + + tx_a.send(test_meta_entry("obj-a")).await.unwrap(); + tx_b.send(test_meta_entry("obj-b")).await.unwrap(); + drop(tx_a); + drop(tx_b); + drop(out_rx); + + let result = timeout( + Duration::from_secs(1), + merge_entry_channels(CancellationToken::new(), vec![rx_a, rx_b], out_tx, 1), + ) + .await + .expect("merge should stop promptly when its consumer disconnects"); + + assert!(result.is_ok(), "consumer disconnect must not surface as a merge worker error"); + } + #[test] fn walk_ascending_versions_contract_reverses_newest_first_metadata() { // Documents the invariant the walk `versions_sort` handling relies on: