Skip to main content

base64_ng/stream/
encoder_reader.rs

1use super::{
2    EncoderDriver, OutputQueue, redacted_inner_state, stream_encoder_failed_error,
3    wrapped_reader_overreported_error,
4};
5use crate::{Alphabet, Engine};
6use std::io::{self, Read};
7
8/// A streaming Base64 encoder for `std::io::Read`.
9pub struct EncoderReader<R, A, const PAD: bool>
10where
11    A: Alphabet,
12{
13    inner: Option<R>,
14    engine: Engine<A, PAD>,
15    driver: EncoderDriver,
16    output: OutputQueue<1024>,
17    finished: bool,
18    failed: bool,
19}
20
21impl<R, A, const PAD: bool> EncoderReader<R, A, PAD>
22where
23    A: Alphabet,
24{
25    /// Creates a new streaming encoder reader.
26    #[must_use]
27    pub fn new(inner: R, engine: Engine<A, PAD>) -> Self {
28        Self {
29            inner: Some(inner),
30            engine,
31            driver: EncoderDriver::new::<A, PAD>(),
32            output: OutputQueue::new(),
33            finished: false,
34            failed: false,
35        }
36    }
37
38    /// Returns a shared reference to the wrapped reader.
39    #[must_use]
40    pub fn get_ref(&self) -> &R {
41        self.inner_ref()
42    }
43
44    /// Returns a mutable reference to the wrapped reader.
45    pub fn get_mut(&mut self) -> &mut R {
46        self.inner_mut()
47    }
48
49    /// Returns the Base64 engine used by this adapter.
50    #[must_use]
51    pub const fn engine(&self) -> Engine<A, PAD> {
52        self.engine
53    }
54
55    /// Returns whether this adapter uses padded Base64.
56    #[must_use]
57    pub const fn is_padded(&self) -> bool {
58        PAD
59    }
60
61    /// Returns the number of raw input bytes currently buffered until a
62    /// complete 3-byte Base64 encode quantum is available.
63    #[must_use]
64    pub const fn pending_len(&self) -> usize {
65        self.driver.pending_input_len()
66    }
67
68    /// Returns whether this encoder reader currently holds a partial input
69    /// quantum.
70    #[must_use]
71    pub const fn has_pending_input(&self) -> bool {
72        self.pending_len() != 0
73    }
74
75    /// Returns how many additional raw input bytes are needed to complete
76    /// the currently buffered encode quantum.
77    ///
78    /// Returns `0` when no partial input quantum is buffered.
79    #[must_use]
80    pub const fn pending_input_needed_len(&self) -> usize {
81        if self.has_pending_input() {
82            3 - self.pending_len()
83        } else {
84            0
85        }
86    }
87
88    /// Returns the number of encoded bytes currently buffered and ready to
89    /// be read before this adapter polls the wrapped reader again.
90    #[must_use]
91    pub const fn buffered_output_len(&self) -> usize {
92        self.output.len()
93    }
94
95    /// Returns the maximum number of encoded bytes this adapter can buffer
96    /// before returning bytes to the caller.
97    #[must_use]
98    pub const fn buffered_output_capacity(&self) -> usize {
99        self.output.capacity()
100    }
101
102    /// Returns how many more encoded bytes can be buffered before this
103    /// adapter must return bytes to the caller.
104    #[must_use]
105    pub const fn buffered_output_remaining_capacity(&self) -> usize {
106        self.output.available_capacity()
107    }
108
109    /// Returns whether this encoder reader currently has encoded output
110    /// waiting in its internal queue.
111    #[must_use]
112    pub const fn has_buffered_output(&self) -> bool {
113        !self.output.is_empty()
114    }
115
116    /// Returns whether this encoder reader has reached EOF in the wrapped
117    /// reader.
118    ///
119    /// This may become `true` before [`Self::is_finished`] when encoded
120    /// output is still buffered for the caller.
121    #[must_use]
122    pub const fn has_finished_input(&self) -> bool {
123        self.finished
124    }
125
126    /// Returns whether this reader has reached EOF and has no encoded
127    /// output buffered for the caller.
128    #[must_use]
129    pub const fn is_finished(&self) -> bool {
130        self.finished && self.output.is_empty()
131    }
132
133    /// Returns whether this adapter has failed closed after an internal
134    /// stream error.
135    #[must_use]
136    pub const fn is_failed(&self) -> bool {
137        self.failed
138    }
139
140    /// Returns whether [`Self::try_into_inner`] can recover the wrapped
141    /// reader without discarding pending input or buffered encoded output.
142    #[must_use]
143    pub const fn can_into_inner(&self) -> bool {
144        self.is_finished() && !self.failed
145    }
146
147    /// Consumes the encoder reader and returns the wrapped reader.
148    #[must_use]
149    pub fn into_inner(mut self) -> R {
150        self.take_inner()
151    }
152
153    /// Consumes the encoder reader only after the encoded stream is fully
154    /// drained.
155    ///
156    /// This is a checked alternative to [`Self::into_inner`] for callers
157    /// that want to avoid accidentally discarding pending input or encoded
158    /// output buffered inside the adapter.
159    #[allow(clippy::result_large_err)]
160    pub fn try_into_inner(mut self) -> Result<R, Self> {
161        if !self.can_into_inner() {
162            return Err(self);
163        }
164        Ok(self.take_inner())
165    }
166
167    fn inner_ref(&self) -> &R {
168        match &self.inner {
169            Some(inner) => inner,
170            None => unreachable!("stream encoder reader inner reader was already taken"),
171        }
172    }
173
174    fn inner_mut(&mut self) -> &mut R {
175        match &mut self.inner {
176            Some(inner) => inner,
177            None => unreachable!("stream encoder reader inner reader was already taken"),
178        }
179    }
180
181    fn take_inner(&mut self) -> R {
182        match self.inner.take() {
183            Some(inner) => inner,
184            None => unreachable!("stream encoder reader inner reader was already taken"),
185        }
186    }
187
188    fn clear_pending(&mut self) {
189        self.driver.wipe();
190    }
191}
192
193impl<R, A, const PAD: bool> Drop for EncoderReader<R, A, PAD>
194where
195    A: Alphabet,
196{
197    fn drop(&mut self) {
198        self.clear_pending();
199        self.output.clear_all();
200    }
201}
202
203impl<R, A, const PAD: bool> core::fmt::Debug for EncoderReader<R, A, PAD>
204where
205    A: Alphabet,
206{
207    fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
208        formatter
209            .debug_struct("EncoderReader")
210            .field("inner", &redacted_inner_state(self.inner.is_some()))
211            .field("engine", &self.engine)
212            .field("driver", &"<redacted>")
213            .field("pending", &"<redacted>")
214            .field("pending_len", &self.pending_len())
215            .field("pending_input_needed_len", &self.pending_input_needed_len())
216            .field("buffered_output_len", &self.output.len())
217            .field("buffered_output_capacity", &self.output.capacity())
218            .field(
219                "buffered_output_remaining_capacity",
220                &self.output.available_capacity(),
221            )
222            .field("can_into_inner", &self.can_into_inner())
223            .field("finished", &self.finished)
224            .field("failed", &self.failed)
225            .finish()
226    }
227}
228
229impl<R, A, const PAD: bool> Read for EncoderReader<R, A, PAD>
230where
231    R: Read,
232    A: Alphabet,
233{
234    fn read(&mut self, output: &mut [u8]) -> io::Result<usize> {
235        if self.failed {
236            return Err(stream_encoder_failed_error());
237        }
238
239        if output.is_empty() {
240            return Ok(0);
241        }
242
243        while self.output.is_empty() && !self.finished {
244            self.fill_output()?;
245        }
246
247        Ok(self.output.pop_slice(output))
248    }
249}
250
251impl<R, A, const PAD: bool> EncoderReader<R, A, PAD>
252where
253    R: Read,
254    A: Alphabet,
255{
256    fn fill_output(&mut self) -> io::Result<()> {
257        let mut input = [0u8; 768];
258        let available = input.len();
259        let read = match self.inner_mut().read(&mut input[..available]) {
260            Ok(read) => read,
261            Err(err) => {
262                crate::wipe_bytes(&mut input);
263                return Err(err);
264            }
265        };
266        if read > available {
267            crate::wipe_bytes(&mut input);
268            self.clear_pending();
269            self.failed = true;
270            return Err(wrapped_reader_overreported_error());
271        }
272        if read == 0 {
273            crate::wipe_bytes(&mut input);
274            self.finished = true;
275            if let Err(err) = self.finish_driver() {
276                self.failed = true;
277                return Err(err);
278            }
279            return Ok(());
280        }
281
282        let result = self.update_driver(&input[..read]);
283        crate::wipe_bytes(&mut input);
284        if let Err(err) = result {
285            self.failed = true;
286            return Err(err);
287        }
288        Ok(())
289    }
290
291    fn update_driver(&mut self, input: &[u8]) -> io::Result<()> {
292        let mut encoded = [0u8; 1024];
293        let step = match self.driver.update(input, &mut encoded) {
294            Ok(step) => step,
295            Err(err) => {
296                crate::wipe_bytes(&mut encoded);
297                self.clear_pending();
298                return Err(err);
299            }
300        };
301        let progress = step.progress();
302        let result = if progress.input_consumed() == input.len() {
303            self.output
304                .push_slice(&encoded[..progress.output_produced()])
305        } else {
306            Err(io::Error::other(
307                "base64 stream encoder did not accept bounded reader input",
308            ))
309        };
310        crate::wipe_bytes(&mut encoded);
311        result
312    }
313
314    fn finish_driver(&mut self) -> io::Result<()> {
315        let mut encoded = [0u8; 4];
316        let step = match self.driver.finish(&mut encoded) {
317            Ok(step) => step,
318            Err(err) => {
319                crate::wipe_bytes(&mut encoded);
320                self.clear_pending();
321                return Err(err);
322            }
323        };
324        let result = self
325            .output
326            .push_slice(&encoded[..step.progress().output_produced()]);
327        crate::wipe_bytes(&mut encoded);
328        result
329    }
330}