use std::cell::{Cell, RefCell};
use std::collections::BTreeMap;
use std::rc::Rc;
use crate::tree::Hash;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct RootAdvance {
pub shard_id: usize,
#[serde(with = "hash_serde")]
pub prior_root: Hash,
#[serde(with = "hash_serde")]
pub new_root: Hash,
pub advance_gen: u64,
}
#[derive(Debug, Default)]
struct ShardEmit {
assigned_gen: u64,
last_told_gen: u64,
}
struct Subscriber {
id: u64,
callback: Rc<dyn Fn(RootAdvance)>,
}
#[derive(Default)]
struct SeamInner {
next_id: u64,
subscribers: Vec<Subscriber>,
shards: BTreeMap<usize, ShardEmit>,
}
#[derive(Clone, Default)]
pub(super) struct BrowserRootAdvanceSeam {
inner: Rc<RefCell<SeamInner>>,
}
thread_local! {
static IN_EMISSION: Cell<bool> = const { Cell::new(false) };
}
pub(super) fn in_emission() -> bool {
IN_EMISSION.with(Cell::get)
}
struct InEmissionGuard {
previous: bool,
}
impl InEmissionGuard {
fn enter() -> Self {
Self {
previous: IN_EMISSION.with(|flag| flag.replace(true)),
}
}
}
impl Drop for InEmissionGuard {
fn drop(&mut self) {
IN_EMISSION.with(|flag| flag.set(self.previous));
}
}
impl BrowserRootAdvanceSeam {
pub(super) fn new() -> Self {
Self::default()
}
fn subscribe(&self, callback: Rc<dyn Fn(RootAdvance)>) -> u64 {
let mut inner = self.inner.borrow_mut();
let id = inner.next_id;
inner.next_id = inner.next_id.wrapping_add(1);
inner.subscribers.push(Subscriber { id, callback });
id
}
fn cancel(&self, id: u64) {
self.inner
.borrow_mut()
.subscribers
.retain(|subscriber| subscriber.id != id);
}
pub(super) fn emit_advance(&self, shard_id: usize, prior_root: Hash, new_root: Hash) {
let (advance_gen, subscribers) = {
let mut inner = self.inner.borrow_mut();
let shard = inner.shards.entry(shard_id).or_default();
shard.assigned_gen = shard.assigned_gen.wrapping_add(1);
let advance_gen = shard.assigned_gen;
shard.last_told_gen = advance_gen;
let subscribers: Vec<Rc<dyn Fn(RootAdvance)>> = inner
.subscribers
.iter()
.map(|subscriber| Rc::clone(&subscriber.callback))
.collect();
(advance_gen, subscribers)
};
if subscribers.is_empty() {
return;
}
let event = RootAdvance {
shard_id,
prior_root,
new_root,
advance_gen,
};
let _wall = InEmissionGuard::enter();
for callback in &subscribers {
callback(event);
}
}
}
#[must_use = "dropping the subscription immediately unsubscribes"]
pub struct BrowserRootAdvanceSubscription {
seam: BrowserRootAdvanceSeam,
id: u64,
active: bool,
}
impl BrowserRootAdvanceSubscription {
pub(super) fn new(seam: &BrowserRootAdvanceSeam, callback: Rc<dyn Fn(RootAdvance)>) -> Self {
let id = seam.subscribe(callback);
Self {
seam: seam.clone(),
id,
active: true,
}
}
pub fn cancel(mut self) {
self.deactivate();
}
fn deactivate(&mut self) {
if self.active {
self.active = false;
self.seam.cancel(self.id);
}
}
}
impl Drop for BrowserRootAdvanceSubscription {
fn drop(&mut self) {
self.deactivate();
}
}
mod hash_serde {
use crate::tree::Hash;
use serde::{Deserialize, Deserializer, Serialize, Serializer};
pub(super) fn serialize<S: Serializer>(hash: &Hash, serializer: S) -> Result<S::Ok, S::Error> {
hash.as_bytes().serialize(serializer)
}
pub(super) fn deserialize<'de, D: Deserializer<'de>>(
deserializer: D,
) -> Result<Hash, D::Error> {
let bytes = <[u8; crate::tree::node::HASH_SIZE]>::deserialize(deserializer)?;
Ok(Hash::from_bytes(bytes))
}
}