#[cfg(feature = "metrics")]
use std::time::Duration;
use std::{
fmt::Debug,
marker::PhantomData,
panic::Location,
sync::{
Arc, Mutex, Weak,
atomic::{AtomicBool, Ordering},
},
};
#[cfg(feature = "metrics")]
use arc_swap::ArcSwap;
use dashmap::DashMap;
use rustc_hash::FxHashMap;
use uuid::Uuid;
#[cfg(feature = "metrics")]
use crate::metrics::CellMetrics;
use crate::{
signal::Signal,
subscription::SubscriptionGuard,
traits::{CellValue, DepNode, Gettable, Mutable, Watchable, WatchableResult},
};
#[cfg(feature = "metrics")]
#[derive(Debug, Clone)]
pub struct SlowSubscriberAlert {
pub subscriber_id: Uuid,
pub duration_ns: u64,
pub threshold_ns: u64,
}
#[cfg(feature = "metrics")]
type SlowSubscriberCallback = Arc<dyn Fn(SlowSubscriberAlert) + Send + Sync>;
#[derive(Debug, Clone)]
pub struct CellMutable;
#[derive(Debug, Clone)]
pub struct CellImmutable;
pub(crate) struct CellInner<T> {
pub(crate) id: Uuid,
pub(crate) subscribers: parking_lot::Mutex<SubscriberRegistry<Subscriber<T>>>,
pub(crate) result_subscribers: parking_lot::Mutex<SubscriberRegistry<ResultSubscriber<T>>>,
pub(crate) value: Mutex<Arc<T>>,
pub(crate) name: Mutex<Option<Arc<str>>>,
pub(crate) owned: DashMap<Uuid, SubscriptionGuard>,
pub(crate) completed: AtomicBool,
pub(crate) errored: AtomicBool,
pub(crate) error: Mutex<Option<Arc<anyhow::Error>>>,
#[cfg(feature = "scheduler")]
pub(crate) height_cache: std::sync::atomic::AtomicU64,
#[cfg(feature = "scheduler")]
pub(crate) no_coalesce: AtomicBool,
#[cfg(feature = "scheduler")]
pub(crate) height_epoch: std::sync::atomic::AtomicU64,
#[cfg(feature = "scheduler")]
pub(crate) height_dependents: Mutex<Vec<std::sync::Weak<dyn HeightInvalidate>>>,
#[cfg(feature = "metrics")]
pub(crate) metrics: Option<Arc<CellMetrics>>,
#[cfg(feature = "metrics")]
pub(crate) slow_subscriber_threshold_ns: ArcSwap<Option<u64>>,
#[cfg(feature = "metrics")]
pub(crate) slow_subscriber_callback: ArcSwap<Option<SlowSubscriberCallback>>,
#[allow(dead_code)]
pub(crate) caller: &'static Location<'static>,
}
pub struct Cell<T, M> {
pub(crate) inner: Arc<CellInner<T>>,
pub(crate) _marker: PhantomData<M>,
}
pub struct WeakCell<T, M> {
inner: Weak<CellInner<T>>,
_marker: PhantomData<M>,
}
impl<T, M> WeakCell<T, M> {
pub fn upgrade(&self) -> Option<Cell<T, M>> {
self.inner.upgrade().map(|inner| Cell {
inner,
_marker: PhantomData,
})
}
pub fn is_alive(&self) -> bool {
self.inner.strong_count() > 0
}
}
impl<T, M> Clone for WeakCell<T, M> {
fn clone(&self) -> Self {
WeakCell {
inner: self.inner.clone(),
_marker: PhantomData,
}
}
}
pub(crate) enum SubSnapshot<S> {
Zero,
One((Uuid, Arc<S>)),
Many(Arc<Vec<(Uuid, Arc<S>)>>),
}
impl<S> Clone for SubSnapshot<S> {
fn clone(&self) -> Self {
match self {
SubSnapshot::Zero => SubSnapshot::Zero,
SubSnapshot::One(pair) => SubSnapshot::One(pair.clone()),
SubSnapshot::Many(subs) => SubSnapshot::Many(subs.clone()),
}
}
}
impl<S> SubSnapshot<S> {
pub(crate) fn as_slice(&self) -> &[(Uuid, Arc<S>)] {
match self {
SubSnapshot::Zero => &[],
SubSnapshot::One(pair) => std::slice::from_ref(pair),
SubSnapshot::Many(subs) => subs.as_slice(),
}
}
}
enum SubIndex<S> {
Zero,
One(Uuid, Arc<S>),
Many(FxHashMap<Uuid, Arc<S>>),
}
impl<S> SubIndex<S> {
fn len(&self) -> usize {
match self {
SubIndex::Zero => 0,
SubIndex::One(..) => 1,
SubIndex::Many(map) => map.len(),
}
}
fn insert(&mut self, id: Uuid, sub: Arc<S>) -> Option<Arc<S>> {
match self {
SubIndex::Zero => {
*self = SubIndex::One(id, sub);
return None;
}
SubIndex::One(existing_id, existing_sub) => {
if *existing_id == id {
return Some(std::mem::replace(existing_sub, sub));
}
}
SubIndex::Many(map) => {
return map.insert(id, sub);
}
}
let (old_id, old_sub) = match std::mem::replace(self, SubIndex::Zero) {
SubIndex::One(old_id, old_sub) => (old_id, old_sub),
_ => unreachable!("promotion is entered only from the One arm"),
};
let mut map = FxHashMap::default();
map.insert(old_id, old_sub);
map.insert(id, sub);
*self = SubIndex::Many(map);
None
}
fn remove(&mut self, id: &Uuid) -> Option<Arc<S>> {
match self {
SubIndex::Zero => None,
SubIndex::One(existing_id, _) => {
if *existing_id != *id {
return None;
}
match std::mem::replace(self, SubIndex::Zero) {
SubIndex::One(_, sub) => Some(sub),
_ => unreachable!("just matched One"),
}
}
SubIndex::Many(map) => map.remove(id),
}
}
}
pub(crate) struct SubscriberRegistry<S> {
index: SubIndex<S>,
snapshot: SubSnapshot<S>,
dirty: bool,
}
impl<S> SubscriberRegistry<S> {
fn new() -> Self {
Self {
index: SubIndex::Zero,
snapshot: SubSnapshot::Zero,
dirty: false,
}
}
#[must_use = "displaced subscriber must be dropped outside the lock"]
fn insert(&mut self, id: Uuid, sub: Arc<S>) -> Option<Arc<S>> {
self.dirty = true;
self.index.insert(id, sub)
}
#[must_use = "removed subscriber and stale snapshot must be dropped outside the lock"]
fn remove(&mut self, id: &Uuid) -> (Option<Arc<S>>, Option<SubSnapshot<S>>) {
let removed = self.index.remove(id);
if removed.is_some() {
self.dirty = true;
let stale = std::mem::replace(&mut self.snapshot, SubSnapshot::Zero);
(removed, Some(stale))
} else {
(removed, None)
}
}
pub(crate) fn len(&self) -> usize {
self.index.len()
}
#[must_use = "displaced snapshot must be dropped outside the lock"]
pub(crate) fn snapshot(&mut self) -> (SubSnapshot<S>, Option<SubSnapshot<S>>) {
if self.dirty {
let next = match &self.index {
SubIndex::Zero => SubSnapshot::Zero,
SubIndex::One(id, sub) => SubSnapshot::One((*id, sub.clone())),
SubIndex::Many(map) => SubSnapshot::Many(Arc::new(
map.iter().map(|(id, sub)| (*id, sub.clone())).collect(),
)),
};
let old = std::mem::replace(&mut self.snapshot, next);
self.dirty = false;
(self.snapshot.clone(), Some(old))
} else {
(self.snapshot.clone(), None)
}
}
}
pub(crate) type SubscriberCallback<T> = Arc<dyn Fn(&Signal<T>) + Send + Sync>;
pub(crate) struct Subscriber<T> {
pub(crate) callback: SubscriberCallback<T>,
}
impl<T> Subscriber<T> {
pub(crate) fn new(callback: impl Fn(&Signal<T>) + Send + Sync + 'static) -> Self {
Self {
callback: Arc::new(callback),
}
}
}
pub(crate) type ResultSubscriberCallback<T> =
Arc<dyn Fn(&Signal<T>) -> Result<(), String> + Send + Sync>;
pub(crate) struct ResultSubscriber<T> {
pub(crate) callback: ResultSubscriberCallback<T>,
}
impl<T> ResultSubscriber<T> {
pub(crate) fn new(
callback: impl Fn(&Signal<T>) -> Result<(), String> + Send + Sync + 'static,
) -> Self {
Self {
callback: Arc::new(callback),
}
}
}
impl<T: CellValue> Cell<T, CellMutable> {
#[track_caller]
pub fn new(initial_value: T) -> Self {
let inner = Arc::new(CellInner {
id: Uuid::new_v4(),
subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
result_subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
value: Mutex::new(Arc::new(initial_value)),
name: Mutex::new(None),
owned: DashMap::new(),
completed: AtomicBool::new(false),
errored: AtomicBool::new(false),
error: Mutex::new(None),
#[cfg(feature = "scheduler")]
height_cache: std::sync::atomic::AtomicU64::new(0),
#[cfg(feature = "scheduler")]
no_coalesce: AtomicBool::new(crate::scheduler::birth_no_coalesce()),
#[cfg(feature = "scheduler")]
height_epoch: std::sync::atomic::AtomicU64::new(1),
#[cfg(feature = "scheduler")]
height_dependents: Mutex::new(Vec::new()),
#[cfg(feature = "metrics")]
metrics: default_metrics(),
#[cfg(feature = "metrics")]
slow_subscriber_threshold_ns: ArcSwap::from_pointee(None),
#[cfg(feature = "metrics")]
slow_subscriber_callback: ArcSwap::from_pointee(None),
caller: Location::caller(),
});
#[cfg(feature = "inspector")]
crate::registry::registry().register(inner.id, Arc::downgrade(&inner) as Weak<dyn DepNode>);
#[cfg(feature = "trace")]
crate::tracing::register_cell(inner.id, Some(Location::caller().to_string()));
Self {
inner,
_marker: PhantomData,
}
}
#[cfg(feature = "metrics")]
#[track_caller]
pub fn with_metrics(initial_value: T) -> Self {
let inner = Arc::new(CellInner {
id: Uuid::new_v4(),
subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
result_subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
value: Mutex::new(Arc::new(initial_value)),
name: Mutex::new(None),
owned: DashMap::new(),
completed: AtomicBool::new(false),
errored: AtomicBool::new(false),
error: Mutex::new(None),
#[cfg(feature = "scheduler")]
height_cache: std::sync::atomic::AtomicU64::new(0),
#[cfg(feature = "scheduler")]
no_coalesce: AtomicBool::new(crate::scheduler::birth_no_coalesce()),
#[cfg(feature = "scheduler")]
height_epoch: std::sync::atomic::AtomicU64::new(1),
#[cfg(feature = "scheduler")]
height_dependents: Mutex::new(Vec::new()),
metrics: Some(Arc::new(CellMetrics::new())),
slow_subscriber_threshold_ns: ArcSwap::from_pointee(None),
slow_subscriber_callback: ArcSwap::from_pointee(None),
caller: Location::caller(),
});
#[cfg(feature = "inspector")]
crate::registry::registry().register(inner.id, Arc::downgrade(&inner) as Weak<dyn DepNode>);
#[cfg(feature = "trace")]
crate::tracing::register_cell(inner.id, Some(Location::caller().to_string()));
Self {
inner,
_marker: PhantomData,
}
}
#[cfg(feature = "metrics")]
pub fn on_slow_subscriber<F>(&self, threshold: Duration, callback: F)
where
F: Fn(SlowSubscriberAlert) + Send + Sync + 'static,
{
self.inner
.slow_subscriber_threshold_ns
.store(Arc::new(Some(threshold.as_nanos() as u64)));
self.inner
.slow_subscriber_callback
.store(Arc::new(Some(Arc::new(callback))));
}
pub fn lock(self) -> Cell<T, CellImmutable> {
Cell {
inner: self.inner,
_marker: PhantomData,
}
}
pub fn with_name(self, name: impl Into<Arc<str>>) -> Self {
let name = name.into();
*self.inner.name.lock().expect("cell name poisoned") = Some(name.clone());
#[cfg(feature = "trace")]
crate::tracing::update_name(self.inner.id, name.to_string());
self
}
#[cfg(feature = "metrics")]
pub fn is_backed_up(&self) -> bool {
self.is_backed_up_threshold(std::time::Duration::from_millis(1))
}
#[cfg(feature = "metrics")]
pub fn is_backed_up_threshold(&self, threshold: std::time::Duration) -> bool {
self.inner
.metrics
.as_ref()
.map(|m| m.last_notify_time_ns() > threshold.as_nanos() as u64)
.unwrap_or(false)
}
#[cfg(feature = "metrics")]
pub fn try_set(&self, value: T) -> Result<(), T> {
if self.is_backed_up() {
Err(value)
} else {
self.set(value);
Ok(())
}
}
#[cfg(feature = "metrics")]
pub fn try_set_threshold(&self, value: T, threshold: std::time::Duration) -> Result<(), T> {
if self.is_backed_up_threshold(threshold) {
Err(value)
} else {
self.set(value);
Ok(())
}
}
}
impl<T, M> Clone for Cell<T, M> {
fn clone(&self) -> Self {
Cell {
inner: Arc::clone(&self.inner),
_marker: PhantomData,
}
}
}
impl<T, M> Cell<T, M> {
pub fn downgrade(&self) -> WeakCell<T, M> {
WeakCell {
inner: Arc::downgrade(&self.inner),
_marker: PhantomData,
}
}
#[cfg(feature = "metrics")]
pub fn metrics(&self) -> Option<&CellMetrics> {
self.inner.metrics.as_ref().map(|m| m.as_ref())
}
pub fn own(&self, guard: SubscriptionGuard)
where
T: Send + Sync + 'static,
{
#[cfg(feature = "inspector")]
crate::registry::registry().mark_owned(guard.source().id(), self.inner.id);
#[cfg(feature = "scheduler")]
{
let dep: std::sync::Weak<dyn HeightInvalidate> =
Arc::downgrade(&(self.inner.clone() as Arc<dyn HeightInvalidate>));
guard.source().add_height_dependent(dep);
}
self.inner.owned.insert(Uuid::new_v4(), guard);
#[cfg(feature = "scheduler")]
invalidate_height_cone(self.inner.as_ref());
#[cfg(feature = "trace")]
crate::tracing::update_owned_count(self.inner.id, self.inner.owned.len());
}
pub fn own_keyed(&self, key: Uuid, guard: SubscriptionGuard)
where
T: Send + Sync + 'static,
{
#[cfg(feature = "inspector")]
{
if let Some((_, old_guard)) = self.inner.owned.remove(&key) {
crate::registry::registry().unmark_owned(old_guard.source().id());
}
crate::registry::registry().mark_owned(guard.source().id(), self.inner.id);
}
#[cfg(feature = "scheduler")]
{
let dep: std::sync::Weak<dyn HeightInvalidate> =
Arc::downgrade(&(self.inner.clone() as Arc<dyn HeightInvalidate>));
guard.source().add_height_dependent(dep);
}
self.inner.owned.insert(key, guard);
#[cfg(feature = "scheduler")]
invalidate_height_cone(self.inner.as_ref());
#[cfg(feature = "trace")]
crate::tracing::update_owned_count(self.inner.id, self.inner.owned.len());
}
}
#[cfg(feature = "scheduler")]
impl<T, M> Cell<T, M> {
pub fn no_coalesce(self) -> Self {
self.inner.no_coalesce.store(true, Ordering::Relaxed);
self
}
}
#[doc(hidden)]
#[cfg(feature = "scheduler")]
pub trait HeightInvalidate: Send + Sync {
fn hi_id(&self) -> Uuid;
fn hi_bump_epoch(&self);
fn hi_dependents(&self) -> Vec<std::sync::Weak<dyn HeightInvalidate>>;
}
#[cfg(feature = "scheduler")]
impl<T: Send + Sync> HeightInvalidate for CellInner<T> {
fn hi_id(&self) -> Uuid {
self.id
}
fn hi_bump_epoch(&self) {
self.height_epoch.fetch_add(1, Ordering::Relaxed);
}
fn hi_dependents(&self) -> Vec<std::sync::Weak<dyn HeightInvalidate>> {
let mut g = self
.height_dependents
.lock()
.unwrap_or_else(|p| p.into_inner());
g.retain(|w| w.strong_count() > 0);
g.clone()
}
}
#[cfg(feature = "scheduler")]
pub(crate) fn invalidate_height_cone(start: &dyn HeightInvalidate) {
start.hi_bump_epoch();
let mut visited = std::collections::HashSet::new();
visited.insert(start.hi_id());
let mut stack = start.hi_dependents();
while let Some(weak) = stack.pop() {
let Some(node) = weak.upgrade() else { continue };
if visited.insert(node.hi_id()) {
node.hi_bump_epoch();
stack.extend(node.hi_dependents());
}
}
}
impl<T: Send + Sync, M: Send + Sync> DepNode for Cell<T, M> {
fn id(&self) -> Uuid {
self.inner.id
}
fn name(&self) -> Option<String> {
self.inner
.name
.lock()
.expect("cell name poisoned")
.as_ref()
.map(|s| s.to_string())
}
fn deps(&self) -> Vec<Arc<dyn DepNode>> {
let mut seen = std::collections::HashSet::new();
self.inner
.owned
.iter()
.filter_map(|entry| {
let source = entry.value().source();
let id = source.id();
if seen.insert(id) {
Some(Arc::clone(source))
} else {
None
}
})
.collect()
}
#[cfg(feature = "scheduler")]
fn height_cache(&self) -> Option<&std::sync::atomic::AtomicU64> {
Some(&self.inner.height_cache)
}
#[cfg(feature = "scheduler")]
fn height_epoch(&self) -> Option<&std::sync::atomic::AtomicU64> {
Some(&self.inner.height_epoch)
}
#[cfg(feature = "scheduler")]
fn add_height_dependent(&self, dep: std::sync::Weak<dyn HeightInvalidate>) {
self.inner
.height_dependents
.lock()
.unwrap_or_else(|p| p.into_inner())
.push(dep);
}
#[cfg(feature = "scheduler")]
fn no_coalesce(&self) -> bool {
self.inner.no_coalesce.load(Ordering::Relaxed)
}
fn subscriber_count(&self) -> usize {
self.inner.subscribers.lock().len() + self.inner.result_subscribers.lock().len()
}
fn owned_count(&self) -> usize {
self.inner.owned.len()
}
}
impl<T: CellValue> Cell<T, CellImmutable> {
pub fn with_name(self, name: impl Into<Arc<str>>) -> Self {
let name = name.into();
*self.inner.name.lock().expect("cell name poisoned") = Some(name.clone());
#[cfg(feature = "trace")]
crate::tracing::update_name(self.inner.id, name.to_string());
self
}
}
impl<T: CellValue, M: Send + Sync + 'static> Cell<T, M> {
#[doc(hidden)]
#[cfg_attr(feature = "profiling", inline(never))]
pub fn notify(&self, signal: Signal<T>) {
if self.inner.completed.load(Ordering::SeqCst) || self.inner.errored.load(Ordering::SeqCst)
{
return;
}
#[cfg(feature = "scheduler")]
if crate::scheduler::tick_active() {
let cell = self.clone();
let terminal = !matches!(signal, Signal::Value(_));
let signal = signal.clone();
crate::scheduler::enqueue(
self.inner.id,
self as &dyn crate::traits::DepNode,
terminal,
Box::new(move || {
cell.write_value(&signal);
cell.fanout(&signal);
}),
);
return;
}
self.write_value(&signal);
self.fanout(&signal);
}
#[cfg_attr(feature = "profiling", inline(never))]
fn write_value(&self, signal: &Signal<T>) {
match signal {
Signal::Value(arc_value) => {
*self.inner.value.lock().expect("cell value poisoned") = arc_value.clone();
}
Signal::Complete => {
self.inner.completed.store(true, Ordering::SeqCst);
}
Signal::Error(err) => {
self.inner.errored.store(true, Ordering::SeqCst);
*self.inner.error.lock().expect("cell error poisoned") = Some(err.clone());
}
}
}
#[cfg_attr(feature = "profiling", inline(never))]
fn fanout(&self, signal: &Signal<T>) {
#[cfg(feature = "profiling")]
crate::profiling::record_fire(self.inner.id);
#[cfg(feature = "profiling")]
let _fanout_span = {
let name = self.inner.name.lock().expect("cell name poisoned").clone();
::tracing::trace_span!(
"hyphae.fanout",
cell.id = %self.inner.id,
cell.name = name.as_deref().unwrap_or(""),
)
.entered()
};
#[cfg(feature = "metrics")]
let notify_start = self
.inner
.metrics
.as_ref()
.map(|_| crate::platform::Instant::now());
let subs = {
let (subs, old_snapshot) = self.inner.subscribers.lock().snapshot();
drop(old_snapshot);
subs
};
#[cfg(feature = "metrics")]
let metrics = &self.inner.metrics;
#[cfg(feature = "metrics")]
let (slow_threshold, slow_callback) = if metrics.is_some() {
(
**self.inner.slow_subscriber_threshold_ns.load(),
(**self.inner.slow_subscriber_callback.load()).clone(),
)
} else {
(None, None)
};
for (_subscriber_id, sub) in subs.as_slice() {
#[cfg(feature = "metrics")]
let sub_start = metrics.as_ref().map(|_| crate::platform::Instant::now());
(sub.callback)(signal);
#[cfg(feature = "metrics")]
if let (Some(m), Some(start)) = (metrics, sub_start) {
let elapsed = start.elapsed().as_nanos() as u64;
m.update_slowest_subscriber(elapsed);
if let (Some(threshold), Some(cb)) = (&slow_threshold, &slow_callback)
&& elapsed > *threshold
{
let alert = SlowSubscriberAlert {
subscriber_id: *_subscriber_id,
duration_ns: elapsed,
threshold_ns: *threshold,
};
cb(alert);
}
}
}
let result_subs = {
let (result_subs, old_snapshot) = self.inner.result_subscribers.lock().snapshot();
drop(old_snapshot);
result_subs
};
for (subscriber_id, sub) in result_subs.as_slice() {
#[cfg(feature = "metrics")]
let sub_start = metrics.as_ref().map(|_| crate::platform::Instant::now());
if let Err(err) = (sub.callback)(signal) {
log::error!(
"hyphae: fallible subscriber {} on cell {} returned error: {}",
subscriber_id,
self.inner.id,
err
);
}
#[cfg(feature = "metrics")]
if let (Some(m), Some(start)) = (metrics, sub_start) {
let elapsed = start.elapsed().as_nanos() as u64;
m.update_slowest_subscriber(elapsed);
if let (Some(threshold), Some(cb)) = (&slow_threshold, &slow_callback)
&& elapsed > *threshold
{
let alert = SlowSubscriberAlert {
subscriber_id: *subscriber_id,
duration_ns: elapsed,
threshold_ns: *threshold,
};
cb(alert);
}
}
}
#[cfg(feature = "metrics")]
if let (Some(metrics), Some(start)) = (&self.inner.metrics, notify_start) {
let duration_ns = start.elapsed().as_nanos() as u64;
metrics.record_notify(duration_ns);
#[cfg(feature = "trace")]
crate::tracing::record_notify(
self.inner.id,
duration_ns,
subs.as_slice().len() + result_subs.as_slice().len(),
self.inner.owned.len(),
metrics.slowest_subscriber_ns(),
);
}
}
}
impl<T: CellValue, U: Send + Sync + 'static> Gettable<T> for Cell<T, U> {
fn get(&self) -> T {
let arc = self
.inner
.value
.lock()
.expect("cell value poisoned")
.clone();
(*arc).clone()
}
}
impl<T: CellValue, U: Send + Sync + 'static> Watchable<T> for Cell<T, U> {
fn subscribe(
&self,
callback: impl Fn(&Signal<T>) + Send + Sync + 'static,
) -> SubscriptionGuard {
let id = Uuid::new_v4();
let sub = Arc::new(Subscriber::new(callback));
let displaced = self.inner.subscribers.lock().insert(id, sub.clone());
drop(displaced);
let current = self
.inner
.value
.lock()
.expect("cell value poisoned")
.clone();
(sub.callback)(&Signal::Value(current));
if self.is_complete() {
(sub.callback)(&Signal::Complete);
} else if self.is_error()
&& let Some(err) = self.error()
{
(sub.callback)(&Signal::Error(err));
}
#[cfg(feature = "metrics")]
if let Some(metrics) = &self.inner.metrics {
metrics.record_subscriber_added();
}
#[cfg(feature = "trace")]
{
let subs_len = self.inner.subscribers.lock().len();
let result_len = self.inner.result_subscribers.lock().len();
crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
}
let source: Arc<dyn DepNode> = Arc::new(self.clone());
let cell = self.clone();
#[cfg(feature = "metrics")]
let metrics = self.inner.metrics.clone();
SubscriptionGuard::new(id, source, move || {
let (removed_sub, stale_snap) = cell.inner.subscribers.lock().remove(&id);
let removed = removed_sub.is_some();
drop(removed_sub);
drop(stale_snap);
#[cfg(feature = "metrics")]
if removed && let Some(m) = &metrics {
m.record_subscriber_removed();
}
#[cfg(not(feature = "metrics"))]
let _ = removed;
#[cfg(feature = "trace")]
{
let subs_len = cell.inner.subscribers.lock().len();
let result_len = cell.inner.result_subscribers.lock().len();
crate::tracing::update_subscriber_count(cell.inner.id, subs_len + result_len);
}
})
}
fn unsubscribe(&self, id: Uuid) {
let (removed_sub, stale_snap) = self.inner.subscribers.lock().remove(&id);
let removed_from_subs = removed_sub.is_some();
drop(removed_sub);
drop(stale_snap);
let removed_from_result = if removed_from_subs {
false
} else {
let (removed, stale_result_snap) = self.inner.result_subscribers.lock().remove(&id);
let did = removed.is_some();
drop(removed);
drop(stale_result_snap);
did
};
if removed_from_subs || removed_from_result {
#[cfg(feature = "metrics")]
if let Some(metrics) = &self.inner.metrics {
metrics.record_subscriber_removed();
}
#[cfg(feature = "trace")]
{
let subs_len = self.inner.subscribers.lock().len();
let result_len = self.inner.result_subscribers.lock().len();
crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
}
}
}
fn is_complete(&self) -> bool {
self.inner.completed.load(Ordering::SeqCst)
}
fn is_error(&self) -> bool {
self.inner.errored.load(Ordering::SeqCst)
}
fn error(&self) -> Option<Arc<anyhow::Error>> {
self.inner
.error
.lock()
.expect("cell error poisoned")
.clone()
}
}
impl<T: CellValue, U: Send + Sync + 'static> WatchableResult<T> for Cell<T, U> {
fn subscribe_result(
&self,
callback: impl Fn(&Signal<T>) -> Result<(), String> + Send + Sync + 'static,
) -> SubscriptionGuard {
let cell_id = self.inner.id;
let log_err = |id: &Uuid, err: &str| {
log::error!(
"hyphae: fallible subscriber {} on cell {} returned error: {}",
id,
cell_id,
err
);
};
let id = Uuid::new_v4();
let current = self
.inner
.value
.lock()
.expect("cell value poisoned")
.clone();
if let Err(err) = callback(&Signal::Value(current)) {
log_err(&id, &err);
}
if self.inner.completed.load(Ordering::SeqCst) {
if let Err(err) = callback(&Signal::Complete) {
log_err(&id, &err);
}
} else if self.inner.errored.load(Ordering::SeqCst)
&& let Some(e) = self
.inner
.error
.lock()
.expect("cell error poisoned")
.clone()
&& let Err(err) = callback(&Signal::Error(e))
{
log_err(&id, &err);
}
let sub = Arc::new(ResultSubscriber::new(callback));
let displaced = self.inner.result_subscribers.lock().insert(id, sub);
drop(displaced);
#[cfg(feature = "metrics")]
if let Some(metrics) = &self.inner.metrics {
metrics.record_subscriber_added();
}
#[cfg(feature = "trace")]
{
let subs_len = self.inner.subscribers.lock().len();
let result_len = self.inner.result_subscribers.lock().len();
crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
}
let source: Arc<dyn DepNode> = Arc::new(self.clone());
let cell = self.clone();
#[cfg(feature = "metrics")]
let metrics = self.inner.metrics.clone();
SubscriptionGuard::new(id, source, move || {
let (removed_sub, stale_snap) = cell.inner.result_subscribers.lock().remove(&id);
let removed = removed_sub.is_some();
drop(removed_sub);
drop(stale_snap);
#[cfg(feature = "metrics")]
if removed && let Some(m) = &metrics {
m.record_subscriber_removed();
}
#[cfg(not(feature = "metrics"))]
let _ = removed;
#[cfg(feature = "trace")]
{
let subs_len = cell.inner.subscribers.lock().len();
let result_len = cell.inner.result_subscribers.lock().len();
crate::tracing::update_subscriber_count(cell.inner.id, subs_len + result_len);
}
})
}
}
impl<T: CellValue> Mutable<T> for Cell<T, CellMutable> {
fn set(&self, value: T) {
self.notify(Signal::value(value)); }
fn complete(&self) {
self.notify(Signal::Complete);
}
fn fail(&self, error: impl Into<anyhow::Error>) {
self.notify(Signal::error(error));
}
}
#[cfg(feature = "inspector")]
impl<T: CellValue> DepNode for CellInner<T> {
fn id(&self) -> Uuid {
self.id
}
fn name(&self) -> Option<String> {
self.name
.lock()
.expect("cell name poisoned")
.as_ref()
.map(|s| s.to_string())
}
fn deps(&self) -> Vec<Arc<dyn DepNode>> {
let mut seen = std::collections::HashSet::new();
self.owned
.iter()
.filter_map(|entry| {
let source = entry.value().source();
let id = source.id();
if seen.insert(id) {
Some(Arc::clone(source))
} else {
None
}
})
.collect()
}
fn subscriber_count(&self) -> usize {
self.subscribers.lock().len() + self.result_subscribers.lock().len()
}
fn owned_count(&self) -> usize {
self.owned.len()
}
fn value_debug(&self) -> Option<String> {
let arc = self.value.lock().expect("cell value poisoned").clone();
Some(format!("{:?}", *arc))
}
fn caller(&self) -> Option<&'static Location<'static>> {
Some(self.caller)
}
}
impl<T> Drop for CellInner<T> {
fn drop(&mut self) {
#[cfg(feature = "trace")]
crate::tracing::deregister_cell(&self.id);
#[cfg(feature = "inspector")]
crate::registry::registry().deregister(&self.id);
}
}
#[cfg(all(feature = "metrics", feature = "trace"))]
fn default_metrics() -> Option<Arc<CellMetrics>> {
Some(Arc::new(CellMetrics::new()))
}
#[cfg(all(feature = "metrics", not(feature = "trace")))]
fn default_metrics() -> Option<Arc<CellMetrics>> {
None
}
#[cfg(test)]
mod sub_index_tests {
use std::sync::Arc;
use uuid::Uuid;
use super::{SubIndex, SubSnapshot};
fn sub(v: i32) -> Arc<i32> {
Arc::new(v)
}
#[test]
fn zero_and_one_stay_inline() {
let mut idx: SubIndex<i32> = SubIndex::Zero;
assert!(matches!(idx, SubIndex::Zero));
assert_eq!(idx.len(), 0);
assert!(idx.insert(Uuid::new_v4(), sub(1)).is_none());
assert!(matches!(idx, SubIndex::One(..)));
assert_eq!(idx.len(), 1);
}
#[test]
fn same_id_insert_overwrites_and_returns_old() {
let id = Uuid::new_v4();
let mut idx: SubIndex<i32> = SubIndex::Zero;
let first = sub(1);
assert!(idx.insert(id, first.clone()).is_none());
let displaced = idx.insert(id, sub(2)).expect("old sub returned");
assert!(Arc::ptr_eq(&displaced, &first));
assert!(matches!(idx, SubIndex::One(..)));
assert_eq!(idx.len(), 1);
}
#[test]
fn second_distinct_id_promotes_to_many_keeping_both() {
let (a, b) = (Uuid::new_v4(), Uuid::new_v4());
let mut idx: SubIndex<i32> = SubIndex::Zero;
assert!(idx.insert(a, sub(1)).is_none());
assert!(idx.insert(b, sub(2)).is_none());
assert!(matches!(idx, SubIndex::Many(_)));
assert_eq!(idx.len(), 2);
let snap = build_snapshot(&idx);
let ids: Vec<Uuid> = snap.as_slice().iter().map(|(id, _)| *id).collect();
assert!(ids.contains(&a) && ids.contains(&b));
}
#[test]
fn remove_from_one_returns_to_zero() {
let id = Uuid::new_v4();
let mut idx: SubIndex<i32> = SubIndex::Zero;
let s = sub(7);
let _ = idx.insert(id, s.clone());
let removed = idx.remove(&id).expect("present");
assert!(Arc::ptr_eq(&removed, &s));
assert!(matches!(idx, SubIndex::Zero));
assert_eq!(idx.len(), 0);
assert!(idx.remove(&Uuid::new_v4()).is_none());
}
#[test]
fn remove_wrong_id_from_one_is_noop() {
let mut idx: SubIndex<i32> = SubIndex::Zero;
let _ = idx.insert(Uuid::new_v4(), sub(1));
assert!(idx.remove(&Uuid::new_v4()).is_none());
assert!(matches!(idx, SubIndex::One(..)));
assert_eq!(idx.len(), 1);
}
#[test]
fn many_does_not_demote_when_shrinking() {
let (a, b) = (Uuid::new_v4(), Uuid::new_v4());
let mut idx: SubIndex<i32> = SubIndex::Zero;
let _ = idx.insert(a, sub(1));
let _ = idx.insert(b, sub(2));
assert!(matches!(idx, SubIndex::Many(_)));
let _ = idx.remove(&a);
assert_eq!(idx.len(), 1);
assert!(matches!(idx, SubIndex::Many(_)));
}
fn build_snapshot(idx: &SubIndex<i32>) -> SubSnapshot<i32> {
match idx {
SubIndex::Zero => SubSnapshot::Zero,
SubIndex::One(id, s) => SubSnapshot::One((*id, s.clone())),
SubIndex::Many(map) => SubSnapshot::Many(Arc::new(
map.iter().map(|(id, s)| (*id, s.clone())).collect(),
)),
}
}
}