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}