Skip to main content

pamoja_ladder/
lib.rs

1//! Cost-aware transport ladder for the pamoja SDK.
2//!
3//! A field node usually has more than one way to reach the wider network, and
4//! those links differ wildly in cost, range, and availability: a local mesh hop is
5//! nearly free, long-range radio is cheap but slow, cellular is metered, and
6//! satellite is expensive. [`TransportLadder`] models that hierarchy. It holds a
7//! set of [`Transport`] rungs ordered cheapest-first and,
8//! on each send, uses the first rung that accepts the message. When no rung is
9//! reachable, the message is buffered in a durable [`Store`]
10//! and replayed later, so connectivity degrades gracefully instead of failing.
11//!
12//! This is the offline-first behavior the target deployments need on day one: an
13//! irrigation node or a fridge alarm keeps recording while every link is down and
14//! loses nothing once one returns.
15//!
16//! # Ordering and the buffer
17//!
18//! Delivery is in order. Once anything is buffered, later sends are buffered too
19//! rather than jumping ahead of the backlog over a recovered link;
20//! [`flush`](TransportLadder::flush) drains the backlog oldest-first, removing each
21//! record only after a rung accepts it. The pattern is to call
22//! [`flush`](TransportLadder::flush) when a link event suggests connectivity may
23//! have returned, and [`send`](TransportLadder::send) for new data.
24//!
25//! # Examples
26//!
27//! ```
28//! use pamoja_ladder::{Delivery, TransportLadder};
29//! use pamoja_loopback::{LoopbackBroker, LoopbackTransport};
30//! use pamoja_sync::MemoryStore;
31//!
32//! # async fn run() -> pamoja_core::Result<()> {
33//! let broker = LoopbackBroker::new();
34//! let mut ladder =
35//!     TransportLadder::new(MemoryStore::new()).rung(LoopbackTransport::new(broker.clone()));
36//! ladder.connect().await?;
37//!
38//! match ladder.send("sensors/1/temperature", b"21.5").await? {
39//!     Delivery::Sent => println!("delivered over a live link"),
40//!     Delivery::Buffered => println!("no link, buffered for later"),
41//! }
42//! # Ok(())
43//! # }
44//! ```
45
46use core::future::Future;
47use core::pin::Pin;
48
49use pamoja_core::{Error, Result, Store, Transport};
50
51/// The outcome of a [`TransportLadder::send`].
52#[derive(Clone, Copy, Debug, PartialEq, Eq)]
53pub enum Delivery {
54    /// The message was delivered immediately over one of the ladder's rungs.
55    Sent,
56    /// No rung accepted the message, so it was buffered for a later
57    /// [`flush`](TransportLadder::flush).
58    Buffered,
59}
60
61/// Object-safe erasure of [`Transport`] so a ladder can hold heterogeneous rungs.
62///
63/// The core [`Transport`] trait uses `async fn`, which is not dyn-compatible; this
64/// wrapper boxes the returned futures so transports of different concrete types can
65/// live together in one ordered list.
66///
67/// The boxed futures are `Send` so a ladder can be driven from a multi-threaded
68/// runtime, which is where one usually lives: a gateway ticks it from a task
69/// rather than blocking a thread on it.
70trait DynTransport: Send {
71    /// Connects the underlying transport.
72    fn connect(&mut self) -> Pin<Box<dyn Future<Output = Result<()>> + Send + '_>>;
73
74    /// Sends a payload to a topic over the underlying transport.
75    fn send<'a>(
76        &'a mut self,
77        topic: &'a str,
78        payload: &'a [u8],
79    ) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>>;
80}
81
82/// Newtype that carries one concrete transport behind the object-safe
83/// [`DynTransport`]. Erasing through a dedicated wrapper, rather than a blanket
84/// impl over every `T: Transport`, keeps these boxed-future methods off the
85/// transports themselves so their own `connect`/`send` stay unambiguous.
86struct Erased<T>(T);
87
88impl<T: Transport + Send> DynTransport for Erased<T> {
89    fn connect(&mut self) -> Pin<Box<dyn Future<Output = Result<()>> + Send + '_>> {
90        Box::pin(Transport::connect(&mut self.0))
91    }
92
93    fn send<'a>(
94        &'a mut self,
95        topic: &'a str,
96        payload: &'a [u8],
97    ) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'a>> {
98        Box::pin(Transport::send(&mut self.0, topic, payload))
99    }
100}
101
102/// An ordered set of transports backed by an offline buffer.
103///
104/// Rungs are tried in the order they are added, so the cheapest, most-preferred
105/// link is added first. A send that no rung accepts is buffered in the
106/// [`Store`] and replayed by [`flush`](Self::flush).
107pub struct TransportLadder<S> {
108    rungs: Vec<Box<dyn DynTransport>>,
109    buffer: S,
110}
111
112impl<S: Store> TransportLadder<S> {
113    /// Creates an empty ladder that buffers into `buffer`.
114    ///
115    /// # Arguments
116    ///
117    /// * `buffer` - the durable queue that holds messages while no rung is
118    ///   reachable.
119    ///
120    /// # Returns
121    ///
122    /// A ladder with no rungs; add them with [`rung`](Self::rung).
123    pub fn new(buffer: S) -> Self {
124        Self {
125            rungs: Vec::new(),
126            buffer,
127        }
128    }
129
130    /// Adds a rung, lowest-cost first.
131    ///
132    /// # Arguments
133    ///
134    /// * `transport` - a transport to try. Rungs added earlier are preferred, so
135    ///   add the cheapest link first and the costliest fallback last.
136    ///
137    /// # Returns
138    ///
139    /// The ladder, for chaining.
140    pub fn rung(mut self, transport: impl Transport + Send + 'static) -> Self {
141        self.rungs.push(Box::new(Erased(transport)));
142        self
143    }
144
145    /// Connects every rung, best-effort.
146    ///
147    /// A rung that fails to connect is left unreachable rather than failing the
148    /// whole ladder; sends simply fall through to the next rung or the buffer.
149    ///
150    /// # Returns
151    ///
152    /// `Ok(())` once every rung has been given the chance to connect.
153    ///
154    /// # Errors
155    ///
156    /// This call is best-effort and currently always returns `Ok(())`.
157    pub async fn connect(&mut self) -> Result<()> {
158        for rung in self.rungs.iter_mut() {
159            let _ = rung.connect().await;
160        }
161        Ok(())
162    }
163
164    /// Sends a payload, falling back down the rungs and then to the buffer.
165    ///
166    /// If the buffer is empty, each rung is tried in order and the first to accept
167    /// the message delivers it. If every rung fails, or the buffer already holds a
168    /// backlog, the message is buffered to preserve order.
169    ///
170    /// # Arguments
171    ///
172    /// * `topic` - the destination topic.
173    /// * `payload` - the bytes to send.
174    ///
175    /// # Returns
176    ///
177    /// [`Delivery::Sent`] if a rung delivered the message, or [`Delivery::Buffered`]
178    /// if it was queued for a later [`flush`](Self::flush).
179    ///
180    /// # Errors
181    ///
182    /// Returns [`Error::Io`] if the message must be buffered
183    /// but the store cannot be written.
184    pub async fn send(&mut self, topic: &str, payload: &[u8]) -> Result<Delivery> {
185        if self.buffer.is_empty().await? && Self::deliver(&mut self.rungs, topic, payload).await {
186            return Ok(Delivery::Sent);
187        }
188        self.buffer.append(&frame(topic, payload)).await?;
189        Ok(Delivery::Buffered)
190    }
191
192    /// Drains the buffer across the rungs, oldest record first.
193    ///
194    /// Each record is sent before it is removed, so the first record no rung can
195    /// deliver halts the drain and leaves it, and everything after it, buffered in
196    /// order for a later retry.
197    ///
198    /// # Returns
199    ///
200    /// The number of records forwarded before the buffer emptied or a rung refused
201    /// one.
202    ///
203    /// # Errors
204    ///
205    /// Returns [`Error::Io`] if the store cannot be read or
206    /// written, or [`Error::Codec`] if a buffered record
207    /// cannot be decoded.
208    pub async fn flush(&mut self) -> Result<usize> {
209        let mut forwarded = 0;
210        while let Some(record) = self.buffer.peek().await? {
211            let (topic, payload) = unframe(&record)?;
212            if !Self::deliver(&mut self.rungs, &topic, &payload).await {
213                break;
214            }
215            self.buffer.pop().await?;
216            forwarded += 1;
217        }
218        Ok(forwarded)
219    }
220
221    /// Returns how many messages are currently buffered.
222    ///
223    /// Takes the ladder mutably, like the rest of its surface. Reading through a
224    /// shared borrow would hold one across the await, which would in turn oblige
225    /// every rung to be `Sync` rather than only `Send`, and that is a heavier
226    /// requirement than a transport should have to meet.
227    ///
228    /// # Returns
229    ///
230    /// The number of records waiting for a [`flush`](Self::flush).
231    ///
232    /// # Errors
233    ///
234    /// Returns [`Error::Io`] if the store length cannot be
235    /// read.
236    pub async fn buffered(&mut self) -> Result<usize> {
237        self.buffer.len().await
238    }
239
240    /// Tries each rung in order, returning whether any accepted the message.
241    async fn deliver(rungs: &mut [Box<dyn DynTransport>], topic: &str, payload: &[u8]) -> bool {
242        for rung in rungs.iter_mut() {
243            if rung.send(topic, payload).await.is_ok() {
244                return true;
245            }
246        }
247        false
248    }
249}
250
251/// Frames a topic and payload into one record for the buffer.
252///
253/// The layout is a four-byte big-endian topic length, the topic bytes, then the
254/// payload, so [`unframe`] can split them back apart.
255fn frame(topic: &str, payload: &[u8]) -> Vec<u8> {
256    let mut record = Vec::with_capacity(4 + topic.len() + payload.len());
257    record.extend_from_slice(&(topic.len() as u32).to_be_bytes());
258    record.extend_from_slice(topic.as_bytes());
259    record.extend_from_slice(payload);
260    record
261}
262
263/// Splits a buffered record back into its topic and payload.
264fn unframe(record: &[u8]) -> Result<(String, Vec<u8>)> {
265    let header: [u8; 4] = record
266        .get(..4)
267        .ok_or_else(|| Error::Codec("ladder record is missing its length header".to_owned()))?
268        .try_into()
269        .expect("a four-byte slice");
270    let topic_len = u32::from_be_bytes(header) as usize;
271    let topic_bytes = record
272        .get(4..4 + topic_len)
273        .ok_or_else(|| Error::Codec("ladder record topic is truncated".to_owned()))?;
274    let topic =
275        String::from_utf8(topic_bytes.to_vec()).map_err(|err| Error::Codec(err.to_string()))?;
276    let payload = record[4 + topic_len..].to_vec();
277    Ok((topic, payload))
278}
279
280#[cfg(test)]
281mod tests {
282    use super::*;
283
284    use std::time::Duration;
285
286    use pamoja_loopback::{Faulty, LoopbackBroker, LoopbackTransport};
287    use pamoja_sync::MemoryStore;
288
289    /// Accepts anything that can move between threads.
290    fn assert_send<T: Send>(_value: T) {}
291
292    #[test]
293    fn a_ladder_can_be_driven_from_a_spawned_task() {
294        // A gateway ticks its ladder from a task on a threaded runtime, so the
295        // futures have to be Send. This does not run them; it fails to compile
296        // if a rung ever stops promising it.
297        let mut ladder = TransportLadder::new(MemoryStore::new())
298            .rung(LoopbackTransport::new(LoopbackBroker::new()));
299        assert_send(ladder.connect());
300        assert_send(ladder.send("sensors/1", b"21.5"));
301        assert_send(ladder.flush());
302        assert_send(ladder.buffered());
303    }
304
305    /// Subscribes a gateway to everything on a broker so the test can observe it.
306    async fn gateway(broker: &LoopbackBroker) -> LoopbackTransport {
307        let mut gateway = LoopbackTransport::new(broker.clone());
308        gateway.connect().await.expect("connect gateway");
309        gateway.subscribe("#").await.expect("subscribe gateway");
310        gateway
311    }
312
313    #[test]
314    fn frame_round_trips_topic_and_payload() {
315        let record = frame("sensors/1/temperature", b"21.5");
316        let (topic, payload) = unframe(&record).expect("unframe");
317        assert_eq!(topic, "sensors/1/temperature");
318        assert_eq!(payload, b"21.5");
319    }
320
321    #[test]
322    fn unframe_rejects_a_truncated_record() {
323        assert!(matches!(unframe(&[0, 0]), Err(Error::Codec(_))));
324        // Claims a four-byte topic but carries only one.
325        assert!(matches!(unframe(&[0, 0, 0, 4, b'a']), Err(Error::Codec(_))));
326    }
327
328    #[tokio::test]
329    async fn send_delivers_over_the_first_working_rung() {
330        let broker = LoopbackBroker::new();
331        let mut observer = gateway(&broker).await;
332
333        let mut ladder =
334            TransportLadder::new(MemoryStore::new()).rung(LoopbackTransport::new(broker.clone()));
335        ladder.connect().await.expect("connect");
336
337        let delivery = ladder
338            .send("sensors/1/temperature", b"21.5")
339            .await
340            .expect("send");
341        assert_eq!(delivery, Delivery::Sent);
342        assert_eq!(ladder.buffered().await.expect("buffered"), 0);
343
344        let message = observer.recv().await.expect("recv").expect("a message");
345        assert_eq!(message.topic, "sensors/1/temperature");
346        assert_eq!(message.payload, b"21.5");
347    }
348
349    #[tokio::test]
350    async fn send_falls_over_to_a_cheaper_rung_that_is_down() {
351        // The preferred rung publishes to its own broker but is broken; the
352        // fallback rung publishes to a second broker and works.
353        let preferred_broker = LoopbackBroker::new();
354        let fallback_broker = LoopbackBroker::new();
355        let mut preferred_observer = gateway(&preferred_broker).await;
356        let mut fallback_observer = gateway(&fallback_broker).await;
357
358        let preferred = Faulty::new(LoopbackTransport::new(preferred_broker.clone()), 1);
359        let fallback = LoopbackTransport::new(fallback_broker.clone());
360        let mut ladder = TransportLadder::new(MemoryStore::new())
361            .rung(preferred)
362            .rung(fallback);
363        ladder.connect().await.expect("connect");
364
365        let delivery = ladder.send("t", b"x").await.expect("send");
366        assert_eq!(delivery, Delivery::Sent);
367
368        let message = fallback_observer
369            .recv()
370            .await
371            .expect("recv")
372            .expect("a message");
373        assert_eq!(message.payload, b"x");
374        // The broken rung delivered nothing: its observer never sees a message.
375        let starved =
376            tokio::time::timeout(Duration::from_millis(50), preferred_observer.recv()).await;
377        assert!(
378            starved.is_err(),
379            "the broken rung must not deliver anything"
380        );
381    }
382
383    #[tokio::test]
384    async fn buffers_when_every_rung_is_down_then_flushes_in_order() {
385        let broker = LoopbackBroker::new();
386        let mut observer = gateway(&broker).await;
387
388        // One simulated outage on the only rung, then it recovers.
389        let rung = Faulty::new(LoopbackTransport::new(broker.clone()), 1);
390        let mut ladder = TransportLadder::new(MemoryStore::new()).rung(rung);
391        ladder.connect().await.expect("connect");
392
393        // First send hits the outage and buffers; the next two preserve order by
394        // buffering behind it rather than racing ahead.
395        assert_eq!(
396            ladder.send("out", b"a").await.expect("send"),
397            Delivery::Buffered
398        );
399        assert_eq!(
400            ladder.send("out", b"b").await.expect("send"),
401            Delivery::Buffered
402        );
403        assert_eq!(
404            ladder.send("out", b"c").await.expect("send"),
405            Delivery::Buffered
406        );
407        assert_eq!(ladder.buffered().await.expect("buffered"), 3);
408
409        // The link is back: drain everything in the order it was accepted.
410        let forwarded = ladder.flush().await.expect("flush");
411        assert_eq!(forwarded, 3);
412        assert_eq!(ladder.buffered().await.expect("buffered"), 0);
413
414        for expected in [b"a", b"b", b"c"] {
415            let message = observer.recv().await.expect("recv").expect("a message");
416            assert_eq!(message.topic, "out");
417            assert_eq!(message.payload, expected);
418        }
419    }
420
421    #[tokio::test]
422    async fn flush_stops_at_the_first_record_no_rung_accepts() {
423        let broker = LoopbackBroker::new();
424
425        // Three outages: the first buffers, the next two buffer behind it, and the
426        // flush attempt spends the remaining outages without draining anything.
427        let rung = Faulty::new(LoopbackTransport::new(broker.clone()), 3);
428        let mut ladder = TransportLadder::new(MemoryStore::new()).rung(rung);
429        ladder.connect().await.expect("connect");
430
431        ladder.send("out", b"a").await.expect("send");
432        ladder.send("out", b"b").await.expect("send");
433        ladder.send("out", b"c").await.expect("send");
434
435        // The link is still down on the first drained record, so nothing forwards
436        // and the backlog stays intact and ordered.
437        assert_eq!(ladder.flush().await.expect("flush"), 0);
438        assert_eq!(ladder.buffered().await.expect("buffered"), 3);
439    }
440}