radiate_engines/events/
subscription.rs1use 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}