use std::cell::RefCell;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, Sender, channel};
use std::sync::{Mutex, OnceLock};
use guinea_core::actor::{UiDispatcher, UiTask, set_ui_dispatcher};
use iced::futures::StreamExt;
use iced::futures::channel::mpsc::{UnboundedReceiver, UnboundedSender, unbounded};
use iced::futures::stream::{Stream, pending};
use crate::Envelope;
struct ChannelDispatcher(Sender<UiTask>);
impl UiDispatcher for ChannelDispatcher {
fn init(&self) {}
fn dispatch(&self, task: UiTask) {
if self.0.send(task).is_ok() {
wake();
}
}
}
thread_local! {
static TASKS: RefCell<Option<Receiver<UiTask>>> = const { RefCell::new(None) };
}
static PENDING_WAKE: AtomicBool = AtomicBool::new(false);
fn wakeups_out() -> &'static Mutex<Option<UnboundedSender<()>>> {
static SENDER: OnceLock<Mutex<Option<UnboundedSender<()>>>> = OnceLock::new();
SENDER.get_or_init(|| Mutex::new(None))
}
fn wakeups_in() -> &'static Mutex<Option<UnboundedReceiver<()>>> {
static RECEIVER: OnceLock<Mutex<Option<UnboundedReceiver<()>>>> = OnceLock::new();
RECEIVER.get_or_init(|| Mutex::new(None))
}
pub(crate) fn install() {
static INIT: std::sync::Once = std::sync::Once::new();
INIT.call_once(|| {
let (tx, rx) = channel();
TASKS.with(|slot| *slot.borrow_mut() = Some(rx));
let (wake_tx, wake_rx) = unbounded();
if let Ok(mut slot) = wakeups_out().lock() {
*slot = Some(wake_tx);
}
if let Ok(mut slot) = wakeups_in().lock() {
*slot = Some(wake_rx);
}
set_ui_dispatcher(ChannelDispatcher(tx));
});
}
fn wake() {
if PENDING_WAKE.swap(true, Ordering::SeqCst) {
return;
}
if let Ok(sender) = wakeups_out().lock()
&& let Some(sender) = sender.as_ref()
{
let _ = sender.unbounded_send(());
}
}
pub(crate) fn wakeups() -> impl Stream<Item = Envelope> + Send + 'static {
let taken = wakeups_in().lock().ok().and_then(|mut slot| slot.take());
match taken {
Some(receiver) => receiver.map(|()| Envelope::settled()).left_stream(),
None => pending().right_stream(),
}
}
pub(crate) fn drain() {
PENDING_WAKE.store(false, Ordering::SeqCst);
let ready: Vec<UiTask> = TASKS.with(|slot| match slot.borrow().as_ref() {
Some(rx) => rx.try_iter().collect(),
None => Vec::new(),
});
for task in ready {
task();
}
}
pub(crate) fn close() {
if let Ok(mut sender) = wakeups_out().lock() {
*sender = None;
}
}