use crate::types::wire::{self, Delivery, Reply as MachReply, Request, service_name};
use bevy::tasks::{IoTaskPool, TaskPool};
use futures_lite::StreamExt;
use std::sync::Arc;
use std::thread;
use tracing::{error, warn};
use crate::errors::Result;
use crate::events::{Event, EventSender, Reply};
pub struct CommandReader {
events: EventSender,
}
impl CommandReader {
#[must_use]
pub fn new(events: EventSender) -> Self {
CommandReader { events }
}
pub fn start(self) -> Result<()> {
let receiver = wire::bind::<Request>(&service_name()).inspect_err(|_| {
error!(
"can not register a Mach port - maybe another Paneru instance is already running?"
);
})?;
thread::spawn(move || {
futures_lite::future::block_on(async move {
let mut requests = std::pin::pin!(receiver);
while let Some(delivery) = requests.next().await {
match delivery {
Ok(delivery) => self.dispatch(delivery),
Err(err) => warn!("reading request: {err}"),
}
}
});
});
Ok(())
}
fn dispatch(&self, delivery: Delivery<Request>) {
let events = self.events.clone();
let Delivery {
value,
reply,
subscriber,
} = delivery;
match value {
Request::Command(command) => {
send(&events, Event::Command { command });
}
Request::WindowSetApply(ops) => {
send(
&events,
Event::Command {
command: crate::commands::Command::Layout(ops),
},
);
}
Request::Query(kind) => {
answer(events, reply, "state query", move |respond_to| {
Event::StateQuery { kind, respond_to }
});
}
Request::WindowSet => {
answer(events, reply, "window set query", |respond_to| {
Event::WindowSetQuery { respond_to }
});
}
Request::ScriptState(request) => {
answer(events, reply, "script state request", move |respond_to| {
Event::ScriptState {
request,
respond_to,
}
});
}
Request::Subscribe => {
if let Some(subscriber) = subscriber {
send(
&events,
Event::StateSubscribe {
subscriber: Arc::new(subscriber),
},
);
} else {
warn!("subscribe request carried no event channel");
}
}
}
}
}
fn send(events: &EventSender, event: Event) {
_ = events
.send(event)
.inspect_err(|err| error!("sending event: {err}"));
}
fn answer(
events: EventSender,
reply: Option<MachReply>,
what: &'static str,
request: impl FnOnce(Reply) -> Event + Send + 'static,
) {
let Some(reply) = reply else {
warn!("{what} arrived without a reply channel");
return;
};
let (tx, rx) = async_channel::bounded(1);
if events
.send(request(tx))
.inspect_err(|err| error!("sending {what}: {err}"))
.is_err()
{
return;
}
IoTaskPool::get_or_init(TaskPool::default)
.spawn(async move {
match rx.recv().await {
Ok(response) => {
if let Err(err) = reply.send(&response) {
warn!("answering {what}: {err}");
}
}
Err(err) => error!("waiting for {what} response: {err}"),
}
})
.detach();
}