use anyhow::{Context as _, Result, bail};
use bytes::Bytes;
use std::marker::PhantomData;
use tokio::sync::mpsc::{self, UnboundedReceiver, UnboundedSender};
#[cfg(all(feature = "uring", target_os = "linux"))]
pub mod common;
pub mod stream;
#[cfg(all(feature = "uring", target_os = "linux"))]
mod tcp_client;
#[cfg(all(feature = "uring", target_os = "linux"))]
mod unix_client;
#[derive(Debug, Clone)]
pub struct Sender<T> {
tx: UnboundedSender<Bytes>,
_phantom: PhantomData<T>,
}
impl<T> Sender<T>
where
T: senax_encoder::Encoder,
{
pub fn send(&self, data: &T) -> Result<()> {
let bytes = senax_encoder::encode(data)?;
self.tx.send(bytes)?;
Ok(())
}
}
#[derive(Debug)]
pub struct Receiver<T> {
rx: UnboundedReceiver<Bytes>,
_phantom: PhantomData<T>,
}
impl<T> Receiver<T>
where
T: senax_encoder::Decoder,
{
pub async fn recv(&mut self) -> Option<Option<Result<T>>> {
match self.rx.recv().await {
Some(mut v) => {
if v.is_empty() {
Some(None)
} else {
Some(Some(senax_encoder::decode(&mut v).context("parse error")))
}
}
None => None,
}
}
}
pub fn link<T>(
stream_id: u64,
port: &str,
pw: &str,
exit_tx: mpsc::Sender<i32>,
send_only: bool,
) -> Result<(Sender<T>, Receiver<T>)>
where
T: senax_encoder::Encoder + senax_encoder::Decoder,
{
let (to_linker, from_linker) = LinkerClient::start(port, stream_id, pw, exit_tx, send_only)?;
Ok((
Sender {
tx: to_linker,
_phantom: Default::default(),
},
Receiver {
rx: from_linker,
_phantom: Default::default(),
},
))
}
pub struct LinkerClient;
#[cfg(all(feature = "uring", target_os = "linux"))]
#[allow(clippy::type_complexity)]
impl LinkerClient {
pub fn start(
port: &str,
stream_id: u64,
pw: &str,
exit_tx: mpsc::Sender<i32>,
send_only: bool,
) -> Result<(UnboundedSender<Bytes>, UnboundedReceiver<Bytes>)> {
let (to_linker, from_local) = mpsc::unbounded_channel();
let (to_local, from_linker) = mpsc::unbounded_channel();
if port.starts_with('/') {
match unix_client::run(
port,
stream_id,
from_local,
to_local,
pw.to_string(),
exit_tx,
send_only,
) {
Ok(_) => {
return Ok((to_linker, from_linker));
}
Err(e) => {
log::warn!("{}", e);
}
}
} else {
match tcp_client::run(
port,
stream_id,
from_local,
to_local,
pw.to_string(),
exit_tx,
send_only,
) {
Ok(_) => {
return Ok((to_linker, from_linker));
}
Err(e) => {
log::warn!("{}", e);
}
}
}
bail!("linker connection failed");
}
}
#[cfg(not(all(feature = "uring", target_os = "linux")))]
#[allow(unused_variables)]
#[allow(clippy::type_complexity)]
impl LinkerClient {
pub fn start(
_port: &str,
_stream_id: u64,
_pw: &str,
_exit_tx: mpsc::Sender<i32>,
_send_only: bool,
) -> Result<(UnboundedSender<Bytes>, UnboundedReceiver<Bytes>)> {
bail!("linker is not supported");
}
}