Skip to main content

radiate_engines/events/
subscription.rs

1use radiate_core::{Expr, RadiateError, error::RadiateResult, radiate_err};
2use radiate_error::radiate_bail;
3use radiate_utils::sentry_id;
4use std::sync::atomic::AtomicUsize;
5use std::{
6    fmt::Debug,
7    sync::{
8        Arc, RwLock,
9        atomic::{AtomicBool, Ordering},
10    },
11};
12
13sentry_id!(SubscriptionId);
14
15#[derive(Clone)]
16pub struct Subscription {
17    pub(crate) id: SubscriptionId,
18    pub(crate) schedule: Arc<RwLock<Option<Expr>>>,
19    pub(crate) permits: Arc<AtomicUsize>,
20    pub(crate) alive: Arc<AtomicBool>,
21}
22
23impl Subscription {
24    pub(super) fn new() -> Self {
25        Subscription {
26            id: SubscriptionId::new(),
27            schedule: Arc::new(RwLock::new(None)),
28            permits: Arc::new(AtomicUsize::new(0)),
29            alive: Arc::new(AtomicBool::new(true)),
30        }
31    }
32
33    pub fn is_alive(&self) -> bool {
34        self.alive.load(Ordering::Acquire)
35    }
36
37    pub fn id(&self) -> SubscriptionId {
38        self.id
39    }
40
41    pub fn schedule(&self, schedule: impl Into<Expr>) -> Result<bool, RadiateError> {
42        let maybe_expr = schedule.into();
43
44        match maybe_expr.clone().into_schedule() {
45            Some(expr) => {
46                *self.schedule.write().unwrap() = Some(expr);
47                Ok(true)
48            }
49            _ => {
50                radiate_bail!(Expr: format!("Invalid schedule expression: {:?}", maybe_expr))
51            }
52        }
53    }
54
55    pub fn unsubscribe(&self) {
56        self.alive.store(false, Ordering::Release);
57        self.permits.store(0, Ordering::Release);
58    }
59
60    pub(super) fn reserve(&self) -> RadiateResult<bool> {
61        if !self.is_alive() {
62            return Ok(false);
63        }
64
65        if !self.try_schedule()? {
66            return Ok(false);
67        }
68
69        self.permits.fetch_add(1, Ordering::Release);
70
71        Ok(true)
72    }
73
74    pub(super) fn take_permit(&self) -> bool {
75        let mut current = self.permits.load(Ordering::Acquire);
76
77        loop {
78            if current == 0 {
79                return false;
80            }
81
82            match self.permits.compare_exchange_weak(
83                current,
84                current - 1,
85                Ordering::AcqRel,
86                Ordering::Acquire,
87            ) {
88                Ok(_) => return true,
89                Err(next) => current = next,
90            }
91        }
92    }
93
94    fn try_schedule(&self) -> RadiateResult<bool> {
95        let mut guard = self.schedule.write().unwrap();
96        if let Some(expr) = &mut *guard {
97            return expr
98                .trigger()?
99                .extract_bool()
100                .ok_or_else(|| radiate_err!(Expr: "Failed to compute schedule as bool"));
101        }
102
103        Ok(true)
104    }
105}