io: add lines example for StreamReader (#5145)

This commit is contained in:
Alice Ryhl
2022-10-31 20:40:52 +01:00
committed by GitHub
parent c2210dfe37
commit a9d5eb2fc7
+177 -103
View File
@@ -1,113 +1,162 @@
use bytes::Buf; use bytes::Buf;
use futures_core::stream::Stream; use futures_core::stream::Stream;
use pin_project_lite::pin_project;
use std::io; use std::io;
use std::pin::Pin; use std::pin::Pin;
use std::task::{Context, Poll}; use std::task::{Context, Poll};
use tokio::io::{AsyncBufRead, AsyncRead, ReadBuf}; use tokio::io::{AsyncBufRead, AsyncRead, ReadBuf};
pin_project! { /// Convert a [`Stream`] of byte chunks into an [`AsyncRead`].
/// Convert a [`Stream`] of byte chunks into an [`AsyncRead`]. ///
/// /// This type performs the inverse operation of [`ReaderStream`].
/// This type performs the inverse operation of [`ReaderStream`]. ///
/// /// This type also implements the [`AsyncBufRead`] trait, so you can use it
/// # Example /// to read a `Stream` of byte chunks line-by-line. See the examples below.
/// ///
/// ``` /// # Example
/// use bytes::Bytes; ///
/// use tokio::io::{AsyncReadExt, Result}; /// ```
/// use tokio_util::io::StreamReader; /// use bytes::Bytes;
/// # #[tokio::main] /// use tokio::io::{AsyncReadExt, Result};
/// # async fn main() -> std::io::Result<()> { /// use tokio_util::io::StreamReader;
/// /// # #[tokio::main(flavor = "current_thread")]
/// // Create a stream from an iterator. /// # async fn main() -> std::io::Result<()> {
/// let stream = tokio_stream::iter(vec![ ///
/// Result::Ok(Bytes::from_static(&[0, 1, 2, 3])), /// // Create a stream from an iterator.
/// Result::Ok(Bytes::from_static(&[4, 5, 6, 7])), /// let stream = tokio_stream::iter(vec![
/// Result::Ok(Bytes::from_static(&[8, 9, 10, 11])), /// Result::Ok(Bytes::from_static(&[0, 1, 2, 3])),
/// ]); /// Result::Ok(Bytes::from_static(&[4, 5, 6, 7])),
/// /// Result::Ok(Bytes::from_static(&[8, 9, 10, 11])),
/// // Convert it to an AsyncRead. /// ]);
/// let mut read = StreamReader::new(stream); ///
/// /// // Convert it to an AsyncRead.
/// // Read five bytes from the stream. /// let mut read = StreamReader::new(stream);
/// let mut buf = [0; 5]; ///
/// read.read_exact(&mut buf).await?; /// // Read five bytes from the stream.
/// assert_eq!(buf, [0, 1, 2, 3, 4]); /// let mut buf = [0; 5];
/// /// read.read_exact(&mut buf).await?;
/// // Read the rest of the current chunk. /// assert_eq!(buf, [0, 1, 2, 3, 4]);
/// assert_eq!(read.read(&mut buf).await?, 3); ///
/// assert_eq!(&buf[..3], [5, 6, 7]); /// // Read the rest of the current chunk.
/// /// assert_eq!(read.read(&mut buf).await?, 3);
/// // Read the next chunk. /// assert_eq!(&buf[..3], [5, 6, 7]);
/// assert_eq!(read.read(&mut buf).await?, 4); ///
/// assert_eq!(&buf[..4], [8, 9, 10, 11]); /// // Read the next chunk.
/// /// assert_eq!(read.read(&mut buf).await?, 4);
/// // We have now reached the end. /// assert_eq!(&buf[..4], [8, 9, 10, 11]);
/// assert_eq!(read.read(&mut buf).await?, 0); ///
/// /// // We have now reached the end.
/// # Ok(()) /// assert_eq!(read.read(&mut buf).await?, 0);
/// # } ///
/// ``` /// # Ok(())
/// /// # }
/// If the stream produces errors which are not [std::io::Error], /// ```
/// the errors can be converted using [`StreamExt`] to map each ///
/// element. /// If the stream produces errors which are not [`std::io::Error`],
/// /// the errors can be converted using [`StreamExt`] to map each
/// ``` /// element.
/// use bytes::Bytes; ///
/// use tokio::io::AsyncReadExt; /// ```
/// use tokio_util::io::StreamReader; /// use bytes::Bytes;
/// use tokio_stream::StreamExt; /// use tokio::io::AsyncReadExt;
/// # #[tokio::main] /// use tokio_util::io::StreamReader;
/// # async fn main() -> std::io::Result<()> { /// use tokio_stream::StreamExt;
/// /// # #[tokio::main(flavor = "current_thread")]
/// // Create a stream from an iterator, including an error. /// # async fn main() -> std::io::Result<()> {
/// let stream = tokio_stream::iter(vec![ ///
/// Result::Ok(Bytes::from_static(&[0, 1, 2, 3])), /// // Create a stream from an iterator, including an error.
/// Result::Ok(Bytes::from_static(&[4, 5, 6, 7])), /// let stream = tokio_stream::iter(vec![
/// Result::Err("Something bad happened!") /// Result::Ok(Bytes::from_static(&[0, 1, 2, 3])),
/// ]); /// Result::Ok(Bytes::from_static(&[4, 5, 6, 7])),
/// /// Result::Err("Something bad happened!")
/// // Use StreamExt to map the stream and error to a std::io::Error /// ]);
/// let stream = stream.map(|result| result.map_err(|err| { ///
/// std::io::Error::new(std::io::ErrorKind::Other, err) /// // Use StreamExt to map the stream and error to a std::io::Error
/// })); /// let stream = stream.map(|result| result.map_err(|err| {
/// /// std::io::Error::new(std::io::ErrorKind::Other, err)
/// // Convert it to an AsyncRead. /// }));
/// let mut read = StreamReader::new(stream); ///
/// /// // Convert it to an AsyncRead.
/// // Read five bytes from the stream. /// let mut read = StreamReader::new(stream);
/// let mut buf = [0; 5]; ///
/// read.read_exact(&mut buf).await?; /// // Read five bytes from the stream.
/// assert_eq!(buf, [0, 1, 2, 3, 4]); /// let mut buf = [0; 5];
/// /// read.read_exact(&mut buf).await?;
/// // Read the rest of the current chunk. /// assert_eq!(buf, [0, 1, 2, 3, 4]);
/// assert_eq!(read.read(&mut buf).await?, 3); ///
/// assert_eq!(&buf[..3], [5, 6, 7]); /// // Read the rest of the current chunk.
/// /// assert_eq!(read.read(&mut buf).await?, 3);
/// // Reading the next chunk will produce an error /// assert_eq!(&buf[..3], [5, 6, 7]);
/// let error = read.read(&mut buf).await.unwrap_err(); ///
/// assert_eq!(error.kind(), std::io::ErrorKind::Other); /// // Reading the next chunk will produce an error
/// assert_eq!(error.into_inner().unwrap().to_string(), "Something bad happened!"); /// let error = read.read(&mut buf).await.unwrap_err();
/// /// assert_eq!(error.kind(), std::io::ErrorKind::Other);
/// // We have now reached the end. /// assert_eq!(error.into_inner().unwrap().to_string(), "Something bad happened!");
/// assert_eq!(read.read(&mut buf).await?, 0); ///
/// /// // We have now reached the end.
/// # Ok(()) /// assert_eq!(read.read(&mut buf).await?, 0);
/// # } ///
/// ``` /// # Ok(())
/// /// # }
/// [`AsyncRead`]: tokio::io::AsyncRead /// ```
/// [`Stream`]: futures_core::Stream ///
/// [`ReaderStream`]: crate::io::ReaderStream /// Using the [`AsyncBufRead`] impl, you can read a `Stream` of byte chunks
/// [`StreamExt`]: tokio_stream::StreamExt /// line-by-line. Note that you will usually also need to convert the error
#[derive(Debug)] /// type when doing this. See the second example for an explanation of how
pub struct StreamReader<S, B> { /// to do this.
#[pin] ///
inner: S, /// ```
chunk: Option<B>, /// use tokio::io::{Result, AsyncBufReadExt};
} /// use tokio_util::io::StreamReader;
/// # #[tokio::main(flavor = "current_thread")]
/// # async fn main() -> std::io::Result<()> {
///
/// // Create a stream of byte chunks.
/// let stream = tokio_stream::iter(vec![
/// Result::Ok(b"The first line.\n".as_slice()),
/// Result::Ok(b"The second line.".as_slice()),
/// Result::Ok(b"\nThe third".as_slice()),
/// Result::Ok(b" line.\nThe fourth line.\nThe fifth line.\n".as_slice()),
/// ]);
///
/// // Convert it to an AsyncRead.
/// let mut read = StreamReader::new(stream);
///
/// // Loop through the lines from the `StreamReader`.
/// let mut line = String::new();
/// let mut lines = Vec::new();
/// loop {
/// line.clear();
/// let len = read.read_line(&mut line).await?;
/// if len == 0 { break; }
/// lines.push(line.clone());
/// }
///
/// // Verify that we got the lines we expected.
/// assert_eq!(
/// lines,
/// vec![
/// "The first line.\n",
/// "The second line.\n",
/// "The third line.\n",
/// "The fourth line.\n",
/// "The fifth line.\n",
/// ]
/// );
/// # Ok(())
/// # }
/// ```
///
/// [`AsyncRead`]: tokio::io::AsyncRead
/// [`AsyncBufRead`]: tokio::io::AsyncBufRead
/// [`Stream`]: futures_core::Stream
/// [`ReaderStream`]: crate::io::ReaderStream
/// [`StreamExt`]: https://docs.rs/tokio-stream/latest/tokio_stream/trait.StreamExt.html
#[derive(Debug)]
pub struct StreamReader<S, B> {
// This field is pinned.
inner: S,
// This field is not pinned.
chunk: Option<B>,
} }
impl<S, B, E> StreamReader<S, B> impl<S, B, E> StreamReader<S, B>
@@ -250,3 +299,28 @@ where
} }
} }
} }
// The code below is a manual expansion of the code that pin-project-lite would
// generate. This is done because pin-project-lite fails by hitting the recusion
// limit on this struct. (Every line of documentation is handled recursively by
// the macro.)
impl<S: Unpin, B> Unpin for StreamReader<S, B> {}
struct StreamReaderProject<'a, S, B> {
inner: Pin<&'a mut S>,
chunk: &'a mut Option<B>,
}
impl<S, B> StreamReader<S, B> {
#[inline]
fn project(self: Pin<&mut Self>) -> StreamReaderProject<'_, S, B> {
// SAFETY: We define that only `inner` should be pinned when `Self` is
// and have an appropriate `impl Unpin` for this.
let me = unsafe { Pin::into_inner_unchecked(self) };
StreamReaderProject {
inner: unsafe { Pin::new_unchecked(&mut me.inner) },
chunk: &mut me.chunk,
}
}
}