reinhardt_streaming/
in_memory.rs1use crate::{StreamingBackend, StreamingError};
2use async_trait::async_trait;
3use std::{
4 collections::{HashMap, VecDeque},
5 sync::Mutex,
6};
7
8pub struct InMemoryStreamingBackend {
10 queues: Mutex<HashMap<String, VecDeque<Vec<u8>>>>,
11}
12
13impl InMemoryStreamingBackend {
14 pub fn new() -> Self {
16 Self {
17 queues: Mutex::new(HashMap::new()),
18 }
19 }
20}
21
22impl Default for InMemoryStreamingBackend {
23 fn default() -> Self {
24 Self::new()
25 }
26}
27
28#[async_trait]
29impl StreamingBackend for InMemoryStreamingBackend {
30 async fn publish(&self, topic: &str, payload: Vec<u8>) -> Result<(), StreamingError> {
31 self.queues
32 .lock()
33 .map_err(|e| StreamingError::Backend(e.to_string()))?
34 .entry(topic.to_owned())
35 .or_default()
36 .push_back(payload);
37 Ok(())
38 }
39
40 async fn poll(&self, topic: &str) -> Result<Option<Vec<u8>>, StreamingError> {
41 Ok(self
42 .queues
43 .lock()
44 .map_err(|e| StreamingError::Backend(e.to_string()))?
45 .entry(topic.to_owned())
46 .or_default()
47 .pop_front())
48 }
49}
50
51#[cfg(test)]
52mod tests {
53 use super::*;
54 use rstest::*;
55
56 #[fixture]
57 fn backend() -> InMemoryStreamingBackend {
58 InMemoryStreamingBackend::new()
59 }
60
61 #[rstest]
62 #[tokio::test]
63 async fn publish_then_poll_returns_message(backend: InMemoryStreamingBackend) {
64 let payload = b"hello".to_vec();
66
67 backend
69 .publish("test-topic", payload.clone())
70 .await
71 .unwrap();
72 let result = backend.poll("test-topic").await.unwrap();
73
74 assert_eq!(result, Some(payload));
76 }
77
78 #[rstest]
79 #[tokio::test]
80 async fn poll_empty_topic_returns_none(backend: InMemoryStreamingBackend) {
81 let result = backend.poll("empty").await.unwrap();
82 assert_eq!(result, None);
83 }
84
85 #[rstest]
86 #[tokio::test]
87 async fn fifo_order_preserved(backend: InMemoryStreamingBackend) {
88 backend.publish("t", b"first".to_vec()).await.unwrap();
90 backend.publish("t", b"second".to_vec()).await.unwrap();
91
92 assert_eq!(backend.poll("t").await.unwrap(), Some(b"first".to_vec()));
94 assert_eq!(backend.poll("t").await.unwrap(), Some(b"second".to_vec()));
95 assert_eq!(backend.poll("t").await.unwrap(), None);
96 }
97
98 #[rstest]
99 #[tokio::test]
100 async fn independent_topics_do_not_interfere(backend: InMemoryStreamingBackend) {
101 backend.publish("a", b"msg-a".to_vec()).await.unwrap();
102 backend.publish("b", b"msg-b".to_vec()).await.unwrap();
103
104 assert_eq!(backend.poll("a").await.unwrap(), Some(b"msg-a".to_vec()));
105 assert_eq!(backend.poll("b").await.unwrap(), Some(b"msg-b".to_vec()));
106 }
107}