use std::cell::Cell;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use crate::CrdtDoc;
pub struct Watcher {
revision: Arc<AtomicU64>,
last_seen: Cell<u64>,
}
impl Watcher {
pub(crate) fn new(revision: Arc<AtomicU64>) -> Self {
let current = revision.load(Ordering::Relaxed);
Self {
revision,
last_seen: Cell::new(current),
}
}
pub fn drain_changed(&self) -> bool {
let current = self.revision.load(Ordering::Relaxed);
if current != self.last_seen.get() {
self.last_seen.set(current);
true
} else {
false
}
}
pub fn revision(&self) -> u64 {
self.revision.load(Ordering::Relaxed)
}
}
pub struct ReactiveQuery<R, F> {
watcher: Watcher,
query: F,
last: R,
}
impl<R, F> ReactiveQuery<R, F>
where
R: Clone + PartialEq,
F: Fn(&CrdtDoc) -> R,
{
pub fn new(doc: &CrdtDoc, query: F) -> Self {
let watcher = doc.watch();
let last = query(doc);
Self {
watcher,
query,
last,
}
}
pub fn poll(&mut self, doc: &CrdtDoc) -> Option<R> {
if !self.watcher.drain_changed() {
return None;
}
let next = (self.query)(doc);
if next != self.last {
self.last = next.clone();
Some(next)
} else {
None
}
}
pub fn current(&self) -> &R {
&self.last
}
}