Skip to main content

deser_core/stream/
streamed.rs

1use alloc::borrow::Cow;
2use alloc::vec::Vec;
3use core::fmt;
4use core::ops::{Deref, DerefMut};
5
6use crate::State;
7use crate::de::{Deserialize, OwnedSink, Sink, SinkHandle};
8use crate::error::Error;
9use crate::event::{Atom, ContainerShape};
10use crate::ser::{Begin, Chunk, Describe, Serialize};
11
12/// A sequence whose elements can be handed out while it's read.
13///
14/// `Streamed<T>` behaves like a `Vec<T>`: it serializes and deserializes
15/// as a sequence and the elements are collected.  But when a value
16/// containing it is read from a stream with `Reader::read_next` of
17/// `deser::io` (or the equivalent of an async runtime), the
18/// elements are handed out as soon as they are complete instead.  This allows processing large (or unbounded) sequences
19/// within a value while the stream is read:
20///
21/// ```
22/// # fn example() -> Result<(), deser::Error> {
23/// # #[cfg(all(feature = "derive", feature = "io"))] {
24/// use deser::Deserialize;
25/// use deser::io::Reader;
26/// use deser::stream::{Part, Streamed};
27/// # use deser::de::{DeserializeDriver, Frame, StreamDeserializer};
28/// # use deser::{Error, Event};
29/// # /// Numbers on a line of their own form a page (with a sequence of the numbers).
30/// # struct Pages;
31/// # impl StreamDeserializer for Pages {
32/// #     fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
33/// #         Ok(match input.iter().position(|&b| b == b'\n') {
34/// #             Some(end) => Frame::Value { start: 0, end, consumed: end + 1 },
35/// #             None if eof && input.is_empty() => Frame::End,
36/// #             None if eof => Frame::Value { start: 0, end: input.len(), consumed: input.len() },
37/// #             None => Frame::Incomplete { consumed: 0 },
38/// #         })
39/// #     }
40/// #     fn drive_frame<'de>(&mut self, frame: &'de [u8], driver: &mut DeserializeDriver<'_, 'de>) -> Result<(), Error> {
41/// #         driver.emit(Event::map_start())?;
42/// #         driver.emit("items")?;
43/// #         driver.emit(Event::seq_start())?;
44/// #         for number in std::str::from_utf8(frame).unwrap().split(' ') {
45/// #             driver.emit(number.parse::<u64>().unwrap())?;
46/// #         }
47/// #         driver.emit(Event::SeqEnd)?;
48/// #         driver.emit(Event::MapEnd)
49/// #     }
50/// # }
51///
52/// #[derive(Deserialize)]
53/// struct Page {
54///     items: Streamed<u32>,
55/// }
56///
57/// // `Pages` is the stream deserializer of a format with pages of numbers
58/// let mut reader = Reader::new(&b"1 2 3"[..], Pages);
59/// let mut items = Vec::new();
60/// while let Some(next) = reader.read_next::<Page, u32>()? {
61///     match next {
62///         Part::Element(item) => items.push(item),
63///         // the elements were handed out, they are not in the page
64///         Part::Done(page) => assert!(page.items.is_empty()),
65///     }
66/// }
67/// assert_eq!(items, [1, 2, 3]);
68/// # } Ok(()) } example().unwrap();
69/// ```
70///
71/// If the format can deserialize values while their input arrives (see
72/// [`StreamDeserializer::feed`](crate::de::StreamDeserializer::feed)), the
73/// memory used does not depend on the length of the sequence.  The
74/// elements have to be owned (they cannot borrow from the input).  Without
75/// IO, the elements are handed out by an
76/// [`ElementReader`](crate::stream::ElementReader).
77#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
78pub struct Streamed<T> {
79    items: Vec<T>,
80}
81
82impl<T> Streamed<T> {
83    /// Creates an empty sequence.
84    pub fn new() -> Streamed<T> {
85        Streamed { items: Vec::new() }
86    }
87
88    /// Returns the collected elements.
89    pub fn into_vec(self) -> Vec<T> {
90        self.items
91    }
92}
93
94impl<T> Default for Streamed<T> {
95    fn default() -> Streamed<T> {
96        Streamed::new()
97    }
98}
99
100impl<T: fmt::Debug> fmt::Debug for Streamed<T> {
101    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
102        self.items.fmt(f)
103    }
104}
105
106impl<T> From<Vec<T>> for Streamed<T> {
107    fn from(items: Vec<T>) -> Streamed<T> {
108        Streamed { items }
109    }
110}
111
112impl<T> From<Streamed<T>> for Vec<T> {
113    fn from(value: Streamed<T>) -> Vec<T> {
114        value.items
115    }
116}
117
118impl<T> FromIterator<T> for Streamed<T> {
119    fn from_iter<I: IntoIterator<Item = T>>(iter: I) -> Streamed<T> {
120        Streamed {
121            items: iter.into_iter().collect(),
122        }
123    }
124}
125
126impl<T> IntoIterator for Streamed<T> {
127    type Item = T;
128    type IntoIter = alloc::vec::IntoIter<T>;
129
130    fn into_iter(self) -> Self::IntoIter {
131        self.items.into_iter()
132    }
133}
134
135impl<'a, T> IntoIterator for &'a Streamed<T> {
136    type Item = &'a T;
137    type IntoIter = core::slice::Iter<'a, T>;
138
139    fn into_iter(self) -> Self::IntoIter {
140        self.items.iter()
141    }
142}
143
144impl<T> Deref for Streamed<T> {
145    type Target = Vec<T>;
146
147    fn deref(&self) -> &Vec<T> {
148        &self.items
149    }
150}
151
152impl<T> DerefMut for Streamed<T> {
153    fn deref_mut(&mut self) -> &mut Vec<T> {
154        &mut self.items
155    }
156}
157
158impl<T: Serialize> Serialize for Streamed<T> {
159    fn serialize(&self, state: &mut State) -> Result<Chunk<'_>, Error> {
160        self.items.serialize(state)
161    }
162
163    fn finish(&self, state: &mut State) -> Result<(), Error> {
164        self.items.finish(state)
165    }
166
167    #[inline]
168    fn __private_begin(&self, state: &mut State) -> Result<Begin<'_>, Error> {
169        self.items.__private_begin(state)
170    }
171
172    fn container_shape(&self) -> ContainerShape {
173        self.items.container_shape()
174    }
175
176    fn describe(&self, d: &mut dyn Describe) {
177        self.items.describe(d)
178    }
179}
180
181/// Hands out a complete element or collects it.
182fn complete<T: Send + 'static>(value: T, items: &mut Vec<T>, state: &State) {
183    let Err(value) = crate::stream::elements::hand_out(value, state) else {
184        return;
185    };
186    items.push(value);
187}
188
189impl<'de, T: Deserialize<'de> + 'static> Deserialize<'de> for Streamed<T> {
190    fn deserialize_into<'out>(
191        out: &'out mut Option<Self>,
192        state: &mut State,
193    ) -> SinkHandle<'out, 'de> {
194        SinkHandle::arena(
195            StreamedSink {
196                out,
197                items: Vec::new(),
198            },
199            state,
200        )
201    }
202}
203
204struct StreamedSink<'a, T> {
205    out: &'a mut Option<Streamed<T>>,
206    items: Vec<T>,
207}
208
209impl<'a, 'de, T: Deserialize<'de> + 'static> Sink<'de> for StreamedSink<'a, T> {
210    fn expecting(&self) -> Cow<'_, str> {
211        Cow::Borrowed("sequence")
212    }
213
214    fn seq(&mut self, _state: &mut State) -> Result<(), Error> {
215        Ok(())
216    }
217
218    fn next_value(&mut self, state: &mut State) -> Result<SinkHandle<'_, 'de>, Error> {
219        Ok(SinkHandle::arena(
220            ElementSink {
221                sink: OwnedSink::deserialize(state),
222                items: &mut self.items,
223            },
224            state,
225        ))
226    }
227
228    fn __private_value_atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
229        let mut value = None;
230        T::__private_atom_into(&mut value, atom, state)?;
231        if let Some(value) = value {
232            complete(value, &mut self.items, state);
233        }
234        Ok(())
235    }
236
237    fn __private_borrowed_value_atom(
238        &mut self,
239        atom: Atom<'de>,
240        state: &mut State,
241    ) -> Result<(), Error> {
242        let mut value = None;
243        T::__private_borrowed_atom_into(&mut value, atom, state)?;
244        if let Some(value) = value {
245            complete(value, &mut self.items, state);
246        }
247        Ok(())
248    }
249
250    fn finish(&mut self, _state: &mut State) -> Result<(), Error> {
251        *self.out = Some(Streamed {
252            items: core::mem::take(&mut self.items),
253        });
254        Ok(())
255    }
256}
257
258/// Deserializes an element and hands it out once it's complete.
259struct ElementSink<'a, 'de, T> {
260    sink: OwnedSink<'de, T>,
261    items: &'a mut Vec<T>,
262}
263
264impl<'a, 'de, T: Deserialize<'de> + 'static> Sink<'de> for ElementSink<'a, 'de, T> {
265    fn atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
266        self.sink.borrow_mut().atom(atom, state)
267    }
268
269    fn borrowed_atom(&mut self, atom: Atom<'de>, state: &mut State) -> Result<(), Error> {
270        self.sink.borrow_mut().borrowed_atom(atom, state)
271    }
272
273    fn map(&mut self, state: &mut State) -> Result<(), Error> {
274        self.sink.borrow_mut().map(state)
275    }
276
277    fn seq(&mut self, state: &mut State) -> Result<(), Error> {
278        self.sink.borrow_mut().seq(state)
279    }
280
281    fn next_key(&mut self, state: &mut State) -> Result<SinkHandle<'_, 'de>, Error> {
282        self.sink.borrow_mut().next_key(state)
283    }
284
285    fn next_value(&mut self, state: &mut State) -> Result<SinkHandle<'_, 'de>, Error> {
286        self.sink.borrow_mut().next_value(state)
287    }
288
289    fn __private_key_atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
290        self.sink.borrow_mut().__private_key_atom(atom, state)
291    }
292
293    fn __private_value_atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
294        self.sink.borrow_mut().__private_value_atom(atom, state)
295    }
296
297    fn __private_borrowed_key_atom(
298        &mut self,
299        atom: Atom<'de>,
300        state: &mut State,
301    ) -> Result<(), Error> {
302        self.sink
303            .borrow_mut()
304            .__private_borrowed_key_atom(atom, state)
305    }
306
307    fn __private_borrowed_value_atom(
308        &mut self,
309        atom: Atom<'de>,
310        state: &mut State,
311    ) -> Result<(), Error> {
312        self.sink
313            .borrow_mut()
314            .__private_borrowed_value_atom(atom, state)
315    }
316
317    fn value_for_key(
318        &mut self,
319        key: &str,
320        state: &mut State,
321    ) -> Result<Option<SinkHandle<'_, 'de>>, Error> {
322        self.sink.borrow_mut().value_for_key(key, state)
323    }
324
325    fn recover(&mut self, err: Error, state: &mut State) -> Result<(), Error> {
326        self.sink.borrow_mut().recover(err, state)
327    }
328
329    fn finish(&mut self, state: &mut State) -> Result<(), Error> {
330        self.sink.borrow_mut().finish(state)?;
331        if let Some(value) = self.sink.take() {
332            complete(value, self.items, state);
333        }
334        Ok(())
335    }
336
337    fn expecting(&self) -> Cow<'_, str> {
338        self.sink.borrow().expecting()
339    }
340}