mod conn;
mod ctrlstream;
pub(crate) mod nbistream;
pub(crate) mod tagstream;
use crate::client::ConnectingStreamError;
pub(crate) use conn::ConecConn;
pub use conn::ConecConnError;
pub(crate) use ctrlstream::CtrlStream;
pub use ctrlstream::{ControlMsg, CtrlStreamError};
use futures::prelude::*;
use quinn::{RecvStream, SendStream};
use serde::{Deserialize, Serialize};
use std::net::SocketAddr;
use tokio_serde::{formats::SymmetricalBincode, SymmetricallyFramed};
use tokio_util::codec::{FramedRead, FramedWrite, LengthDelimitedCodec};
pub type InStream = FramedRead<RecvStream, LengthDelimitedCodec>;
pub type OutStream = FramedWrite<SendStream, LengthDelimitedCodec>;
pub(crate) async fn outstream_init(
send: SendStream,
from: String,
sid: u64,
) -> Result<OutStream, ConnectingStreamError> {
let mut write_stream = SymmetricallyFramed::new(
FramedWrite::new(send, LengthDelimitedCodec::new()),
SymmetricalBincode::<(String, u64)>::default(),
);
write_stream.send((from, sid)).await?;
write_stream.flush().await?;
Ok(write_stream.into_inner())
}
#[derive(Clone, Debug)]
pub(crate) enum ConecConnAddr {
Portnum(u16),
Sockaddr(SocketAddr),
}
impl From<SocketAddr> for ConecConnAddr {
fn from(addr: SocketAddr) -> Self {
ConecConnAddr::Sockaddr(addr)
}
}
impl From<u16> for ConecConnAddr {
fn from(port: u16) -> Self {
ConecConnAddr::Portnum(port)
}
}
impl ConecConnAddr {
pub(crate) fn is_sockaddr(&self) -> bool {
match self {
Self::Portnum(_) => false,
Self::Sockaddr(_) => true,
}
}
pub(crate) fn get_port(&self) -> Option<u16> {
match self {
Self::Portnum(p) => Some(*p),
Self::Sockaddr(_) => None,
}
}
pub(crate) fn get_addr(&self) -> Option<&SocketAddr> {
match self {
Self::Portnum(_) => None,
Self::Sockaddr(ref s) => Some(s),
}
}
}
#[derive(Serialize, Deserialize, Debug, Clone)]
pub(crate) enum StreamTo {
Broadcast(u64),
Client(u64),
}