use std::marker::PhantomData;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::watch;
pub trait ReactiveRepr: Copy {
fn to_usize(self) -> usize;
fn from_usize(value: usize) -> Self;
}
impl ReactiveRepr for usize {
#[inline]
fn to_usize(self) -> usize {
self
}
#[inline]
fn from_usize(value: usize) -> Self {
value
}
}
#[derive(Debug)]
pub struct Reactive<T> {
value: AtomicUsize,
signal: watch::Sender<usize>,
_repr: PhantomData<fn() -> T>,
}
impl<T: ReactiveRepr> Reactive<T> {
#[must_use]
pub fn new(value: T) -> Self {
let bits = value.to_usize();
let (signal, _) = watch::channel(bits);
Self {
value: AtomicUsize::new(bits),
signal,
_repr: PhantomData,
}
}
#[must_use]
pub fn get(&self) -> T {
T::from_usize(self.value.load(Ordering::Acquire))
}
pub fn set(&self, value: T) {
let bits = value.to_usize();
self.value.store(bits, Ordering::Release);
let _unused = self.signal.send(bits);
}
#[must_use]
pub fn watch(&self) -> Changed<T> {
Changed {
rx: self.signal.subscribe(),
_repr: PhantomData,
}
}
}
impl<T: ReactiveRepr + Default> Default for Reactive<T> {
fn default() -> Self {
Self::new(T::default())
}
}
#[derive(Debug, Clone)]
pub struct Changed<T> {
rx: watch::Receiver<usize>,
_repr: PhantomData<fn() -> T>,
}
impl<T: ReactiveRepr> Changed<T> {
pub async fn changed(&mut self) -> Option<T> {
self.rx.changed().await.ok()?;
Some(T::from_usize(*self.rx.borrow_and_update()))
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
#[test]
fn get_set_roundtrip() {
let r = Reactive::<usize>::new(3);
assert_eq!(r.get(), 3);
r.set(7);
assert_eq!(r.get(), 7);
}
#[tokio::test]
async fn changed_yields_the_new_value() {
let r = Arc::new(Reactive::<usize>::new(0));
let mut w = r.watch();
let handle = tokio::spawn(async move { w.changed().await });
tokio::task::yield_now().await;
r.set(42);
assert_eq!(handle.await.unwrap(), Some(42));
}
#[tokio::test]
async fn changed_returns_none_once_source_dropped() {
let r = Reactive::<usize>::new(0);
let mut w = r.watch();
drop(r);
assert_eq!(
w.changed().await,
None,
"no source left: should report closed"
);
}
#[tokio::test]
async fn set_without_watchers_is_a_noop_send() {
let r = Reactive::<usize>::new(1);
r.set(2);
assert_eq!(r.get(), 2);
let mut w = r.watch();
r.set(3);
assert_eq!(w.changed().await, Some(3));
}
}