Skip to main content

deser_tokio/
lib.rs

1//! Read and write [deser](https://docs.rs/deser) values with
2//! [tokio](https://tokio.rs).
3//!
4//! This crate connects the stream serializers and stream deserializers of
5//! the data formats (which implement [`StreamSerializer`] and
6//! [`StreamDeserializer`], see [`deser::stream`](deser_core::stream)) to
7//! tokio's [`AsyncRead`] and [`AsyncWrite`].  It works with every format.  Values
8//! of formats which support it (like JSON and CBOR) are deserialized while
9//! their input arrives, other values are buffered until they are complete,
10//! so streams of values (like JSON Lines, CBOR sequences or YAML documents)
11//! can be read from sockets with bounded memory:
12//!
13//! ```
14//! # #[tokio::main(flavor = "current_thread")]
15//! # async fn main() -> Result<(), deser::Error> {
16//! use deser::{Deserialize, Serialize};
17//! use deser_json::{DeserializerConfig, Serializer, SerializerConfig, StreamDeserializer, Trailing};
18//! use deser_tokio::{Reader, Writer};
19//!
20//! const READ_LINES: DeserializerConfig =
21//!     DeserializerConfig::builder().trailing(Trailing::Newline).build();
22//! const WRITE_LINES: SerializerConfig =
23//!     SerializerConfig::builder().trailing(Trailing::Newline).build();
24//!
25//! #[derive(Debug, Serialize, Deserialize)]
26//! struct Request {
27//!     id: u64,
28//!     method: String,
29//! }
30//!
31//! # let (client, server) = tokio::io::duplex(1024);
32//! # let client = tokio::spawn(async move {
33//! #     let (input, output) = tokio::io::split(client);
34//! #     let mut requests = Writer::new(output, Serializer::with_config(WRITE_LINES));
35//! #     requests.write(&Request { id: 1, method: "ping".into() }).await.unwrap();
36//! #     requests.shutdown().await.unwrap();
37//! #     let mut responses = Reader::new(input, StreamDeserializer::with_config(READ_LINES));
38//! #     assert_eq!(responses.read::<u64>().await.unwrap(), Some(1));
39//! # });
40//! let (input, output) = tokio::io::split(server);
41//!
42//! // JSON Lines in, JSON Lines out
43//! let mut requests = Reader::new(input, StreamDeserializer::with_config(READ_LINES));
44//! let mut responses = Writer::new(output, Serializer::with_config(WRITE_LINES));
45//! while let Some(request) = requests.read::<Request>().await? {
46//!     responses.write(&request.id).await?;
47//! }
48//! # client.await.unwrap();
49//! # Ok(()) }
50//! ```
51//!
52//! Single values are read with [`from_reader`] and written with
53//! [`to_writer`].  With the `codec` feature, [`Codec`] implements the codec
54//! traits of [`tokio-util`](https://docs.rs/tokio-util) for use with
55//! `FramedRead`, `FramedWrite` and `Framed`.
56//!
57//! # Multi-Threaded Runtimes
58//!
59//! The futures are `Send` if the reader or writer and the stream
60//! deserializer or serializer are, so they can be spawned on multi-threaded
61//! runtimes.
62//!
63//! # Large Values
64//!
65//! Formats which support it (like JSON and CBOR) serialize values in
66//! parts: once the output of a value exceeds the
67//! [buffer limit](Writer::set_buffer_limit), what was serialized so far is
68//! written and the serialization continues after that.  The memory used for
69//! writing does not depend on the size of the values either.
70//!
71//! # Cancellation
72//!
73//! Reading is cancellation safe: if a future that reads a value is
74//! dropped, the data read so far stays in the buffer of the [`Reader`] and
75//! the next read continues with it.  This allows reading in
76//! `tokio::select!`.  Writing is not cancellation safe, a value might have
77//! been written partially.  A [`Writer`] refuses to write more values after
78//! a value was abandoned after a part of it was written.
79#![doc(html_logo_url = "https://raw.githubusercontent.com/mitsuhiko/deser/main/artwork/logo.svg")]
80#![cfg_attr(docsrs, feature(doc_cfg))]
81
82use std::any::Any;
83use std::future::poll_fn;
84use std::marker::PhantomData;
85use std::pin::Pin;
86use std::task::{Context, Poll, ready};
87
88use deser_core::de::{
89    Deserialize, DeserializeDriver, DeserializeOwned, OwnedDriver, StreamDeserializer,
90};
91use deser_core::ser::{Serialize, SerializeDriver, StreamSerializer};
92use deser_core::stream::{
93    DEFAULT_BUFFER_LIMIT, ElementReader, ElementStatus, InputBuffer, Part, Status,
94};
95use deser_core::{Error, ErrorKind};
96use futures_core::Stream;
97use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt, ReadBuf};
98
99#[cfg(feature = "codec")]
100mod codec;
101
102#[cfg(feature = "codec")]
103pub use self::codec::Codec;
104
105/// Reads values from an [`AsyncRead`].
106///
107/// The values are split and deserialized with a [`StreamDeserializer`] (for
108/// instance `deser_json::StreamDeserializer`).  The reader buffers the
109/// input so it does not need to be buffered.  If the format supports it
110/// (see [`StreamDeserializer::supports_partial`]), values are deserialized
111/// while their input arrives which means that only incomplete tokens are
112/// buffered.
113pub struct Reader<R, D: StreamDeserializer> {
114    reader: R,
115    buffer: InputBuffer<D>,
116    // a value that is being deserialized while its input arrives (an
117    // `OwnedDriver<'static, T>`), kept when a read is cancelled.
118    pending: Option<Box<dyn Any + Send>>,
119}
120
121// the reader is never pinned structurally
122impl<R, D: StreamDeserializer> Unpin for Reader<R, D> {}
123
124impl<R: AsyncRead + Unpin, D: StreamDeserializer> Reader<R, D> {
125    /// Creates a reader.
126    ///
127    /// To continue a stream whose context is known (for instance the
128    /// names of the columns of a CSV file), create the stream deserializer
129    /// with that context.
130    pub fn new(reader: R, deserializer: D) -> Reader<R, D> {
131        Reader {
132            reader,
133            buffer: InputBuffer::new(deserializer),
134            pending: None,
135        }
136    }
137
138    /// Sets the context the values are deserialized in.
139    ///
140    /// This replaces the context of the stream deserializer (see
141    /// [`StreamDeserializer::context`](deser_core::de::StreamDeserializer::context)),
142    /// which is the one of the configuration it was created with.  The
143    /// values of the context are the defaults of the extension values of
144    /// the state (see [`Context`](deser_core::Context)).  A context set by
145    /// the callback of [`read_with`](Self::read_with) takes precedence.
146    pub fn set_context(&mut self, context: deser_core::Context) {
147        self.buffer.set_context(context);
148    }
149
150    /// Returns the context the values are deserialized in.
151    pub fn context(&self) -> &deser_core::Context {
152        self.buffer.context()
153    }
154
155    /// Reads more input into the buffer.
156    fn poll_read_more(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Error>> {
157        let mut buf = ReadBuf::new(self.buffer.read_buf());
158        ready!(Pin::new(&mut self.reader).poll_read(cx, &mut buf))?;
159        let read = buf.filled().len();
160        if read == 0 {
161            self.buffer.set_eof();
162        } else {
163            self.buffer.filled(read);
164        }
165        Poll::Ready(Ok(()))
166    }
167
168    /// Reads until the frame of the next value is complete.
169    ///
170    /// Resolves to `false` if there are no more values.
171    fn poll_fill(&mut self, cx: &mut Context<'_>) -> Poll<Result<bool, Error>> {
172        loop {
173            match self.buffer.poll()? {
174                Status::Ready => return Poll::Ready(Ok(true)),
175                Status::End => return Poll::Ready(Ok(false)),
176                Status::NeedInput => ready!(self.poll_read_more(cx))?,
177            }
178        }
179    }
180
181    /// Reads the next value, the setup is invoked with its driver.
182    fn poll_read_setup<T, F>(
183        &mut self,
184        cx: &mut Context<'_>,
185        setup: &mut Option<F>,
186    ) -> Poll<Result<Option<T>, Error>>
187    where
188        T: DeserializeOwned + 'static,
189        F: FnOnce(&mut DeserializeDriver<'_, '_>),
190    {
191        if !self.buffer.supports_partial() {
192            if !ready!(self.poll_fill(cx))? {
193                return Poll::Ready(Ok(None));
194            }
195            let setup = setup.take();
196            return Poll::Ready(
197                self.buffer
198                    .deserialize_with(|driver| {
199                        if let Some(setup) = setup {
200                            setup(driver);
201                        }
202                    })
203                    .map(Some),
204            );
205        }
206
207        // continue a value that is being read or start a new one
208        let mut driver = match self.pending.take() {
209            Some(pending) => match pending.downcast::<OwnedDriver<'static, T>>() {
210                Ok(driver) => *driver,
211                Err(pending) => {
212                    self.pending = Some(pending);
213                    return Poll::Ready(Err(Error::new(
214                        ErrorKind::InvalidState,
215                        "a value of another type is being read",
216                    )));
217                }
218            },
219            None => {
220                let mut driver = OwnedDriver::<'static, T>::new();
221                if let Some(setup) = setup.take() {
222                    driver.with(|driver| setup(driver));
223                }
224                driver
225            }
226        };
227        loop {
228            match driver.with(|driver| self.buffer.drive_partial(driver))? {
229                Status::Ready => return Poll::Ready(driver.finish().map(Some)),
230                Status::End => return Poll::Ready(Ok(None)),
231                Status::NeedInput => match self.poll_read_more(cx) {
232                    Poll::Ready(Ok(())) => {}
233                    rv => {
234                        // the value continues with the next read
235                        self.pending = Some(Box::new(driver));
236                        return rv.map(|rv| rv.map(|_| None));
237                    }
238                },
239            }
240        }
241    }
242
243    /// Polls for the next value.
244    ///
245    /// This is the poll based version of [`read`](Self::read).
246    pub fn poll_read<T: DeserializeOwned + 'static>(
247        &mut self,
248        cx: &mut Context<'_>,
249    ) -> Poll<Result<Option<T>, Error>> {
250        self.poll_read_setup(cx, &mut None::<fn(&mut DeserializeDriver<'_, '_>)>)
251    }
252
253    /// Reads the next value.
254    ///
255    /// Resolves to `None` if there are no more values.  If the future is
256    /// dropped before it resolves, the next read continues where it
257    /// stopped (a value of another type cannot be read then).  Whether
258    /// reading can continue after an error depends on the format (for
259    /// instance with JSON Lines it continues with the next line).
260    pub async fn read<T: DeserializeOwned + 'static>(&mut self) -> Result<Option<T>, Error> {
261        poll_fn(|cx| self.poll_read(cx)).await
262    }
263
264    /// Reads the next value with a configured driver.
265    ///
266    /// The callback is invoked with the driver before the value is
267    /// deserialized, for instance to add [`Layer`](deser_core::de::Layer)s.
268    pub async fn read_with<T, F>(&mut self, setup: F) -> Result<Option<T>, Error>
269    where
270        T: DeserializeOwned + 'static,
271        F: FnOnce(&mut DeserializeDriver<'_, '_>),
272    {
273        let mut setup = Some(setup);
274        poll_fn(|cx| self.poll_read_setup(cx, &mut setup)).await
275    }
276
277    /// Polls for the next element of the [`Streamed`](deser_core::stream::Streamed) sequence of a value or
278    /// the value.
279    ///
280    /// This is the poll based version of [`read_next`](Self::read_next).
281    pub fn poll_read_next<T, E>(
282        &mut self,
283        cx: &mut Context<'_>,
284    ) -> Poll<Result<Option<Part<E, T>>, Error>>
285    where
286        T: DeserializeOwned + 'static,
287        E: Send + 'static,
288    {
289        let mut reader = match self.pending.take() {
290            Some(pending) => match pending.downcast::<ElementReader<T, E>>() {
291                Ok(reader) => reader,
292                Err(pending) => {
293                    self.pending = Some(pending);
294                    return Poll::Ready(Err(Error::new(
295                        ErrorKind::InvalidState,
296                        "a value of another type is being read",
297                    )));
298                }
299            },
300            None => Box::new(ElementReader::<T, E>::new()),
301        };
302        loop {
303            match reader.poll(&mut self.buffer)? {
304                ElementStatus::Ready(next) => {
305                    if reader.is_reading() {
306                        self.pending = Some(reader);
307                    }
308                    return Poll::Ready(Ok(Some(next)));
309                }
310                ElementStatus::End => return Poll::Ready(Ok(None)),
311                ElementStatus::NeedInput => match self.poll_read_more(cx) {
312                    Poll::Ready(Ok(())) => {}
313                    rv => {
314                        // the value continues with the next read
315                        self.pending = Some(reader);
316                        return rv.map(|rv| rv.map(|_| None));
317                    }
318                },
319            }
320        }
321    }
322
323    /// Reads the next element of the [`Streamed`](deser_core::stream::Streamed) sequence of a value or
324    /// the value.
325    ///
326    /// `T` is the type of the value and `E` the type of the elements of a
327    /// [`Streamed<E>`](deser_core::stream::Streamed) sequence within it.  The elements are
328    /// handed out as they are read ([`Part::Element`]), the value once it's
329    /// complete ([`Part::Done`]).  The next call continues with the next
330    /// value.  Resolves to `None` if there are no more values.  This is
331    /// cancellation safe, until the value is complete the reader can only
332    /// be used to read the value with the same types.
333    pub async fn read_next<T, E>(&mut self) -> Result<Option<Part<E, T>>, Error>
334    where
335        T: DeserializeOwned + 'static,
336        E: Send + 'static,
337    {
338        poll_fn(|cx| self.poll_read_next(cx)).await
339    }
340
341    /// Converts the reader into a [`Stream`] of the elements of the
342    /// [`Streamed`](deser_core::stream::Streamed) sequence of values and the values.
343    ///
344    /// See [`read_next`](Self::read_next).  The stream ends after the first
345    /// error.
346    pub fn into_element_stream<T, E>(self) -> ElementStream<R, D, T, E>
347    where
348        T: DeserializeOwned + 'static,
349        E: Send + 'static,
350    {
351        ElementStream {
352            reader: self,
353            failed: false,
354            _marker: PhantomData,
355        }
356    }
357
358    /// Reads the next value which can borrow from the reader's buffer.
359    ///
360    /// The complete value is buffered first.
361    pub async fn read_borrowed<'a, T: Deserialize<'a>>(&'a mut self) -> Result<Option<T>, Error> {
362        if !poll_fn(|cx| self.poll_fill(cx)).await? {
363            return Ok(None);
364        }
365        self.buffer.deserialize().map(Some)
366    }
367
368    /// Returns `true` if there are no more values.
369    ///
370    /// This reads until the start of the next value or the end of the
371    /// stream (see
372    /// [`deser::io::Reader::is_end`](https://docs.rs/deser/latest/deser/io/struct.Reader.html#method.is_end)).
373    pub async fn is_end(&mut self) -> Result<bool, Error> {
374        poll_fn(|cx| {
375            loop {
376                match self.buffer.peek()? {
377                    Status::Ready => return Poll::Ready(Ok(false)),
378                    Status::End => return Poll::Ready(Ok(true)),
379                    Status::NeedInput => ready!(self.poll_read_more(cx))?,
380                }
381            }
382        })
383        .await
384    }
385
386    /// Checks that there are no more values.
387    ///
388    /// Fails if another value follows (or if the data that follows is not
389    /// valid).
390    pub async fn end(&mut self) -> Result<(), Error> {
391        if poll_fn(|cx| self.poll_fill(cx)).await? {
392            return Err(self.buffer.trailing_error());
393        }
394        Ok(())
395    }
396
397    /// Converts the reader into a [`Stream`] of values.
398    ///
399    /// The stream ends after the first error.
400    pub fn into_stream<T: DeserializeOwned + 'static>(self) -> ReaderStream<R, D, T> {
401        ReaderStream {
402            reader: self,
403            failed: false,
404            _marker: PhantomData,
405        }
406    }
407
408    /// Returns the stream deserializer.
409    ///
410    /// This gives access to what the stream established so far, for
411    /// instance the names of the columns of a CSV file.
412    pub fn deserializer(&self) -> &D {
413        self.buffer.deserializer()
414    }
415
416    /// Returns a reference to the underlying reader.
417    pub fn get_ref(&self) -> &R {
418        &self.reader
419    }
420
421    /// Returns a mutable reference to the underlying reader.
422    ///
423    /// Reading from it directly is likely to corrupt the stream as the
424    /// reader buffers data.
425    pub fn get_mut(&mut self) -> &mut R {
426        &mut self.reader
427    }
428
429    /// Returns the underlying reader.
430    ///
431    /// Data that was read into the buffer but not deserialized yet is lost.
432    pub fn into_inner(self) -> R {
433        self.reader
434    }
435}
436
437/// A [`Stream`] of the values of a [`Reader`].
438///
439/// Created with [`Reader::into_stream`].
440pub struct ReaderStream<R, D: StreamDeserializer, T> {
441    reader: Reader<R, D>,
442    failed: bool,
443    _marker: PhantomData<fn() -> T>,
444}
445
446impl<R, D: StreamDeserializer, T> Unpin for ReaderStream<R, D, T> {}
447
448impl<R, D: StreamDeserializer, T> ReaderStream<R, D, T> {
449    /// Returns the reader.
450    pub fn into_inner(self) -> Reader<R, D> {
451        self.reader
452    }
453}
454
455impl<R: AsyncRead + Unpin, D: StreamDeserializer, T: DeserializeOwned + 'static> Stream
456    for ReaderStream<R, D, T>
457{
458    type Item = Result<T, Error>;
459
460    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
461        if self.failed {
462            return Poll::Ready(None);
463        }
464        match ready!(self.reader.poll_read(cx)) {
465            Ok(value) => Poll::Ready(value.map(Ok)),
466            Err(err) => {
467                self.failed = true;
468                Poll::Ready(Some(Err(err)))
469            }
470        }
471    }
472}
473
474/// A [`Stream`] of the elements of the [`Streamed`](deser_core::stream::Streamed) sequence of values and
475/// the values.
476///
477/// Created with [`Reader::into_element_stream`].
478pub struct ElementStream<R, D: StreamDeserializer, T, E> {
479    reader: Reader<R, D>,
480    failed: bool,
481    _marker: PhantomData<fn() -> (T, E)>,
482}
483
484impl<R, D: StreamDeserializer, T, E> Unpin for ElementStream<R, D, T, E> {}
485
486impl<R, D: StreamDeserializer, T, E> ElementStream<R, D, T, E> {
487    /// Returns the reader.
488    pub fn into_inner(self) -> Reader<R, D> {
489        self.reader
490    }
491}
492
493impl<R, D, T, E> Stream for ElementStream<R, D, T, E>
494where
495    R: AsyncRead + Unpin,
496    D: StreamDeserializer,
497    T: DeserializeOwned + 'static,
498    E: Send + 'static,
499{
500    type Item = Result<Part<E, T>, Error>;
501
502    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
503        if self.failed {
504            return Poll::Ready(None);
505        }
506        match ready!(self.reader.poll_read_next(cx)) {
507            Ok(next) => Poll::Ready(next.map(Ok)),
508            Err(err) => {
509                self.failed = true;
510                Poll::Ready(Some(Err(err)))
511            }
512        }
513    }
514}
515
516/// Writes values to an [`AsyncWrite`].
517///
518/// The values are serialized with a [`StreamSerializer`] (the serializer
519/// of a data format, for instance `deser_json::Serializer`) and its output
520/// is written with [`write_all`](tokio::io::AsyncWriteExt::write_all), wrap
521/// the writer in a [`BufWriter`](tokio::io::BufWriter) when writing many
522/// small values (and [`flush`](Self::flush) it).  If the format supports it
523/// (see [`StreamSerializer::supports_partial`]), the output of large values
524/// is written in parts while they are serialized, so the memory used does
525/// not depend on the size of the values (see
526/// [`set_buffer_limit`](Self::set_buffer_limit)).
527pub struct Writer<W, S: StreamSerializer> {
528    writer: W,
529    serializer: S,
530    limit: usize,
531    // output is being written, if the future is dropped meanwhile it's
532    // unknown what was written
533    writing: bool,
534    // the context the values are serialized in
535    context: deser_core::Context,
536}
537
538impl<W: AsyncWrite + Unpin, S: StreamSerializer> Writer<W, S> {
539    /// Creates a writer.
540    ///
541    /// To append to a stream that was written before, create the
542    /// serializer with the state of the stream (for instance the number of
543    /// values that were written).
544    pub fn new(writer: W, serializer: S) -> Writer<W, S> {
545        Writer {
546            writer,
547            serializer,
548            limit: DEFAULT_BUFFER_LIMIT,
549            writing: false,
550            context: deser_core::Context::new(),
551        }
552    }
553
554    /// Sets the context the values are serialized in.
555    ///
556    /// The values of the context are the defaults of the extension values
557    /// of the state (see [`Context`](deser_core::Context)).  It takes
558    /// precedence over the context of the serializer (which is the one of
559    /// the configuration it was created with), a context set by the
560    /// callback of [`write_with`](Self::write_with) takes precedence over
561    /// both.
562    pub fn set_context(&mut self, context: deser_core::Context) {
563        self.context = context;
564    }
565
566    /// Returns the context the values are serialized in.
567    pub fn context(&self) -> &deser_core::Context {
568        &self.context
569    }
570
571    /// Sets how much output of a value is buffered before it's written.
572    ///
573    /// If the format supports it, the output of a value is written once it
574    /// exceeds the limit, the serialization continues after that.  The
575    /// default is [`DEFAULT_BUFFER_LIMIT`] (8 KiB).  With `usize::MAX`
576    /// every value is serialized completely before it's written.  See
577    /// [`deser::io::Writer::set_buffer_limit`](https://docs.rs/deser/latest/deser/io/struct.Writer.html#method.set_buffer_limit).
578    pub fn set_buffer_limit(&mut self, limit: usize) {
579        self.limit = limit;
580    }
581
582    /// Returns how much output of a value is buffered before it's written.
583    pub fn buffer_limit(&self) -> usize {
584        self.limit
585    }
586
587    /// Serializes a value and writes it.
588    ///
589    /// If the value fails to serialize before any of it was written (which
590    /// is always the case for values whose output is below the
591    /// [buffer limit](Self::set_buffer_limit)), nothing is written and the
592    /// next value can be written.  If a value is abandoned after a part of
593    /// it was written, because it fails to serialize, a write fails or the
594    /// future is dropped, the stream holds an incomplete value and the
595    /// writer refuses to write more values.
596    pub async fn write<T: Serialize + ?Sized>(&mut self, value: &T) -> Result<(), Error> {
597        self.write_with(value, |_| {}).await
598    }
599
600    /// Serializes a value with a configured driver and writes it.
601    ///
602    /// The callback is invoked with the driver before the value is
603    /// serialized, for instance to add [`Layer`](deser_core::ser::Layer)s.
604    pub async fn write_with<T, F>(&mut self, value: &T, setup: F) -> Result<(), Error>
605    where
606        T: Serialize + ?Sized,
607        F: FnOnce(&mut SerializeDriver<'_>),
608    {
609        if self.writing || self.serializer.in_progress() {
610            return Err(Error::new(
611                ErrorKind::InvalidState,
612                "a value was only partially written, the stream cannot continue",
613            ));
614        }
615        let mut driver = SerializeDriver::new(&value);
616        setup(&mut driver);
617        if !self.context.is_empty() {
618            driver.set_default_context(self.context.clone());
619        }
620        // output that was not written (for instance of values serialized
621        // before the serializer was given to the writer) comes first
622        self.write_output().await?;
623        let limit = match self.serializer.supports_partial() {
624            true => self.limit.max(1),
625            false => usize::MAX,
626        };
627        loop {
628            let done = self.serializer.drive_partial(&mut driver, limit)?;
629            self.write_output().await?;
630            if done {
631                return Ok(());
632            }
633        }
634    }
635
636    /// Writes the output of the serializer and clears it.
637    async fn write_output(&mut self) -> Result<(), Error> {
638        if self.serializer.output().is_empty() {
639            return Ok(());
640        }
641        // stays set if the future is dropped or the write fails
642        self.writing = true;
643        self.writer.write_all(self.serializer.output()).await?;
644        self.serializer.clear_output();
645        self.writing = false;
646        Ok(())
647    }
648
649    /// Flushes the underlying writer.
650    pub async fn flush(&mut self) -> Result<(), Error> {
651        self.writer.flush().await?;
652        Ok(())
653    }
654
655    /// Shuts down the underlying writer.
656    pub async fn shutdown(&mut self) -> Result<(), Error> {
657        self.writer.shutdown().await?;
658        Ok(())
659    }
660
661    /// Returns the stream serializer.
662    ///
663    /// This gives access to the state of the stream, for instance the
664    /// number of values that were written.
665    pub fn serializer(&self) -> &S {
666        &self.serializer
667    }
668
669    /// Returns a reference to the underlying writer.
670    pub fn get_ref(&self) -> &W {
671        &self.writer
672    }
673
674    /// Returns a mutable reference to the underlying writer.
675    pub fn get_mut(&mut self) -> &mut W {
676        &mut self.writer
677    }
678
679    /// Returns the underlying writer.
680    pub fn into_inner(self) -> W {
681        self.writer
682    }
683
684    /// Returns the underlying writer and the stream serializer.
685    pub fn into_parts(self) -> (W, S) {
686        (self.writer, self.serializer)
687    }
688}
689
690/// Reads a single value from an [`AsyncRead`].
691///
692/// Fails if there is no value or if another value follows it.  How the
693/// value is read depends on the format, for instance with JSON the reader
694/// is read to the end.
695///
696/// ```
697/// # #[tokio::main(flavor = "current_thread")]
698/// # async fn main() {
699/// let input = &b"[1, 2, 3]"[..];
700/// let de = deser_json::StreamDeserializer::new();
701/// let value: Vec<u32> = deser_tokio::from_reader(input, de).await.unwrap();
702/// assert_eq!(value, [1, 2, 3]);
703/// # }
704/// ```
705pub async fn from_reader<T, R, D>(reader: R, deserializer: D) -> Result<T, Error>
706where
707    T: DeserializeOwned + 'static,
708    R: AsyncRead + Unpin,
709    D: StreamDeserializer,
710{
711    let mut reader = Reader::new(reader, deserializer);
712    let value = reader
713        .read()
714        .await?
715        .ok_or_else(|| Error::new(ErrorKind::EndOfFile, "empty input"))?;
716    reader.end().await?;
717    Ok(value)
718}
719
720/// Writes a single value to an [`AsyncWrite`] and flushes it.
721///
722/// ```
723/// # #[tokio::main(flavor = "current_thread")]
724/// # async fn main() {
725/// let mut out = Vec::new();
726/// let ser = deser_json::Serializer::new();
727/// deser_tokio::to_writer(&mut out, ser, &vec![1, 2]).await.unwrap();
728/// assert_eq!(out, b"[1,2]");
729/// # }
730/// ```
731pub async fn to_writer<W, S, T>(writer: W, serializer: S, value: &T) -> Result<(), Error>
732where
733    W: AsyncWrite + Unpin,
734    S: StreamSerializer,
735    T: Serialize + ?Sized,
736{
737    let mut writer = Writer::new(writer, serializer);
738    writer.write(value).await?;
739    writer.flush().await
740}