Skip to main content

zstd/stream/read/
mod.rs

1//! Implement pull-based [`Read`] trait for both compressing and decompressing.
2use std::io::{self, BufRead, BufReader, Read};
3
4use crate::dict::{DecoderDictionary, EncoderDictionary};
5use crate::stream::{raw, zio};
6use zstd_safe;
7
8#[cfg(test)]
9mod tests;
10
11/// A decoder that decompress input data from another `Read`.
12///
13/// This allows to read a stream of compressed data
14/// (good for files or heavy network stream).
15pub struct Decoder<'a, R> {
16    reader: zio::Reader<R, raw::Decoder<'a>>,
17}
18
19/// An encoder that compress input data from another `Read`.
20pub struct Encoder<'a, R> {
21    reader: zio::Reader<R, raw::Encoder<'a>>,
22}
23
24impl<R: Read> Decoder<'static, BufReader<R>> {
25    /// Creates a new decoder.
26    pub fn new(reader: R) -> io::Result<Self> {
27        let buffer_size = zstd_safe::DCtx::in_size();
28
29        Self::with_buffer(BufReader::with_capacity(buffer_size, reader))
30    }
31}
32
33impl<R: BufRead> Decoder<'static, R> {
34    /// Creates a new decoder around a `BufRead`.
35    pub fn with_buffer(reader: R) -> io::Result<Self> {
36        Self::with_dictionary(reader, &[])
37    }
38    /// Creates a new decoder, using an existing dictionary.
39    ///
40    /// The dictionary must be the same as the one used during compression.
41    pub fn with_dictionary(reader: R, dictionary: &[u8]) -> io::Result<Self> {
42        let decoder = raw::Decoder::with_dictionary(dictionary)?;
43        let reader = zio::Reader::new(reader, decoder);
44
45        Ok(Decoder { reader })
46    }
47}
48impl<'a, R: BufRead> Decoder<'a, R> {
49    /// Creates a new decoder which employs the provided context for deserialization.
50    pub fn with_context(
51        reader: R,
52        context: &'a mut zstd_safe::DCtx<'static>,
53    ) -> Self {
54        Self {
55            reader: zio::Reader::new(
56                reader,
57                raw::Decoder::with_context(context),
58            ),
59        }
60    }
61
62    /// Sets this `Decoder` to stop after the first frame.
63    ///
64    /// By default, it keeps concatenating frames until EOF is reached.
65    #[must_use]
66    pub fn single_frame(mut self) -> Self {
67        self.reader.set_single_frame();
68        self
69    }
70
71    /// Creates a new decoder, using an existing `DecoderDictionary`.
72    ///
73    /// The dictionary must be the same as the one used during compression.
74    pub fn with_prepared_dictionary<'b>(
75        reader: R,
76        dictionary: &'a DecoderDictionary<'b>,
77    ) -> io::Result<Self>
78    where
79        'b: 'a,
80    {
81        let decoder = raw::Decoder::with_prepared_dictionary(dictionary)?;
82        let reader = zio::Reader::new(reader, decoder);
83
84        Ok(Decoder { reader })
85    }
86
87    /// Creates a new decoder, using a ref prefix.
88    ///
89    /// The prefix must be the same as the one used during compression.
90    pub fn with_ref_prefix<'b>(
91        reader: R,
92        ref_prefix: &'b [u8],
93    ) -> io::Result<Self>
94    where
95        'b: 'a,
96    {
97        let decoder = raw::Decoder::with_ref_prefix(ref_prefix)?;
98        let reader = zio::Reader::new(reader, decoder);
99
100        Ok(Decoder { reader })
101    }
102
103    /// Recommendation for the size of the output buffer.
104    pub fn recommended_output_size() -> usize {
105        zstd_safe::DCtx::out_size()
106    }
107
108    /// Acquire a reference to the underlying reader.
109    pub fn get_ref(&self) -> &R {
110        self.reader.reader()
111    }
112
113    /// Acquire a mutable reference to the underlying reader.
114    ///
115    /// Note that mutation of the reader may result in surprising results if
116    /// this decoder is continued to be used.
117    pub fn get_mut(&mut self) -> &mut R {
118        self.reader.reader_mut()
119    }
120
121    /// Consume the rest of the current frame, so the underlying reader is
122    /// left pointing just after it.
123    ///
124    /// zstd can hand out the last of the decoded data before it has read the
125    /// frame epilogue, so a reader that stops as soon as it has the bytes it
126    /// wanted leaves the input somewhere inside the frame. Call this to line
127    /// the reader back up with the end of the frame - for instance to carry on
128    /// reading whatever follows the compressed section.
129    ///
130    /// This does not touch the underlying reader if the frame is already
131    /// complete, and never starts decoding the next frame.
132    pub fn finish_frame(&mut self) -> io::Result<()> {
133        self.reader.finish_frame()
134    }
135
136    /// Return the inner `Read`.
137    ///
138    /// Calling `finish()` is not *required* after reading a stream -
139    /// just use it if you need to get the `Read` back.
140    ///
141    /// This consumes the rest of the current frame first, so the reader is
142    /// returned pointing just after it - see [`Self::finish_frame()`]. That
143    /// may read from the underlying reader, and any error doing so is
144    /// ignored; call `finish_frame()` directly if you need to see it.
145    pub fn finish(mut self) -> R {
146        let _ = self.finish_frame();
147        self.reader.into_inner()
148    }
149
150    crate::decoder_common!(reader);
151}
152
153impl<R: BufRead> Read for Decoder<'_, R> {
154    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
155        self.reader.read(buf)
156    }
157}
158
159impl<R: Read> Encoder<'static, BufReader<R>> {
160    /// Creates a new encoder.
161    pub fn new(reader: R, level: i32) -> io::Result<Self> {
162        let buffer_size = zstd_safe::CCtx::in_size();
163
164        Self::with_buffer(BufReader::with_capacity(buffer_size, reader), level)
165    }
166}
167
168impl<R: BufRead> Encoder<'static, R> {
169    /// Creates a new encoder around a `BufRead`.
170    pub fn with_buffer(reader: R, level: i32) -> io::Result<Self> {
171        Self::with_dictionary(reader, level, &[])
172    }
173
174    /// Creates a new encoder, using an existing dictionary.
175    ///
176    /// The dictionary must be the same as the one used during compression.
177    pub fn with_dictionary(
178        reader: R,
179        level: i32,
180        dictionary: &[u8],
181    ) -> io::Result<Self> {
182        let encoder = raw::Encoder::with_dictionary(level, dictionary)?;
183        let reader = zio::Reader::new(reader, encoder);
184
185        Ok(Encoder { reader })
186    }
187}
188
189impl<'a, R: BufRead> Encoder<'a, R> {
190    /// Creates a new encoder which employs the provided context for serialization.
191    pub fn with_context(
192        reader: R,
193        context: &'a mut zstd_safe::CCtx<'static>,
194    ) -> Self {
195        Self {
196            reader: zio::Reader::new(
197                reader,
198                raw::Encoder::with_context(context),
199            ),
200        }
201    }
202
203    /// Creates a new encoder, using an existing `EncoderDictionary`.
204    ///
205    /// The dictionary must be the same as the one used during compression.
206    pub fn with_prepared_dictionary<'b>(
207        reader: R,
208        dictionary: &'a EncoderDictionary<'b>,
209    ) -> io::Result<Self>
210    where
211        'b: 'a,
212    {
213        let encoder = raw::Encoder::with_prepared_dictionary(dictionary)?;
214        let reader = zio::Reader::new(reader, encoder);
215
216        Ok(Encoder { reader })
217    }
218
219    /// Recommendation for the size of the output buffer.
220    pub fn recommended_output_size() -> usize {
221        zstd_safe::CCtx::out_size()
222    }
223
224    /// Acquire a reference to the underlying reader.
225    pub fn get_ref(&self) -> &R {
226        self.reader.reader()
227    }
228
229    /// Acquire a mutable reference to the underlying reader.
230    ///
231    /// Note that mutation of the reader may result in surprising results if
232    /// this encoder is continued to be used.
233    pub fn get_mut(&mut self) -> &mut R {
234        self.reader.reader_mut()
235    }
236
237    /// Flush any internal buffer.
238    ///
239    /// This ensures all input consumed so far is compressed.
240    ///
241    /// Since it prevents bundling currently buffered data with future input,
242    /// it may affect compression ratio.
243    ///
244    /// * Returns the number of bytes written to `out`.
245    /// * Returns `Ok(0)` when everything has been flushed.
246    pub fn flush(&mut self, out: &mut [u8]) -> io::Result<usize> {
247        self.reader.flush(out)
248    }
249
250    /// Return the inner `Read`.
251    ///
252    /// Calling `finish()` is not *required* after reading a stream -
253    /// just use it if you need to get the `Read` back.
254    pub fn finish(self) -> R {
255        self.reader.into_inner()
256    }
257
258    crate::encoder_common!(reader);
259}
260
261impl<R: BufRead> Read for Encoder<'_, R> {
262    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
263        self.reader.read(buf)
264    }
265}
266
267fn _assert_traits() {
268    use std::io::Cursor;
269
270    fn _assert_send<T: Send>(_: T) {}
271
272    _assert_send(Decoder::new(Cursor::new(Vec::new())));
273    _assert_send(Encoder::new(Cursor::new(Vec::new()), 1));
274}