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}