Skip to main content

deser_core/io/
mod.rs

1//! Reading and writing values from and to streams.
2//!
3//! This module requires the `io` feature (which is enabled by default).
4//! For values in memory, the formats have deserializers (which read values
5//! from slices) and serializers (which write values into buffers) that do
6//! not need this module.
7//!
8//! Data formats parse complete inputs (slices) and serialize into complete
9//! outputs.  For streams they have stream serializers and stream
10//! deserializers which do not do IO themselves (see
11//! [`deser::stream`](crate::stream)).  This module connects them to
12//! [`Read`] and [`Write`], such as files, sockets or pipes.  A [`Reader`]
13//! reads values with a
14//! [`StreamDeserializer`], a [`Writer`]
15//! writes values with a [`StreamSerializer`].
16//! The configurations of the formats create them (`config.reader(input)`
17//! and `config.writer(output)`):
18//!
19//! ```
20//! # fn example() -> Result<(), deser::Error> {
21//! use deser::io::{Reader, Writer};
22//! # use deser::de::{DeserializeDriver, Frame, StreamDeserializer};
23//! # use deser::ser::{SerializeDriver, Serializer, StreamSerializer};
24//! # use deser::Error;
25//! # /// A format with a number per line.
26//! # #[derive(Default)]
27//! # struct Lines(Vec<u8>);
28//! # impl StreamDeserializer for Lines {
29//! #     fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
30//! #         Ok(match input.iter().position(|&b| b == b'\n') {
31//! #             Some(end) => Frame::Value { start: 0, end, consumed: end + 1 },
32//! #             None if eof && input.is_empty() => Frame::End,
33//! #             None if eof => Frame::Value { start: 0, end: input.len(), consumed: input.len() },
34//! #             None => Frame::Incomplete { consumed: 0 },
35//! #         })
36//! #     }
37//! #     fn drive_frame<'de>(&mut self, frame: &'de [u8], driver: &mut DeserializeDriver<'_, 'de>) -> Result<(), Error> {
38//! #         let value: u64 = std::str::from_utf8(frame).unwrap().parse().unwrap();
39//! #         driver.emit(value)
40//! #     }
41//! # }
42//! # impl Serializer for Lines {
43//! #     fn drive(&mut self, driver: &mut SerializeDriver<'_>) -> Result<(), Error> {
44//! #         driver.drive(|event, _| {
45//! #             if let deser::Event::Atom(deser::Atom::U64(v)) = event {
46//! #                 self.0.extend_from_slice(format!("{v}\n").as_bytes());
47//! #             }
48//! #             Ok(())
49//! #         })
50//! #     }
51//! # }
52//! # impl StreamSerializer for Lines {
53//! #     fn output(&self) -> &[u8] { &self.0 }
54//! #     fn clear_output(&mut self) { self.0.clear() }
55//! # }
56//! // `Lines` is the stream serializer and deserializer of a format with a
57//! // number per line
58//! let mut reader = Reader::new(&b"1\n2\n3\n"[..], Lines::default());
59//! let mut writer = Writer::new(Vec::new(), Lines::default());
60//! while let Some(value) = reader.read::<u64>()? {
61//!     writer.write(&(value * 2))?;
62//! }
63//! assert_eq!(writer.into_inner(), b"2\n4\n6\n");
64//! # Ok(()) } example().unwrap();
65//! ```
66//!
67//! Readers deserialize values while their input arrives if the format
68//! supports it (see [`StreamDeserializer::drive_partial`]),
69//! otherwise the complete value is buffered first.  Values which borrow
70//! from the reader's buffer are read with [`Reader::read_borrowed`].
71//!
72//! # Generic Code
73//!
74//! A [`Writer`] is a [`Serializer`] and a [`Reader`]
75//! is a [`Deserializer`], so code which is generic
76//! over serializers and deserializers (like `deser-transcode`) can write
77//! to and read from streams.
78//!
79//! # Large Sequences
80//!
81//! Values which contain a large (or unbounded) sequence can be processed
82//! while they are read: a [`Streamed`](crate::stream::Streamed) sequence hands out
83//! its elements as they are read with [`Reader::read_next`] (and behaves
84//! like a `Vec` otherwise).
85//!
86//! # Writing Large Values
87//!
88//! Formats whose output can be written before the value is complete (like
89//! JSON and CBOR) serialize values in parts (see
90//! [`StreamSerializer::drive_partial`]).
91//! [`Writer::write`] uses this if possible: once the output of a value
92//! exceeds the [buffer limit](Writer::set_buffer_limit), what was
93//! serialized so far is written and the serialization continues, which
94//! means that the memory used does not depend on the size of the values.
95//! Values below the limit are written at once.  Output that can still
96//! change (for instance the header of a container whose length is not
97//! known upfront) is held back until it's final.
98//!
99//! # Errors
100//!
101//! Errors refer to positions in the stream: the offsets, lines and columns
102//! of errors are relative to the start of the stream, not to the start of
103//! the frame.  Failed reads and writes are errors of the kind
104//! [`ErrorKind::Io`] with the IO error as source.
105use core::any::Any;
106use core::marker::PhantomData;
107use std::io::{Read, Write};
108
109use crate::Context;
110use crate::de::{
111    Deserialize, DeserializeDriver, DeserializeOwned, Deserializer, StreamDeserializer,
112};
113use crate::error::{Error, ErrorKind};
114use crate::ser::{Serialize, SerializeDriver, Serializer, StreamSerializer};
115use crate::stream::{
116    DEFAULT_BUFFER_LIMIT, ElementReader, ElementStatus, InputBuffer, Part, Status,
117};
118
119/// Reads values from a [`Read`].
120///
121/// The values are split and deserialized with a [`StreamDeserializer`]
122/// (for instance `deser_json::StreamDeserializer`, which the configuration
123/// of the format creates with `config.reader(input)`).  The reader buffers
124/// the input so it does not need to be buffered.
125///
126/// Readers implement [`Deserializer`]: every call to
127/// [`drive`](Deserializer::drive) reads the next value.  Values read this
128/// way cannot borrow from the reader, use
129/// [`read_borrowed`](Self::read_borrowed) for that.
130pub struct Reader<R, D: StreamDeserializer> {
131    reader: R,
132    buffer: InputBuffer<D>,
133    // the value that is read with `read_next` (an `ElementReader`)
134    pending: Option<Box<dyn Any + Send>>,
135}
136
137impl<R: Read, D: StreamDeserializer> Reader<R, D> {
138    /// Creates a reader.
139    ///
140    /// To continue a stream whose context is known (for instance the
141    /// names of the columns of a CSV file), create the stream deserializer
142    /// with that context.
143    pub fn new(reader: R, deserializer: D) -> Reader<R, D> {
144        Reader {
145            reader,
146            buffer: InputBuffer::new(deserializer),
147            pending: None,
148        }
149    }
150
151    /// Sets the context the values are deserialized in.
152    ///
153    /// This replaces the context of the stream deserializer (see
154    /// [`StreamDeserializer::context`]), which is the one of the
155    /// configuration it was created with.  The values of the context are
156    /// the defaults of the extension values of the state (see
157    /// [`Context`]).  A context set by the callback of
158    /// [`read_with`](Self::read_with) takes precedence.
159    pub fn set_context(&mut self, context: Context) {
160        self.buffer.set_context(context);
161    }
162
163    /// Returns the context the values are deserialized in.
164    pub fn context(&self) -> &Context {
165        self.buffer.context()
166    }
167
168    /// Fails if a value is being read with [`read_next`](Self::read_next).
169    fn ensure_idle(&self) -> Result<(), Error> {
170        match self.pending {
171            Some(_) => Err(Error::new(
172                ErrorKind::InvalidState,
173                "a value is being read with read_next",
174            )),
175            None => Ok(()),
176        }
177    }
178
179    /// Reads more input into the buffer.
180    fn read_more(&mut self) -> Result<(), Error> {
181        let buf = self.buffer.read_buf();
182        let read = loop {
183            match self.reader.read(buf) {
184                Ok(read) => break read,
185                Err(err) if err.kind() == std::io::ErrorKind::Interrupted => {}
186                Err(err) => return Err(err.into()),
187            }
188        };
189        if read == 0 {
190            self.buffer.set_eof();
191        } else {
192            self.buffer.filled(read);
193        }
194        Ok(())
195    }
196
197    /// Reads until the frame of the next value is complete.
198    ///
199    /// Returns `false` if there are no more values.
200    fn fill(&mut self) -> Result<bool, Error> {
201        loop {
202            match self.buffer.poll()? {
203                Status::Ready => return Ok(true),
204                Status::End => return Ok(false),
205                Status::NeedInput => self.read_more()?,
206            }
207        }
208    }
209
210    /// Feeds the next value into a driver.
211    ///
212    /// Returns `false` if there are no more values.
213    fn drive_partial(&mut self, driver: &mut DeserializeDriver<'_, '_>) -> Result<bool, Error> {
214        loop {
215            match self.buffer.drive_partial(driver)? {
216                Status::Ready => return Ok(true),
217                Status::End => return Ok(false),
218                Status::NeedInput => self.read_more()?,
219            }
220        }
221    }
222
223    /// Reads the next value.
224    ///
225    /// Returns `None` if there are no more values.  If the format supports
226    /// it (see [`StreamDeserializer::supports_partial`]), the value is
227    /// deserialized while the input is read which means that only
228    /// incomplete tokens are buffered.  Otherwise the complete value is
229    /// buffered first.  Whether reading can continue after an error depends
230    /// on the format (for instance with JSON Lines it continues with the
231    /// next line).
232    pub fn read<T: DeserializeOwned>(&mut self) -> Result<Option<T>, Error> {
233        self.read_with(|_| {})
234    }
235
236    /// Reads the next value with a configured driver.
237    ///
238    /// The callback is invoked with the driver before the value is
239    /// deserialized, for instance to add [`Layer`](crate::de::Layer)s.
240    pub fn read_with<T, F>(&mut self, setup: F) -> Result<Option<T>, Error>
241    where
242        T: DeserializeOwned,
243        F: FnOnce(&mut DeserializeDriver<'_, '_>),
244    {
245        self.ensure_idle()?;
246        if !self.buffer.supports_partial() {
247            if !self.fill()? {
248                return Ok(None);
249            }
250            return self
251                .buffer
252                .deserialize_with(|driver| setup(driver))
253                .map(Some);
254        }
255
256        let mut out = None::<T>;
257        {
258            let mut driver = DeserializeDriver::<'_, 'static>::new(&mut out);
259            setup(&mut driver);
260            if !self.drive_partial(&mut driver)? {
261                return Ok(None);
262            }
263        }
264        out.ok_or_else(|| Error::new(ErrorKind::EndOfFile, "empty input"))
265            .map(Some)
266    }
267
268    /// Reads the next value which can borrow from the reader's buffer.
269    ///
270    /// The complete value is buffered first.
271    ///
272    /// ```
273    /// # use deser::de::{DeserializeDriver, Frame, StreamDeserializer};
274    /// # use deser::Error;
275    /// # struct Lines;
276    /// # impl StreamDeserializer for Lines {
277    /// #     fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
278    /// #         Ok(match input.iter().position(|&b| b == b'\n') {
279    /// #             Some(end) => Frame::Value { start: 0, end, consumed: end + 1 },
280    /// #             None if eof && input.is_empty() => Frame::End,
281    /// #             None if eof => Frame::Value { start: 0, end: input.len(), consumed: input.len() },
282    /// #             None => Frame::Incomplete { consumed: 0 },
283    /// #         })
284    /// #     }
285    /// #     fn drive_frame<'de>(&mut self, frame: &'de [u8], driver: &mut DeserializeDriver<'_, 'de>) -> Result<(), Error> {
286    /// #         driver.emit_borrowed(std::str::from_utf8(frame).unwrap())
287    /// #     }
288    /// # }
289    /// use deser::io::Reader;
290    ///
291    /// // `Lines` is the stream deserializer of a format with a string per
292    /// // line
293    /// let mut reader = Reader::new(&b"hello\nworld\n"[..], Lines);
294    /// let value: &str = reader.read_borrowed().unwrap().unwrap();
295    /// assert_eq!(value, "hello");
296    /// ```
297    pub fn read_borrowed<'a, T: Deserialize<'a>>(&'a mut self) -> Result<Option<T>, Error> {
298        self.ensure_idle()?;
299        if !self.fill()? {
300            return Ok(None);
301        }
302        self.buffer.deserialize().map(Some)
303    }
304
305    /// Reads the next element of the [`Streamed`](crate::stream::Streamed) sequence of a value or
306    /// the value.
307    ///
308    /// `T` is the type of the value and `E` the type of the elements of a
309    /// [`Streamed<E>`](crate::stream::Streamed) sequence within it.  The elements are
310    /// handed out as they are read ([`Part::Element`]), the value once it's
311    /// complete ([`Part::Done`]).  The next call continues with the next
312    /// value.  Returns `None` if there are no more values.  See [`Streamed`](crate::stream::Streamed)
313    /// for an example.
314    ///
315    /// Until the value is complete, the reader can only be used to read the
316    /// value with the same types.
317    pub fn read_next<T, E>(&mut self) -> Result<Option<Part<E, T>>, Error>
318    where
319        T: DeserializeOwned + 'static,
320        E: Send + 'static,
321    {
322        let mut reader = match self.pending.take() {
323            Some(pending) => match pending.downcast::<ElementReader<T, E>>() {
324                Ok(reader) => reader,
325                Err(pending) => {
326                    self.pending = Some(pending);
327                    return Err(Error::new(
328                        ErrorKind::InvalidState,
329                        "a value of another type is being read",
330                    ));
331                }
332            },
333            None => Box::new(ElementReader::<T, E>::new()),
334        };
335        loop {
336            match reader.poll(&mut self.buffer)? {
337                ElementStatus::Ready(next) => {
338                    if reader.is_reading() {
339                        self.pending = Some(reader);
340                    }
341                    return Ok(Some(next));
342                }
343                ElementStatus::End => return Ok(None),
344                ElementStatus::NeedInput => {
345                    if let Err(err) = self.read_more() {
346                        self.pending = Some(reader);
347                        return Err(err);
348                    }
349                }
350            }
351        }
352    }
353
354    /// Returns `true` if there are no more values.
355    ///
356    /// This reads until the start of the next value or the end of the
357    /// stream.  If the format cannot find the start of a value on its own
358    /// (see [`StreamDeserializer::peek`]), the next value is buffered
359    /// completely.  This is useful to read values with the reader's
360    /// [`Deserializer`] implementation:
361    ///
362    /// ```
363    /// # use deser::de::{DeserializeDriver, Frame, StreamDeserializer};
364    /// # use deser::Error;
365    /// # struct Lines;
366    /// # impl StreamDeserializer for Lines {
367    /// #     fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
368    /// #         Ok(match input.iter().position(|&b| b == b'\n') {
369    /// #             Some(end) => Frame::Value { start: 0, end, consumed: end + 1 },
370    /// #             None if eof && input.is_empty() => Frame::End,
371    /// #             None if eof => Frame::Value { start: 0, end: input.len(), consumed: input.len() },
372    /// #             None => Frame::Incomplete { consumed: 0 },
373    /// #         })
374    /// #     }
375    /// #     fn drive_frame<'de>(&mut self, frame: &'de [u8], driver: &mut DeserializeDriver<'_, 'de>) -> Result<(), Error> {
376    /// #         driver.emit_borrowed(std::str::from_utf8(frame).unwrap())
377    /// #     }
378    /// # }
379    /// use deser::de::Deserializer;
380    /// use deser::io::Reader;
381    ///
382    /// // `Lines` is the stream deserializer of a format with a string per
383    /// // line
384    /// let mut reader = Reader::new(&b"hello\nworld\n"[..], Lines);
385    /// let mut values = Vec::new();
386    /// while !reader.is_end().unwrap() {
387    ///     values.push(reader.deserialize::<String>().unwrap());
388    /// }
389    /// assert_eq!(values, ["hello", "world"]);
390    /// ```
391    pub fn is_end(&mut self) -> Result<bool, Error> {
392        self.ensure_idle()?;
393        loop {
394            match self.buffer.peek()? {
395                Status::Ready => return Ok(false),
396                Status::End => return Ok(true),
397                Status::NeedInput => self.read_more()?,
398            }
399        }
400    }
401
402    /// Checks that there are no more values.
403    ///
404    /// Fails if another value follows (or if the data that follows is not
405    /// valid).
406    pub fn end(&mut self) -> Result<(), Error> {
407        self.ensure_idle()?;
408        if self.fill()? {
409            return Err(self.buffer.trailing_error());
410        }
411        Ok(())
412    }
413
414    /// Returns an iterator over the remaining values.
415    ///
416    /// The iterator stops after the first error.
417    pub fn iter<T: DeserializeOwned>(&mut self) -> Iter<'_, R, D, T> {
418        Iter {
419            reader: self,
420            failed: false,
421            _marker: PhantomData,
422        }
423    }
424
425    /// Returns the stream deserializer.
426    ///
427    /// This gives access to what the stream established so far, for
428    /// instance the names of the columns of a CSV file.
429    pub fn deserializer(&self) -> &D {
430        self.buffer.deserializer()
431    }
432
433    /// Returns a reference to the underlying reader.
434    pub fn get_ref(&self) -> &R {
435        &self.reader
436    }
437
438    /// Returns a mutable reference to the underlying reader.
439    ///
440    /// Reading from it directly is likely to corrupt the stream as the
441    /// reader buffers data.
442    pub fn get_mut(&mut self) -> &mut R {
443        &mut self.reader
444    }
445
446    /// Returns the underlying reader.
447    ///
448    /// Data that was read into the buffer but not deserialized yet is lost.
449    pub fn into_inner(self) -> R {
450        self.reader
451    }
452
453    /// Returns the underlying reader and the stream deserializer.
454    ///
455    /// Data that was read into the buffer but not deserialized yet is lost.
456    pub fn into_parts(self) -> (R, D) {
457        (self.reader, self.buffer.into_parts().0)
458    }
459}
460
461/// Reads the next value.
462///
463/// Values cannot borrow from the reader: borrowed data is passed on like
464/// data that is only valid for the call (see
465/// [`DeserializeDriver::transient`]).  If there are no more values this
466/// fails (see [`Reader::is_end`]).
467impl<'de, R: Read, D: StreamDeserializer> Deserializer<'de> for Reader<R, D> {
468    fn drive(&mut self, driver: &mut DeserializeDriver<'_, 'de>) -> Result<(), Error> {
469        self.ensure_idle()?;
470        let found = if self.buffer.supports_partial() {
471            self.drive_partial(driver)?
472        } else if self.fill()? {
473            self.buffer.drive_transient(driver)?;
474            true
475        } else {
476            false
477        };
478        match found {
479            true => Ok(()),
480            false => Err(Error::new(ErrorKind::EndOfFile, "empty input")),
481        }
482    }
483}
484
485/// An iterator over the values of a [`Reader`].
486pub struct Iter<'r, R, D: StreamDeserializer, T> {
487    reader: &'r mut Reader<R, D>,
488    failed: bool,
489    _marker: PhantomData<fn() -> T>,
490}
491
492impl<R: Read, D: StreamDeserializer, T: DeserializeOwned> Iterator for Iter<'_, R, D, T> {
493    type Item = Result<T, Error>;
494
495    fn next(&mut self) -> Option<Self::Item> {
496        if self.failed {
497            return None;
498        }
499        match self.reader.read() {
500            Ok(Some(value)) => Some(Ok(value)),
501            Ok(None) => None,
502            Err(err) => {
503                self.failed = true;
504                Some(Err(err))
505            }
506        }
507    }
508}
509
510/// Writes values to a [`Write`].
511///
512/// The values are serialized with a [`StreamSerializer`] (the serializer
513/// of a format, which the configuration of the format creates with
514/// `config.writer(output)`).  The output of a value is written with
515/// [`write_all`](Write::write_all), wrap the writer in a
516/// [`BufWriter`](std::io::BufWriter) when writing many small values.  If
517/// the format supports it, the output of large values is written in parts
518/// while they are serialized (see
519/// [`set_buffer_limit`](Self::set_buffer_limit)).
520///
521/// Writers implement [`Serializer`]: [`serialize`](Serializer::serialize)
522/// is the same as [`write`](Self::write).
523pub struct Writer<W, S: StreamSerializer> {
524    writer: W,
525    serializer: S,
526    limit: usize,
527    context: Context,
528}
529
530impl<W: Write, S: StreamSerializer> Writer<W, S> {
531    /// Creates a writer.
532    ///
533    /// To append to a stream that was written before, create the
534    /// serializer with the state of the stream (for instance the number of
535    /// values that were written).
536    pub fn new(writer: W, serializer: S) -> Writer<W, S> {
537        Writer {
538            writer,
539            serializer,
540            limit: DEFAULT_BUFFER_LIMIT,
541            context: Context::new(),
542        }
543    }
544
545    /// Sets the context the values are serialized in.
546    ///
547    /// The values of the context are the defaults of the extension values
548    /// of the state (see [`Context`]).  It takes precedence over the
549    /// context of the serializer (which is the one of the configuration it
550    /// was created with), a context set by the callback of
551    /// [`write_with`](Self::write_with) takes precedence over both.
552    pub fn set_context(&mut self, context: Context) {
553        self.context = context;
554    }
555
556    /// Returns the context the values are serialized in.
557    pub fn context(&self) -> &Context {
558        &self.context
559    }
560
561    /// Sets how much output of a value is buffered before it's written.
562    ///
563    /// If the format supports it (see
564    /// [`StreamSerializer::supports_partial`]), the output of a value is
565    /// written once it exceeds the limit, the serialization continues
566    /// after that.  This way the memory used does not depend on the size of
567    /// the values.  The default is [`DEFAULT_BUFFER_LIMIT`] (8 KiB).  With
568    /// `usize::MAX` every value is serialized completely before it's
569    /// written.
570    ///
571    /// ```
572    /// use deser::io::Writer;
573    /// # use deser::ser::{EventSink, SerializeDriver, SerializeRef, Serializer, StreamSerializer};
574    /// # use deser::{Error, Event, State};
575    /// # /// Writes `x` for every event, can stop between values.
576    /// # #[derive(Default)]
577    /// # struct Xs { out: Vec<u8>, partial: bool }
578    /// # struct Sink<'a>(&'a mut Vec<u8>, usize);
579    /// # impl EventSink for Sink<'_> {
580    /// #     fn event(&mut self, _: Event<'_>, _: SerializeRef<'_>, _: &mut State) -> Result<(), Error> {
581    /// #         self.0.push(b'x');
582    /// #         Ok(())
583    /// #     }
584    /// #     fn pause(&mut self) -> bool {
585    /// #         self.0.len() >= self.1
586    /// #     }
587    /// # }
588    /// # impl Serializer for Xs {
589    /// #     fn drive(&mut self, driver: &mut SerializeDriver<'_>) -> Result<(), Error> {
590    /// #         driver.drive(|_, _| Ok(self.out.push(b'x')))
591    /// #     }
592    /// # }
593    /// # impl StreamSerializer for Xs {
594    /// #     fn output(&self) -> &[u8] { &self.out }
595    /// #     fn clear_output(&mut self) { self.out.clear() }
596    /// #     fn supports_partial(&self) -> bool { true }
597    /// #     fn drive_partial(&mut self, driver: &mut SerializeDriver<'_>, limit: usize) -> Result<bool, Error> {
598    /// #         let done = driver.drive_until(&mut Sink(&mut self.out, limit))?;
599    /// #         self.partial = !done;
600    /// #         Ok(done)
601    /// #     }
602    /// #     fn in_progress(&self) -> bool { self.partial }
603    /// # }
604    /// /// Counts the writes.
605    /// struct Counter(usize);
606    ///
607    /// impl std::io::Write for Counter {
608    ///     fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
609    ///         self.0 += 1;
610    ///         Ok(buf.len())
611    ///     }
612    ///
613    ///     fn flush(&mut self) -> std::io::Result<()> {
614    ///         Ok(())
615    ///     }
616    /// }
617    ///
618    /// // `Xs` is a format which writes an `x` for every event
619    /// let mut writer = Writer::new(Counter(0), Xs::default());
620    /// writer.set_buffer_limit(100);
621    /// writer.write(&vec![0; 20_000]).unwrap();
622    /// assert!(writer.get_ref().0 > 50);
623    /// ```
624    pub fn set_buffer_limit(&mut self, limit: usize) {
625        self.limit = limit;
626    }
627
628    /// Returns how much output of a value is buffered before it's written.
629    pub fn buffer_limit(&self) -> usize {
630        self.limit
631    }
632
633    /// Serializes a value and writes it.
634    ///
635    /// If the value fails to serialize before any of it was written (which
636    /// is always the case for values whose output is below the
637    /// [buffer limit](Self::set_buffer_limit)), nothing is written and the
638    /// next value can be written.  If a value is abandoned after a part of
639    /// it was written, because it fails to serialize or a write fails, the
640    /// stream holds an incomplete value and the writer refuses to write
641    /// more values (see [`StreamSerializer::in_progress`]).
642    pub fn write<T: Serialize + ?Sized>(&mut self, value: &T) -> Result<(), Error> {
643        self.write_driver(&mut SerializeDriver::new(&value))
644    }
645
646    /// Serializes a value with a configured driver and writes it.
647    ///
648    /// The callback is invoked with the driver before the value is
649    /// serialized, for instance to add [`Layer`](crate::ser::Layer)s.
650    pub fn write_with<T, F>(&mut self, value: &T, setup: F) -> Result<(), Error>
651    where
652        T: Serialize + ?Sized,
653        F: FnOnce(&mut SerializeDriver<'_>),
654    {
655        let mut driver = SerializeDriver::new(&value);
656        setup(&mut driver);
657        self.write_driver(&mut driver)
658    }
659
660    /// Serializes the value of a driver and writes it.
661    fn write_driver(&mut self, driver: &mut SerializeDriver<'_>) -> Result<(), Error> {
662        if !self.context.is_empty() {
663            driver.set_default_context(self.context.clone());
664        }
665        if self.serializer.in_progress() {
666            return Err(Error::in_progress());
667        }
668        // output that was not written (for instance of values serialized
669        // before the serializer was given to the writer) comes first
670        write_output(&mut self.writer, &mut self.serializer)?;
671        let limit = match self.serializer.supports_partial() {
672            true => self.limit.max(1),
673            false => usize::MAX,
674        };
675        loop {
676            let done = self.serializer.drive_partial(driver, limit)?;
677            write_output(&mut self.writer, &mut self.serializer)?;
678            if done {
679                return Ok(());
680            }
681        }
682    }
683
684    /// Flushes the underlying writer.
685    pub fn flush(&mut self) -> Result<(), Error> {
686        self.writer.flush()?;
687        Ok(())
688    }
689
690    /// Returns the stream serializer.
691    ///
692    /// This gives access to the state of the stream, for instance the
693    /// number of values that were written.
694    pub fn serializer(&self) -> &S {
695        &self.serializer
696    }
697
698    /// Returns a reference to the underlying writer.
699    pub fn get_ref(&self) -> &W {
700        &self.writer
701    }
702
703    /// Returns a mutable reference to the underlying writer.
704    pub fn get_mut(&mut self) -> &mut W {
705        &mut self.writer
706    }
707
708    /// Returns the underlying writer.
709    pub fn into_inner(self) -> W {
710        self.writer
711    }
712
713    /// Returns the underlying writer and the stream serializer.
714    pub fn into_parts(self) -> (W, S) {
715        (self.writer, self.serializer)
716    }
717}
718
719/// Writes the output of a serializer and clears it.
720fn write_output<W: Write, S: StreamSerializer>(
721    writer: &mut W,
722    serializer: &mut S,
723) -> Result<(), Error> {
724    let output = serializer.output();
725    if output.is_empty() {
726        return Ok(());
727    }
728    // the output is gone even if the write fails, a part of it might have
729    // been written
730    let rv = writer.write_all(output);
731    serializer.clear_output();
732    rv.map_err(Error::from)
733}
734
735/// Serializes values and writes them, like [`Writer::write`].
736impl<W: Write, S: StreamSerializer> Serializer for Writer<W, S> {
737    fn drive(&mut self, driver: &mut SerializeDriver<'_>) -> Result<(), Error> {
738        self.write_driver(driver)
739    }
740}
741
742/// Reads a single value from a [`Read`].
743///
744/// Fails if there is no value or if another value follows it.
745pub fn from_reader<T, R, D>(reader: R, deserializer: D) -> Result<T, Error>
746where
747    T: DeserializeOwned,
748    R: Read,
749    D: StreamDeserializer,
750{
751    let mut reader = Reader::new(reader, deserializer);
752    let value = reader
753        .read()?
754        .ok_or_else(|| Error::new(ErrorKind::EndOfFile, "empty input"))?;
755    reader.end()?;
756    Ok(value)
757}
758
759/// Writes a single value to a [`Write`].
760pub fn to_writer<W, S, T>(writer: W, serializer: S, value: &T) -> Result<(), Error>
761where
762    W: Write,
763    S: StreamSerializer,
764    T: Serialize + ?Sized,
765{
766    Writer::new(writer, serializer).write(value)
767}