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, Describe, Emit, 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::drive_partial`](crate::de::StreamDeserializer::drive_partial)),
73/// the 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<'a>(value: &'a Self, state: &mut State) -> Result<Emit<'a>, Error> {
160        Vec::<T>::serialize(&value.items, state)
161    }
162
163    fn finish(value: &Self, state: &mut State) -> Result<(), Error> {
164        Vec::<T>::finish(&value.items, state)
165    }
166
167    #[inline]
168    fn __private_begin<'a>(value: &'a Self, state: &mut State) -> Result<Begin<'a>, Error> {
169        Vec::<T>::__private_begin(&value.items, state)
170    }
171
172    fn container_shape(value: &Self) -> ContainerShape {
173        Vec::<T>::container_shape(&value.items)
174    }
175
176    fn describe(value: &Self, d: &mut dyn Describe) {
177        Vec::<T>::describe(&value.items, 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
189/// What streamed values expect.
190const STREAMED_NAME: &str = "sequence";
191
192impl<'de, T: Deserialize<'de> + 'static> Deserialize<'de> for Streamed<T> {
193    fn deserialize_into<'out>(
194        out: &'out mut Option<Self>,
195        state: &mut State,
196    ) -> SinkHandle<'out, 'de> {
197        SinkHandle::arena(
198            StreamedSink {
199                out,
200                items: Vec::new(),
201            },
202            state,
203        )
204    }
205
206    fn expecting() -> Cow<'static, str> {
207        Cow::Borrowed(STREAMED_NAME)
208    }
209}
210
211struct StreamedSink<'a, T> {
212    out: &'a mut Option<Streamed<T>>,
213    items: Vec<T>,
214}
215
216impl<'a, 'de, T: Deserialize<'de> + 'static> Sink<'de> for StreamedSink<'a, T> {
217    fn expecting(&self) -> Cow<'_, str> {
218        Cow::Borrowed(STREAMED_NAME)
219    }
220
221    fn seq(&mut self, _state: &mut State) -> Result<(), Error> {
222        Ok(())
223    }
224
225    fn next_value(&mut self, state: &mut State) -> Result<SinkHandle<'_, 'de>, Error> {
226        Ok(SinkHandle::arena(
227            ElementSink {
228                sink: OwnedSink::deserialize(state),
229                items: &mut self.items,
230            },
231            state,
232        ))
233    }
234
235    fn __private_value_atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
236        let mut value = None;
237        T::__private_atom_into(&mut value, atom, state)?;
238        if let Some(value) = value {
239            complete(value, &mut self.items, state);
240        }
241        Ok(())
242    }
243
244    fn __private_borrowed_value_atom(
245        &mut self,
246        atom: Atom<'de>,
247        state: &mut State,
248    ) -> Result<(), Error> {
249        let mut value = None;
250        T::__private_borrowed_atom_into(&mut value, atom, state)?;
251        if let Some(value) = value {
252            complete(value, &mut self.items, state);
253        }
254        Ok(())
255    }
256
257    fn finish(&mut self, _state: &mut State) -> Result<(), Error> {
258        *self.out = Some(Streamed {
259            items: core::mem::take(&mut self.items),
260        });
261        Ok(())
262    }
263}
264
265/// Deserializes an element and hands it out once it's complete.
266struct ElementSink<'a, 'de, T> {
267    sink: OwnedSink<'de, T>,
268    items: &'a mut Vec<T>,
269}
270
271impl<'a, 'de, T: Deserialize<'de> + 'static> Sink<'de> for ElementSink<'a, 'de, T> {
272    fn atom(&mut self, atom: Atom, state: &mut State) -> Result<(), Error> {
273        self.sink.get_mut().atom(atom, state)
274    }
275
276    fn borrowed_atom(&mut self, atom: Atom<'de>, state: &mut State) -> Result<(), Error> {
277        self.sink.get_mut().borrowed_atom(atom, state)
278    }
279
280    fn map(&mut self, state: &mut State) -> Result<(), Error> {
281        self.sink.get_mut().map(state)
282    }
283
284    fn seq(&mut self, state: &mut State) -> Result<(), Error> {
285        self.sink.get_mut().seq(state)
286    }
287
288    forward_to_owned!(sink);
289
290    fn finish(&mut self, state: &mut State) -> Result<(), Error> {
291        self.sink.get_mut().finish(state)?;
292        if let Some(value) = self.sink.take() {
293            complete(value, self.items, state);
294        }
295        Ok(())
296    }
297}