use std::{
collections::HashMap,
sync::{
Arc, Mutex, MutexGuard, OnceLock, PoisonError,
atomic::{AtomicU64, Ordering},
},
};
use cranpose_core::{EventStream, rememberEventStream};
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct HostMessage {
pub channel: String,
pub payload: String,
}
impl HostMessage {
pub fn new(channel: impl Into<String>, payload: impl Into<String>) -> Self {
Self {
channel: channel.into(),
payload: payload.into(),
}
}
}
type Observer = Arc<dyn Fn(String) + Send + Sync>;
type Outbox = Arc<dyn Fn(HostMessage) + Send + Sync>;
struct Mailbox {
observers: Vec<(u64, String, Observer)>,
latest: HashMap<String, String>,
}
impl Mailbox {
fn new() -> Self {
Self {
observers: Vec::new(),
latest: HashMap::new(),
}
}
fn observe(&mut self, id: u64, channel: &str, observer: Observer) -> Option<String> {
self.observers.push((id, channel.to_owned(), observer));
self.latest.get(channel).cloned()
}
fn publish(&mut self, message: &HostMessage) -> Vec<Observer> {
self.latest
.insert(message.channel.clone(), message.payload.clone());
self.observers
.iter()
.filter(|(_, channel, _)| *channel == message.channel)
.map(|(_, _, observer)| Arc::clone(observer))
.collect()
}
fn remove_observer(&mut self, id: u64) {
self.observers.retain(|(existing, _, _)| *existing != id);
}
fn clear(&mut self) {
self.latest.clear();
}
}
fn mailbox() -> MutexGuard<'static, Mailbox> {
static MAILBOX: OnceLock<Mutex<Mailbox>> = OnceLock::new();
MAILBOX
.get_or_init(|| Mutex::new(Mailbox::new()))
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
fn outbox() -> MutexGuard<'static, Option<Outbox>> {
static OUTBOX: OnceLock<Mutex<Option<Outbox>>> = OnceLock::new();
OUTBOX
.get_or_init(|| Mutex::new(None))
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
static NEXT_ID: AtomicU64 = AtomicU64::new(1);
pub struct HostMessageObserver {
id: u64,
}
impl Drop for HostMessageObserver {
fn drop(&mut self) {
mailbox().remove_observer(self.id);
}
}
pub fn observe_host_messages(
channel: &str,
observer: impl Fn(String) + Send + Sync + 'static,
) -> HostMessageObserver {
let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
let observer: Observer = Arc::new(observer);
let replay = mailbox().observe(id, channel, Arc::clone(&observer));
if let Some(payload) = replay {
observer(payload);
}
HostMessageObserver { id }
}
pub fn publish_host_message(message: HostMessage) {
let observers = mailbox().publish(&message);
for observer in observers {
observer(message.payload.clone());
}
}
pub fn clear_host_messages() {
mailbox().clear();
}
pub fn install_host_outbox(outbox_fn: impl Fn(HostMessage) + Send + Sync + 'static) {
*outbox() = Some(Arc::new(outbox_fn));
}
pub fn clear_host_outbox() {
*outbox() = None;
}
pub fn send_to_host(channel: &str, payload: &str) -> bool {
let Some(outbox_fn) = outbox().clone() else {
log::debug!("host message on {channel} dropped: no host is attached");
return false;
};
outbox_fn(HostMessage::new(channel, payload));
true
}
#[expect(non_snake_case)]
#[track_caller]
pub fn rememberHostMessages(channel: &str) -> EventStream<String> {
let channel = channel.to_owned();
rememberEventStream(channel.clone(), move |sender| {
observe_host_messages(&channel, move |payload| sender.send(payload))
})
}
#[cfg(test)]
#[path = "tests/host_messages_tests.rs"]
mod tests;