use std::cell::{Cell, RefCell};
use std::future::Future;
use std::rc::Rc;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use async_channel::{Sender, bounded};
use super::worker::{Shared, StoreRequest, WorldRequest};
use crate::ecs::state::PaneruQueryState;
use crate::types::script_state::{ScriptState, ScriptStateWrite, WriteOutcome};
use crate::types::windowset::WindowSet;
const SHUTTING_DOWN: &str = "the window manager is shutting down";
const NO_DISPATCH: &str = "only available inside a paneru.on handler or a paneru.bind callback";
#[derive(Clone)]
pub(super) struct WorldAccess {
world: Sender<WorldRequest>,
store: Sender<StoreRequest>,
revision: Arc<AtomicU64>,
}
async fn ask<R, T>(channel: &Sender<R>, request: impl FnOnce(Sender<T>) -> R) -> Result<T, String> {
let (reply, answer) = bounded(1);
channel
.send(request(reply))
.await
.map_err(|_| SHUTTING_DOWN.to_string())?;
answer.recv().await.map_err(|_| SHUTTING_DOWN.to_string())
}
impl WorldAccess {
pub(super) fn new(
world: Sender<WorldRequest>,
store: Sender<StoreRequest>,
revision: Arc<AtomicU64>,
) -> Self {
Self {
world,
store,
revision,
}
}
async fn state(&self) -> Shared<PaneruQueryState> {
ask(&self.world, |reply| WorldRequest::State { reply })
.await
.unwrap_or_else(Err)
}
async fn window_set(&self) -> Shared<WindowSet> {
ask(&self.world, |reply| WorldRequest::WindowSet { reply })
.await
.unwrap_or_else(Err)
}
async fn script_state(&self) -> Result<ScriptState, String> {
ask(&self.store, |reply| StoreRequest::Read { reply })
.await
.unwrap_or_else(Err)
}
async fn write_script_state(&self, write: &ScriptStateWrite) -> Result<WriteOutcome, String> {
ask(&self.store, |reply| StoreRequest::Write {
write: write.clone(),
reply,
})
.await
.unwrap_or_else(Err)
}
}
struct SharedRead<T> {
cached: RefCell<Option<Arc<T>>>,
waiting: RefCell<Option<Vec<Sender<Shared<T>>>>>,
}
impl<T> SharedRead<T> {
fn new() -> Self {
Self {
cached: RefCell::new(None),
waiting: RefCell::new(None),
}
}
fn clear(&self) {
self.cached.borrow_mut().take();
}
async fn get<F>(&self, read: impl FnOnce() -> F) -> Shared<T>
where
F: Future<Output = Shared<T>>,
{
let cached = self.cached.borrow().clone();
if let Some(cached) = cached {
return Ok(cached);
}
let joined = {
let mut waiting = self.waiting.borrow_mut();
if let Some(queue) = waiting.as_mut() {
let (tell, told) = bounded(1);
queue.push(tell);
Some(told)
} else {
*waiting = Some(Vec::new());
None
}
};
if let Some(told) = joined {
return told
.recv()
.await
.unwrap_or_else(|_| Err(SHUTTING_DOWN.to_string()));
}
let answer = read().await;
if let Ok(value) = &answer {
*self.cached.borrow_mut() = Some(Arc::clone(value));
}
let queued = self.waiting.borrow_mut().take().unwrap_or_default();
for waiter in queued {
let _ = waiter.try_send(answer.clone());
}
answer
}
}
pub(super) struct DispatchWorld {
access: WorldAccess,
in_flight: Cell<usize>,
state: SharedRead<PaneruQueryState>,
window_set: SharedRead<WindowSet>,
script_state: RefCell<Option<(u64, ScriptState)>>,
}
impl DispatchWorld {
pub(super) fn new(access: WorldAccess) -> Rc<Self> {
Rc::new(Self {
access,
in_flight: Cell::new(0),
state: SharedRead::new(),
window_set: SharedRead::new(),
script_state: RefCell::new(None),
})
}
pub(super) fn enter(self: &Rc<Self>) -> Dispatch {
self.in_flight.set(self.in_flight.get() + 1);
Dispatch {
world: Rc::clone(self),
}
}
fn available(&self, call: &str) -> Result<(), String> {
if self.in_flight.get() == 0 {
return Err(format!("{call} is {NO_DISPATCH}"));
}
Ok(())
}
pub(super) async fn query_state(&self) -> Result<Arc<PaneruQueryState>, String> {
self.available("paneru.query")?;
self.state.get(|| self.access.state()).await
}
pub(super) async fn layout(&self) -> Result<Arc<WindowSet>, String> {
self.available("the window set")?;
self.window_set.get(|| self.access.window_set()).await
}
pub(super) async fn script_state(&self) -> Result<ScriptState, String> {
self.available("paneru.state")?;
let revision = self.access.revision.load(Ordering::Acquire);
let cached = self.script_state.borrow().clone();
match cached {
Some((stamp, state)) if stamp == revision => Ok(state),
_ => {
let fresh = self.access.script_state().await?;
*self.script_state.borrow_mut() = Some((revision, fresh.clone()));
Ok(fresh)
}
}
}
pub(super) async fn write_script_state(
&self,
write: &ScriptStateWrite,
) -> Result<WriteOutcome, String> {
self.available("paneru.state")?;
self.access.write_script_state(write).await
}
}
pub(super) struct Dispatch {
world: Rc<DispatchWorld>,
}
impl Drop for Dispatch {
fn drop(&mut self) {
let remaining = self.world.in_flight.get().saturating_sub(1);
self.world.in_flight.set(remaining);
if remaining == 0 {
self.world.state.clear();
self.world.window_set.clear();
}
}
}