diff --git a/src/io/mod.rs b/src/io/mod.rs index 3e44a6b4e..6d1341c05 100644 --- a/src/io/mod.rs +++ b/src/io/mod.rs @@ -36,6 +36,7 @@ mod flush; mod read_exact; mod read_to_end; mod read; +mod read_until; mod split; mod window; mod write_all; @@ -44,6 +45,7 @@ pub use self::flush::{flush, Flush}; pub use self::read_exact::{read_exact, ReadExact}; pub use self::read_to_end::{read_to_end, ReadToEnd}; pub use self::read::read; +pub use self::read_until::{read_until, ReadUntil}; pub use self::split::{ReadHalf, WriteHalf}; pub use self::window::Window; pub use self::write_all::{write_all, WriteAll}; diff --git a/src/io/read_until.rs b/src/io/read_until.rs new file mode 100644 index 000000000..bb8c0073c --- /dev/null +++ b/src/io/read_until.rs @@ -0,0 +1,70 @@ +use std::io::{self, Read, BufRead}; +use std::mem; + +use futures::{Poll, Future}; + +/// A future which can be used to easily read the contents of a stream into a +/// vector until the delimiter is reached. +/// +/// Created by the [`read_until`] function. +/// +/// [`read_until`]: fn.read_until.html +pub struct ReadUntil { + state: State, +} + +enum State { + Reading { + a: A, + byte: u8, + buf: Vec, + }, + Empty, +} + +/// Creates a future which will read all the bytes associated with the I/O +/// object `A` into the buffer provided until the delimiter `byte` is reached. +/// This method is the async equivalent to [`BufRead::read_until`]. +/// +/// In case of an error the buffer and the object will be discarded, with +/// the error yielded. In the case of success the object will be destroyed and +/// the buffer will be returned, with all bytes up to, and including, the delimiter +/// (if found). +/// +/// [`ButRead::read_until`]: https://doc.rust-lang.org/std/io/trait.BufRead.html#method.read_until +pub fn read_until(a: A, byte: u8, buf: Vec) -> ReadUntil + where A: BufRead +{ + ReadUntil { + state: State::Reading { + a: a, + byte: byte, + buf: buf, + } + } +} + +impl Future for ReadUntil + where A: Read + BufRead +{ + type Item = (A, Vec); + type Error = io::Error; + + fn poll(&mut self) -> Poll<(A, Vec), io::Error> { + match self.state { + State::Reading { ref mut a, byte, ref mut buf } => { + // If we get `Ok(n)`, then we know the stream hit EOF or the delimiter. + // and just return it, as we are finished. + // If we hit "would block" then all the read data so far + // is in our buffer, and otherwise we propagate errors. + try_nb!(a.read_until(byte, buf)); + }, + State::Empty => panic!("poll ReadUntil after it's done"), + } + + match mem::replace(&mut self.state, State::Empty) { + State::Reading { a, byte: _, buf } => Ok((a, buf).into()), + State::Empty => unreachable!(), + } + } +}