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}