Skip to main content

reinhardt_streaming/
in_memory.rs

1use crate::{StreamingBackend, StreamingError};
2use async_trait::async_trait;
3use std::{
4	collections::{HashMap, VecDeque},
5	sync::Mutex,
6};
7
8/// In-memory streaming backend for unit tests (no Kafka required).
9pub struct InMemoryStreamingBackend {
10	queues: Mutex<HashMap<String, VecDeque<Vec<u8>>>>,
11}
12
13impl InMemoryStreamingBackend {
14	/// Create an empty in-memory backend.
15	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		// Arrange
65		let payload = b"hello".to_vec();
66
67		// Act
68		backend
69			.publish("test-topic", payload.clone())
70			.await
71			.unwrap();
72		let result = backend.poll("test-topic").await.unwrap();
73
74		// Assert
75		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		// Arrange
89		backend.publish("t", b"first".to_vec()).await.unwrap();
90		backend.publish("t", b"second".to_vec()).await.unwrap();
91
92		// Act & Assert
93		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}