diff --git a/crates/async-compression/src/generic/bufread/decoder.rs b/crates/async-compression/src/generic/bufread/decoder.rs index 9cad7380..35dc0cb0 100644 --- a/crates/async-compression/src/generic/bufread/decoder.rs +++ b/crates/async-compression/src/generic/bufread/decoder.rs @@ -2,7 +2,6 @@ use crate::{ codecs::DecodeV2, core::util::{PartialBuffer, WriteBuffer}, }; - use std::{io::Result, ops::ControlFlow}; #[derive(Debug)] @@ -11,6 +10,7 @@ enum State { Flushing, Done, Next, + Error(std::io::Error), } #[derive(Debug)] @@ -54,7 +54,14 @@ impl Decoder { Ok(true) => State::Flushing, // ignore the first error, occurs when input is empty // but we need to run decode to flush - Err(err) if !first => return ControlFlow::Break(Err(err)), + Err(err) if !first => { + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } + } // poll for more data for the next decode _ => break, } @@ -66,7 +73,12 @@ impl Decoder { Ok(true) => { if self.multiple_members { if let Err(err) = decoder.reinit() { - return ControlFlow::Break(Err(err)); + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } } // The decode stage might consume all the input, @@ -78,7 +90,14 @@ impl Decoder { } } Ok(false) => State::Flushing, - Err(err) => return ControlFlow::Break(Err(err)), + Err(err) => { + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } + } } } @@ -95,6 +114,13 @@ impl Decoder { State::Decoding } } + + State::Error(_) => { + let State::Error(err) = std::mem::replace(&mut self.state, State::Done) else { + unreachable!() + }; + return ControlFlow::Break(Err(err)); + } }; if output.has_no_spare_space() { diff --git a/crates/async-compression/src/generic/bufread/encoder.rs b/crates/async-compression/src/generic/bufread/encoder.rs index 4ce90df1..dde539bb 100644 --- a/crates/async-compression/src/generic/bufread/encoder.rs +++ b/crates/async-compression/src/generic/bufread/encoder.rs @@ -10,6 +10,7 @@ enum State { Flushing, Finishing, Done, + Error(std::io::Error), } #[derive(Debug)] @@ -53,7 +54,12 @@ impl Encoder { State::Finishing } else { if let Err(err) = encoder.encode(input, output) { - return ControlFlow::Break(Err(err)); + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } } *read += input.written().len(); @@ -72,16 +78,37 @@ impl Encoder { break; } Ok(false) => State::Flushing, - Err(err) => return ControlFlow::Break(Err(err)), + Err(err) => { + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } + } }, State::Finishing => match encoder.finish(output) { Ok(true) => State::Done, Ok(false) => State::Finishing, - Err(err) => return ControlFlow::Break(Err(err)), + Err(err) => { + self.state = State::Error(err); + if output.written_len() > 0 { + return ControlFlow::Break(Ok(())); + } else { + continue; + } + } }, State::Done => return ControlFlow::Break(Ok(())), + + State::Error(_) => { + let State::Error(err) = std::mem::replace(&mut self.state, State::Done) else { + unreachable!() + }; + return ControlFlow::Break(Err(err)); + } }; if output.has_no_spare_space() { diff --git a/crates/async-compression/tests/gzip.rs b/crates/async-compression/tests/gzip.rs index 3ece5cd7..e976d97b 100644 --- a/crates/async-compression/tests/gzip.rs +++ b/crates/async-compression/tests/gzip.rs @@ -119,3 +119,24 @@ fn gzip_bufread_chunks_compress_flushes_when_reader_pending() { chunks_received.load(Ordering::Relaxed) ); } + +#[test] +#[ntest::timeout(1000)] +#[cfg(feature = "futures-io")] +fn gzip_bufread_chunks_decompress_without_footer_emits_all_payload() { + use flate2::bufread::GzDecoder; + use std::io::Read; + + let mut bytes = compress_with_header(&[1, 2, 3, 4, 5, 6]); + + // Remove the footer. + bytes.truncate(bytes.len() - 8); + + let mut decoder = GzDecoder::new(bytes.as_slice()); + + let mut output = vec![]; + let result = decoder.read_to_end(&mut output); + + assert!(result.is_err()); + assert_eq!(output, &[1, 2, 3, 4, 5, 6][..]); +}