Skip to main content

macula_rust/station_link/
pubsub.rs

1//! PubSub, as macula 12's link does it. A PUBLISH carries a publication
2//! signed with the link's identity key; SUBSCRIBE and UNSUBSCRIBE are control
3//! frames naming the realm, the topic and this node, with no
4//! acknowledgement; an EVENT carries a publication that is verified before it
5//! is delivered, and delivered once however many copies arrive, until it
6//! expires.
7
8use 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
20/// How many events a subscription holds that its reader has not taken; an
21/// event arriving at a full subscription is dropped and counted.
22const EVENT_BUFFER: usize = 64;
23
24/// What [`Link::publish`] sends: the realm and topic, the payload, and its
25/// time to live, `None` for macula's 10 minutes.
26#[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/// A publication a subscription heard, verified: who published it (the key
35/// id its signature verified under), where, its seq and time, the payload,
36/// and how it arrived (direct or plumtree).
37#[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/// Numbers one publisher's publications: the first is the wall clock in
49/// microseconds, and each after is one more than the last, or the clock,
50/// whichever is later, as macula_publication_seq does. Links of one identity
51/// key share one, so their seqs never repeat.
52#[derive(Debug, Default)]
53pub struct PublicationSeq {
54    last: Mutex<u64>,
55}
56
57impl PublicationSeq {
58    /// The seq for the next publication.
59    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/// Remembers the publications a node delivered, by publication hash, until
71/// each expires, so an event heard on several links, or twice on one, is
72/// delivered once. The links of one node share one.
73#[derive(Debug, Default)]
74pub struct EventDedup {
75    seen: Mutex<HashMap<[u8; 48], u64>>,
76}
77
78impl EventDedup {
79    /// Whether `hash` is new at `now_ms`, remembering it until `expires_at`.
80    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/// A PUBLISH signed once, to be sent on several links of one node: every copy
92/// is the same publication, so a subscriber hearing it on several links
93/// delivers it once.
94#[derive(Debug, Clone, PartialEq)]
95pub struct SignedPublication {
96    frame: Value,
97}
98
99impl SignedPublication {
100    /// Signs `p` with `key` at the next seq of `seq`.
101    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
123/// One subscriber of a realm and topic on a link.
124pub(super) struct SubscriberSlot {
125    id: u64,
126    events: mpsc::Sender<Event>,
127}
128
129/// One subscription to a realm and topic on a link, until
130/// [`Subscription::unsubscribe`] or the link ends.
131pub 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    /// The next event, or `None` once the subscription has ended.
141    pub async fn recv(&mut self) -> Option<Event> {
142        self.events.recv().await
143    }
144
145    /// Ends the subscription, and sends UNSUBSCRIBE once no other
146    /// subscription on the link holds its realm and topic.
147    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    /// A subscription dropped without unsubscribing unsubscribes as it goes.
163    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
186/// Removes a subscriber, and whether it was the last on its realm and topic.
187fn 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    /// Signs `p` as a PUBLISH and sends it.
202    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    /// Sends a publication signed once for several links, as it is.
208    pub async fn publish_signed(&self, p: &SignedPublication) -> Result<(), LinkError> {
209        self.inner.write_control(&p.frame).await
210    }
211
212    /// Subscribes to `topic` in `realm`: the first subscription to a realm
213    /// and topic on the link sends SUBSCRIBE.
214    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
258/// An EVENT's publication, once it verifies, delivered to every subscription
259/// on its realm and topic, unless it was delivered before.
260pub(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}