Skip to main content

pamoja_sync/
forward.rs

1//! Forwarding buffered records onto a transport when a link appears.
2
3use pamoja_core::{Result, Store, Transport};
4
5/// Drains `store` onto `transport`, publishing each record to `topic`, oldest first.
6///
7/// Each record is sent before it is removed from the store, so a send failure
8/// leaves that record and every record after it buffered in order to retry later:
9/// nothing is lost and nothing is reordered. Delivery is at-least-once, since a
10/// crash between a successful send and the record's removal redelivers it on the
11/// next run.
12///
13/// This is the "forward" half of store-and-forward: buffer with a
14/// [`Store`](pamoja_core::Store) while offline, then call this when a link
15/// appears.
16///
17/// # Arguments
18///
19/// * `store` - the queue to drain.
20/// * `transport` - a connected transport to publish on.
21/// * `topic` - the topic every record is published to.
22///
23/// # Returns
24///
25/// The number of records forwarded once the store is drained empty.
26///
27/// # Errors
28///
29/// Returns the transport's error if a send fails, leaving the unsent records
30/// buffered in order, or [`Error::Io`](pamoja_core::Error::Io) if the store
31/// cannot be read.
32///
33/// # Examples
34///
35/// ```no_run
36/// use pamoja_core::{Store, Transport};
37/// use pamoja_sync::{drain_to, MemoryStore};
38///
39/// # async fn run(transport: &mut impl Transport) -> pamoja_core::Result<()> {
40/// let mut outbox = MemoryStore::new();
41/// outbox.append(b"21.5").await?;
42/// let forwarded = drain_to(&mut outbox, transport, "sensors/1/temperature").await?;
43/// assert_eq!(forwarded, 1);
44/// # Ok(())
45/// # }
46/// ```
47pub async fn drain_to<S, T>(store: &mut S, transport: &mut T, topic: &str) -> Result<usize>
48where
49    S: Store,
50    T: Transport,
51{
52    let mut forwarded = 0;
53    while let Some(record) = store.peek().await? {
54        transport.send(topic, &record).await?;
55        store.pop().await?;
56        forwarded += 1;
57    }
58    Ok(forwarded)
59}
60
61#[cfg(test)]
62mod tests {
63    use super::*;
64    use crate::MemoryStore;
65    use pamoja_core::Error;
66
67    /// A transport that records what it sends and can fail after a set count.
68    #[derive(Default)]
69    struct MockTransport {
70        sent: Vec<(String, Vec<u8>)>,
71        fail_after: Option<usize>,
72    }
73
74    impl Transport for MockTransport {
75        async fn connect(&mut self) -> Result<()> {
76            Ok(())
77        }
78
79        async fn send(&mut self, topic: &str, payload: &[u8]) -> Result<()> {
80            if self
81                .fail_after
82                .is_some_and(|limit| self.sent.len() >= limit)
83            {
84                return Err(Error::Transport("link down".to_owned()));
85            }
86            self.sent.push((topic.to_owned(), payload.to_vec()));
87            Ok(())
88        }
89
90        async fn subscribe(&mut self, _topic: &str) -> Result<()> {
91            Ok(())
92        }
93    }
94
95    #[tokio::test]
96    async fn drains_every_record_in_order() {
97        let mut store = MemoryStore::new();
98        for record in [b"a", b"b", b"c"] {
99            store.append(record).await.expect("append");
100        }
101        let mut transport = MockTransport::default();
102
103        let forwarded = drain_to(&mut store, &mut transport, "out")
104            .await
105            .expect("drain");
106
107        assert_eq!(forwarded, 3);
108        assert!(store.is_empty().await.expect("is_empty"));
109        assert_eq!(
110            transport.sent,
111            vec![
112                ("out".to_owned(), b"a".to_vec()),
113                ("out".to_owned(), b"b".to_vec()),
114                ("out".to_owned(), b"c".to_vec()),
115            ]
116        );
117    }
118
119    #[tokio::test]
120    async fn send_failure_preserves_remaining_records_in_order() {
121        let mut store = MemoryStore::new();
122        for record in [b"a", b"b", b"c"] {
123            store.append(record).await.expect("append");
124        }
125        let mut transport = MockTransport {
126            fail_after: Some(1),
127            ..Default::default()
128        };
129
130        let result = drain_to(&mut store, &mut transport, "out").await;
131
132        assert!(matches!(result, Err(Error::Transport(_))));
133        // Exactly the first record was forwarded and removed; the rest stay queued
134        // in their original order for a later retry.
135        assert_eq!(transport.sent, vec![("out".to_owned(), b"a".to_vec())]);
136        assert_eq!(store.len().await.expect("len"), 2);
137        assert_eq!(store.pop().await.expect("pop"), Some(b"b".to_vec()));
138        assert_eq!(store.pop().await.expect("pop"), Some(b"c".to_vec()));
139    }
140}