canary-rs/apps/magpie/src/ipc.rs

220 lines
6.8 KiB
Rust
Raw Normal View History

2022-10-29 05:22:36 +00:00
use std::io::Read;
2022-10-28 04:52:42 +00:00
use std::ops::{Deref, DerefMut};
use std::path::{Path, PathBuf};
2022-10-29 05:22:36 +00:00
use std::str::from_utf8;
2022-10-28 04:52:42 +00:00
use std::sync::mpsc::{channel, Receiver, Sender};
use std::time::Duration;
2022-10-28 04:23:03 +00:00
2022-10-29 05:22:36 +00:00
use mio::net::{UnixListener, UnixStream};
2022-10-28 04:52:42 +00:00
use mio::{Events, Interest, Poll, Token, Waker};
2022-10-28 02:15:47 +00:00
use mio_signals::{Signal, Signals};
use slab::Slab;
2022-10-28 04:23:03 +00:00
use crate::window::{WindowMessage, WindowMessageSender};
2022-10-28 02:15:47 +00:00
const SOCK_NAME: &str = "magpie.sock";
2022-10-28 04:23:03 +00:00
pub enum IpcMessage {}
2022-10-28 04:52:42 +00:00
pub struct IpcMessageSender {
waker: Waker,
sender: Sender<IpcMessage>,
}
2022-10-28 02:15:47 +00:00
/// Wraps [mio::net::UnixListener] with automatic file deletion on drop.
pub struct Listener {
pub uds: UnixListener,
pub path: PathBuf,
}
impl Drop for Listener {
fn drop(&mut self) {
match std::fs::remove_file(&self.path) {
Ok(_) => {}
Err(e) => eprintln!("Could not delete UnixListener {:?}", e),
}
}
}
impl Deref for Listener {
type Target = UnixListener;
fn deref(&self) -> &UnixListener {
&self.uds
}
}
impl DerefMut for Listener {
fn deref_mut(&mut self) -> &mut UnixListener {
&mut self.uds
}
}
2022-10-29 05:22:36 +00:00
pub struct Client {
connection: UnixStream,
}
impl Client {
pub fn new(connection: UnixStream) -> Self {
Self { connection }
}
pub fn on_readable(&mut self) -> std::io::Result<bool> {
let mut connection_closed = false;
let mut received_data = vec![0; 4096];
let mut bytes_read = 0;
loop {
match self.connection.read(&mut received_data[bytes_read..]) {
Ok(0) => {
connection_closed = true;
break;
}
Ok(n) => {
bytes_read += n;
if bytes_read >= received_data.len() {
received_data.resize(received_data.len() + 1024, 0);
}
}
Err(ref err) if err.kind() == std::io::ErrorKind::WouldBlock => break,
Err(ref err) if err.kind() == std::io::ErrorKind::Interrupted => continue,
Err(err) => return Err(err),
}
}
if bytes_read > 0 {
let received_data = &received_data[..bytes_read];
if let Ok(str_buf) = from_utf8(received_data) {
println!("Received data: {}", str_buf.trim_end());
} else {
println!("Received (non-UTF-8) data: {:?}", received_data);
}
}
Ok(connection_closed)
}
}
2022-10-28 02:15:47 +00:00
pub struct Ipc {
2022-10-28 04:52:42 +00:00
pub message_recv: Receiver<IpcMessage>,
2022-10-28 04:23:03 +00:00
pub window_sender: WindowMessageSender,
2022-10-28 02:15:47 +00:00
pub poll: Poll,
pub events: Events,
pub quit: bool,
pub listener: Listener,
pub signals: Signals,
pub listener_token: Token,
pub signals_token: Token,
2022-10-28 04:23:03 +00:00
pub message_recv_token: Token,
2022-10-28 02:15:47 +00:00
pub clients: Slab<Client>,
}
impl Ipc {
2022-10-28 04:52:42 +00:00
pub fn new(window_sender: WindowMessageSender) -> std::io::Result<(Self, IpcMessageSender)> {
2022-10-28 02:15:47 +00:00
let sock_dir = std::env::var("XDG_RUNTIME_DIR").expect("XDG_RUNTIME_DIR not set");
let sock_dir = Path::new(&sock_dir);
let sock_path = sock_dir.join(SOCK_NAME);
eprintln!("Making socket at: {:?}", sock_path);
let mut listener = Listener {
uds: UnixListener::bind(&sock_path)?,
path: sock_path.to_path_buf(),
};
let mut signals = Signals::new(Signal::Interrupt | Signal::Quit)?;
let events = Events::with_capacity(128);
let poll = Poll::new()?;
let listener_token = Token(usize::MAX);
let signals_token = Token(listener_token.0 - 1);
2022-10-28 04:23:03 +00:00
let message_recv_token = Token(signals_token.0 - 1);
2022-10-28 02:15:47 +00:00
let registry = poll.registry();
let interest = Interest::READABLE;
registry.register(&mut listener.uds, listener_token, interest)?;
registry.register(&mut signals, signals_token, interest)?;
2022-10-28 04:52:42 +00:00
let (sender, message_recv) = channel();
let sender = IpcMessageSender {
waker: Waker::new(registry, message_recv_token)?,
sender,
};
let ipc = Self {
2022-10-28 04:23:03 +00:00
message_recv,
window_sender,
2022-10-28 02:15:47 +00:00
poll,
events,
quit: false,
listener,
signals,
listener_token,
signals_token,
2022-10-28 04:23:03 +00:00
message_recv_token,
2022-10-28 02:15:47 +00:00
clients: Default::default(),
2022-10-28 04:52:42 +00:00
};
Ok((ipc, sender))
2022-10-28 02:15:47 +00:00
}
pub fn poll(&mut self, timeout: Option<Duration>) -> std::io::Result<()> {
self.poll.poll(&mut self.events, timeout)?;
for event in self.events.iter() {
if event.token() == self.listener_token {
loop {
match self.listener.accept() {
Ok((connection, address)) => {
2022-10-29 05:22:36 +00:00
let token = Token(self.clients.vacant_key());
println!(
"Accepting connection (Client #{}) from {:?}",
token.0, address
);
let mut client = Client::new(connection);
let interest = Interest::READABLE;
self.poll.registry().register(
&mut client.connection,
token,
interest,
)?;
self.clients.insert(client);
2022-10-28 02:15:47 +00:00
}
Err(ref err) if err.kind() == std::io::ErrorKind::WouldBlock => break,
Err(err) => return Err(err),
}
}
} else if event.token() == self.signals_token {
while let Some(received) = self.signals.receive()? {
eprintln!("Received {:?} signal; exiting...", received);
2022-10-28 04:23:03 +00:00
let _ = self.window_sender.send_event(WindowMessage::Quit);
2022-10-28 02:15:47 +00:00
self.quit = true;
}
2022-10-28 04:52:42 +00:00
} else if let Some(client) = self.clients.get_mut(event.token().0) {
2022-10-29 05:22:36 +00:00
let disconnected = client.on_readable()?;
if disconnected {
println!("Client #{} disconnected", event.token().0);
let mut client = self.clients.remove(event.token().0);
self.poll.registry().deregister(&mut client.connection)?;
}
2022-10-28 02:15:47 +00:00
} else {
panic!("Unrecognized event token: {:?}", event);
}
}
Ok(())
}
pub fn run(mut self) {
while !self.quit {
let wait = Duration::from_millis(100);
match self.poll(Some(wait)) {
Ok(_) => {}
Err(e) => {
eprintln!("IPC poll error: {:?}", e);
}
}
}
}
}