buf: Add BufStreamExt trait and add a core feature (#897)

This change adds an extension trait to `BufStream` and puts the core
trait behind a feature flag for optional use.

This mainly adds the additional functions in an extension trait to
allow the user to select if they want just the core trait or the fully
featured version. Now the user can add the core feature to _not_
include the extension trait. By deafult, this feature is disabled.
This commit is contained in:
Lucio Franco
2019-02-19 13:30:37 -08:00
committed by Carl Lerche
parent c08e73c8d4
commit dd66096ea0
17 changed files with 562 additions and 531 deletions
+145 -2
View File
@@ -10,11 +10,154 @@
//! `Buf` (i.e, byte collections).
extern crate bytes;
#[cfg(feature = "ext")]
extern crate either;
#[allow(unused)]
#[macro_use]
extern crate futures;
pub mod buf_stream;
#[cfg(feature = "ext")]
pub mod ext;
pub mod errors;
mod size_hint;
mod str;
#[doc(inline)]
pub use buf_stream::BufStream;
#[cfg(feature = "ext")]
pub use ext::BufStreamExt;
pub use self::size_hint::SizeHint;
use futures::Poll;
use bytes::{Buf, Bytes, BytesMut};
use std::io;
use errors::internal::Never;
/// An asynchronous stream of bytes.
///
/// `BufStream` asynchronously yields values implementing `Buf`, i.e. byte
/// buffers.
pub trait BufStream {
/// Values yielded by the `BufStream`.
///
/// Each item is a sequence of bytes representing a chunk of the total
/// `ByteStream`.
type Item: Buf;
/// The error type this `BufStream` might generate.
type Error;
/// Attempt to pull out the next buffer of this stream, registering the
/// current task for wakeup if the value is not yet available, and returning
/// `None` if the stream is exhausted.
///
/// # Return value
///
/// There are several possible return values, each indicating a distinct
/// stream state:
///
/// - `Ok(Async::NotReady)` means that this stream's next value is not ready
/// yet. Implementations will ensure that the current task will be notified
/// when the next value may be ready.
///
/// - `Ok(Async::Ready(Some(buf)))` means that the stream has successfully
/// produced a value, `buf`, and may produce further values on subsequent
/// `poll_buf` calls.
///
/// - `Ok(Async::Ready(None))` means that the stream has terminated, and
/// `poll_buf` should not be invoked again.
///
/// # Panics
///
/// Once a stream is finished, i.e. `Ready(None)` has been returned, further
/// calls to `poll_buf` may result in a panic or other "bad behavior".
fn poll_buf(&mut self) -> Poll<Option<Self::Item>, Self::Error>;
/// Returns the bounds on the remaining length of the stream.
///
/// The size hint allows the caller to perform certain optimizations that
/// are dependent on the byte stream size. For example, `collect` uses the
/// size hint to pre-allocate enough capacity to store the entirety of the
/// data received from the byte stream.
///
/// When `SizeHint::upper()` returns `Some` with a value equal to
/// `SizeHint::lower()`, this represents the exact number of bytes that will
/// be yielded by the `BufStream`.
///
/// # Implementation notes
///
/// While not enforced, implementations are expected to respect the values
/// returned from `SizeHint`. Any deviation is considered an implementation
/// bug. Consumers may rely on correctness in order to use the value as part
/// of protocol impelmentations. For example, an HTTP library may use the
/// size hint to set the `content-length` header.
///
/// However, `size_hint` must not be trusted to omit bounds checks in unsafe
/// code. An incorrect implementation of `size_hint()` must not lead to
/// memory safety violations.
fn size_hint(&self) -> SizeHint {
SizeHint::default()
}
}
impl BufStream for Vec<u8> {
type Item = io::Cursor<Vec<u8>>;
type Error = Never;
fn poll_buf(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
if self.is_empty() {
return Ok(None.into());
}
poll_bytes(self)
}
}
impl BufStream for &'static [u8] {
type Item = io::Cursor<&'static [u8]>;
type Error = Never;
fn poll_buf(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
if self.is_empty() {
return Ok(None.into());
}
poll_bytes(self)
}
}
impl BufStream for Bytes {
type Item = io::Cursor<Bytes>;
type Error = Never;
fn poll_buf(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
if self.is_empty() {
return Ok(None.into());
}
poll_bytes(self)
}
}
impl BufStream for BytesMut {
type Item = io::Cursor<BytesMut>;
type Error = Never;
fn poll_buf(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
if self.is_empty() {
return Ok(None.into());
}
poll_bytes(self)
}
}
fn poll_bytes<T: Default>(buf: &mut T)
-> Poll<Option<io::Cursor<T>>, Never>
{
use std::mem;
let bytes = mem::replace(buf, Default::default());
let buf = io::Cursor::new(bytes);
Ok(Some(buf).into())
}