Files
tokio/tokio-util/src/udp/frame.rs
T

231 lines
7.5 KiB
Rust
Raw Normal View History

2019-10-22 10:13:49 -07:00
use crate::codec::{Decoder, Encoder};
2019-08-16 07:26:10 -07:00
2020-11-06 10:59:15 -05:00
use tokio::{io::ReadBuf, net::UdpSocket, stream::Stream};
2019-08-16 07:26:10 -07:00
2019-02-21 11:56:15 -08:00
use bytes::{BufMut, BytesMut};
2019-12-18 22:57:22 +03:00
use futures_core::ready;
use futures_sink::Sink;
2019-05-14 10:27:36 -07:00
use std::net::{Ipv4Addr, SocketAddr, SocketAddrV4};
use std::pin::Pin;
2019-10-22 10:13:49 -07:00
use std::task::{Context, Poll};
2020-11-06 10:59:15 -05:00
use std::{io, mem::MaybeUninit};
2016-11-22 11:48:09 -08:00
2020-11-06 10:59:15 -05:00
/// A unified [`Stream`] and [`Sink`] interface to an underlying `UdpSocket`, using
2018-02-07 01:41:31 +04:00
/// the `Encoder` and `Decoder` traits to encode and decode frames.
2016-11-22 11:48:09 -08:00
///
2018-02-07 10:42:27 -08:00
/// Raw UDP sockets work with datagrams, but higher-level code usually wants to
/// batch these into meaningful chunks, called "frames". This method layers
/// framing on top of this socket by using the `Encoder` and `Decoder` traits to
/// handle encoding and decoding of messages frames. Note that the incoming and
/// outgoing frame types may be distinct.
///
2020-11-06 10:59:15 -05:00
/// This function returns a *single* object that is both [`Stream`] and [`Sink`];
2018-02-07 10:42:27 -08:00
/// grouping this into a single object is often useful for layering things which
/// require both read and write access to the underlying object.
///
/// If you want to work more directly with the streams and sink, consider
2020-11-06 10:59:15 -05:00
/// calling [`split`] on the `UdpFramed` returned by this method, which will break
2018-02-07 10:42:27 -08:00
/// them into separate objects, allowing them to interact more easily.
2020-11-06 10:59:15 -05:00
///
/// [`Stream`]: tokio::stream::Stream
/// [`Sink`]: futures_sink::Sink
/// [`split`]: https://docs.rs/futures/0.3/futures/stream/trait.StreamExt.html#method.split
2017-10-25 18:03:31 -07:00
#[must_use = "sinks do nothing unless polled"]
2019-12-21 20:04:30 +03:00
#[cfg_attr(docsrs, doc(all(feature = "codec", feature = "udp")))]
2017-12-06 17:19:21 +01:00
#[derive(Debug)]
2016-11-22 11:48:09 -08:00
pub struct UdpFramed<C> {
socket: UdpSocket,
codec: C,
2018-02-07 01:41:31 +04:00
rd: BytesMut,
wr: BytesMut,
2016-11-22 11:48:09 -08:00
out_addr: SocketAddr,
2017-09-11 15:56:41 +02:00
flushed: bool,
is_readable: bool,
current_addr: Option<SocketAddr>,
2016-11-22 11:48:09 -08:00
}
2020-11-06 10:59:15 -05:00
const INITIAL_RD_CAPACITY: usize = 64 * 1024;
const INITIAL_WR_CAPACITY: usize = 8 * 1024;
impl<C: Decoder + Unpin> Stream for UdpFramed<C> {
type Item = Result<(C::Item, SocketAddr), C::Error>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let pin = self.get_mut();
2016-11-22 11:48:09 -08:00
pin.rd.reserve(INITIAL_RD_CAPACITY);
2018-02-07 01:41:31 +04:00
loop {
// Are there are still bytes left in the read buffer to decode?
if pin.is_readable {
if let Some(frame) = pin.codec.decode_eof(&mut pin.rd)? {
let current_addr = pin
.current_addr
.expect("will always be set before this line is called");
return Poll::Ready(Some(Ok((frame, current_addr))));
}
// if this line has been reached then decode has returned `None`.
pin.is_readable = false;
pin.rd.clear();
}
2019-08-27 17:53:57 -07:00
// We're out of data. Try and fetch more data to decode
let addr = unsafe {
// Convert `&mut [MaybeUnit<u8>]` to `&mut [u8]` because we will be
// writing to it via `poll_recv_from` and therefore initializing the memory.
2020-11-06 10:59:15 -05:00
let buf = &mut *(pin.rd.bytes_mut() as *mut _ as *mut [MaybeUninit<u8>]);
let mut read = ReadBuf::uninit(buf);
let ptr = read.filled().as_ptr();
let res = ready!(Pin::new(&mut pin.socket).poll_recv_from(cx, &mut read));
assert_eq!(ptr, read.filled().as_ptr());
let addr = res?;
pin.rd.advance_mut(read.filled().len());
addr
};
pin.current_addr = Some(addr);
pin.is_readable = true;
}
2016-11-22 11:48:09 -08:00
}
}
2020-03-04 15:54:41 -05:00
impl<I, C: Encoder<I> + Unpin> Sink<(I, SocketAddr)> for UdpFramed<C> {
type Error = C::Error;
2017-09-11 15:56:41 +02:00
fn poll_ready(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
2017-09-11 15:56:41 +02:00
if !self.flushed {
match self.poll_flush(cx)? {
Poll::Ready(()) => {}
Poll::Pending => return Poll::Pending,
2016-11-22 11:48:09 -08:00
}
}
Poll::Ready(Ok(()))
}
2020-03-04 15:54:41 -05:00
fn start_send(self: Pin<&mut Self>, item: (I, SocketAddr)) -> Result<(), Self::Error> {
2018-02-07 01:41:31 +04:00
let (frame, out_addr) = item;
2017-09-11 15:56:41 +02:00
let pin = self.get_mut();
pin.codec.encode(frame, &mut pin.wr)?;
pin.out_addr = out_addr;
pin.flushed = false;
Ok(())
2016-11-22 11:48:09 -08:00
}
fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
2017-09-11 15:56:41 +02:00
if self.flushed {
return Poll::Ready(Ok(()));
2016-11-22 11:48:09 -08:00
}
let Self {
ref mut socket,
ref mut out_addr,
ref mut wr,
..
} = *self;
2019-10-22 10:13:49 -07:00
let n = ready!(socket.poll_send_to(cx, &wr, &out_addr))?;
2017-09-11 15:56:41 +02:00
2016-11-22 11:48:09 -08:00
let wrote_all = n == self.wr.len();
self.wr.clear();
2017-09-11 15:56:41 +02:00
self.flushed = true;
let res = if wrote_all {
Ok(())
2016-11-22 11:48:09 -08:00
} else {
2019-02-21 11:56:15 -08:00
Err(io::Error::new(
io::ErrorKind::Other,
"failed to write entire datagram to socket",
)
.into())
};
Poll::Ready(res)
2016-11-22 11:48:09 -08:00
}
2017-02-05 17:06:57 -08:00
fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
ready!(self.poll_flush(cx))?;
Poll::Ready(Ok(()))
2017-02-05 17:06:57 -08:00
}
2016-11-22 11:48:09 -08:00
}
2018-02-07 10:42:27 -08:00
impl<C> UdpFramed<C> {
/// Create a new `UdpFramed` backed by the given socket and codec.
///
2018-05-08 14:44:17 -04:00
/// See struct level documentation for more details.
2018-02-07 10:42:27 -08:00
pub fn new(socket: UdpSocket, codec: C) -> UdpFramed<C> {
2020-11-06 10:59:15 -05:00
Self {
socket,
codec,
2018-02-07 10:42:27 -08:00
out_addr: SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 0)),
rd: BytesMut::with_capacity(INITIAL_RD_CAPACITY),
wr: BytesMut::with_capacity(INITIAL_WR_CAPACITY),
flushed: true,
is_readable: false,
current_addr: None,
2018-02-07 10:42:27 -08:00
}
2016-11-22 11:48:09 -08:00
}
/// Returns a reference to the underlying I/O stream wrapped by `Framed`.
///
2017-12-05 16:55:25 +01:00
/// # Note
///
/// Care should be taken to not tamper with the underlying stream of data
/// coming in as it may corrupt the stream of frames otherwise being worked
/// with.
2016-11-22 11:48:09 -08:00
pub fn get_ref(&self) -> &UdpSocket {
&self.socket
}
/// Returns a mutable reference to the underlying I/O stream wrapped by
/// `Framed`.
///
2017-12-05 16:55:25 +01:00
/// # Note
///
/// Care should be taken to not tamper with the underlying stream of data
/// coming in as it may corrupt the stream of frames otherwise being worked
/// with.
2016-11-22 11:48:09 -08:00
pub fn get_mut(&mut self) -> &mut UdpSocket {
&mut self.socket
}
/// Consumes the `Framed`, returning its underlying I/O stream.
pub fn into_inner(self) -> UdpSocket {
self.socket
}
2020-11-06 10:59:15 -05:00
/// Returns a reference to the underlying codec wrapped by
/// `Framed`.
///
/// Note that care should be taken to not tamper with the underlying codec
/// as it may corrupt the stream of frames otherwise being worked with.
pub fn codec(&self) -> &C {
&self.codec
}
/// Returns a mutable reference to the underlying codec wrapped by
/// `UdpFramed`.
///
/// Note that care should be taken to not tamper with the underlying codec
/// as it may corrupt the stream of frames otherwise being worked with.
pub fn codec_mut(&mut self) -> &mut C {
&mut self.codec
}
/// Returns a reference to the read buffer.
pub fn read_buffer(&self) -> &BytesMut {
&self.rd
}
/// Returns a mutable reference to the read buffer.
pub fn read_buffer_mut(&mut self) -> &mut BytesMut {
&mut self.rd
}
2016-11-22 11:48:09 -08:00
}