Improve naming of ConnectionHandle variable bindings

This commit is contained in:
Dirkjan Ochtman
2019-01-12 14:38:05 +01:00
parent e73fadaa6f
commit f95e89e3bd
3 changed files with 266 additions and 268 deletions
+92 -93
View File
@@ -86,11 +86,11 @@ impl Endpoint {
/// Get an application-facing event
pub fn poll(&mut self) -> Option<(ConnectionHandle, Event)> {
while let Some(&conn) = self.eventful_conns.iter().next() {
if let Some(e) = self.connections[conn.0].poll() {
return Some((conn, e));
while let Some(&ch) = self.eventful_conns.iter().next() {
if let Some(e) = self.connections[ch.0].poll() {
return Some((ch, e));
}
self.eventful_conns.remove(&conn);
self.eventful_conns.remove(&ch);
}
None
}
@@ -101,9 +101,9 @@ impl Endpoint {
return Some(x);
}
loop {
let &conn = self.dirty_conns.iter().next()?;
let &ch = self.dirty_conns.iter().next()?;
loop {
if let Some(io) = self.connections[conn.0].poll_io(now) {
if let Some(io) = self.connections[ch.0].poll_io(now) {
return Some(match io {
connection::Io::Transmit {
destination,
@@ -115,20 +115,20 @@ impl Endpoint {
packet,
},
connection::Io::TimerUpdate { timer, update } => Io::TimerUpdate {
connection: conn,
connection: ch,
timer,
update,
},
connection::Io::RetireConnectionId { connection_id } => {
self.connection_ids.remove(&connection_id);
let new_cid = self.new_cid();
self.connection_ids.insert(new_cid, conn);
self.connections[conn.0].issue_cid(new_cid);
self.connection_ids.insert(new_cid, ch);
self.connections[ch.0].issue_cid(new_cid);
continue;
}
});
} else {
self.dirty_conns.remove(&conn);
self.dirty_conns.remove(&ch);
break;
}
}
@@ -186,20 +186,20 @@ impl Endpoint {
}
/// Connection is either ready to accept data or failed.
fn conn_ready(&mut self, conn: ConnectionHandle) {
if self.connections[conn.0].side().is_server() {
fn conn_ready(&mut self, ch: ConnectionHandle) {
if self.connections[ch.0].side().is_server() {
self.incoming_handshakes -= 1;
self.incoming.push_back(conn);
self.incoming.push_back(ch);
}
if self.config.local_cid_len != 0 && !self.connections[conn.0].is_closed() {
if self.config.local_cid_len != 0 && !self.connections[ch.0].is_closed() {
/// Draft 17 §5.1.1: endpoints SHOULD provide and maintain at least eight
/// connection IDs
const LOCAL_CID_COUNT: usize = 8;
// We've already issued one CID as part of the normal handshake process.
for _ in 1..LOCAL_CID_COUNT {
let cid = self.new_cid();
self.connection_ids.insert(cid, conn);
self.connections[conn.0].issue_cid(cid);
self.connection_ids.insert(cid, ch);
self.connections[ch.0].issue_cid(cid);
}
}
}
@@ -217,13 +217,13 @@ impl Endpoint {
//
let dst_cid = partial_decode.dst_cid();
let conn = {
let conn = if self.config.local_cid_len > 0 {
let known_ch = {
let ch = if self.config.local_cid_len > 0 {
self.connection_ids.get(&dst_cid)
} else {
None
};
conn.or_else(|| self.connection_ids_initial.get(&dst_cid))
ch.or_else(|| self.connection_ids_initial.get(&dst_cid))
.or_else(|| {
// If CIDs are in use, only stateless resets (which use short headers) will
// legitimately have unknown CIDs.
@@ -235,17 +235,16 @@ impl Endpoint {
})
.cloned()
};
if let Some(conn_id) = conn {
let had_1rtt = self.connections[conn_id.0].has_1rtt();
self.connections[conn_id.0].handle_decode(now, remote, ecn, partial_decode);
if let Some(ch) = known_ch {
let had_1rtt = self.connections[ch.0].has_1rtt();
self.connections[ch.0].handle_decode(now, remote, ecn, partial_decode);
if !had_1rtt
&& (self.connections[conn_id.0].has_1rtt()
|| !self.connections[conn_id.0].is_handshaking())
&& (self.connections[ch.0].has_1rtt() || !self.connections[ch.0].is_handshaking())
{
self.conn_ready(conn_id);
self.conn_ready(ch);
}
self.dirty_conns.insert(conn_id);
self.eventful_conns.insert(conn_id);
self.dirty_conns.insert(ch);
self.eventful_conns.insert(ch);
return;
}
@@ -357,7 +356,7 @@ impl Endpoint {
) -> Result<ConnectionHandle, ConnectError> {
let remote_id = ConnectionId::random(&mut self.rng, MAX_CID_SIZE);
trace!(self.log, "initial dcid"; "value" => %remote_id);
let conn = self.add_connection(
let ch = self.add_connection(
remote_id,
remote_id,
remote,
@@ -366,8 +365,8 @@ impl Endpoint {
server_name: server_name.into(),
}),
)?;
self.dirty_conns.insert(conn);
Ok(conn)
self.dirty_conns.insert(ch);
Ok(ch)
}
fn new_cid(&mut self) -> ConnectionId {
@@ -416,7 +415,7 @@ impl Endpoint {
let remote_validated = self.server_config.as_ref().map_or(false, |cfg| {
cfg.use_stateless_retry && client_config.is_none()
});
let conn = self.connections.insert(Connection::new(
let id = self.connections.insert(Connection::new(
self.log.new(o!("connection" => local_id)),
Arc::clone(&self.config),
initial_id,
@@ -427,13 +426,13 @@ impl Endpoint {
tls,
remote_validated,
));
let conn = ConnectionHandle(conn);
let ch = ConnectionHandle(id);
if self.config.local_cid_len > 0 {
self.connection_ids.insert(local_id, conn);
self.connection_ids.insert(local_id, ch);
}
self.connection_remotes.insert(remote, conn);
Ok(conn)
self.connection_remotes.insert(remote, ch);
Ok(ch)
}
fn handle_initial(
@@ -560,7 +559,7 @@ impl Endpoint {
}
}
let conn = self
let ch = self
.add_connection(
dst_cid,
src_cid,
@@ -571,19 +570,19 @@ impl Endpoint {
)
.unwrap();
if dst_cid.len() != 0 {
self.connection_ids_initial.insert(dst_cid, conn);
self.connection_ids_initial.insert(dst_cid, ch);
}
match self.connections[conn.0].handle_initial(now, ecn, packet_number as u64, packet) {
match self.connections[ch.0].handle_initial(now, ecn, packet_number as u64, packet) {
Ok(()) => {
self.incoming_handshakes += 1;
self.dirty_conns.insert(conn);
if self.connections[conn.0].has_1rtt() {
self.conn_ready(conn);
self.dirty_conns.insert(ch);
if self.connections[ch.0].has_1rtt() {
self.conn_ready(ch);
}
}
Err(e) => {
debug!(self.log, "handshake failed"; "reason" => %e);
self.forget(conn);
self.forget(ch);
self.io.push_back(Io::Transmit {
destination: remote,
ecn: None,
@@ -593,33 +592,33 @@ impl Endpoint {
}
}
fn forget(&mut self, conn: ConnectionHandle) {
if self.connections[conn.0].side().is_server() {
fn forget(&mut self, ch: ConnectionHandle) {
if self.connections[ch.0].side().is_server() {
self.connection_ids_initial
.remove(&self.connections[conn.0].init_cid);
.remove(&self.connections[ch.0].init_cid);
}
if self.config.local_cid_len > 0 {
for cid in self.connections[conn.0].loc_cids() {
for cid in self.connections[ch.0].loc_cids() {
self.connection_ids.remove(cid);
}
}
self.connection_remotes
.remove(&self.connections[conn.0].remote());
self.dirty_conns.remove(&conn);
self.eventful_conns.remove(&conn);
self.connections.remove(conn.0);
.remove(&self.connections[ch.0].remote());
self.dirty_conns.remove(&ch);
self.eventful_conns.remove(&ch);
self.connections.remove(ch.0);
}
/// Handle a timer expiring
pub fn timeout(&mut self, now: u64, conn: ConnectionHandle, timer: Timer) {
if self.connections[conn.0].timeout(now, timer) {
self.forget(conn);
pub fn timeout(&mut self, now: u64, ch: ConnectionHandle, timer: Timer) {
if self.connections[ch.0].timeout(now, timer) {
self.forget(ch);
return;
}
if let Timer::Idle = timer {
self.eventful_conns.insert(conn);
self.eventful_conns.insert(ch);
}
self.dirty_conns.insert(conn);
self.dirty_conns.insert(ch);
}
/// Transmit data on a stream
@@ -630,12 +629,12 @@ impl Endpoint {
/// - when applied to a stream that does not have an active outgoing channel
pub fn write(
&mut self,
conn: ConnectionHandle,
ch: ConnectionHandle,
stream: StreamId,
data: &[u8],
) -> Result<usize, WriteError> {
let result = self.connections[conn.0].write(stream, data);
self.dirty_conns.insert(conn);
let result = self.connections[ch.0].write(stream, data);
self.dirty_conns.insert(ch);
result
}
@@ -646,9 +645,9 @@ impl Endpoint {
///
/// # Panics
/// - when applied to a stream that does not have an active outgoing channel
pub fn finish(&mut self, conn: ConnectionHandle, stream: StreamId) {
self.connections[conn.0].finish(stream);
self.dirty_conns.insert(conn);
pub fn finish(&mut self, ch: ConnectionHandle, stream: StreamId) {
self.connections[ch.0].finish(stream);
self.dirty_conns.insert(ch);
}
/// Read data from a stream
@@ -660,14 +659,14 @@ impl Endpoint {
/// - when applied to a stream that does not have an active incoming channel
pub fn read(
&mut self,
conn: ConnectionHandle,
ch: ConnectionHandle,
stream: StreamId,
buf: &mut [u8],
) -> Result<usize, ReadError> {
self.dirty_conns.insert(conn); // May need to send flow control frames after reading
match self.connections[conn.0].read(stream, buf) {
self.dirty_conns.insert(ch); // May need to send flow control frames after reading
match self.connections[ch.0].read(stream, buf) {
x @ Err(ReadError::Finished) | x @ Err(ReadError::Reset { .. }) => {
self.connections[conn.0].maybe_cleanup(stream);
self.connections[ch.0].maybe_cleanup(stream);
x
}
x => x,
@@ -688,13 +687,13 @@ impl Endpoint {
/// - when applied to a stream that does not have an active incoming channel
pub fn read_unordered(
&mut self,
conn: ConnectionHandle,
ch: ConnectionHandle,
stream: StreamId,
) -> Result<(Bytes, u64), ReadError> {
self.dirty_conns.insert(conn); // May need to send flow control frames after reading
match self.connections[conn.0].read_unordered(stream) {
self.dirty_conns.insert(ch); // May need to send flow control frames after reading
match self.connections[ch.0].read_unordered(stream) {
x @ Err(ReadError::Finished) | x @ Err(ReadError::Reset { .. }) => {
self.connections[conn.0].maybe_cleanup(stream);
self.connections[ch.0].maybe_cleanup(stream);
x
}
x => x,
@@ -705,67 +704,67 @@ impl Endpoint {
///
/// # Panics
/// - when applied to a receive stream or an unopened send stream
pub fn reset(&mut self, conn: ConnectionHandle, stream: StreamId, error_code: u16) {
self.connections[conn.0].reset(stream, error_code);
self.dirty_conns.insert(conn);
pub fn reset(&mut self, ch: ConnectionHandle, stream: StreamId, error_code: u16) {
self.connections[ch.0].reset(stream, error_code);
self.dirty_conns.insert(ch);
}
/// Instruct the peer to abandon transmitting data on a stream
///
/// # Panics
/// - when applied to a stream that has not begun receiving data
pub fn stop_sending(&mut self, conn: ConnectionHandle, stream: StreamId, error_code: u16) {
self.connections[conn.0].stop_sending(stream, error_code);
self.dirty_conns.insert(conn);
pub fn stop_sending(&mut self, ch: ConnectionHandle, stream: StreamId, error_code: u16) {
self.connections[ch.0].stop_sending(stream, error_code);
self.dirty_conns.insert(ch);
}
/// Create a new stream
///
/// Returns `None` if the maximum number of streams currently permitted by the remote endpoint
/// are already open.
pub fn open(&mut self, conn: ConnectionHandle, direction: Directionality) -> Option<StreamId> {
self.connections[conn.0].open(direction)
pub fn open(&mut self, ch: ConnectionHandle, direction: Directionality) -> Option<StreamId> {
self.connections[ch.0].open(direction)
}
/// Ping the remote endpoint
///
/// Useful for preventing an otherwise idle connection from timing out.
pub fn ping(&mut self, conn: ConnectionHandle) {
self.connections[conn.0].ping();
self.dirty_conns.insert(conn);
pub fn ping(&mut self, ch: ConnectionHandle) {
self.connections[ch.0].ping();
self.dirty_conns.insert(ch);
}
/// Close a connection immediately
///
/// This does not ensure delivery of outstanding data. It is the application's responsibility
/// to call this only when all important communications have been completed.
pub fn close(&mut self, now: u64, conn: ConnectionHandle, error_code: u16, reason: Bytes) {
if self.connections[conn.0].is_drained() {
self.forget(conn);
pub fn close(&mut self, now: u64, ch: ConnectionHandle, error_code: u16, reason: Bytes) {
if self.connections[ch.0].is_drained() {
self.forget(ch);
return;
}
self.connections[conn.0].close(now, error_code, reason);
self.dirty_conns.insert(conn);
self.connections[ch.0].close(now, error_code, reason);
self.dirty_conns.insert(ch);
}
pub fn accept(&mut self) -> Option<ConnectionHandle> {
self.incoming.pop_front()
}
pub fn accept_stream(&mut self, conn: ConnectionHandle) -> Option<StreamId> {
let id = self.connections[conn.0].accept()?;
self.dirty_conns.insert(conn);
pub fn accept_stream(&mut self, ch: ConnectionHandle) -> Option<StreamId> {
let id = self.connections[ch.0].accept()?;
self.dirty_conns.insert(ch);
Some(id)
}
#[doc(hidden)]
pub fn force_key_update(&mut self, conn: ConnectionHandle) {
self.connections[conn.0].force_key_update();
self.ping(conn);
pub fn force_key_update(&mut self, ch: ConnectionHandle) {
self.connections[ch.0].force_key_update();
self.ping(ch);
}
pub fn connection(&self, handle: ConnectionHandle) -> &Connection {
&self.connections[handle.0]
pub fn connection(&self, ch: ConnectionHandle) -> &Connection {
&self.connections[ch.0]
}
}
+116 -119
View File
@@ -216,18 +216,18 @@ impl Pair {
fn connect(&mut self) -> (ConnectionHandle, ConnectionHandle) {
info!(self.log, "connecting");
let client_conn = self
let client_ch = self
.client
.connect(self.server.addr, &client_config(), "localhost")
.unwrap();
self.drive();
let server_conn = if let Some(c) = self.server.accept() {
let server_ch = if let Some(c) = self.server.accept() {
c
} else {
panic!("server didn't connect");
};
assert_matches!(self.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_conn);
(client_conn, server_conn)
assert_matches!(self.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_ch);
(client_ch, server_ch)
}
}
@@ -407,14 +407,14 @@ fn version_negotiate() {
#[test]
fn lifecycle() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
assert_matches!(pair.client.poll(), None);
assert!(pair.client.connection(client_conn).using_ecn());
assert!(pair.server.connection(server_conn).using_ecn());
assert!(pair.client.connection(client_ch).using_ecn());
assert!(pair.server.connection(server_ch).using_ecn());
const REASON: &[u8] = b"whee";
info!(pair.log, "closing");
pair.client.close(pair.time, client_conn, 42, REASON.into());
pair.client.close(pair.time, client_ch, 42, REASON.into());
pair.drive();
assert!(pair.spins > 0);
assert_matches!(pair.server.poll(),
@@ -452,7 +452,7 @@ fn server_stateless_reset() {
};
let mut pair = Pair::new(server, Config::default(), server_config());
let (client_conn, _) = pair.connect();
let (client_ch, _) = pair.connect();
pair.server.endpoint = Endpoint::new(
pair.log.new(o!("side" => "Server")),
Config {
@@ -464,10 +464,10 @@ fn server_stateless_reset() {
.unwrap();
// Send something big enough to allow room for a smaller stateless reset.
pair.client
.close(pair.time, client_conn, 42, (&[0xab; 128][..]).into());
.close(pair.time, client_ch, 42, (&[0xab; 128][..]).into());
info!(pair.log, "resetting");
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::ConnectionLost { reason: ConnectionError::Reset })) if conn == client_conn);
assert_matches!(pair.client.poll(), Some((conn, Event::ConnectionLost { reason: ConnectionError::Reset })) if conn == client_ch);
}
#[test]
@@ -485,7 +485,7 @@ fn client_stateless_reset() {
};
let mut pair = Pair::new(Config::default(), client, server_config());
let (_, server_conn) = pair.connect();
let (_, server_ch) = pair.connect();
pair.client.endpoint = Endpoint::new(
pair.log.new(o!("side" => "Client")),
Config {
@@ -497,32 +497,32 @@ fn client_stateless_reset() {
.unwrap();
// Send something big enough to allow room for a smaller stateless reset.
pair.server
.close(pair.time, server_conn, 42, (&[0xab; 128][..]).into());
.close(pair.time, server_ch, 42, (&[0xab; 128][..]).into());
info!(pair.log, "resetting");
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::ConnectionLost { reason: ConnectionError::Reset })) if conn == server_conn);
assert_matches!(pair.server.poll(), Some((conn, Event::ConnectionLost { reason: ConnectionError::Reset })) if conn == server_ch);
}
#[test]
fn finish_stream() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair.client.open(client_conn, Directionality::Uni).unwrap();
let s = pair.client.open(client_ch, Directionality::Uni).unwrap();
const MSG: &[u8] = b"hello";
pair.client.write(client_conn, s, MSG).unwrap();
pair.client.finish(client_conn, s);
pair.client.write(client_ch, s, MSG).unwrap();
pair.client.finish(client_ch, s);
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::StreamFinished { stream })) if conn == client_conn && stream == s);
assert_matches!(pair.client.poll(), Some((conn, Event::StreamFinished { stream })) if conn == client_ch && stream == s);
assert_matches!(pair.client.poll(), None);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), None);
assert_matches!(pair.server.read_unordered(server_conn, s), Ok((ref data, 0)) if data == MSG);
assert_matches!(pair.server.read_unordered(server_ch, s), Ok((ref data, 0)) if data == MSG);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Err(ReadError::Finished)
);
}
@@ -530,24 +530,24 @@ fn finish_stream() {
#[test]
fn reset_stream() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair.client.open(client_conn, Directionality::Uni).unwrap();
let s = pair.client.open(client_ch, Directionality::Uni).unwrap();
const MSG: &[u8] = b"hello";
pair.client.write(client_conn, s, MSG).unwrap();
pair.client.write(client_ch, s, MSG).unwrap();
pair.drive();
info!(pair.log, "resetting stream");
const ERROR: u16 = 42;
pair.client.reset(client_conn, s, ERROR);
pair.client.reset(client_ch, s, ERROR);
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_conn && stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_ch && stream == s);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Err(ReadError::Reset { error_code: ERROR })
);
assert_matches!(pair.client.poll(), None);
@@ -556,28 +556,28 @@ fn reset_stream() {
#[test]
fn stop_stream() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair.client.open(client_conn, Directionality::Uni).unwrap();
let s = pair.client.open(client_ch, Directionality::Uni).unwrap();
const MSG: &[u8] = b"hello";
pair.client.write(client_conn, s, MSG).unwrap();
pair.client.write(client_ch, s, MSG).unwrap();
pair.drive();
info!(pair.log, "stopping stream");
const ERROR: u16 = 42;
pair.server.stop_sending(server_conn, s, ERROR);
pair.server.stop_sending(server_ch, s, ERROR);
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_conn && stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_ch && stream == s);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Err(ReadError::Reset { error_code: ERROR })
);
assert_matches!(
pair.client.write(client_conn, s, b"foo"),
pair.client.write(client_ch, s, b"foo"),
Err(WriteError::Stopped { error_code: ERROR })
);
}
@@ -590,7 +590,7 @@ fn reject_self_signed_cert() {
let mut pair = Pair::default();
info!(pair.log, "connecting");
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &Arc::new(client_config), "localhost")
.unwrap();
@@ -598,18 +598,18 @@ fn reject_self_signed_cert() {
assert_matches!(pair.client.poll(),
Some((conn, Event::ConnectionLost { reason: ConnectionError::TransportError {
error_code
}})) if conn == client_conn && error_code == TransportError::crypto(AlertDescription::BadCertificate));
}})) if conn == client_ch && error_code == TransportError::crypto(AlertDescription::BadCertificate));
}
#[test]
fn congestion() {
let mut pair = Pair::default();
let (client_conn, _) = pair.connect();
let (client_ch, _) = pair.connect();
let initial_congestion_state = pair.client.connection(client_conn).congestion_state();
let s = pair.client.open(client_conn, Directionality::Uni).unwrap();
let initial_congestion_state = pair.client.connection(client_ch).congestion_state();
let s = pair.client.open(client_ch, Directionality::Uni).unwrap();
loop {
match pair.client.write(client_conn, s, &[42; 1024]) {
match pair.client.write(client_ch, s, &[42; 1024]) {
Ok(n) => {
assert!(n <= 1024);
pair.drive_client();
@@ -623,19 +623,19 @@ fn congestion() {
}
}
pair.drive();
assert!(pair.client.connection(client_conn).congestion_state() >= initial_congestion_state);
pair.client.write(client_conn, s, &[42; 1024]).unwrap();
assert!(pair.client.connection(client_ch).congestion_state() >= initial_congestion_state);
pair.client.write(client_ch, s, &[42; 1024]).unwrap();
}
#[test]
fn high_latency_handshake() {
let mut pair = Pair::default();
pair.latency = 200 * 1000;
let (client_conn, server_conn) = pair.connect();
assert_eq!(pair.client.connection(client_conn).bytes_in_flight(), 0);
assert_eq!(pair.server.connection(server_conn).bytes_in_flight(), 0);
assert!(pair.client.connection(client_conn).using_ecn());
assert!(pair.server.connection(server_conn).using_ecn());
let (client_ch, server_ch) = pair.connect();
assert_eq!(pair.client.connection(client_ch).bytes_in_flight(), 0);
assert_eq!(pair.server.connection(server_ch).bytes_in_flight(), 0);
assert!(pair.client.connection(client_ch).using_ecn());
assert!(pair.server.connection(server_ch).using_ecn());
}
#[test]
@@ -644,13 +644,13 @@ fn zero_rtt() {
let config = client_config();
// Establish normal connection
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &config, "localhost")
.unwrap();
pair.drive();
pair.server.accept().unwrap();
pair.client.close(pair.time, client_conn, 0, [][..].into());
pair.client.close(pair.time, client_ch, 0, [][..].into());
pair.drive();
pair.client.addr = SocketAddr::new(
@@ -658,23 +658,23 @@ fn zero_rtt() {
CLIENT_PORTS.lock().unwrap().next().unwrap(),
);
info!(pair.log, "resuming session");
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &config, "localhost")
.unwrap();
assert!(pair.client.connection(client_conn).has_0rtt());
let s = pair.client.open(client_conn, Directionality::Uni).unwrap();
assert!(pair.client.connection(client_ch).has_0rtt());
let s = pair.client.open(client_ch, Directionality::Uni).unwrap();
const MSG: &[u8] = b"Hello, 0-RTT!";
pair.client.write(client_conn, s, MSG).unwrap();
pair.client.write(client_ch, s, MSG).unwrap();
pair.drive();
assert!(pair.client.connection(client_conn).accepted_0rtt());
let server_conn = if let Some(c) = pair.server.accept() {
assert!(pair.client.connection(client_ch).accepted_0rtt());
let server_ch = if let Some(c) = pair.server.accept() {
c
} else {
panic!("server didn't connect");
};
assert_matches!(pair.server.read_unordered(server_conn, s), Ok((ref data, 0)) if data == MSG);
assert_eq!(pair.client.connection(client_conn).lost_packets(), 0);
assert_matches!(pair.server.read_unordered(server_ch, s), Ok((ref data, 0)) if data == MSG);
assert_eq!(pair.client.connection(client_ch).lost_packets(), 0);
}
#[test]
@@ -737,46 +737,46 @@ fn stream_id_backpressure() {
..Config::default()
};
let mut pair = Pair::new(server, Default::default(), server_config());
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair
.client
.open(client_conn, Directionality::Uni)
.open(client_ch, Directionality::Uni)
.expect("couldn't open first stream");
assert_eq!(
pair.client.open(client_conn, Directionality::Uni),
pair.client.open(client_ch, Directionality::Uni),
None,
"only one stream is permitted at a time"
);
// Close the first stream to make room for the second
pair.client.finish(client_conn, s);
pair.client.finish(client_ch, s);
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::StreamFinished { stream })) if conn == client_conn && stream == s);
assert_matches!(pair.client.poll(), Some((conn, Event::StreamFinished { stream })) if conn == client_ch && stream == s);
assert_matches!(pair.client.poll(), None);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Err(ReadError::Finished)
);
// Server will only send MAX_STREAM_ID now that the application's been notified
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::StreamAvailable { directionality: Directionality::Uni })) if conn == client_conn);
assert_matches!(pair.client.poll(), Some((conn, Event::StreamAvailable { directionality: Directionality::Uni })) if conn == client_ch);
assert_matches!(pair.client.poll(), None);
// Try opening the second stream again, now that we've made room
let s = pair
.client
.open(client_conn, Directionality::Uni)
.open(client_ch, Directionality::Uni)
.expect("didn't get stream id budget");
pair.client.finish(client_conn, s);
pair.client.finish(client_ch, s);
pair.drive();
// Make sure the server actually processes data on the newly-available stream
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), None);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Err(ReadError::Finished)
);
}
@@ -784,34 +784,34 @@ fn stream_id_backpressure() {
#[test]
fn key_update() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair
.client
.open(client_conn, Directionality::Bi)
.open(client_ch, Directionality::Bi)
.expect("couldn't open first stream");
const MSG1: &[u8] = b"hello1";
pair.client.write(client_conn, s, MSG1).unwrap();
pair.client.write(client_ch, s, MSG1).unwrap();
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), None);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Ok((ref data, 0)) if data == MSG1
);
pair.client.connections[client_conn.0].force_key_update();
pair.client.connections[client_ch.0].force_key_update();
const MSG2: &[u8] = b"hello2";
pair.client.write(client_conn, s, MSG2).unwrap();
pair.client.write(client_ch, s, MSG2).unwrap();
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_conn && stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_ch && stream == s);
assert_matches!(pair.server.poll(), None);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Ok((ref data, 6)) if data == MSG2
);
}
@@ -819,65 +819,65 @@ fn key_update() {
#[test]
fn key_update_reordered() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
let s = pair
.client
.open(client_conn, Directionality::Bi)
.open(client_ch, Directionality::Bi)
.expect("couldn't open first stream");
const MSG1: &[u8] = b"1";
pair.client.write(client_conn, s, MSG1).unwrap();
pair.client.write(client_ch, s, MSG1).unwrap();
pair.client.drive(&pair.log, pair.time, pair.server.addr);
assert!(!pair.client.outbound.is_empty());
pair.client.delay_outbound();
pair.client.connections[client_conn.0].force_key_update();
pair.client.connections[client_ch.0].force_key_update();
info!(pair.log, "updated keys");
const MSG2: &[u8] = b"two";
pair.client.write(client_conn, s, MSG2).unwrap();
pair.client.write(client_ch, s, MSG2).unwrap();
pair.client.drive(&pair.log, pair.time, pair.server.addr);
pair.client.finish_delay();
pair.drive();
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_conn);
assert_matches!(pair.server.accept_stream(server_conn), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_conn && stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamOpened)) if conn == server_ch);
assert_matches!(pair.server.accept_stream(server_ch), Some(stream) if stream == s);
assert_matches!(pair.server.poll(), Some((conn, Event::StreamReadable { stream })) if conn == server_ch && stream == s);
assert_matches!(pair.server.poll(), None);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Ok((ref data, 1)) if data == MSG2
);
assert_matches!(
pair.server.read_unordered(server_conn, s),
pair.server.read_unordered(server_ch, s),
Ok((ref data, 0)) if data == MSG1
);
assert_eq!(pair.client.connection(client_conn).lost_packets(), 0);
assert_eq!(pair.client.connection(client_ch).lost_packets(), 0);
}
#[test]
fn initial_retransmit() {
let mut pair = Pair::default();
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &client_config(), "localhost")
.unwrap();
pair.client.drive(&pair.log, pair.time, pair.server.addr);
pair.client.outbound.clear(); // Drop initial
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_conn);
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_ch);
}
#[test]
fn instant_close() {
let mut pair = Pair::default();
info!(pair.log, "connecting");
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &client_config(), "localhost")
.unwrap();
pair.client.close(pair.time, client_conn, 0, Bytes::new());
pair.client.close(pair.time, client_ch, 0, Bytes::new());
pair.drive();
assert_matches!(pair.client.poll(), None);
assert_matches!(pair.server.poll(), None);
@@ -887,13 +887,13 @@ fn instant_close() {
fn instant_close_2() {
let mut pair = Pair::default();
info!(pair.log, "connecting");
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &client_config(), "localhost")
.unwrap();
// Unlike `instant_close`, the server sees a valid Initial packet first.
pair.drive_client();
pair.client.close(pair.time, client_conn, 42, Bytes::new());
pair.client.close(pair.time, client_ch, 42, Bytes::new());
pair.drive();
assert_matches!(pair.client.poll(), None);
assert_matches!(pair.server.poll(), Some((_, Event::ConnectionLost { reason: ConnectionError::ApplicationClosed {
@@ -904,10 +904,10 @@ fn instant_close_2() {
#[test]
fn idle_timeout() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
pair.client.ping(client_conn);
while !pair.client.connection(client_conn).is_closed()
|| !pair.server.connection(server_conn).is_closed()
let (client_ch, server_ch) = pair.connect();
pair.client.ping(client_ch);
while !pair.client.connection(client_ch).is_closed()
|| !pair.server.connection(server_ch).is_closed()
{
pair.step();
pair.client.inbound.clear();
@@ -970,7 +970,7 @@ fn server_busy() {
#[test]
fn server_hs_retransmit() {
let mut pair = Pair::default();
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &client_config(), "localhost")
.unwrap();
@@ -993,7 +993,7 @@ fn server_hs_retransmit() {
pair.client.inbound.drain(..);
}
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_conn);
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_ch);
}
#[test]
@@ -1001,7 +1001,7 @@ fn decode_coalesced() {
// We can't currently generate coalesced packets natively, but we must support decoding
// them. Hack around the problem by manually concatenating the server's first flight.
let mut pair = Pair::default();
let client_conn = pair
let client_ch = pair
.client
.connect(pair.server.addr, &client_config(), "localhost")
.unwrap();
@@ -1018,25 +1018,22 @@ fn decode_coalesced() {
.inbound
.push_back((pair.time, Some(EcnCodepoint::ECT0), coalesced.into()));
pair.drive();
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_conn);
assert_eq!(pair.client.connection(client_conn).lost_packets(), 0);
assert_matches!(pair.client.poll(), Some((conn, Event::Connected { .. })) if conn == client_ch);
assert_eq!(pair.client.connection(client_ch).lost_packets(), 0);
}
#[test]
fn migration() {
let mut pair = Pair::default();
let (client_conn, server_conn) = pair.connect();
let (client_ch, server_ch) = pair.connect();
pair.client.addr = SocketAddr::new(
Ipv4Addr::new(127, 0, 0, 1).into(),
CLIENT_PORTS.lock().unwrap().next().unwrap(),
);
pair.client.ping(client_conn);
pair.client.ping(client_ch);
pair.drive();
assert_matches!(pair.client.poll(), None);
assert_eq!(
pair.server.connection(server_conn).remote(),
pair.client.addr
);
assert_eq!(pair.server.connection(server_ch).remote(), pair.client.addr);
}
fn test_flow_control(config: Config, window_size: usize) {
+58 -56
View File
@@ -291,7 +291,7 @@ impl Endpoint {
};
let conn = ConnectionInner {
endpoint: self.inner.clone(),
conn: handle,
handle,
side: Side::Client,
};
Ok((recv, conn))
@@ -310,7 +310,7 @@ impl NewConnection {
fn new(endpoint: Rc<RefCell<EndpointInner>>, handle: quinn::ConnectionHandle) -> Self {
let conn = Rc::new(ConnectionInner {
endpoint,
conn: handle,
handle,
side: Side::Server,
});
NewConnection {
@@ -365,13 +365,13 @@ impl Future for Driver {
}
}
}
while let Some((connection, event)) = endpoint.inner.poll() {
while let Some((ch, event)) = endpoint.inner.poll() {
use crate::quinn::Event::*;
match event {
Connected { .. } => {
let _ = endpoint
.pending
.get_mut(&connection)
.get_mut(&ch)
.unwrap()
.connecting
.take()
@@ -381,14 +381,14 @@ impl Future for Driver {
ConnectionLost { reason } => {
// HACK HACK HACK: Handshake currently emits ConnectionLost, which means we might not know about
// this connection yet. This should probably be made more consistent.
if let Some(x) = endpoint.pending.get_mut(&connection) {
if let Some(x) = endpoint.pending.get_mut(&ch) {
x.fail(reason);
}
}
StreamWritable { stream } => {
if let Some(writer) = endpoint
.pending
.get_mut(&connection)
.get_mut(&ch)
.unwrap()
.blocked_writers
.remove(&stream)
@@ -397,28 +397,28 @@ impl Future for Driver {
}
}
StreamOpened => {
let pending = endpoint.pending.get_mut(&connection).unwrap();
let pending = endpoint.pending.get_mut(&ch).unwrap();
if let Some(x) = pending.incoming_streams_reader.take() {
x.notify();
}
}
StreamReadable { stream } => {
let pending = endpoint.pending.get_mut(&connection).unwrap();
let pending = endpoint.pending.get_mut(&ch).unwrap();
if let Some(reader) = pending.blocked_readers.remove(&stream) {
reader.notify();
}
}
StreamAvailable { directionality } => {
let pending = endpoint.pending.get_mut(&connection).unwrap();
let pending = endpoint.pending.get_mut(&ch).unwrap();
let queue = match directionality {
Directionality::Uni => &mut pending.uni_opening,
Directionality::Bi => &mut pending.bi_opening,
};
while let Some(ch) = queue.pop_front() {
if let Some(id) = endpoint.inner.open(connection, directionality) {
let _ = ch.send(Ok(id));
while let Some(connection) = queue.pop_front() {
if let Some(id) = endpoint.inner.open(ch, directionality) {
let _ = connection.send(Ok(id));
} else {
queue.push_front(ch);
queue.push_front(connection);
break;
}
}
@@ -426,7 +426,7 @@ impl Future for Driver {
StreamFinished { stream } => {
let _ = endpoint
.pending
.get_mut(&connection)
.get_mut(&ch)
.unwrap()
.finishing
.remove(&stream)
@@ -483,27 +483,27 @@ impl Future for Driver {
}
}
TimerUpdate {
connection,
connection: ch,
timer: timer @ quinn::Timer::Close,
update: quinn::TimerUpdate::Start(time),
} => {
let instant = endpoint.epoch + duration_micros(time);
endpoint.timers.push(Timer {
conn: connection,
ch,
ty: timer,
delay: Delay::new(instant),
cancel: None,
});
}
TimerUpdate {
connection,
connection: ch,
timer,
update: quinn::TimerUpdate::Start(time),
} => {
// Loss detection and idle timers start before the connection is established
let pending = endpoint
.pending
.entry(connection)
.entry(ch)
.or_insert_with(|| Pending::new(None));
let cancel = &mut pending.cancel_timers[timer as usize];
let instant = endpoint.epoch + duration_micros(time);
@@ -514,20 +514,20 @@ impl Future for Driver {
*cancel = Some(send);
trace!(endpoint.log, "timer start"; "timer" => ?timer, "time" => ?duration_micros(time));
endpoint.timers.push(Timer {
conn: connection,
ch,
ty: timer,
delay: Delay::new(instant),
cancel: Some(recv),
});
}
TimerUpdate {
connection,
connection: ch,
timer,
update: quinn::TimerUpdate::Stop,
} => {
trace!(endpoint.log, "timer stop"; "timer" => ?timer);
// If a connection was lost, we already canceled its loss/idle timers.
if let Some(pending) = endpoint.pending.get_mut(&connection) {
if let Some(pending) = endpoint.pending.get_mut(&ch) {
if let Some(x) = pending.cancel_timers[timer as usize].take() {
let _ = x.send(());
}
@@ -552,12 +552,12 @@ impl Future for Driver {
let mut fired = false;
loop {
match endpoint.timers.poll() {
Ok(Async::Ready(Some(Some((conn, timer))))) => {
Ok(Async::Ready(Some(Some((ch, timer))))) => {
trace!(endpoint.log, "timeout"; "timer" => ?timer);
endpoint.inner.timeout(now, conn, timer);
endpoint.inner.timeout(now, ch, timer);
if timer == quinn::Timer::Close {
// Connection drained
if let Some(x) = endpoint.pending.get_mut(&conn).and_then(|p| {
if let Some(x) = endpoint.pending.get_mut(&ch).and_then(|p| {
p.drained = true;
p.draining.take()
}) {
@@ -584,8 +584,8 @@ impl Future for Driver {
impl Drop for Driver {
fn drop(&mut self) {
let mut endpoint = self.0.borrow_mut();
for connection in endpoint.pending.values_mut() {
connection.fail(ConnectionError::TransportError {
for ch in endpoint.pending.values_mut() {
ch.fail(ConnectionError::TransportError {
error_code: quinn::TransportError::INTERNAL_ERROR,
});
}
@@ -601,7 +601,7 @@ fn micros_from(x: Duration) -> u64 {
struct ConnectionInner {
endpoint: Rc<RefCell<EndpointInner>>,
conn: ConnectionHandle,
handle: ConnectionHandle,
side: Side,
}
@@ -620,10 +620,10 @@ impl Connection {
let (send, recv) = oneshot::channel();
{
let mut endpoint = self.0.endpoint.borrow_mut();
if let Some(x) = endpoint.inner.open(self.0.conn, Directionality::Uni) {
if let Some(x) = endpoint.inner.open(self.0.handle, Directionality::Uni) {
let _ = send.send(Ok(x));
} else {
let pending = endpoint.pending.get_mut(&self.0.conn).unwrap();
let pending = endpoint.pending.get_mut(&self.0.handle).unwrap();
pending.uni_opening.push_back(send);
// We don't notify the driver here because there's no way to ask the peer for more streams
}
@@ -639,10 +639,10 @@ impl Connection {
let (send, recv) = oneshot::channel();
{
let mut endpoint = self.0.endpoint.borrow_mut();
if let Some(x) = endpoint.inner.open(self.0.conn, Directionality::Bi) {
if let Some(x) = endpoint.inner.open(self.0.handle, Directionality::Bi) {
let _ = send.send(Ok(x));
} else {
let pending = endpoint.pending.get_mut(&self.0.conn).unwrap();
let pending = endpoint.pending.get_mut(&self.0.handle).unwrap();
pending.bi_opening.push_back(send);
// We don't notify the driver here because there's no way to ask the peer for more streams
}
@@ -671,7 +671,7 @@ impl Connection {
{
let endpoint = &mut *self.0.endpoint.borrow_mut();
let pending = endpoint.pending.get_mut(&self.0.conn).unwrap();
let pending = endpoint.pending.get_mut(&self.0.handle).unwrap();
assert!(
pending.draining.is_none(),
"a connection can only be closed once"
@@ -680,7 +680,7 @@ impl Connection {
endpoint.inner.close(
micros_from(endpoint.epoch.elapsed()),
self.0.conn,
self.0.handle,
error_code,
reason.into(),
);
@@ -699,7 +699,7 @@ impl Connection {
.endpoint
.borrow()
.inner
.connection(self.0.conn)
.connection(self.0.handle)
.remote()
}
@@ -709,7 +709,7 @@ impl Connection {
.endpoint
.borrow()
.inner
.connection(self.0.conn)
.connection(self.0.handle)
.loc_cids()
.cloned()
.collect::<Vec<_>>()
@@ -721,7 +721,7 @@ impl Connection {
.endpoint
.borrow()
.inner
.connection(self.0.conn)
.connection(self.0.handle)
.rem_cid()
}
@@ -731,7 +731,7 @@ impl Connection {
.endpoint
.borrow()
.inner
.connection(self.0.conn)
.connection(self.0.handle)
.protocol()
.map(|x| x.into())
}
@@ -743,18 +743,18 @@ impl Connection {
.endpoint
.borrow_mut()
.inner
.force_key_update(self.0.conn)
.force_key_update(self.0.handle)
}
}
impl Drop for ConnectionInner {
fn drop(&mut self) {
let endpoint = &mut *self.endpoint.borrow_mut();
if let hash_map::Entry::Occupied(pending) = endpoint.pending.entry(self.conn) {
if let hash_map::Entry::Occupied(pending) = endpoint.pending.entry(self.handle) {
if pending.get().draining.is_none() && !pending.get().drained {
endpoint.inner.close(
micros_from(endpoint.epoch.elapsed()),
self.conn,
self.handle,
0,
(&[][..]).into(),
);
@@ -856,10 +856,10 @@ impl Write for BiStream {
fn poll_write(&mut self, buf: &[u8]) -> Poll<usize, WriteError> {
let mut endpoint = self.conn.endpoint.borrow_mut();
use crate::quinn::WriteError::*;
let n = match endpoint.inner.write(self.conn.conn, self.stream, buf) {
let n = match endpoint.inner.write(self.conn.handle, self.stream, buf) {
Ok(n) => n,
Err(Blocked) => {
let pending = endpoint.pending.get_mut(&self.conn.conn).unwrap();
let pending = endpoint.pending.get_mut(&self.conn.handle).unwrap();
if let Some(ref x) = pending.error {
return Err(WriteError::ConnectionClosed(x.clone()));
}
@@ -877,12 +877,12 @@ impl Write for BiStream {
fn poll_finish(&mut self) -> Poll<(), ConnectionError> {
let mut endpoint = self.conn.endpoint.borrow_mut();
if self.finishing.is_none() {
endpoint.inner.finish(self.conn.conn, self.stream);
endpoint.inner.finish(self.conn.handle, self.stream);
let (send, recv) = oneshot::channel();
self.finishing = Some(recv);
endpoint
.pending
.get_mut(&self.conn.conn)
.get_mut(&self.conn.handle)
.unwrap()
.finishing
.insert(self.stream, send);
@@ -903,7 +903,7 @@ impl Write for BiStream {
let endpoint = &mut *self.conn.endpoint.borrow_mut();
endpoint
.inner
.reset(self.conn.conn, self.stream, error_code);
.reset(self.conn.handle, self.stream, error_code);
endpoint.notify();
}
}
@@ -912,8 +912,8 @@ impl Read for BiStream {
fn poll_read_unordered(&mut self) -> Poll<(Bytes, u64), ReadError> {
let endpoint = &mut *self.conn.endpoint.borrow_mut();
use crate::quinn::ReadError::*;
let pending = endpoint.pending.get_mut(&self.conn.conn).unwrap();
match endpoint.inner.read_unordered(self.conn.conn, self.stream) {
let pending = endpoint.pending.get_mut(&self.conn.handle).unwrap();
match endpoint.inner.read_unordered(self.conn.handle, self.stream) {
Ok((bytes, offset)) => Ok(Async::Ready((bytes, offset))),
Err(Blocked) => {
if let Some(ref x) = pending.error {
@@ -936,8 +936,8 @@ impl Read for BiStream {
fn poll_read(&mut self, buf: &mut [u8]) -> Poll<usize, ReadError> {
let endpoint = &mut *self.conn.endpoint.borrow_mut();
use crate::quinn::ReadError::*;
let pending = endpoint.pending.get_mut(&self.conn.conn).unwrap();
match endpoint.inner.read(self.conn.conn, self.stream, buf) {
let pending = endpoint.pending.get_mut(&self.conn.handle).unwrap();
match endpoint.inner.read(self.conn.handle, self.stream, buf) {
Ok(n) => Ok(Async::Ready(n)),
Err(Blocked) => {
if let Some(ref x) = pending.error {
@@ -961,7 +961,7 @@ impl Read for BiStream {
let endpoint = &mut *self.conn.endpoint.borrow_mut();
endpoint
.inner
.stop_sending(self.conn.conn, self.stream, error_code);
.stop_sending(self.conn.handle, self.stream, error_code);
endpoint.notify();
self.recvd = true;
}
@@ -1008,10 +1008,12 @@ impl Drop for BiStream {
Directionality::Uni => (ours, !ours),
};
if send && !self.finished {
endpoint.inner.reset(self.conn.conn, self.stream, 0);
endpoint.inner.reset(self.conn.handle, self.stream, 0);
}
if recv && !self.recvd {
endpoint.inner.stop_sending(self.conn.conn, self.stream, 0);
endpoint
.inner
.stop_sending(self.conn.handle, self.stream, 0);
}
endpoint.notify();
}
@@ -1131,7 +1133,7 @@ pub enum ReadError {
}
struct Timer {
conn: ConnectionHandle,
ch: ConnectionHandle,
ty: quinn::Timer,
delay: Delay,
cancel: Option<oneshot::Receiver<()>>,
@@ -1150,7 +1152,7 @@ impl Future for Timer {
match self.delay.poll() {
Err(e) => panic!("unexpected timer error: {}", e),
Ok(Async::NotReady) => Ok(Async::NotReady),
Ok(Async::Ready(())) => Ok(Async::Ready(Some((self.conn, self.ty)))),
Ok(Async::Ready(())) => Ok(Async::Ready(Some((self.ch, self.ty)))),
}
}
}
@@ -1171,7 +1173,7 @@ impl FuturesStream for IncomingStreams {
type Error = ConnectionError;
fn poll(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
let mut endpoint = self.0.endpoint.borrow_mut();
if let Some(x) = endpoint.inner.accept_stream(self.0.conn) {
if let Some(x) = endpoint.inner.accept_stream(self.0.handle) {
let stream = BiStream::new(self.0.clone(), x);
let stream = if x.directionality() == Directionality::Uni {
NewStream::Uni(RecvStream(stream))
@@ -1180,7 +1182,7 @@ impl FuturesStream for IncomingStreams {
};
return Ok(Async::Ready(Some(stream)));
}
let pending = endpoint.pending.get_mut(&self.0.conn).unwrap();
let pending = endpoint.pending.get_mut(&self.0.handle).unwrap();
if let Some(ref x) = pending.error {
Err(x.clone())
} else {