use super::interprocess::name_onto;
use crate::message::Message;
use crate::{Error, Result};
use interprocess::local_socket::{LocalSocketListener, LocalSocketStream};
use serde::de::DeserializeOwned;
use serde::Serialize;
use std::io::{Read, Write};
pub struct Listener {
internal: Box<dyn ListenerImpl>,
closed: bool,
}
impl Listener {
pub const fn new(internal: Box<dyn ListenerImpl>) -> Self {
Self {
internal,
closed: false,
}
}
pub fn listen_as_socket<'a, S>(name: S, global: bool) -> Result<Self>
where
S: AsRef<str>,
{
let bound = name_onto!(LocalSocketListener::bind; name, global)?;
Ok(Self::new(Box::new(bound)))
}
pub fn accept(&mut self) -> Result<Connection> {
if self.closed {
return Err(Error::Closed(false));
}
self.internal.accept()
}
pub fn close(&mut self) -> Result<()> {
if self.closed {
return Err(Error::Closed(false));
}
self.closed = true; self.internal.close()
}
pub fn is_closed(&self) -> bool {
self.closed
}
}
impl Drop for Listener {
fn drop(&mut self) {
let _ = self.close();
}
}
pub struct Connection {
internal: Box<dyn ConnectionImpl>,
closed: bool,
}
impl Connection {
pub const fn new(internal: Box<dyn ConnectionImpl>) -> Self {
Self {
internal,
closed: false,
}
}
pub fn connect_to_socket<S>(name: S, global: bool) -> Result<Self>
where
S: AsRef<str>,
{
let bound = name_onto!(LocalSocketStream::connect; name, global)?;
Ok(Self::new(Box::new(bound)))
}
fn _send<T>(&mut self, message: Message<T>) -> Result<()>
where
T: Serialize,
{
message.write_to(&mut self.internal)
}
fn _receive<T>(&mut self) -> Result<Message<T>>
where
T: DeserializeOwned,
{
Message::<T>::read_from(&mut self.internal)
}
pub fn send<T>(&mut self, message_data: &T) -> Result<()>
where
T: Serialize,
{
if self.closed {
return Err(Error::Closed(false));
}
let message = Message::Data(message_data);
self._send(message)
}
pub fn receive<T>(&mut self) -> Result<T>
where
T: DeserializeOwned,
{
if self.closed {
return Err(Error::Closed(false));
}
let message = self._receive()?;
match message {
Message::ClosingConnection => {
self._close();
Err(Error::Closed(true))
}
Message::Data(data) => Ok(data),
}
}
pub fn send_and_receive<A, B>(&mut self, data: &A) -> Result<B>
where
A: Serialize,
B: DeserializeOwned,
{
self.send(data)?;
self.receive()
}
fn _close(&mut self) {
self.internal.close();
self.closed = true;
}
pub fn close(&mut self) {
if self.closed {
return;
}
let _ = self._send::<()>(Message::ClosingConnection);
self._close();
}
pub fn is_closed(&self) -> bool {
self.closed
}
}
impl Drop for Connection {
fn drop(&mut self) {
self.close();
}
}
pub trait ListenerImpl {
fn accept(&mut self) -> Result<Connection>;
fn close(&mut self) -> Result<()>;
}
impl ListenerImpl for LocalSocketListener {
fn accept(&mut self) -> Result<Connection> {
Ok(Connection::from(LocalSocketListener::accept(self)?))
}
fn close(&mut self) -> Result<()> {
Ok(())
}
}
impl From<LocalSocketListener> for Listener {
fn from(value: LocalSocketListener) -> Self {
Self::new(Box::new(value))
}
}
pub trait ConnectionImpl: Read + Write {
fn close(&mut self);
}
impl ConnectionImpl for LocalSocketStream {
fn close(&mut self) {
let _ = self.flush();
}
}
impl From<LocalSocketStream> for Connection {
fn from(value: LocalSocketStream) -> Self {
Connection::new(Box::new(value))
}
}