#![cfg(not(target_arch = "wasm32"))]
#![allow(clippy::allow_attributes, missing_docs, reason = "// TODO(mingwei)")]
use std::collections::VecDeque;
use std::pin::Pin;
use byteorder::{NetworkEndian, WriteBytesExt};
use futures::{Sink, StreamExt};
use tokio::net::tcp::{OwnedReadHalf, OwnedWriteHalf};
use tokio::net::TcpStream;
use tokio_util::codec::{FramedRead, FramedWrite, LengthDelimitedCodec};
use super::graph::Hydroflow;
use super::graph_ext::GraphExt;
use super::handoff::VecHandoff;
use super::port::{RecvPort, SendPort};
pub mod network_vertex;
const ADDRESS_LEN: usize = 4;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Message {
pub address: u32,
pub batch: bytes::Bytes,
}
impl Message {
fn encode(&self, v: &mut Vec<u8>) {
v.write_u32::<NetworkEndian>(self.address).unwrap();
v.extend(self.batch.iter());
}
pub fn decode(v: bytes::Bytes) -> Self {
let address = u32::from_be_bytes(v[0..ADDRESS_LEN].try_into().unwrap());
let batch = v.slice(ADDRESS_LEN..);
Message { address, batch }
}
}
impl Hydroflow<'_> {
fn register_read_tcp_stream(&mut self, reader: OwnedReadHalf) -> RecvPort<VecHandoff<Message>> {
let reader = FramedRead::new(reader, LengthDelimitedCodec::new());
let (send_port, recv_port) = self.make_edge("tcp ingress handoff");
self.add_input_from_stream(
"tcp ingress",
send_port,
reader.map(|buf| Some(<Message>::decode(buf.unwrap().into()))),
);
recv_port
}
fn register_write_tcp_stream(
&mut self,
writer: OwnedWriteHalf,
) -> SendPort<VecHandoff<Message>> {
let mut writer = FramedWrite::new(writer, LengthDelimitedCodec::new());
let mut message_queue = VecDeque::new();
let (input_port, output_port) =
self.make_edge::<_, VecHandoff<Message>>("tcp egress handoff");
self.add_subgraph_sink("tcp egress", output_port, move |ctx, recv| {
let waker = ctx.waker();
let mut cx = std::task::Context::from_waker(&waker);
message_queue.extend(recv.take_inner());
while !message_queue.is_empty() {
if let std::task::Poll::Ready(Ok(())) = Pin::new(&mut writer).poll_ready(&mut cx) {
let v = message_queue.pop_front().unwrap();
let mut buf = Vec::new();
v.encode(&mut buf);
Pin::new(&mut writer).start_send(buf.into()).unwrap();
}
}
let _ = Pin::new(&mut writer).poll_flush(&mut cx);
});
input_port
}
pub fn add_write_tcp_stream(&mut self, stream: TcpStream) -> SendPort<VecHandoff<Message>> {
let (_, writer) = stream.into_split();
self.register_write_tcp_stream(writer)
}
pub fn add_read_tcp_stream(&mut self, stream: TcpStream) -> RecvPort<VecHandoff<Message>> {
let (reader, _) = stream.into_split();
self.register_read_tcp_stream(reader)
}
pub fn add_tcp_stream(
&mut self,
stream: TcpStream,
) -> (SendPort<VecHandoff<Message>>, RecvPort<VecHandoff<Message>>) {
let (reader, writer) = stream.into_split();
(
self.register_write_tcp_stream(writer),
self.register_read_tcp_stream(reader),
)
}
}