macula_rust/station_link/
pubsub.rs1use std::collections::HashMap;
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex, Weak};
11
12use tokio::sync::mpsc;
13
14use crate::cbor::Value;
15use crate::frame::{self, PublicationSpec};
16use crate::node_key::NodeKey;
17
18use super::{now_ms, Inner, Link, LinkError};
19
20const EVENT_BUFFER: usize = 64;
23
24#[derive(Debug, Clone, PartialEq)]
27pub struct Publication {
28 pub realm: [u8; 32],
29 pub topic: String,
30 pub payload: Value,
31 pub ttl_ms: Option<u64>,
32}
33
34#[derive(Debug, Clone, PartialEq)]
38pub struct Event {
39 pub publisher: [u8; 32],
40 pub realm: [u8; 32],
41 pub topic: String,
42 pub seq: u64,
43 pub published_at: u64,
44 pub payload: Value,
45 pub delivered_via: String,
46}
47
48#[derive(Debug, Default)]
53pub struct PublicationSeq {
54 last: Mutex<u64>,
55}
56
57impl PublicationSeq {
58 pub fn next(&self) -> u64 {
60 let now = std::time::SystemTime::now()
61 .duration_since(std::time::UNIX_EPOCH)
62 .map(|d| d.as_micros() as u64)
63 .unwrap_or(0);
64 let mut last = self.last.lock().unwrap_or_else(|p| p.into_inner());
65 *last = now.max(*last + 1);
66 *last
67 }
68}
69
70#[derive(Debug, Default)]
74pub struct EventDedup {
75 seen: Mutex<HashMap<[u8; 48], u64>>,
76}
77
78impl EventDedup {
79 fn first(&self, hash: [u8; 48], expires_at: u64, now_ms: i64) -> bool {
81 let mut seen = self.seen.lock().unwrap_or_else(|p| p.into_inner());
82 if seen.contains_key(&hash) {
83 return false;
84 }
85 seen.retain(|_, until| *until >= now_ms as u64);
86 seen.insert(hash, expires_at);
87 true
88 }
89}
90
91#[derive(Debug, Clone, PartialEq)]
95pub struct SignedPublication {
96 frame: Value,
97}
98
99impl SignedPublication {
100 pub fn sign(
102 key: &NodeKey,
103 seq: &PublicationSeq,
104 p: Publication,
105 ) -> Result<SignedPublication, LinkError> {
106 let frame = frame::sign_publish(
107 &PublicationSpec {
108 realm: p.realm,
109 topic: p.topic,
110 seq: seq.next(),
111 published_at: now_ms() as u64,
112 payload: p.payload,
113 ttl_ms: p.ttl_ms,
114 },
115 key,
116 )?;
117 Ok(SignedPublication { frame })
118 }
119}
120
121static NEXT_SUBSCRIBER: AtomicU64 = AtomicU64::new(1);
122
123pub(super) struct SubscriberSlot {
125 id: u64,
126 events: mpsc::Sender<Event>,
127}
128
129pub struct Subscription {
132 link: Weak<Inner>,
133 key: ([u8; 32], String),
134 id: u64,
135 events: mpsc::Receiver<Event>,
136 unsubscribed: bool,
137}
138
139impl Subscription {
140 pub async fn recv(&mut self) -> Option<Event> {
142 self.events.recv().await
143 }
144
145 pub async fn unsubscribe(mut self) -> Result<(), LinkError> {
148 self.unsubscribed = true;
149 let Some(inner) = self.link.upgrade() else {
150 return Ok(());
151 };
152 if drop_subscriber(&inner, &self.key, self.id) {
153 let frame =
154 frame::unsubscribe_frame(self.key.1.as_bytes(), &self.key.0, &inner.self_id)?;
155 inner.send_control(&frame).await?;
156 }
157 Ok(())
158 }
159}
160
161impl Drop for Subscription {
162 fn drop(&mut self) {
164 if self.unsubscribed {
165 return;
166 }
167 let Some(inner) = self.link.upgrade() else {
168 return;
169 };
170 if !drop_subscriber(&inner, &self.key, self.id) {
171 return;
172 }
173 let key = self.key.clone();
174 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
175 runtime.spawn(async move {
176 if let Ok(frame) =
177 frame::unsubscribe_frame(key.1.as_bytes(), &key.0, &inner.self_id)
178 {
179 let _ = inner.send_control(&frame).await;
180 }
181 });
182 }
183 }
184}
185
186fn drop_subscriber(inner: &Inner, key: &([u8; 32], String), id: u64) -> bool {
188 let mut state = inner.lock();
189 let Some(slots) = state.subs.get_mut(key) else {
190 return false;
191 };
192 slots.retain(|s| s.id != id);
193 if slots.is_empty() {
194 state.subs.remove(key);
195 return true;
196 }
197 false
198}
199
200impl Link {
201 pub async fn publish(&self, p: Publication) -> Result<(), LinkError> {
203 let signed = SignedPublication::sign(&self.inner.key, &self.inner.publication_seq, p)?;
204 self.publish_signed(&signed).await
205 }
206
207 pub async fn publish_signed(&self, p: &SignedPublication) -> Result<(), LinkError> {
209 self.inner.write_control(&p.frame).await
210 }
211
212 pub async fn subscribe(
215 &self,
216 realm: &[u8; 32],
217 topic: &str,
218 ) -> Result<Subscription, LinkError> {
219 let key = (*realm, topic.to_string());
220 let (events_tx, events) = mpsc::channel(EVENT_BUFFER);
221 let id = NEXT_SUBSCRIBER.fetch_add(1, Ordering::Relaxed);
222 let first = {
223 let mut state = self.inner.lock();
224 if let Some(e) = &state.ended {
225 return Err(e.clone());
226 }
227 let slots = state.subs.entry(key.clone()).or_default();
228 slots.push(SubscriberSlot {
229 id,
230 events: events_tx,
231 });
232 slots.len() == 1
233 };
234 let subscription = Subscription {
235 link: Arc::downgrade(&self.inner),
236 key: key.clone(),
237 id,
238 events,
239 unsubscribed: false,
240 };
241 if first {
242 let frame = frame::subscribe_frame(topic.as_bytes(), realm, &self.inner.self_id);
243 let sent = match frame {
244 Ok(frame) => self.inner.send_control(&frame).await,
245 Err(e) => Err(e.into()),
246 };
247 if let Err(e) = sent {
248 drop_subscriber(&self.inner, &key, id);
249 let mut subscription = subscription;
250 subscription.unsubscribed = true;
251 return Err(e);
252 }
253 }
254 Ok(subscription)
255 }
256}
257
258pub(super) fn evented(inner: &Arc<Inner>, v: &Value) {
261 let now = now_ms();
262 let Ok(publication) = frame::verify_publication(v, inner.profile, now) else {
263 inner.count("event_unverified");
264 return;
265 };
266 let delivered_via = match v.get("delivered_via") {
267 Some(Value::Text(t)) => t.clone(),
268 _ => String::new(),
269 };
270 if !inner
271 .dedup
272 .first(publication.publication_hash, publication.expires_at, now)
273 {
274 inner.count("event_duplicate");
275 return;
276 }
277 let event = Event {
278 publisher: publication.publisher,
279 realm: publication.realm,
280 topic: publication.topic.clone(),
281 seq: publication.seq,
282 published_at: publication.published_at,
283 payload: publication.payload,
284 delivered_via,
285 };
286 let mut state = inner.lock();
287 let Some(slots) = state.subs.get(&(publication.realm, publication.topic)) else {
288 *state
289 .unrouted
290 .entry("event_unsubscribed".into())
291 .or_default() += 1;
292 return;
293 };
294 let overflowed = slots
295 .iter()
296 .filter(|s| s.events.try_send(event.clone()).is_err())
297 .count();
298 if overflowed > 0 {
299 *state.unrouted.entry("event_overflow".into()).or_default() += overflowed as u64;
300 }
301}