1#[cfg(feature = "metrics")]
2use std::time::Duration;
3use std::{
4 fmt::Debug,
5 marker::PhantomData,
6 panic::Location,
7 sync::{
8 Arc, Mutex, Weak,
9 atomic::{AtomicBool, Ordering},
10 },
11};
12
13#[cfg(feature = "metrics")]
14use arc_swap::ArcSwap;
15use dashmap::DashMap;
16use rustc_hash::FxHashMap;
17use uuid::Uuid;
18
19#[cfg(feature = "metrics")]
20use crate::metrics::CellMetrics;
21use crate::{
22 signal::Signal,
23 subscription::SubscriptionGuard,
24 traits::{CellValue, DepNode, Gettable, Mutable, Watchable, WatchableResult},
25};
26
27#[cfg(feature = "metrics")]
29#[derive(Debug, Clone)]
30pub struct SlowSubscriberAlert {
31 pub subscriber_id: Uuid,
33 pub duration_ns: u64,
35 pub threshold_ns: u64,
37}
38
39#[cfg(feature = "metrics")]
40type SlowSubscriberCallback = Arc<dyn Fn(SlowSubscriberAlert) + Send + Sync>;
41
42#[derive(Debug, Clone)]
43pub struct CellMutable;
44
45#[derive(Debug, Clone)]
46pub struct CellImmutable;
47
48pub(crate) struct CellInner<T> {
50 pub(crate) id: Uuid,
51 pub(crate) subscribers: parking_lot::Mutex<SubscriberRegistry<Subscriber<T>>>,
56 pub(crate) result_subscribers: parking_lot::Mutex<SubscriberRegistry<ResultSubscriber<T>>>,
59 pub(crate) value: Mutex<Arc<T>>,
66 pub(crate) name: Mutex<Option<Arc<str>>>,
70 pub(crate) owned: DashMap<Uuid, SubscriptionGuard>,
72 pub(crate) completed: AtomicBool,
74 pub(crate) errored: AtomicBool,
76 pub(crate) error: Mutex<Option<Arc<anyhow::Error>>>,
79 #[cfg(feature = "scheduler")]
86 pub(crate) height_cache: std::sync::atomic::AtomicU64,
87 #[cfg(feature = "scheduler")]
96 pub(crate) no_coalesce: AtomicBool,
97 #[cfg(feature = "metrics")]
99 pub(crate) metrics: Option<Arc<CellMetrics>>,
100 #[cfg(feature = "metrics")]
102 pub(crate) slow_subscriber_threshold_ns: ArcSwap<Option<u64>>,
103 #[cfg(feature = "metrics")]
105 pub(crate) slow_subscriber_callback: ArcSwap<Option<SlowSubscriberCallback>>,
106 #[allow(dead_code)]
108 pub(crate) caller: &'static Location<'static>,
109}
110
111pub struct Cell<T, M> {
113 pub(crate) inner: Arc<CellInner<T>>,
114 pub(crate) _marker: PhantomData<M>,
115}
116
117pub struct WeakCell<T, M> {
119 inner: Weak<CellInner<T>>,
120 _marker: PhantomData<M>,
121}
122
123impl<T, M> WeakCell<T, M> {
124 pub fn upgrade(&self) -> Option<Cell<T, M>> {
127 self.inner.upgrade().map(|inner| Cell {
128 inner,
129 _marker: PhantomData,
130 })
131 }
132
133 pub fn is_alive(&self) -> bool {
139 self.inner.strong_count() > 0
140 }
141}
142
143impl<T, M> Clone for WeakCell<T, M> {
144 fn clone(&self) -> Self {
145 WeakCell {
146 inner: self.inner.clone(),
147 _marker: PhantomData,
148 }
149 }
150}
151
152pub(crate) enum SubSnapshot<S> {
185 Zero,
186 One((Uuid, Arc<S>)),
187 Many(Arc<Vec<(Uuid, Arc<S>)>>),
188}
189
190impl<S> Clone for SubSnapshot<S> {
193 fn clone(&self) -> Self {
194 match self {
195 SubSnapshot::Zero => SubSnapshot::Zero,
196 SubSnapshot::One(pair) => SubSnapshot::One(pair.clone()),
197 SubSnapshot::Many(subs) => SubSnapshot::Many(subs.clone()),
198 }
199 }
200}
201
202impl<S> SubSnapshot<S> {
203 pub(crate) fn as_slice(&self) -> &[(Uuid, Arc<S>)] {
208 match self {
209 SubSnapshot::Zero => &[],
210 SubSnapshot::One(pair) => std::slice::from_ref(pair),
211 SubSnapshot::Many(subs) => subs.as_slice(),
212 }
213 }
214}
215
216enum SubIndex<S> {
234 Zero,
235 One(Uuid, Arc<S>),
236 Many(FxHashMap<Uuid, Arc<S>>),
237}
238
239impl<S> SubIndex<S> {
240 fn len(&self) -> usize {
241 match self {
242 SubIndex::Zero => 0,
243 SubIndex::One(..) => 1,
244 SubIndex::Many(map) => map.len(),
245 }
246 }
247
248 fn insert(&mut self, id: Uuid, sub: Arc<S>) -> Option<Arc<S>> {
252 match self {
253 SubIndex::Zero => {
254 *self = SubIndex::One(id, sub);
255 return None;
256 }
257 SubIndex::One(existing_id, existing_sub) => {
258 if *existing_id == id {
259 return Some(std::mem::replace(existing_sub, sub));
260 }
261 }
264 SubIndex::Many(map) => {
265 return map.insert(id, sub);
266 }
267 }
268
269 let (old_id, old_sub) = match std::mem::replace(self, SubIndex::Zero) {
272 SubIndex::One(old_id, old_sub) => (old_id, old_sub),
273 _ => unreachable!("promotion is entered only from the One arm"),
274 };
275 let mut map = FxHashMap::default();
276 map.insert(old_id, old_sub);
277 map.insert(id, sub);
278 *self = SubIndex::Many(map);
279 None
280 }
281
282 fn remove(&mut self, id: &Uuid) -> Option<Arc<S>> {
285 match self {
286 SubIndex::Zero => None,
287 SubIndex::One(existing_id, _) => {
288 if *existing_id != *id {
289 return None;
290 }
291 match std::mem::replace(self, SubIndex::Zero) {
292 SubIndex::One(_, sub) => Some(sub),
293 _ => unreachable!("just matched One"),
294 }
295 }
296 SubIndex::Many(map) => map.remove(id),
297 }
298 }
299}
300
301pub(crate) struct SubscriberRegistry<S> {
302 index: SubIndex<S>,
303 snapshot: SubSnapshot<S>,
304 dirty: bool,
305}
306
307impl<S> SubscriberRegistry<S> {
308 fn new() -> Self {
309 Self {
310 index: SubIndex::Zero,
311 snapshot: SubSnapshot::Zero,
312 dirty: false,
313 }
314 }
315
316 #[must_use = "displaced subscriber must be dropped outside the lock"]
319 fn insert(&mut self, id: Uuid, sub: Arc<S>) -> Option<Arc<S>> {
320 self.dirty = true;
321 self.index.insert(id, sub)
322 }
323
324 #[must_use = "removed subscriber must be dropped outside the lock"]
327 fn remove(&mut self, id: &Uuid) -> Option<Arc<S>> {
328 let removed = self.index.remove(id);
329 if removed.is_some() {
330 self.dirty = true;
331 }
332 removed
333 }
334
335 pub(crate) fn len(&self) -> usize {
336 self.index.len()
337 }
338
339 #[must_use = "displaced snapshot must be dropped outside the lock"]
344 pub(crate) fn snapshot(&mut self) -> (SubSnapshot<S>, Option<SubSnapshot<S>>) {
345 if self.dirty {
346 let next = match &self.index {
351 SubIndex::Zero => SubSnapshot::Zero,
352 SubIndex::One(id, sub) => SubSnapshot::One((*id, sub.clone())),
353 SubIndex::Many(map) => SubSnapshot::Many(Arc::new(
354 map.iter().map(|(id, sub)| (*id, sub.clone())).collect(),
355 )),
356 };
357 let old = std::mem::replace(&mut self.snapshot, next);
358 self.dirty = false;
359 (self.snapshot.clone(), Some(old))
360 } else {
361 (self.snapshot.clone(), None)
362 }
363 }
364}
365
366pub(crate) type SubscriberCallback<T> = Arc<dyn Fn(&Signal<T>) + Send + Sync>;
368
369pub(crate) struct Subscriber<T> {
370 pub(crate) callback: SubscriberCallback<T>,
371}
372
373impl<T> Subscriber<T> {
374 pub(crate) fn new(callback: impl Fn(&Signal<T>) + Send + Sync + 'static) -> Self {
375 Self {
376 callback: Arc::new(callback),
377 }
378 }
379}
380
381pub(crate) type ResultSubscriberCallback<T> =
383 Arc<dyn Fn(&Signal<T>) -> Result<(), String> + Send + Sync>;
384
385pub(crate) struct ResultSubscriber<T> {
386 pub(crate) callback: ResultSubscriberCallback<T>,
387}
388
389impl<T> ResultSubscriber<T> {
390 pub(crate) fn new(
391 callback: impl Fn(&Signal<T>) -> Result<(), String> + Send + Sync + 'static,
392 ) -> Self {
393 Self {
394 callback: Arc::new(callback),
395 }
396 }
397}
398
399impl<T: CellValue> Cell<T, CellMutable> {
400 #[track_caller]
401 pub fn new(initial_value: T) -> Self {
402 let inner = Arc::new(CellInner {
403 id: Uuid::new_v4(),
404 subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
405 result_subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
406 value: Mutex::new(Arc::new(initial_value)),
407 name: Mutex::new(None),
408 owned: DashMap::new(),
409 completed: AtomicBool::new(false),
410 errored: AtomicBool::new(false),
411 error: Mutex::new(None),
412 #[cfg(feature = "scheduler")]
413 height_cache: std::sync::atomic::AtomicU64::new(0),
414 #[cfg(feature = "scheduler")]
415 no_coalesce: AtomicBool::new(crate::scheduler::birth_no_coalesce()),
416 #[cfg(feature = "metrics")]
417 metrics: default_metrics(),
418 #[cfg(feature = "metrics")]
419 slow_subscriber_threshold_ns: ArcSwap::from_pointee(None),
420 #[cfg(feature = "metrics")]
421 slow_subscriber_callback: ArcSwap::from_pointee(None),
422 caller: Location::caller(),
423 });
424 #[cfg(feature = "inspector")]
425 crate::registry::registry().register(inner.id, Arc::downgrade(&inner) as Weak<dyn DepNode>);
426 #[cfg(feature = "trace")]
427 crate::tracing::register_cell(inner.id, Some(Location::caller().to_string()));
428 Self {
429 inner,
430 _marker: PhantomData,
431 }
432 }
433
434 #[cfg(feature = "metrics")]
436 #[track_caller]
437 pub fn with_metrics(initial_value: T) -> Self {
438 let inner = Arc::new(CellInner {
439 id: Uuid::new_v4(),
440 subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
441 result_subscribers: parking_lot::Mutex::new(SubscriberRegistry::new()),
442 value: Mutex::new(Arc::new(initial_value)),
443 name: Mutex::new(None),
444 owned: DashMap::new(),
445 completed: AtomicBool::new(false),
446 errored: AtomicBool::new(false),
447 error: Mutex::new(None),
448 #[cfg(feature = "scheduler")]
449 height_cache: std::sync::atomic::AtomicU64::new(0),
450 #[cfg(feature = "scheduler")]
451 no_coalesce: AtomicBool::new(crate::scheduler::birth_no_coalesce()),
452 metrics: Some(Arc::new(CellMetrics::new())),
453 slow_subscriber_threshold_ns: ArcSwap::from_pointee(None),
454 slow_subscriber_callback: ArcSwap::from_pointee(None),
455 caller: Location::caller(),
456 });
457 #[cfg(feature = "inspector")]
458 crate::registry::registry().register(inner.id, Arc::downgrade(&inner) as Weak<dyn DepNode>);
459 #[cfg(feature = "trace")]
460 crate::tracing::register_cell(inner.id, Some(Location::caller().to_string()));
461 Self {
462 inner,
463 _marker: PhantomData,
464 }
465 }
466 #[cfg(feature = "metrics")]
488 pub fn on_slow_subscriber<F>(&self, threshold: Duration, callback: F)
489 where
490 F: Fn(SlowSubscriberAlert) + Send + Sync + 'static,
491 {
492 self.inner
493 .slow_subscriber_threshold_ns
494 .store(Arc::new(Some(threshold.as_nanos() as u64)));
495 self.inner
496 .slow_subscriber_callback
497 .store(Arc::new(Some(Arc::new(callback))));
498 }
499
500 pub fn lock(self) -> Cell<T, CellImmutable> {
503 Cell {
504 inner: self.inner,
505 _marker: PhantomData,
506 }
507 }
508
509 pub fn with_name(self, name: impl Into<Arc<str>>) -> Self {
510 let name = name.into();
511 *self.inner.name.lock().expect("cell name poisoned") = Some(name.clone());
512 #[cfg(feature = "trace")]
513 crate::tracing::update_name(self.inner.id, name.to_string());
514 self
515 }
516
517 #[cfg(feature = "metrics")]
525 pub fn is_backed_up(&self) -> bool {
526 self.is_backed_up_threshold(std::time::Duration::from_millis(1))
527 }
528
529 #[cfg(feature = "metrics")]
534 pub fn is_backed_up_threshold(&self, threshold: std::time::Duration) -> bool {
535 self.inner
536 .metrics
537 .as_ref()
538 .map(|m| m.last_notify_time_ns() > threshold.as_nanos() as u64)
539 .unwrap_or(false)
540 }
541
542 #[cfg(feature = "metrics")]
548 pub fn try_set(&self, value: T) -> Result<(), T> {
549 if self.is_backed_up() {
550 Err(value)
551 } else {
552 self.set(value);
553 Ok(())
554 }
555 }
556
557 #[cfg(feature = "metrics")]
561 pub fn try_set_threshold(&self, value: T, threshold: std::time::Duration) -> Result<(), T> {
562 if self.is_backed_up_threshold(threshold) {
563 Err(value)
564 } else {
565 self.set(value);
566 Ok(())
567 }
568 }
569}
570
571impl<T, M> Clone for Cell<T, M> {
572 fn clone(&self) -> Self {
573 Cell {
574 inner: Arc::clone(&self.inner),
575 _marker: PhantomData,
576 }
577 }
578}
579
580impl<T, M> Cell<T, M> {
581 pub fn downgrade(&self) -> WeakCell<T, M> {
584 WeakCell {
585 inner: Arc::downgrade(&self.inner),
586 _marker: PhantomData,
587 }
588 }
589
590 #[cfg(feature = "metrics")]
595 pub fn metrics(&self) -> Option<&CellMetrics> {
596 self.inner.metrics.as_ref().map(|m| m.as_ref())
597 }
598
599 pub fn own(&self, guard: SubscriptionGuard) {
601 #[cfg(feature = "inspector")]
602 crate::registry::registry().mark_owned(guard.source().id(), self.inner.id);
603 self.inner.owned.insert(Uuid::new_v4(), guard);
604 #[cfg(feature = "scheduler")]
606 crate::scheduler::bump_topology_epoch();
607 #[cfg(feature = "trace")]
608 crate::tracing::update_owned_count(self.inner.id, self.inner.owned.len());
609 }
610
611 pub fn own_keyed(&self, key: Uuid, guard: SubscriptionGuard) {
617 #[cfg(feature = "inspector")]
618 {
619 if let Some((_, old_guard)) = self.inner.owned.remove(&key) {
621 crate::registry::registry().unmark_owned(old_guard.source().id());
622 }
623 crate::registry::registry().mark_owned(guard.source().id(), self.inner.id);
624 }
625 self.inner.owned.insert(key, guard);
626 #[cfg(feature = "scheduler")]
628 crate::scheduler::bump_topology_epoch();
629 #[cfg(feature = "trace")]
630 crate::tracing::update_owned_count(self.inner.id, self.inner.owned.len());
631 }
632}
633
634#[cfg(feature = "scheduler")]
639impl<T, M> Cell<T, M> {
640 pub fn no_coalesce(self) -> Self {
659 self.inner.no_coalesce.store(true, Ordering::Relaxed);
660 self
661 }
662}
663
664impl<T: Send + Sync, M: Send + Sync> DepNode for Cell<T, M> {
665 fn id(&self) -> Uuid {
666 self.inner.id
667 }
668
669 fn name(&self) -> Option<String> {
670 self.inner
671 .name
672 .lock()
673 .expect("cell name poisoned")
674 .as_ref()
675 .map(|s| s.to_string())
676 }
677
678 fn deps(&self) -> Vec<Arc<dyn DepNode>> {
679 let mut seen = std::collections::HashSet::new();
681 self.inner
682 .owned
683 .iter()
684 .filter_map(|entry| {
685 let source = entry.value().source();
686 let id = source.id();
687 if seen.insert(id) {
688 Some(Arc::clone(source))
689 } else {
690 None
691 }
692 })
693 .collect()
694 }
695
696 #[cfg(feature = "scheduler")]
697 fn height_cache(&self) -> Option<&std::sync::atomic::AtomicU64> {
698 Some(&self.inner.height_cache)
699 }
700
701 #[cfg(feature = "scheduler")]
702 fn no_coalesce(&self) -> bool {
703 self.inner.no_coalesce.load(Ordering::Relaxed)
704 }
705
706 fn subscriber_count(&self) -> usize {
707 self.inner.subscribers.lock().len() + self.inner.result_subscribers.lock().len()
708 }
709
710 fn owned_count(&self) -> usize {
711 self.inner.owned.len()
712 }
713}
714
715impl<T: CellValue> Cell<T, CellImmutable> {
716 pub fn with_name(self, name: impl Into<Arc<str>>) -> Self {
717 let name = name.into();
718 *self.inner.name.lock().expect("cell name poisoned") = Some(name.clone());
719 #[cfg(feature = "trace")]
720 crate::tracing::update_name(self.inner.id, name.to_string());
721 self
722 }
723}
724
725impl<T: CellValue, M: Send + Sync + 'static> Cell<T, M> {
726 #[doc(hidden)]
736 #[cfg_attr(feature = "profiling", inline(never))]
737 pub fn notify(&self, signal: Signal<T>) {
738 if self.inner.completed.load(Ordering::SeqCst) || self.inner.errored.load(Ordering::SeqCst)
740 {
741 return;
742 }
743
744 #[cfg(feature = "scheduler")]
750 if crate::scheduler::tick_active() {
751 let cell = self.clone();
752 let signal = signal.clone();
753 crate::scheduler::enqueue(
754 self.inner.id,
755 self as &dyn crate::traits::DepNode,
756 move || {
757 cell.write_value(&signal);
758 cell.fanout(&signal);
759 },
760 );
761 return;
762 }
763
764 self.write_value(&signal);
771 self.fanout(&signal);
772 }
773
774 #[cfg_attr(feature = "profiling", inline(never))]
777 fn write_value(&self, signal: &Signal<T>) {
778 match signal {
779 Signal::Value(arc_value) => {
780 *self.inner.value.lock().expect("cell value poisoned") = arc_value.clone();
787 }
788 Signal::Complete => {
789 self.inner.completed.store(true, Ordering::SeqCst);
790 }
791 Signal::Error(err) => {
792 self.inner.errored.store(true, Ordering::SeqCst);
793 *self.inner.error.lock().expect("cell error poisoned") = Some(err.clone());
794 }
795 }
796 }
797
798 #[cfg_attr(feature = "profiling", inline(never))]
801 fn fanout(&self, signal: &Signal<T>) {
802 #[cfg(feature = "profiling")]
807 crate::profiling::record_fire(self.inner.id);
808
809 #[cfg(feature = "profiling")]
815 let _fanout_span = {
816 let name = self.inner.name.lock().expect("cell name poisoned").clone();
817 ::tracing::trace_span!(
818 "hyphae.fanout",
819 cell.id = %self.inner.id,
820 cell.name = name.as_deref().unwrap_or(""),
821 )
822 .entered()
823 };
824
825 #[cfg(feature = "metrics")]
827 let notify_start = self
828 .inner
829 .metrics
830 .as_ref()
831 .map(|_| crate::platform::Instant::now());
832
833 let subs = {
843 let (subs, old_snapshot) = self.inner.subscribers.lock().snapshot();
844 drop(old_snapshot);
845 subs
846 };
847
848 #[cfg(feature = "metrics")]
852 let metrics = &self.inner.metrics;
853 #[cfg(feature = "metrics")]
854 let (slow_threshold, slow_callback) = if metrics.is_some() {
855 (
856 **self.inner.slow_subscriber_threshold_ns.load(),
857 (**self.inner.slow_subscriber_callback.load()).clone(),
858 )
859 } else {
860 (None, None)
861 };
862
863 for (_subscriber_id, sub) in subs.as_slice() {
868 #[cfg(feature = "metrics")]
869 let sub_start = metrics.as_ref().map(|_| crate::platform::Instant::now());
870
871 (sub.callback)(signal);
872
873 #[cfg(feature = "metrics")]
874 if let (Some(m), Some(start)) = (metrics, sub_start) {
875 let elapsed = start.elapsed().as_nanos() as u64;
876 m.update_slowest_subscriber(elapsed);
877
878 if let (Some(threshold), Some(cb)) = (&slow_threshold, &slow_callback)
879 && elapsed > *threshold
880 {
881 let alert = SlowSubscriberAlert {
882 subscriber_id: *_subscriber_id,
883 duration_ns: elapsed,
884 threshold_ns: *threshold,
885 };
886 cb(alert);
887 }
888 }
889 }
890
891 let result_subs = {
898 let (result_subs, old_snapshot) = self.inner.result_subscribers.lock().snapshot();
899 drop(old_snapshot);
900 result_subs
901 };
902
903 for (subscriber_id, sub) in result_subs.as_slice() {
904 #[cfg(feature = "metrics")]
905 let sub_start = metrics.as_ref().map(|_| crate::platform::Instant::now());
906
907 if let Err(err) = (sub.callback)(signal) {
908 log::error!(
909 "hyphae: fallible subscriber {} on cell {} returned error: {}",
910 subscriber_id,
911 self.inner.id,
912 err
913 );
914 }
915
916 #[cfg(feature = "metrics")]
917 if let (Some(m), Some(start)) = (metrics, sub_start) {
918 let elapsed = start.elapsed().as_nanos() as u64;
919 m.update_slowest_subscriber(elapsed);
920
921 if let (Some(threshold), Some(cb)) = (&slow_threshold, &slow_callback)
922 && elapsed > *threshold
923 {
924 let alert = SlowSubscriberAlert {
925 subscriber_id: *subscriber_id,
926 duration_ns: elapsed,
927 threshold_ns: *threshold,
928 };
929 cb(alert);
930 }
931 }
932 }
933
934 #[cfg(feature = "metrics")]
936 if let (Some(metrics), Some(start)) = (&self.inner.metrics, notify_start) {
937 let duration_ns = start.elapsed().as_nanos() as u64;
938 metrics.record_notify(duration_ns);
939 #[cfg(feature = "trace")]
940 crate::tracing::record_notify(
941 self.inner.id,
942 duration_ns,
943 subs.as_slice().len() + result_subs.as_slice().len(),
944 self.inner.owned.len(),
945 metrics.slowest_subscriber_ns(),
946 );
947 }
948 }
949}
950
951impl<T: CellValue, U: Send + Sync + 'static> Gettable<T> for Cell<T, U> {
952 fn get(&self) -> T {
953 let arc = self
956 .inner
957 .value
958 .lock()
959 .expect("cell value poisoned")
960 .clone();
961 (*arc).clone()
962 }
963}
964
965impl<T: CellValue, U: Send + Sync + 'static> Watchable<T> for Cell<T, U> {
966 fn subscribe(
967 &self,
968 callback: impl Fn(&Signal<T>) + Send + Sync + 'static,
969 ) -> SubscriptionGuard {
970 let id = Uuid::new_v4();
971 let sub = Arc::new(Subscriber::new(callback));
972
973 let displaced = self.inner.subscribers.lock().insert(id, sub.clone());
993 drop(displaced);
994
995 let current = self
1003 .inner
1004 .value
1005 .lock()
1006 .expect("cell value poisoned")
1007 .clone();
1008 (sub.callback)(&Signal::Value(current));
1009
1010 if self.is_complete() {
1012 (sub.callback)(&Signal::Complete);
1013 } else if self.is_error()
1014 && let Some(err) = self.error()
1015 {
1016 (sub.callback)(&Signal::Error(err));
1017 }
1018
1019 #[cfg(feature = "metrics")]
1021 if let Some(metrics) = &self.inner.metrics {
1022 metrics.record_subscriber_added();
1023 }
1024 #[cfg(feature = "trace")]
1025 {
1026 let subs_len = self.inner.subscribers.lock().len();
1027 let result_len = self.inner.result_subscribers.lock().len();
1028 crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
1029 }
1030
1031 let source: Arc<dyn DepNode> = Arc::new(self.clone());
1032 let cell = self.clone();
1033 #[cfg(feature = "metrics")]
1034 let metrics = self.inner.metrics.clone();
1035 SubscriptionGuard::new(id, source, move || {
1036 let removed_sub = cell.inner.subscribers.lock().remove(&id);
1039 let removed = removed_sub.is_some();
1040 drop(removed_sub);
1041 #[cfg(feature = "metrics")]
1042 if removed && let Some(m) = &metrics {
1043 m.record_subscriber_removed();
1044 }
1045 #[cfg(not(feature = "metrics"))]
1046 let _ = removed;
1047 #[cfg(feature = "trace")]
1048 {
1049 let subs_len = cell.inner.subscribers.lock().len();
1050 let result_len = cell.inner.result_subscribers.lock().len();
1051 crate::tracing::update_subscriber_count(cell.inner.id, subs_len + result_len);
1052 }
1053 })
1054 }
1055
1056 fn unsubscribe(&self, id: Uuid) {
1057 let removed_sub = self.inner.subscribers.lock().remove(&id);
1062 let removed_from_subs = removed_sub.is_some();
1063 drop(removed_sub);
1064 let removed_from_result = if removed_from_subs {
1065 false
1066 } else {
1067 let removed = self.inner.result_subscribers.lock().remove(&id);
1068 let did = removed.is_some();
1069 drop(removed);
1070 did
1071 };
1072 if removed_from_subs || removed_from_result {
1073 #[cfg(feature = "metrics")]
1075 if let Some(metrics) = &self.inner.metrics {
1076 metrics.record_subscriber_removed();
1077 }
1078 #[cfg(feature = "trace")]
1079 {
1080 let subs_len = self.inner.subscribers.lock().len();
1081 let result_len = self.inner.result_subscribers.lock().len();
1082 crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
1083 }
1084 }
1085 }
1086
1087 fn is_complete(&self) -> bool {
1088 self.inner.completed.load(Ordering::SeqCst)
1089 }
1090
1091 fn is_error(&self) -> bool {
1092 self.inner.errored.load(Ordering::SeqCst)
1093 }
1094
1095 fn error(&self) -> Option<Arc<anyhow::Error>> {
1096 self.inner
1097 .error
1098 .lock()
1099 .expect("cell error poisoned")
1100 .clone()
1101 }
1102}
1103
1104impl<T: CellValue, U: Send + Sync + 'static> WatchableResult<T> for Cell<T, U> {
1105 fn subscribe_result(
1106 &self,
1107 callback: impl Fn(&Signal<T>) -> Result<(), String> + Send + Sync + 'static,
1108 ) -> SubscriptionGuard {
1109 let cell_id = self.inner.id;
1110 let log_err = |id: &Uuid, err: &str| {
1111 log::error!(
1112 "hyphae: fallible subscriber {} on cell {} returned error: {}",
1113 id,
1114 cell_id,
1115 err
1116 );
1117 };
1118
1119 let id = Uuid::new_v4();
1120
1121 let current = self
1123 .inner
1124 .value
1125 .lock()
1126 .expect("cell value poisoned")
1127 .clone();
1128 if let Err(err) = callback(&Signal::Value(current)) {
1129 log_err(&id, &err);
1130 }
1131
1132 if self.inner.completed.load(Ordering::SeqCst) {
1134 if let Err(err) = callback(&Signal::Complete) {
1135 log_err(&id, &err);
1136 }
1137 } else if self.inner.errored.load(Ordering::SeqCst)
1138 && let Some(e) = self
1139 .inner
1140 .error
1141 .lock()
1142 .expect("cell error poisoned")
1143 .clone()
1144 && let Err(err) = callback(&Signal::Error(e))
1145 {
1146 log_err(&id, &err);
1147 }
1148
1149 let sub = Arc::new(ResultSubscriber::new(callback));
1150 let displaced = self.inner.result_subscribers.lock().insert(id, sub);
1153 drop(displaced);
1154
1155 #[cfg(feature = "metrics")]
1156 if let Some(metrics) = &self.inner.metrics {
1157 metrics.record_subscriber_added();
1158 }
1159 #[cfg(feature = "trace")]
1160 {
1161 let subs_len = self.inner.subscribers.lock().len();
1162 let result_len = self.inner.result_subscribers.lock().len();
1163 crate::tracing::update_subscriber_count(self.inner.id, subs_len + result_len);
1164 }
1165
1166 let source: Arc<dyn DepNode> = Arc::new(self.clone());
1167 let cell = self.clone();
1168 #[cfg(feature = "metrics")]
1169 let metrics = self.inner.metrics.clone();
1170 SubscriptionGuard::new(id, source, move || {
1171 let removed_sub = cell.inner.result_subscribers.lock().remove(&id);
1174 let removed = removed_sub.is_some();
1175 drop(removed_sub);
1176 #[cfg(feature = "metrics")]
1177 if removed && let Some(m) = &metrics {
1178 m.record_subscriber_removed();
1179 }
1180 #[cfg(not(feature = "metrics"))]
1181 let _ = removed;
1182 #[cfg(feature = "trace")]
1183 {
1184 let subs_len = cell.inner.subscribers.lock().len();
1185 let result_len = cell.inner.result_subscribers.lock().len();
1186 crate::tracing::update_subscriber_count(cell.inner.id, subs_len + result_len);
1187 }
1188 })
1189 }
1190}
1191
1192impl<T: CellValue> Mutable<T> for Cell<T, CellMutable> {
1193 fn set(&self, value: T) {
1194 self.notify(Signal::value(value)); }
1196
1197 fn complete(&self) {
1198 self.notify(Signal::Complete);
1199 }
1200
1201 fn fail(&self, error: impl Into<anyhow::Error>) {
1202 self.notify(Signal::error(error));
1203 }
1204}
1205
1206#[cfg(feature = "inspector")]
1211impl<T: CellValue> DepNode for CellInner<T> {
1212 fn id(&self) -> Uuid {
1213 self.id
1214 }
1215
1216 fn name(&self) -> Option<String> {
1217 self.name
1218 .lock()
1219 .expect("cell name poisoned")
1220 .as_ref()
1221 .map(|s| s.to_string())
1222 }
1223
1224 fn deps(&self) -> Vec<Arc<dyn DepNode>> {
1225 let mut seen = std::collections::HashSet::new();
1226 self.owned
1227 .iter()
1228 .filter_map(|entry| {
1229 let source = entry.value().source();
1230 let id = source.id();
1231 if seen.insert(id) {
1232 Some(Arc::clone(source))
1233 } else {
1234 None
1235 }
1236 })
1237 .collect()
1238 }
1239
1240 fn subscriber_count(&self) -> usize {
1241 self.subscribers.lock().len() + self.result_subscribers.lock().len()
1242 }
1243
1244 fn owned_count(&self) -> usize {
1245 self.owned.len()
1246 }
1247
1248 fn value_debug(&self) -> Option<String> {
1249 let arc = self.value.lock().expect("cell value poisoned").clone();
1250 Some(format!("{:?}", &*arc))
1251 }
1252
1253 fn caller(&self) -> Option<&'static Location<'static>> {
1254 Some(self.caller)
1255 }
1256}
1257
1258impl<T> Drop for CellInner<T> {
1259 fn drop(&mut self) {
1260 #[cfg(feature = "trace")]
1261 crate::tracing::deregister_cell(&self.id);
1262 #[cfg(feature = "inspector")]
1263 crate::registry::registry().deregister(&self.id);
1264 }
1265}
1266
1267#[cfg(all(feature = "metrics", feature = "trace"))]
1268fn default_metrics() -> Option<Arc<CellMetrics>> {
1269 Some(Arc::new(CellMetrics::new()))
1270}
1271
1272#[cfg(all(feature = "metrics", not(feature = "trace")))]
1273fn default_metrics() -> Option<Arc<CellMetrics>> {
1274 None
1275}
1276
1277#[cfg(test)]
1278mod sub_index_tests {
1279 use std::sync::Arc;
1280
1281 use uuid::Uuid;
1282
1283 use super::{SubIndex, SubSnapshot};
1284
1285 fn sub(v: i32) -> Arc<i32> {
1288 Arc::new(v)
1289 }
1290
1291 #[test]
1292 fn zero_and_one_stay_inline() {
1293 let mut idx: SubIndex<i32> = SubIndex::Zero;
1294 assert!(matches!(idx, SubIndex::Zero));
1295 assert_eq!(idx.len(), 0);
1296
1297 assert!(idx.insert(Uuid::new_v4(), sub(1)).is_none());
1299 assert!(matches!(idx, SubIndex::One(..)));
1300 assert_eq!(idx.len(), 1);
1301 }
1302
1303 #[test]
1304 fn same_id_insert_overwrites_and_returns_old() {
1305 let id = Uuid::new_v4();
1306 let mut idx: SubIndex<i32> = SubIndex::Zero;
1307 let first = sub(1);
1308 assert!(idx.insert(id, first.clone()).is_none());
1309
1310 let displaced = idx.insert(id, sub(2)).expect("old sub returned");
1313 assert!(Arc::ptr_eq(&displaced, &first));
1314 assert!(matches!(idx, SubIndex::One(..)));
1315 assert_eq!(idx.len(), 1);
1316 }
1317
1318 #[test]
1319 fn second_distinct_id_promotes_to_many_keeping_both() {
1320 let (a, b) = (Uuid::new_v4(), Uuid::new_v4());
1321 let mut idx: SubIndex<i32> = SubIndex::Zero;
1322 assert!(idx.insert(a, sub(1)).is_none());
1323 assert!(idx.insert(b, sub(2)).is_none());
1325 assert!(matches!(idx, SubIndex::Many(_)));
1326 assert_eq!(idx.len(), 2);
1327
1328 let snap = build_snapshot(&idx);
1330 let ids: Vec<Uuid> = snap.as_slice().iter().map(|(id, _)| *id).collect();
1331 assert!(ids.contains(&a) && ids.contains(&b));
1332 }
1333
1334 #[test]
1335 fn remove_from_one_returns_to_zero() {
1336 let id = Uuid::new_v4();
1337 let mut idx: SubIndex<i32> = SubIndex::Zero;
1338 let s = sub(7);
1339 let _ = idx.insert(id, s.clone());
1340
1341 let removed = idx.remove(&id).expect("present");
1342 assert!(Arc::ptr_eq(&removed, &s));
1343 assert!(matches!(idx, SubIndex::Zero));
1344 assert_eq!(idx.len(), 0);
1345
1346 assert!(idx.remove(&Uuid::new_v4()).is_none());
1348 }
1349
1350 #[test]
1351 fn remove_wrong_id_from_one_is_noop() {
1352 let mut idx: SubIndex<i32> = SubIndex::Zero;
1353 let _ = idx.insert(Uuid::new_v4(), sub(1));
1354 assert!(idx.remove(&Uuid::new_v4()).is_none());
1355 assert!(matches!(idx, SubIndex::One(..)));
1356 assert_eq!(idx.len(), 1);
1357 }
1358
1359 #[test]
1360 fn many_does_not_demote_when_shrinking() {
1361 let (a, b) = (Uuid::new_v4(), Uuid::new_v4());
1362 let mut idx: SubIndex<i32> = SubIndex::Zero;
1363 let _ = idx.insert(a, sub(1));
1364 let _ = idx.insert(b, sub(2));
1365 assert!(matches!(idx, SubIndex::Many(_)));
1366
1367 let _ = idx.remove(&a);
1370 assert_eq!(idx.len(), 1);
1371 assert!(matches!(idx, SubIndex::Many(_)));
1372 }
1373
1374 fn build_snapshot(idx: &SubIndex<i32>) -> SubSnapshot<i32> {
1377 match idx {
1378 SubIndex::Zero => SubSnapshot::Zero,
1379 SubIndex::One(id, s) => SubSnapshot::One((*id, s.clone())),
1380 SubIndex::Many(map) => SubSnapshot::Many(Arc::new(
1381 map.iter().map(|(id, s)| (*id, s.clone())).collect(),
1382 )),
1383 }
1384 }
1385}