#![expect(clippy::let_underscore_must_use, reason = "a missed wake costs a tick, not a result")]
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, SendError, Sender, channel};
use std::sync::{Arc, OnceLock};
use kevy_sys::Waker;
use crate::Argv;
#[repr(align(64))]
#[derive(Debug)]
pub(crate) struct InboxSignal {
pub(crate) waker: OnceLock<Arc<Waker>>,
pub(crate) wake_pending: AtomicBool,
}
#[derive(Clone)]
pub struct SnapshotGate(#[allow(dead_code)] Arc<dyn std::any::Any + Send + Sync>);
impl SnapshotGate {
#[must_use]
pub fn new(inner: Arc<dyn std::any::Any + Send + Sync>) -> Self {
Self(inner)
}
}
impl std::fmt::Debug for SnapshotGate {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("SnapshotGate")
}
}
#[derive(Debug)]
pub enum ReplicaApply {
SnapshotBegin,
SnapshotChunk(Vec<u8>),
SnapshotEnd {
ack_offset: u64,
routed: bool,
gate: Option<SnapshotGate>,
},
Frame {
offset: u64,
argv: Argv,
},
}
#[derive(Debug, Clone)]
pub struct ReplicaInboxSender {
inner: Sender<ReplicaApply>,
signal: Arc<InboxSignal>,
}
impl ReplicaInboxSender {
pub fn send(&self, ev: ReplicaApply) -> Result<(), SendError<ReplicaApply>> {
self.inner.send(ev)?;
if !self.signal.wake_pending.swap(true, Ordering::AcqRel)
&& let Some(w) = self.signal.waker.get()
{
let _ = w.wake();
}
Ok(())
}
}
#[derive(Debug)]
pub struct ReplicaInboxReceiver {
pub(crate) inner: Receiver<ReplicaApply>,
pub(crate) signal: Arc<InboxSignal>,
}
impl ReplicaInboxReceiver {
pub(crate) fn attach_waker(&self, waker: Arc<Waker>) {
let _ = self.signal.waker.set(waker);
}
}
#[must_use]
pub fn replica_inbox_pair() -> (ReplicaInboxSender, ReplicaInboxReceiver) {
let (tx, rx) = channel();
let signal =
Arc::new(InboxSignal { waker: OnceLock::new(), wake_pending: AtomicBool::new(false) });
(
ReplicaInboxSender { inner: tx, signal: Arc::clone(&signal) },
ReplicaInboxReceiver { inner: rx, signal },
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pair_round_trips_one_event() {
let (tx, rx) = replica_inbox_pair();
tx.send(ReplicaApply::SnapshotBegin).unwrap();
match rx.inner.recv().unwrap() {
ReplicaApply::SnapshotBegin => {}
other => panic!("expected SnapshotBegin, got {other:?}"),
}
}
#[test]
fn drop_receiver_makes_send_fail() {
let (tx, rx) = replica_inbox_pair();
drop(rx);
let err = tx.send(ReplicaApply::SnapshotBegin).unwrap_err();
match err.0 {
ReplicaApply::SnapshotBegin => {}
other => panic!("expected payload roundtrip, got {other:?}"),
}
}
#[test]
fn sender_is_clone_send_sync() {
fn assert_traits<T: Clone + Send + Sync>() {}
assert_traits::<ReplicaInboxSender>();
}
}