diff --git a/quinn/src/connection.rs b/quinn/src/connection.rs index 763d09979..62ee123d9 100644 --- a/quinn/src/connection.rs +++ b/quinn/src/connection.rs @@ -1,12 +1,10 @@ use std::collections::VecDeque; +use std::mem; use std::net::SocketAddr; -use std::str; use std::sync::{Arc, Mutex}; use std::time::Instant; -use std::{io, mem}; use bytes::Bytes; -use err_derive::Error; use fnv::FnvHashMap; use futures::sync::{mpsc, oneshot}; use futures::task::{self, Task}; @@ -14,7 +12,6 @@ use futures::Stream as FuturesStream; use futures::{Async, Future, Poll}; use quinn_proto::{self as quinn, ConnectionHandle, Directionality, StreamId, TimerUpdate}; use slog::Logger; -use tokio_io::{AsyncRead, AsyncWrite}; use tokio_timer::Delay; pub use crate::quinn::{ @@ -26,6 +23,7 @@ pub use crate::tls::{Certificate, CertificateChain, PrivateKey}; pub use crate::builders::{ ClientConfigBuilder, EndpointBuilder, EndpointError, ServerConfigBuilder, }; +use crate::streams::{NewStream, RecvStream, SendStream}; use crate::{ConnectionEvent, EndpointEvent}; /// Connecting future @@ -286,14 +284,6 @@ impl FuturesStream for IncomingStreams { } } -/// A stream initiated by a remote peer. -pub enum NewStream { - /// A unidirectional stream. - Uni(RecvStream), - /// A bidirectional stream. - Bi(SendStream, RecvStream), -} - pub struct ConnectionRef(Arc>); impl ConnectionRef { @@ -363,21 +353,21 @@ impl std::ops::Deref for ConnectionRef { pub struct ConnectionInner { log: Logger, epoch: Instant, - inner: quinn::Connection, + pub(crate) inner: quinn::Connection, driver: Option, handle: ConnectionHandle, connected: bool, timers: [Option; quinn::Timer::COUNT], conn_events: mpsc::UnboundedReceiver, endpoint_events: mpsc::UnboundedSender<(ConnectionHandle, EndpointEvent)>, - blocked_writers: FnvHashMap, - blocked_readers: FnvHashMap, + pub(crate) blocked_writers: FnvHashMap, + pub(crate) blocked_readers: FnvHashMap, uni_opening: VecDeque>>, bi_opening: VecDeque>>, incoming_streams_reader: Option, - finishing: FnvHashMap>>, + pub(crate) finishing: FnvHashMap>>, /// Always set to Some before the connection becomes drained - error: Option, + pub(crate) error: Option, /// Number of live handles that can be used to initiate or handle I/O; excludes the driver ref_count: usize, } @@ -527,7 +517,7 @@ impl ConnectionInner { } /// Wake up a blocked `Driver` task to process I/O - fn notify(&self) { + pub(crate) fn notify(&self) { if let Some(x) = self.driver.as_ref() { x.notify(); } @@ -567,7 +557,7 @@ impl ConnectionInner { self.close(0, Bytes::new()); } - fn check_0rtt(&self) -> Result<(), ()> { + pub(crate) fn check_0rtt(&self) -> Result<(), ()> { match self.inner.is_handshaking() || self.inner.accepted_0rtt() { true => Ok(()), false => Err(()), @@ -586,418 +576,3 @@ impl Drop for ConnectionInner { } } } - -/// A stream that can only be used to send data -pub struct SendStream { - conn: ConnectionRef, - stream: StreamId, - is_0rtt: bool, - finishing: Option>>, - finished: bool, -} - -impl SendStream { - fn new(conn: ConnectionRef, stream: StreamId, is_0rtt: bool) -> Self { - Self { - conn, - stream, - is_0rtt, - finishing: None, - finished: false, - } - } - - /// Write bytes to the stream. - /// - /// Returns the number of bytes written on success. Congestion and flow control may cause this - /// to be shorter than `buf.len()`, indicating that only a prefix of `buf` was written. - pub fn poll_write(&mut self, buf: &[u8]) -> Poll { - use crate::quinn::WriteError::*; - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt { - conn.check_0rtt() - .map_err(|()| WriteError::ZeroRttRejected)?; - } - let n = match conn.inner.write(self.stream, buf) { - Ok(n) => n, - Err(Blocked) => { - if let Some(ref x) = conn.error { - return Err(WriteError::ConnectionClosed(x.clone())); - } - conn.blocked_writers.insert(self.stream, task::current()); - return Ok(Async::NotReady); - } - Err(Stopped { error_code }) => { - return Err(WriteError::Stopped { error_code }); - } - Err(UnknownStream) => { - return Err(WriteError::UnknownStream); - } - }; - conn.notify(); - Ok(Async::Ready(n)) - } - - /// Shut down the send stream gracefully. - /// - /// No new data may be written after calling this method. Completes when the peer has - /// acknowledged all sent data, retransmitting data as needed. - pub fn poll_finish(&mut self) -> Poll<(), FinishError> { - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt { - conn.check_0rtt() - .map_err(|()| FinishError::ZeroRttRejected)?; - } - if self.finishing.is_none() { - conn.inner.finish(self.stream); - let (send, recv) = oneshot::channel(); - self.finishing = Some(recv); - conn.finishing.insert(self.stream, send); - conn.notify(); - } - let r = self.finishing.as_mut().unwrap().poll().unwrap(); - match r { - Async::Ready(None) => { - self.finished = true; - Ok(Async::Ready(())) - } - Async::Ready(Some(e)) => Err(FinishError::ConnectionLost(e)), - Async::NotReady => Ok(Async::NotReady), - } - } - - /// Close the send stream immediately. - /// - /// No new data can be written after calling this method. Locally buffered data is dropped, - /// and previously transmitted data will no longer be retransmitted if lost. If `poll_finish` - /// was called previously and all data has already been transmitted at least once, the peer - /// may still receive all written data. - pub fn reset(&mut self, error_code: u16) { - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt && conn.check_0rtt().is_err() { - return; - } - conn.inner.reset(self.stream, error_code); - conn.notify(); - } -} - -impl io::Write for SendStream { - fn write(&mut self, buf: &[u8]) -> io::Result { - match self.poll_write(buf) { - Ok(Async::Ready(n)) => Ok(n), - Ok(Async::NotReady) => Err(io::Error::new(io::ErrorKind::WouldBlock, "stream blocked")), - Err(e) => Err(e.into()), - } - } - - fn flush(&mut self) -> io::Result<()> { - Ok(()) - } -} - -impl AsyncWrite for SendStream { - fn shutdown(&mut self) -> Poll<(), io::Error> { - self.poll_finish().map_err(|e| e.into()) - } -} - -impl Drop for SendStream { - fn drop(&mut self) { - let mut conn = self.conn.lock().unwrap(); - if conn.error.is_some() { - return; - } - if !self.finished { - conn.inner.reset(self.stream, 0); - conn.notify(); - } - } -} - -/// A stream that can only be used to receive data -pub struct RecvStream { - conn: ConnectionRef, - stream: StreamId, - is_0rtt: bool, - recvd: bool, -} - -impl RecvStream { - fn new(conn: ConnectionRef, stream: StreamId, is_0rtt: bool) -> Self { - Self { - conn, - stream, - is_0rtt, - recvd: false, - } - } - - /// Read data contiguously from the stream. - /// - /// Returns the number of bytes read into `buf` on success. - /// - /// Applications involving bulk data transfer should consider using unordered reads for - /// improved performance. - /// - /// # Panics - /// - If called after `poll_read_unordered` was called on the same stream. - /// This is forbidden because an unordered read could consume a segment of data from a - /// location other than the start of the receive buffer, making it impossible for future - pub fn poll_read(&mut self, buf: &mut [u8]) -> Poll { - use crate::quinn::ReadError::*; - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt { - conn.check_0rtt().map_err(|()| ReadError::ZeroRttRejected)?; - } - match conn.inner.read(self.stream, buf) { - Ok(n) => Ok(Async::Ready(n)), - Err(Blocked) => { - if let Some(ref x) = conn.error { - return Err(ReadError::ConnectionClosed(x.clone())); - } - conn.blocked_readers.insert(self.stream, task::current()); - Ok(Async::NotReady) - } - Err(Reset { error_code }) => { - self.recvd = true; - Err(ReadError::Reset { error_code }) - } - Err(Finished) => { - self.recvd = true; - Err(ReadError::Finished) - } - Err(UnknownStream) => Err(ReadError::UnknownStream), - } - } - - /// Read a segment of data from any offset in the stream. - /// - /// Returns a segment of data and their offset in the stream. Segments may be received in any - /// order and may overlap. - /// - /// Unordered reads have reduced overhead and higher throughput, and should therefore be - /// preferred when applicable. - pub fn poll_read_unordered(&mut self) -> Poll<(Bytes, u64), ReadError> { - use crate::quinn::ReadError::*; - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt { - conn.check_0rtt().map_err(|()| ReadError::ZeroRttRejected)?; - } - match conn.inner.read_unordered(self.stream) { - Ok((bytes, offset)) => Ok(Async::Ready((bytes, offset))), - Err(Blocked) => { - if let Some(ref x) = conn.error { - return Err(ReadError::ConnectionClosed(x.clone())); - } - conn.blocked_readers.insert(self.stream, task::current()); - Ok(Async::NotReady) - } - Err(Reset { error_code }) => { - self.recvd = true; - Err(ReadError::Reset { error_code }) - } - Err(Finished) => { - self.recvd = true; - Err(ReadError::Finished) - } - Err(UnknownStream) => Err(ReadError::UnknownStream), - } - } - - /// Uses unordered reads to be more efficient than using `AsyncRead` would allow - pub fn read_to_end(self, size_limit: usize) -> ReadToEnd { - ReadToEnd { - stream: Some(self), - size_limit, - buffer: Vec::new(), - } - } - - /// Close the receive stream immediately. - /// - /// The peer is notified and will cease transmitting on this stream, as if it had reset the - /// stream itself. Further data may still be received on this stream if it was already in - /// flight. Once called, a `ReadError::Reset` should be expected soon, although a peer might - /// manage to finish the stream before it receives the reset, and a misbehaving peer might - /// ignore the request entirely and continue sending until halted by flow control. - /// - /// Has no effect if the incoming stream already finished. - pub fn stop(&mut self, error_code: u16) { - let mut conn = self.conn.lock().unwrap(); - if self.is_0rtt && conn.check_0rtt().is_err() { - return; - } - conn.inner.stop_sending(self.stream, error_code); - conn.notify(); - self.recvd = true; - } -} - -/// Future produced by `read_to_end` -pub struct ReadToEnd { - stream: Option, - buffer: Vec, - size_limit: usize, -} - -impl Future for ReadToEnd { - type Item = (RecvStream, Box<[u8]>); - type Error = ReadError; - fn poll(&mut self) -> Poll { - loop { - match self.stream.as_mut().unwrap().poll_read_unordered() { - Ok(Async::Ready((data, offset))) => { - let len = self.buffer.len().max(offset as usize + data.len()); - if len > self.size_limit { - return Err(ReadError::Finished); - } - self.buffer.resize(len, 0); - self.buffer[offset as usize..offset as usize + data.len()] - .copy_from_slice(&data); - } - Ok(Async::NotReady) => { - return Ok(Async::NotReady); - } - Err(ReadError::Finished) => { - return Ok(Async::Ready(( - self.stream.take().unwrap(), - mem::replace(&mut self.buffer, Vec::new()).into(), - ))); - } - Err(e) => { - return Err(e); - } - } - } - } -} - -impl io::Read for RecvStream { - fn read(&mut self, buf: &mut [u8]) -> io::Result { - match self.poll_read(buf) { - Ok(Async::Ready(n)) => Ok(n), - Err(ReadError::Finished) => Ok(0), - Ok(Async::NotReady) => Err(io::Error::new(io::ErrorKind::WouldBlock, "stream blocked")), - Err(e) => Err(e.into()), - } - } -} - -impl AsyncRead for RecvStream { - unsafe fn prepare_uninitialized_buffer(&self, _: &mut [u8]) -> bool { - false - } -} - -impl Drop for RecvStream { - fn drop(&mut self) { - let mut conn = self.conn.lock().unwrap(); - if conn.error.is_some() { - return; - } - if !self.recvd { - conn.inner.stop_sending(self.stream, 0); - conn.notify(); - } - } -} - -/// Errors that arise from finishing a stream -#[derive(Debug, Error, Clone)] -pub enum FinishError { - /// The connection was lost. - #[error(display = "connection lost: {}", _0)] - ConnectionLost(ConnectionError), - /// This was a 0-RTT stream and the server rejected it. - /// - /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). - #[error(display = "0-RTT rejected")] - ZeroRttRejected, -} - -impl From for io::Error { - fn from(x: FinishError) -> Self { - use self::FinishError::*; - match x { - ConnectionLost(e) => e.into(), - ZeroRttRejected => io::Error::new(io::ErrorKind::ConnectionReset, "0-RTT rejected"), - } - } -} - -/// Errors that arise from reading from a stream. -#[derive(Debug, Error, Clone)] -pub enum ReadError { - /// The peer abandoned transmitting data on this stream. - #[error(display = "stream reset by peer: error {}", error_code)] - Reset { - /// The error code supplied by the peer. - error_code: u16, - }, - /// The data on this stream has been fully delivered and no more will be transmitted. - #[error(display = "the stream has been completely received")] - Finished, - /// The connection was closed. - #[error(display = "connection closed: {}", _0)] - ConnectionClosed(ConnectionError), - /// Unknown stream - #[error(display = "unknown stream")] - UnknownStream, - /// This was a 0-RTT stream and the server rejected it. - /// - /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). - #[error(display = "0-RTT rejected")] - ZeroRttRejected, -} - -impl From for io::Error { - fn from(x: ReadError) -> Self { - use self::ReadError::*; - let kind = match x { - ConnectionClosed(e) => { - return e.into(); - } - Reset { .. } | ZeroRttRejected => io::ErrorKind::ConnectionReset, - Finished => io::ErrorKind::UnexpectedEof, - UnknownStream => io::ErrorKind::NotConnected, - }; - io::Error::new(kind, x) - } -} - -/// Errors that arise from writing to a stream -#[derive(Debug, Error, Clone)] -pub enum WriteError { - /// The peer is no longer accepting data on this stream. - #[error(display = "sending stopped by peer: error {}", error_code)] - Stopped { - /// The error code supplied by the peer. - error_code: u16, - }, - /// The connection was closed. - #[error(display = "connection closed: {}", _0)] - ConnectionClosed(ConnectionError), - /// Unknown stream - #[error(display = "unknown stream")] - UnknownStream, - /// This was a 0-RTT stream and the server rejected it. - /// - /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). - #[error(display = "0-RTT rejected")] - ZeroRttRejected, -} - -impl From for io::Error { - fn from(x: WriteError) -> Self { - use self::WriteError::*; - let kind = match x { - ConnectionClosed(e) => { - return e.into(); - } - Stopped { .. } | ZeroRttRejected => io::ErrorKind::ConnectionReset, - UnknownStream => io::ErrorKind::NotConnected, - }; - io::Error::new(kind, x) - } -} diff --git a/quinn/src/lib.rs b/quinn/src/lib.rs index 0636b6e5a..215e15618 100644 --- a/quinn/src/lib.rs +++ b/quinn/src/lib.rs @@ -72,14 +72,14 @@ pub use crate::builders::{ }; mod connection; -pub use connection::{ - Connecting, Connection, ConnectionDriver, IncomingStreams, NewStream, ReadError, ReadToEnd, - RecvStream, SendStream, WriteError, -}; +pub use connection::{Connecting, Connection, ConnectionDriver, IncomingStreams}; mod endpoint; pub use endpoint::{Endpoint, EndpointDriver, Incoming}; +mod streams; +pub use streams::{NewStream, ReadError, ReadToEnd, RecvStream, SendStream, WriteError}; + #[cfg(test)] mod tests; diff --git a/quinn/src/streams.rs b/quinn/src/streams.rs new file mode 100644 index 000000000..a0b39b35a --- /dev/null +++ b/quinn/src/streams.rs @@ -0,0 +1,444 @@ +use std::str; +use std::{io, mem}; + +use bytes::Bytes; +use err_derive::Error; +use futures::sync::oneshot; +use futures::task; +use futures::{Async, Future, Poll}; +use quinn_proto::StreamId; +use tokio_io::{AsyncRead, AsyncWrite}; + +pub use crate::quinn::{ + ConnectError, ConnectionError, ConnectionId, DatagramEvent, ServerConfig, Transmit, + TransportConfig, ALPN_QUIC_H3, ALPN_QUIC_HTTP, +}; +pub use crate::tls::{Certificate, CertificateChain, PrivateKey}; + +pub use crate::builders::{ + ClientConfigBuilder, EndpointBuilder, EndpointError, ServerConfigBuilder, +}; +use crate::connection::ConnectionRef; + +/// A stream initiated by a remote peer. +pub enum NewStream { + /// A unidirectional stream. + Uni(RecvStream), + /// A bidirectional stream. + Bi(SendStream, RecvStream), +} + +/// A stream that can only be used to send data +pub struct SendStream { + conn: ConnectionRef, + stream: StreamId, + is_0rtt: bool, + finishing: Option>>, + finished: bool, +} + +impl SendStream { + pub(crate) fn new(conn: ConnectionRef, stream: StreamId, is_0rtt: bool) -> Self { + Self { + conn, + stream, + is_0rtt, + finishing: None, + finished: false, + } + } + + /// Write bytes to the stream. + /// + /// Returns the number of bytes written on success. Congestion and flow control may cause this + /// to be shorter than `buf.len()`, indicating that only a prefix of `buf` was written. + pub fn poll_write(&mut self, buf: &[u8]) -> Poll { + use crate::quinn::WriteError::*; + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt { + conn.check_0rtt() + .map_err(|()| WriteError::ZeroRttRejected)?; + } + let n = match conn.inner.write(self.stream, buf) { + Ok(n) => n, + Err(Blocked) => { + if let Some(ref x) = conn.error { + return Err(WriteError::ConnectionClosed(x.clone())); + } + conn.blocked_writers.insert(self.stream, task::current()); + return Ok(Async::NotReady); + } + Err(Stopped { error_code }) => { + return Err(WriteError::Stopped { error_code }); + } + Err(UnknownStream) => { + return Err(WriteError::UnknownStream); + } + }; + conn.notify(); + Ok(Async::Ready(n)) + } + + /// Shut down the send stream gracefully. + /// + /// No new data may be written after calling this method. Completes when the peer has + /// acknowledged all sent data, retransmitting data as needed. + pub fn poll_finish(&mut self) -> Poll<(), FinishError> { + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt { + conn.check_0rtt() + .map_err(|()| FinishError::ZeroRttRejected)?; + } + if self.finishing.is_none() { + conn.inner.finish(self.stream); + let (send, recv) = oneshot::channel(); + self.finishing = Some(recv); + conn.finishing.insert(self.stream, send); + conn.notify(); + } + let r = self.finishing.as_mut().unwrap().poll().unwrap(); + match r { + Async::Ready(None) => { + self.finished = true; + Ok(Async::Ready(())) + } + Async::Ready(Some(e)) => Err(FinishError::ConnectionLost(e)), + Async::NotReady => Ok(Async::NotReady), + } + } + + /// Close the send stream immediately. + /// + /// No new data can be written after calling this method. Locally buffered data is dropped, + /// and previously transmitted data will no longer be retransmitted if lost. If `poll_finish` + /// was called previously and all data has already been transmitted at least once, the peer + /// may still receive all written data. + pub fn reset(&mut self, error_code: u16) { + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt && conn.check_0rtt().is_err() { + return; + } + conn.inner.reset(self.stream, error_code); + conn.notify(); + } +} + +impl io::Write for SendStream { + fn write(&mut self, buf: &[u8]) -> io::Result { + match self.poll_write(buf) { + Ok(Async::Ready(n)) => Ok(n), + Ok(Async::NotReady) => Err(io::Error::new(io::ErrorKind::WouldBlock, "stream blocked")), + Err(e) => Err(e.into()), + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +impl AsyncWrite for SendStream { + fn shutdown(&mut self) -> Poll<(), io::Error> { + self.poll_finish().map_err(|e| e.into()) + } +} + +impl Drop for SendStream { + fn drop(&mut self) { + let mut conn = self.conn.lock().unwrap(); + if conn.error.is_some() { + return; + } + if !self.finished { + conn.inner.reset(self.stream, 0); + conn.notify(); + } + } +} + +/// A stream that can only be used to receive data +pub struct RecvStream { + conn: ConnectionRef, + stream: StreamId, + is_0rtt: bool, + recvd: bool, +} + +impl RecvStream { + pub(crate) fn new(conn: ConnectionRef, stream: StreamId, is_0rtt: bool) -> Self { + Self { + conn, + stream, + is_0rtt, + recvd: false, + } + } + + /// Read data contiguously from the stream. + /// + /// Returns the number of bytes read into `buf` on success. + /// + /// Applications involving bulk data transfer should consider using unordered reads for + /// improved performance. + /// + /// # Panics + /// - If called after `poll_read_unordered` was called on the same stream. + /// This is forbidden because an unordered read could consume a segment of data from a + /// location other than the start of the receive buffer, making it impossible for future + pub fn poll_read(&mut self, buf: &mut [u8]) -> Poll { + use crate::quinn::ReadError::*; + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt { + conn.check_0rtt().map_err(|()| ReadError::ZeroRttRejected)?; + } + match conn.inner.read(self.stream, buf) { + Ok(n) => Ok(Async::Ready(n)), + Err(Blocked) => { + if let Some(ref x) = conn.error { + return Err(ReadError::ConnectionClosed(x.clone())); + } + conn.blocked_readers.insert(self.stream, task::current()); + Ok(Async::NotReady) + } + Err(Reset { error_code }) => { + self.recvd = true; + Err(ReadError::Reset { error_code }) + } + Err(Finished) => { + self.recvd = true; + Err(ReadError::Finished) + } + Err(UnknownStream) => Err(ReadError::UnknownStream), + } + } + + /// Read a segment of data from any offset in the stream. + /// + /// Returns a segment of data and their offset in the stream. Segments may be received in any + /// order and may overlap. + /// + /// Unordered reads have reduced overhead and higher throughput, and should therefore be + /// preferred when applicable. + pub fn poll_read_unordered(&mut self) -> Poll<(Bytes, u64), ReadError> { + use crate::quinn::ReadError::*; + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt { + conn.check_0rtt().map_err(|()| ReadError::ZeroRttRejected)?; + } + match conn.inner.read_unordered(self.stream) { + Ok((bytes, offset)) => Ok(Async::Ready((bytes, offset))), + Err(Blocked) => { + if let Some(ref x) = conn.error { + return Err(ReadError::ConnectionClosed(x.clone())); + } + conn.blocked_readers.insert(self.stream, task::current()); + Ok(Async::NotReady) + } + Err(Reset { error_code }) => { + self.recvd = true; + Err(ReadError::Reset { error_code }) + } + Err(Finished) => { + self.recvd = true; + Err(ReadError::Finished) + } + Err(UnknownStream) => Err(ReadError::UnknownStream), + } + } + + /// Uses unordered reads to be more efficient than using `AsyncRead` would allow + pub fn read_to_end(self, size_limit: usize) -> ReadToEnd { + ReadToEnd { + stream: Some(self), + size_limit, + buffer: Vec::new(), + } + } + + /// Close the receive stream immediately. + /// + /// The peer is notified and will cease transmitting on this stream, as if it had reset the + /// stream itself. Further data may still be received on this stream if it was already in + /// flight. Once called, a `ReadError::Reset` should be expected soon, although a peer might + /// manage to finish the stream before it receives the reset, and a misbehaving peer might + /// ignore the request entirely and continue sending until halted by flow control. + /// + /// Has no effect if the incoming stream already finished. + pub fn stop(&mut self, error_code: u16) { + let mut conn = self.conn.lock().unwrap(); + if self.is_0rtt && conn.check_0rtt().is_err() { + return; + } + conn.inner.stop_sending(self.stream, error_code); + conn.notify(); + self.recvd = true; + } +} + +/// Future produced by `read_to_end` +pub struct ReadToEnd { + stream: Option, + buffer: Vec, + size_limit: usize, +} + +impl Future for ReadToEnd { + type Item = (RecvStream, Box<[u8]>); + type Error = ReadError; + fn poll(&mut self) -> Poll { + loop { + match self.stream.as_mut().unwrap().poll_read_unordered() { + Ok(Async::Ready((data, offset))) => { + let len = self.buffer.len().max(offset as usize + data.len()); + if len > self.size_limit { + return Err(ReadError::Finished); + } + self.buffer.resize(len, 0); + self.buffer[offset as usize..offset as usize + data.len()] + .copy_from_slice(&data); + } + Ok(Async::NotReady) => { + return Ok(Async::NotReady); + } + Err(ReadError::Finished) => { + return Ok(Async::Ready(( + self.stream.take().unwrap(), + mem::replace(&mut self.buffer, Vec::new()).into(), + ))); + } + Err(e) => { + return Err(e); + } + } + } + } +} + +impl io::Read for RecvStream { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + match self.poll_read(buf) { + Ok(Async::Ready(n)) => Ok(n), + Err(ReadError::Finished) => Ok(0), + Ok(Async::NotReady) => Err(io::Error::new(io::ErrorKind::WouldBlock, "stream blocked")), + Err(e) => Err(e.into()), + } + } +} + +impl AsyncRead for RecvStream { + unsafe fn prepare_uninitialized_buffer(&self, _: &mut [u8]) -> bool { + false + } +} + +impl Drop for RecvStream { + fn drop(&mut self) { + let mut conn = self.conn.lock().unwrap(); + if conn.error.is_some() { + return; + } + if !self.recvd { + conn.inner.stop_sending(self.stream, 0); + conn.notify(); + } + } +} + +/// Errors that arise from finishing a stream +#[derive(Debug, Error, Clone)] +pub enum FinishError { + /// The connection was lost. + #[error(display = "connection lost: {}", _0)] + ConnectionLost(ConnectionError), + /// This was a 0-RTT stream and the server rejected it. + /// + /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). + #[error(display = "0-RTT rejected")] + ZeroRttRejected, +} + +impl From for io::Error { + fn from(x: FinishError) -> Self { + use self::FinishError::*; + match x { + ConnectionLost(e) => e.into(), + ZeroRttRejected => io::Error::new(io::ErrorKind::ConnectionReset, "0-RTT rejected"), + } + } +} + +/// Errors that arise from reading from a stream. +#[derive(Debug, Error, Clone)] +pub enum ReadError { + /// The peer abandoned transmitting data on this stream. + #[error(display = "stream reset by peer: error {}", error_code)] + Reset { + /// The error code supplied by the peer. + error_code: u16, + }, + /// The data on this stream has been fully delivered and no more will be transmitted. + #[error(display = "the stream has been completely received")] + Finished, + /// The connection was closed. + #[error(display = "connection closed: {}", _0)] + ConnectionClosed(ConnectionError), + /// Unknown stream + #[error(display = "unknown stream")] + UnknownStream, + /// This was a 0-RTT stream and the server rejected it. + /// + /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). + #[error(display = "0-RTT rejected")] + ZeroRttRejected, +} + +impl From for io::Error { + fn from(x: ReadError) -> Self { + use self::ReadError::*; + let kind = match x { + ConnectionClosed(e) => { + return e.into(); + } + Reset { .. } | ZeroRttRejected => io::ErrorKind::ConnectionReset, + Finished => io::ErrorKind::UnexpectedEof, + UnknownStream => io::ErrorKind::NotConnected, + }; + io::Error::new(kind, x) + } +} + +/// Errors that arise from writing to a stream +#[derive(Debug, Error, Clone)] +pub enum WriteError { + /// The peer is no longer accepting data on this stream. + #[error(display = "sending stopped by peer: error {}", error_code)] + Stopped { + /// The error code supplied by the peer. + error_code: u16, + }, + /// The connection was closed. + #[error(display = "connection closed: {}", _0)] + ConnectionClosed(ConnectionError), + /// Unknown stream + #[error(display = "unknown stream")] + UnknownStream, + /// This was a 0-RTT stream and the server rejected it. + /// + /// Can only occur on clients for 0-RTT streams (opened using `Connecting::into_0rtt()`). + #[error(display = "0-RTT rejected")] + ZeroRttRejected, +} + +impl From for io::Error { + fn from(x: WriteError) -> Self { + use self::WriteError::*; + let kind = match x { + ConnectionClosed(e) => { + return e.into(); + } + Stopped { .. } | ZeroRttRejected => io::ErrorKind::ConnectionReset, + UnknownStream => io::ErrorKind::NotConnected, + }; + io::Error::new(kind, x) + } +}