mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-27 00:00:12 +02:00
chore: enable full CI run (#1399)
* update all tests * fix doc examples * misc API tweaks
This commit is contained in:
@@ -24,11 +24,11 @@ publish = false
|
||||
[dependencies]
|
||||
tokio-io = { version = "0.2.0", path = "../tokio-io" }
|
||||
bytes = "0.4.7"
|
||||
futures-core-preview = "0.3.0-alpha.17"
|
||||
futures-sink-preview = "0.3.0-alpha.17"
|
||||
futures-core-preview = "= 0.3.0-alpha.17"
|
||||
futures-sink-preview = "= 0.3.0-alpha.17"
|
||||
log = "0.4"
|
||||
|
||||
[dev-dependencies]
|
||||
futures-preview = "0.3.0-alpha.17"
|
||||
tokio-current-thread = { version = "0.2.0", path = "../tokio-current-thread" }
|
||||
futures-util-preview = "= 0.3.0-alpha.17"
|
||||
tokio = { version = "0.2.0", path = "../tokio" }
|
||||
tokio-test = { version = "0.2.0", path = "../tokio-test" }
|
||||
|
||||
@@ -41,9 +41,8 @@
|
||||
//! ```
|
||||
//! #![feature(async_await)]
|
||||
//!
|
||||
//! use tokio_io::{AsyncRead, AsyncWrite};
|
||||
//! use tokio_codec::{Framed, LengthDelimitedCodec};
|
||||
//! use futures::SinkExt;
|
||||
//! use tokio::codec::{Framed, LengthDelimitedCodec};
|
||||
//! use tokio::prelude::*;
|
||||
//!
|
||||
//! use bytes::Bytes;
|
||||
//!
|
||||
|
||||
+10
-18
@@ -1,17 +1,15 @@
|
||||
#![feature(async_await)]
|
||||
#![deny(warnings, rust_2018_idioms)]
|
||||
|
||||
use tokio::prelude::*;
|
||||
use tokio_codec::{Decoder, Encoder, Framed, FramedParts};
|
||||
use tokio_test::assert_ok;
|
||||
|
||||
use bytes::{Buf, BufMut, BytesMut, IntoBuf};
|
||||
use std::io::{self, Read};
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
use tokio_codec::{Decoder, Encoder, Framed, FramedParts};
|
||||
use tokio_current_thread::block_on_all;
|
||||
use tokio_io::AsyncRead;
|
||||
|
||||
use bytes::{Buf, BufMut, BytesMut, IntoBuf};
|
||||
use futures::future::FutureExt;
|
||||
use futures::stream::StreamExt;
|
||||
|
||||
const INITIAL_CAPACITY: usize = 8 * 1024;
|
||||
|
||||
/// Encode and decode u32 values.
|
||||
@@ -65,19 +63,13 @@ impl AsyncRead for DontReadIntoThis {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn can_read_from_existing_buf() {
|
||||
#[tokio::test]
|
||||
async fn can_read_from_existing_buf() {
|
||||
let mut parts = FramedParts::new(DontReadIntoThis, U32Codec);
|
||||
parts.read_buf = vec![0, 0, 0, 42].into();
|
||||
|
||||
let framed = Framed::from_parts(parts);
|
||||
|
||||
let num = block_on_all(
|
||||
framed
|
||||
.into_future()
|
||||
.map(|(first_num, _)| first_num.unwrap()),
|
||||
)
|
||||
.unwrap();
|
||||
let mut framed = Framed::from_parts(parts);
|
||||
let num = assert_ok!(framed.next().await.unwrap());
|
||||
|
||||
assert_eq!(num, 42);
|
||||
}
|
||||
|
||||
@@ -1,19 +1,18 @@
|
||||
#![feature(async_await)]
|
||||
#![deny(warnings, rust_2018_idioms)]
|
||||
|
||||
use tokio::prelude::*;
|
||||
use tokio_codec::{Decoder, FramedRead};
|
||||
use tokio_test::assert_ready;
|
||||
use tokio_test::task::MockTask;
|
||||
|
||||
use bytes::{Buf, BytesMut, IntoBuf};
|
||||
use std::collections::VecDeque;
|
||||
use std::io::{self, Read};
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::task::Poll::{Pending, Ready};
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
use bytes::{Buf, BytesMut, IntoBuf};
|
||||
use futures::Stream;
|
||||
|
||||
use tokio_codec::{Decoder, FramedRead};
|
||||
use tokio_io::AsyncRead;
|
||||
use tokio_test::assert_ready;
|
||||
use tokio_test::task::MockTask;
|
||||
|
||||
macro_rules! mock {
|
||||
($($x:expr,)*) => {{
|
||||
let mut v = VecDeque::new();
|
||||
@@ -261,29 +260,23 @@ struct Mock {
|
||||
calls: VecDeque<io::Result<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl Read for Mock {
|
||||
fn read(&mut self, dst: &mut [u8]) -> io::Result<usize> {
|
||||
match self.calls.pop_front() {
|
||||
Some(Ok(data)) => {
|
||||
debug_assert!(dst.len() >= data.len());
|
||||
dst[..data.len()].copy_from_slice(&data[..]);
|
||||
Ok(data.len())
|
||||
}
|
||||
Some(Err(e)) => Err(e),
|
||||
None => Ok(0),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncRead for Mock {
|
||||
fn poll_read(
|
||||
self: Pin<&mut Self>,
|
||||
mut self: Pin<&mut Self>,
|
||||
_cx: &mut Context<'_>,
|
||||
buf: &mut [u8],
|
||||
) -> Poll<io::Result<usize>> {
|
||||
match Pin::get_mut(self).read(buf) {
|
||||
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Pending,
|
||||
other => Ready(other),
|
||||
use io::ErrorKind::WouldBlock;
|
||||
|
||||
match self.calls.pop_front() {
|
||||
Some(Ok(data)) => {
|
||||
debug_assert!(buf.len() >= data.len());
|
||||
buf[..data.len()].copy_from_slice(&data[..]);
|
||||
Ready(Ok(data.len()))
|
||||
}
|
||||
Some(Err(ref e)) if e.kind() == WouldBlock => Pending,
|
||||
Some(Err(e)) => Ready(Err(e)),
|
||||
None => Ready(Ok(0)),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -293,10 +286,10 @@ struct Slice<'a>(&'a [u8]);
|
||||
|
||||
impl<'a> AsyncRead for Slice<'a> {
|
||||
fn poll_read(
|
||||
self: Pin<&mut Self>,
|
||||
_cx: &mut Context<'_>,
|
||||
mut self: Pin<&mut Self>,
|
||||
cx: &mut Context<'_>,
|
||||
buf: &mut [u8],
|
||||
) -> Poll<io::Result<usize>> {
|
||||
Ready(Pin::get_mut(self).0.read(buf))
|
||||
Pin::new(&mut self.0).poll_read(cx, buf)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,761 @@
|
||||
#![deny(warnings, rust_2018_idioms)]
|
||||
|
||||
use tokio::codec::*;
|
||||
use tokio::io::{AsyncRead, AsyncWrite};
|
||||
use tokio::prelude::*;
|
||||
use tokio_test::task::MockTask;
|
||||
use tokio_test::{
|
||||
assert_err, assert_ok, assert_pending, assert_ready, assert_ready_err, assert_ready_ok,
|
||||
};
|
||||
|
||||
use bytes::{BufMut, Bytes, BytesMut};
|
||||
use futures_util::pin_mut;
|
||||
use std::collections::VecDeque;
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::task::Poll::*;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
macro_rules! mock {
|
||||
($($x:expr,)*) => {{
|
||||
let mut v = VecDeque::new();
|
||||
v.extend(vec![$($x),*]);
|
||||
Mock { calls: v }
|
||||
}};
|
||||
}
|
||||
|
||||
macro_rules! assert_next_eq {
|
||||
($io:ident, $expect:expr) => {{
|
||||
MockTask::new().enter(|cx| {
|
||||
let res = assert_ready!($io.as_mut().poll_next(cx));
|
||||
match res {
|
||||
Some(Ok(v)) => assert_eq!(v, $expect.as_ref()),
|
||||
Some(Err(e)) => panic!("error = {:?}", e),
|
||||
None => panic!("none"),
|
||||
}
|
||||
});
|
||||
}};
|
||||
}
|
||||
|
||||
macro_rules! assert_next_pending {
|
||||
($io:ident) => {{
|
||||
MockTask::new().enter(|cx| match $io.as_mut().poll_next(cx) {
|
||||
Ready(Some(Ok(v))) => panic!("value = {:?}", v),
|
||||
Ready(Some(Err(e))) => panic!("error = {:?}", e),
|
||||
Ready(None) => panic!("done"),
|
||||
Pending => {}
|
||||
});
|
||||
}};
|
||||
}
|
||||
|
||||
macro_rules! assert_next_err {
|
||||
($io:ident) => {{
|
||||
MockTask::new().enter(|cx| match $io.as_mut().poll_next(cx) {
|
||||
Ready(Some(Ok(v))) => panic!("value = {:?}", v),
|
||||
Ready(Some(Err(_))) => {}
|
||||
Ready(None) => panic!("done"),
|
||||
Pending => panic!("pending"),
|
||||
});
|
||||
}};
|
||||
}
|
||||
|
||||
macro_rules! assert_done {
|
||||
($io:ident) => {{
|
||||
MockTask::new().enter(|cx| {
|
||||
let res = assert_ready!($io.as_mut().poll_next(cx));
|
||||
match res {
|
||||
Some(Ok(v)) => panic!("value = {:?}", v),
|
||||
Some(Err(e)) => panic!("error = {:?}", e),
|
||||
None => {}
|
||||
}
|
||||
});
|
||||
}};
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_empty_io_yields_nothing() {
|
||||
let io = Box::pin(FramedRead::new(mock!(), LengthDelimitedCodec::new()));
|
||||
pin_mut!(io);
|
||||
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_frame_one_packet() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00\x00\x09abcdefghi"),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_frame_one_packet_little_endian() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.little_endian()
|
||||
.new_read(mock! {
|
||||
data(b"\x09\x00\x00\x00abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_frame_one_packet_native_endian() {
|
||||
let d = if cfg!(target_endian = "big") {
|
||||
b"\x00\x00\x00\x09abcdefghi"
|
||||
} else {
|
||||
b"\x09\x00\x00\x00abcdefghi"
|
||||
};
|
||||
let io = length_delimited::Builder::new()
|
||||
.native_endian()
|
||||
.new_read(mock! {
|
||||
data(d),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_multi_frame_one_packet() {
|
||||
let mut d: Vec<u8> = vec![];
|
||||
d.extend_from_slice(b"\x00\x00\x00\x09abcdefghi");
|
||||
d.extend_from_slice(b"\x00\x00\x00\x03123");
|
||||
d.extend_from_slice(b"\x00\x00\x00\x0bhello world");
|
||||
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(&d),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_eq!(io, b"123");
|
||||
assert_next_eq!(io, b"hello world");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_frame_multi_packet() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00"),
|
||||
data(b"\x00\x09abc"),
|
||||
data(b"defghi"),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_multi_frame_multi_packet() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00"),
|
||||
data(b"\x00\x09abc"),
|
||||
data(b"defghi"),
|
||||
data(b"\x00\x00\x00\x0312"),
|
||||
data(b"3\x00\x00\x00\x0bhello world"),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_eq!(io, b"123");
|
||||
assert_next_eq!(io, b"hello world");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_frame_multi_packet_wait() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00"),
|
||||
Pending,
|
||||
data(b"\x00\x09abc"),
|
||||
Pending,
|
||||
data(b"defghi"),
|
||||
Pending,
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_pending!(io);
|
||||
assert_next_pending!(io);
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_pending!(io);
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_multi_frame_multi_packet_wait() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00"),
|
||||
Pending,
|
||||
data(b"\x00\x09abc"),
|
||||
Pending,
|
||||
data(b"defghi"),
|
||||
Pending,
|
||||
data(b"\x00\x00\x00\x0312"),
|
||||
Pending,
|
||||
data(b"3\x00\x00\x00\x0bhello world"),
|
||||
Pending,
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_pending!(io);
|
||||
assert_next_pending!(io);
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_pending!(io);
|
||||
assert_next_pending!(io);
|
||||
assert_next_eq!(io, b"123");
|
||||
assert_next_eq!(io, b"hello world");
|
||||
assert_next_pending!(io);
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_incomplete_head() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00"),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_incomplete_head_multi() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
Pending,
|
||||
data(b"\x00"),
|
||||
Pending,
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_pending!(io);
|
||||
assert_next_pending!(io);
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_incomplete_payload() {
|
||||
let io = FramedRead::new(
|
||||
mock! {
|
||||
data(b"\x00\x00\x00\x09ab"),
|
||||
Pending,
|
||||
data(b"cd"),
|
||||
Pending,
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_pending!(io);
|
||||
assert_next_pending!(io);
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_max_frame_len() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.max_frame_length(5)
|
||||
.new_read(mock! {
|
||||
data(b"\x00\x00\x00\x09abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_update_max_frame_len_at_rest() {
|
||||
let io = length_delimited::Builder::new().new_read(mock! {
|
||||
data(b"\x00\x00\x00\x09abcdefghi"),
|
||||
data(b"\x00\x00\x00\x09abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
io.decoder_mut().set_max_frame_length(5);
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_update_max_frame_len_in_flight() {
|
||||
let io = length_delimited::Builder::new().new_read(mock! {
|
||||
data(b"\x00\x00\x00\x09abcd"),
|
||||
Pending,
|
||||
data(b"efghi"),
|
||||
data(b"\x00\x00\x00\x09abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_pending!(io);
|
||||
io.decoder_mut().set_max_frame_length(5);
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_err!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_one_byte_length_field() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_field_length(1)
|
||||
.new_read(mock! {
|
||||
data(b"\x09abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_header_offset() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_field_length(2)
|
||||
.length_field_offset(4)
|
||||
.new_read(mock! {
|
||||
data(b"zzzz\x00\x09abcdefghi"),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_multi_frame_one_packet_skip_none_adjusted() {
|
||||
let mut d: Vec<u8> = vec![];
|
||||
d.extend_from_slice(b"xx\x00\x09abcdefghi");
|
||||
d.extend_from_slice(b"yy\x00\x03123");
|
||||
d.extend_from_slice(b"zz\x00\x0bhello world");
|
||||
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_field_length(2)
|
||||
.length_field_offset(2)
|
||||
.num_skip(0)
|
||||
.length_adjustment(4)
|
||||
.new_read(mock! {
|
||||
data(&d),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"xx\x00\x09abcdefghi");
|
||||
assert_next_eq!(io, b"yy\x00\x03123");
|
||||
assert_next_eq!(io, b"zz\x00\x0bhello world");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_single_multi_frame_one_packet_length_includes_head() {
|
||||
let mut d: Vec<u8> = vec![];
|
||||
d.extend_from_slice(b"\x00\x0babcdefghi");
|
||||
d.extend_from_slice(b"\x00\x05123");
|
||||
d.extend_from_slice(b"\x00\x0dhello world");
|
||||
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_field_length(2)
|
||||
.length_adjustment(-2)
|
||||
.new_read(mock! {
|
||||
data(&d),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
assert_next_eq!(io, b"abcdefghi");
|
||||
assert_next_eq!(io, b"123");
|
||||
assert_next_eq!(io, b"hello world");
|
||||
assert_done!(io);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_frame_length_adjusted() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_adjustment(-2)
|
||||
.new_write(mock! {
|
||||
data(b"\x00\x00\x00\x0b"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_nothing_yields_nothing() {
|
||||
let io = FramedWrite::new(mock!(), LengthDelimitedCodec::new());
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.poll_flush(cx));
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_frame_one_packet() {
|
||||
let io = FramedWrite::new(
|
||||
mock! {
|
||||
data(b"\x00\x00\x00\x09"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_multi_frame_one_packet() {
|
||||
let io = FramedWrite::new(
|
||||
mock! {
|
||||
data(b"\x00\x00\x00\x09"),
|
||||
data(b"abcdefghi"),
|
||||
data(b"\x00\x00\x00\x03"),
|
||||
data(b"123"),
|
||||
data(b"\x00\x00\x00\x0b"),
|
||||
data(b"hello world"),
|
||||
flush(),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("123")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("hello world")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_multi_frame_multi_packet() {
|
||||
let io = FramedWrite::new(
|
||||
mock! {
|
||||
data(b"\x00\x00\x00\x09"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
data(b"\x00\x00\x00\x03"),
|
||||
data(b"123"),
|
||||
flush(),
|
||||
data(b"\x00\x00\x00\x0b"),
|
||||
data(b"hello world"),
|
||||
flush(),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("123")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("hello world")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_frame_would_block() {
|
||||
let io = FramedWrite::new(
|
||||
mock! {
|
||||
Pending,
|
||||
data(b"\x00\x00"),
|
||||
Pending,
|
||||
data(b"\x00\x09"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
},
|
||||
LengthDelimitedCodec::new(),
|
||||
);
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
|
||||
assert_pending!(io.as_mut().poll_flush(cx));
|
||||
assert_pending!(io.as_mut().poll_flush(cx));
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_frame_little_endian() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.little_endian()
|
||||
.new_write(mock! {
|
||||
data(b"\x09\x00\x00\x00"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_single_frame_with_short_length_field() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.length_field_length(1)
|
||||
.new_write(mock! {
|
||||
data(b"\x09"),
|
||||
data(b"abcdefghi"),
|
||||
flush(),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdefghi")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_max_frame_len() {
|
||||
let io = length_delimited::Builder::new()
|
||||
.max_frame_length(5)
|
||||
.new_write(mock! {});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_err!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_update_max_frame_len_at_rest() {
|
||||
let io = length_delimited::Builder::new().new_write(mock! {
|
||||
data(b"\x00\x00\x00\x06"),
|
||||
data(b"abcdef"),
|
||||
flush(),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
io.encoder_mut().set_max_frame_length(5);
|
||||
|
||||
assert_err!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_update_max_frame_len_in_flight() {
|
||||
let io = length_delimited::Builder::new().new_write(mock! {
|
||||
data(b"\x00\x00\x00\x06"),
|
||||
data(b"ab"),
|
||||
Pending,
|
||||
data(b"cdef"),
|
||||
flush(),
|
||||
});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
|
||||
assert_pending!(io.as_mut().poll_flush(cx));
|
||||
|
||||
io.encoder_mut().set_max_frame_length(5);
|
||||
|
||||
assert_ready_ok!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert_err!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_zero() {
|
||||
let io = length_delimited::Builder::new().new_write(mock! {});
|
||||
pin_mut!(io);
|
||||
|
||||
MockTask::new().enter(|cx| {
|
||||
assert_ready_ok!(io.as_mut().poll_ready(cx));
|
||||
assert_ok!(io.as_mut().start_send(Bytes::from("abcdef")));
|
||||
|
||||
assert_ready_err!(io.as_mut().poll_flush(cx));
|
||||
|
||||
assert!(io.get_ref().calls.is_empty());
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn encode_overflow() {
|
||||
// Test reproducing tokio-rs/tokio#681.
|
||||
let mut codec = length_delimited::Builder::new().new_codec();
|
||||
let mut buf = BytesMut::with_capacity(1024);
|
||||
|
||||
// Put some data into the buffer without resizing it to hold more.
|
||||
let some_as = std::iter::repeat(b'a').take(1024).collect::<Vec<_>>();
|
||||
buf.put_slice(&some_as[..]);
|
||||
|
||||
// Trying to encode the length header should resize the buffer if it won't fit.
|
||||
codec.encode(Bytes::from("hello"), &mut buf).unwrap();
|
||||
}
|
||||
|
||||
// ===== Test utils =====
|
||||
|
||||
struct Mock {
|
||||
calls: VecDeque<Poll<io::Result<Op>>>,
|
||||
}
|
||||
|
||||
enum Op {
|
||||
Data(Vec<u8>),
|
||||
Flush,
|
||||
}
|
||||
|
||||
use self::Op::*;
|
||||
|
||||
impl AsyncRead for Mock {
|
||||
fn poll_read(
|
||||
mut self: Pin<&mut Self>,
|
||||
_cx: &mut Context<'_>,
|
||||
dst: &mut [u8],
|
||||
) -> Poll<io::Result<usize>> {
|
||||
match self.calls.pop_front() {
|
||||
Some(Ready(Ok(Op::Data(data)))) => {
|
||||
debug_assert!(dst.len() >= data.len());
|
||||
dst[..data.len()].copy_from_slice(&data[..]);
|
||||
Ready(Ok(data.len()))
|
||||
}
|
||||
Some(Ready(Ok(_))) => panic!(),
|
||||
Some(Ready(Err(e))) => Ready(Err(e)),
|
||||
Some(Pending) => Pending,
|
||||
None => Ready(Ok(0)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AsyncWrite for Mock {
|
||||
fn poll_write(
|
||||
mut self: Pin<&mut Self>,
|
||||
_cx: &mut Context<'_>,
|
||||
src: &[u8],
|
||||
) -> Poll<Result<usize, io::Error>> {
|
||||
match self.calls.pop_front() {
|
||||
Some(Ready(Ok(Op::Data(data)))) => {
|
||||
let len = data.len();
|
||||
assert!(src.len() >= len, "expect={:?}; actual={:?}", data, src);
|
||||
assert_eq!(&data[..], &src[..len]);
|
||||
Ready(Ok(len))
|
||||
}
|
||||
Some(Ready(Ok(_))) => panic!(),
|
||||
Some(Ready(Err(e))) => Ready(Err(e)),
|
||||
Some(Pending) => Pending,
|
||||
None => Ready(Ok(0)),
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_flush(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> {
|
||||
match self.calls.pop_front() {
|
||||
Some(Ready(Ok(Op::Flush))) => Ready(Ok(())),
|
||||
Some(Ready(Ok(_))) => panic!(),
|
||||
Some(Ready(Err(e))) => Ready(Err(e)),
|
||||
Some(Pending) => Pending,
|
||||
None => Ready(Ok(())),
|
||||
}
|
||||
}
|
||||
|
||||
fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> {
|
||||
Ready(Ok(()))
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a> From<&'a [u8]> for Op {
|
||||
fn from(src: &'a [u8]) -> Op {
|
||||
Op::Data(src.into())
|
||||
}
|
||||
}
|
||||
|
||||
impl From<Vec<u8>> for Op {
|
||||
fn from(src: Vec<u8>) -> Op {
|
||||
Op::Data(src)
|
||||
}
|
||||
}
|
||||
|
||||
fn data(bytes: &[u8]) -> Poll<io::Result<Op>> {
|
||||
Ready(Ok(bytes.into()))
|
||||
}
|
||||
|
||||
fn flush() -> Poll<io::Result<Op>> {
|
||||
Ready(Ok(Flush))
|
||||
}
|
||||
Reference in New Issue
Block a user