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}