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}