1mod service_event;
16
17pub use service_event::ServiceEvent;
18use crate::name_generator;
19
20use std::{ time, collections::HashMap };
21use regex::Regex;
22use lazy_static::lazy_static;
23use redis::{Commands, Connection, Client};
24use uuid::Uuid;
25
26#[derive(Debug, Eq, PartialEq)]
27pub enum EventQueueError {
28 ConnectionError(String),
29 JSONDumpError(String),
30 JSONParseError(String),
31 EnqueueError(String),
32 DequeueError(String),
33 EmptyQueue,
34 TimeoutExpired
35}
36
37pub type EventQueueResult<T> = Result<T, EventQueueError>;
38pub type Timestamp = u64;
39
40type EventId = String;
41type SerializedEventData = String;
42type EventMap = HashMap<EventId, SerializedEventData>;
43type StreamEntry = HashMap<String, EventMap>;
44type StreamMap = HashMap<String, Vec<StreamEntry>>;
45
46#[derive(Debug, Eq, PartialEq)]
47pub struct TimestampedEvent(Timestamp, ServiceEvent);
48
49impl TimestampedEvent {
50 pub fn timestamp(&self) -> Timestamp {
51 self.0
52 }
53
54 pub fn event(&self) -> &ServiceEvent {
55 &self.1
56 }
57}
58
59pub struct EventQueue {
60 redis_client: Client,
61 message_queue_name: String,
62 event_stream_name: String,
63 response_stream_name: String
64}
65
66impl EventQueue {
67 pub fn new(queue_name: &str, connection_url: &str) -> Self {
68 let redis_client = redis::Client::open(connection_url).unwrap();
69 let message_queue_name = name_generator::generate_message_queue_name(queue_name);
70 let event_stream_name = name_generator::generate_event_stream_name(queue_name);
71 let response_stream_name = name_generator::generate_response_stream_name(queue_name);
72
73 EventQueue {
74 redis_client,
75 message_queue_name,
76 event_stream_name,
77 response_stream_name
78 }
79 }
80
81 fn extract_timestamp_from_event_key(key: &str) -> u64 {
82 lazy_static! {
83 static ref KEY_REGEX: Regex = Regex::new(r"(?P<timestamp>\d+)-\d+").unwrap();
84 }
85
86 let timestamp = match KEY_REGEX.captures(key) {
87 None => panic!("invalid event key passed to function"),
88 Some(captures) => captures["timestamp"].to_string()
89 };
90
91 timestamp.parse::<u64>().unwrap()
92 }
93
94 fn setup_connection(&self) -> EventQueueResult<redis::Connection> {
95 let connection = self.redis_client.get_connection();
96
97 match connection {
98 Err(error) => Err(EventQueueError::ConnectionError(error.to_string())),
99 Ok(connection) => Ok(connection)
100 }
101 }
102
103 fn get_service_event_by_key(&self, connection: &mut Connection, event_key: &str, event_type: &str) -> EventQueueResult<ServiceEvent> {
104 let event_data_list: Vec<StreamEntry> = match connection.xrange_count(
105 &self.event_stream_name,
106 event_key,
107 event_key,
108 1
109 ) {
110 Err(error) => return Err(EventQueueError::DequeueError(error.to_string())),
111 Ok(data) => data
112 };
113
114 let event_data = match event_data_list.into_iter().next() {
115 None => return Err(EventQueueError::DequeueError(String::from("unexpected empty value in stream"))),
116 Some(event_data) => event_data
117 };
118
119 let event = match event_data.get(event_key) {
120 None => return Err(EventQueueError::DequeueError(String::from("expected event map, found None"))),
121 Some(event) => match event.get(event_type) {
122 None => return Err(EventQueueError::DequeueError(String::from("expected event at key \"event\", found None"))),
123 Some(event) => event
124 }
125 };
126
127 let event: ServiceEvent = match serde_json::from_str(event) {
128 Err(error) => return Err(EventQueueError::JSONParseError(error.to_string())),
129 Ok(event) => event
130 };
131
132 Ok(event)
133 }
134
135 fn get_last_response_id(&self, connection: &mut Connection) -> EventQueueResult<String> {
136 let last_response: Vec<StreamEntry> = match connection.xrevrange_count(&self.response_stream_name, "+", "-", 1) {
137 Err(error) => return Err(EventQueueError::DequeueError(error.to_string())),
138 Ok(response) => response
139 };
140
141 if last_response.is_empty() {
142 return Ok(String::from("0-0"));
143 }
144
145 if last_response.len() != 1 {
146 return Err(EventQueueError::DequeueError(String::from("unexpected response length")));
147 }
148
149 let last_response = &last_response[0];
150 let id = last_response.keys().next().unwrap().to_string();
151
152 Ok(id)
153 }
154
155 pub fn enqueue(&mut self, event: &ServiceEvent) -> EventQueueResult<Timestamp> {
156 let mut connection = self.setup_connection()?;
157
158 let event_as_json = match serde_json::to_string(&event) {
159 Err(error) => return Err(EventQueueError::JSONDumpError(error.to_string())),
160 Ok(json) => json
161 };
162
163 let event_key: String = match connection.xadd(
164 &self.event_stream_name,
165 "*",
166 &[("event", &event_as_json)]
167 ) {
168 Err(error) => return Err(EventQueueError::EnqueueError(error.to_string())),
169 Ok(key) => key
170 };
171
172 if let Err(error) = connection.lpush::<_, _, ()>(
173 &self.message_queue_name,
174 &event_key
175 ) {
176 return Err(EventQueueError::EnqueueError(error.to_string()));
177 }
178
179 let timestamp = Self::extract_timestamp_from_event_key(&event_key);
180
181 Ok(timestamp)
182 }
183
184 pub fn dequeue(&mut self) -> EventQueueResult<TimestampedEvent> {
185 let mut connection = self.setup_connection()?;
186
187 let event_key: String = match connection.rpop(&self.message_queue_name, None) {
188 Err(error) => return Err(EventQueueError::DequeueError(error.to_string())),
189 Ok(key) => match key {
190 None => return Err(EventQueueError::EmptyQueue),
191 Some(key) => key
192 }
193 };
194
195 let event = self.get_service_event_by_key(&mut connection, &event_key, "event")?;
196 let timestamp = Self::extract_timestamp_from_event_key(&event_key);
197
198 Ok(TimestampedEvent(timestamp, event))
199 }
200
201 pub fn dequeue_blocking(&mut self, timeout: u16) -> EventQueueResult<TimestampedEvent> {
202 let mut connection = self.setup_connection()?;
203
204 let event_kvp: (String, String) = match connection.brpop(
205 &self.message_queue_name,
206 timeout.into()
207 ) {
208 Err(error) => return Err(EventQueueError::DequeueError(error.to_string())),
209 Ok(key) => match key {
210 None => return Err(EventQueueError::EmptyQueue),
211 Some(kvp) => kvp
212 }
213 };
214
215 let event_key = event_kvp.1;
216
217 let event = self.get_service_event_by_key(&mut connection, &event_key, "event")?;
218 let timestamp = Self::extract_timestamp_from_event_key(&event_key);
219
220 Ok(TimestampedEvent(timestamp, event))
221 }
222
223 pub fn enqueue_response(&mut self, event: &ServiceEvent) -> EventQueueResult<()> {
224 let mut connection = self.setup_connection()?;
225
226 let event_as_json = match serde_json::to_string(&event) {
227 Err(error) => return Err(EventQueueError::JSONDumpError(error.to_string())),
228 Ok(json) => json
229 };
230
231 let uuid_string = Uuid::from_u128(event.uuid()).to_string();
232 let response_key: String = match connection.xadd(
233 &self.event_stream_name,
234 "*",
235 &[("response", &event_as_json)]
236 ) {
237 Err(error) => return Err(EventQueueError::EnqueueError(error.to_string())),
238 Ok(key) => key
239 };
240
241 if let Err(error) = connection.xadd::<_, _, _, _, ()>(&self.response_stream_name, "*", &[(&uuid_string, &response_key)]) {
242 return Err(EventQueueError::EnqueueError(error.to_string()));
243 }
244
245 Ok(())
246 }
247
248 pub fn await_response(&mut self, event: &ServiceEvent) -> EventQueueResult<TimestampedEvent> {
249 let mut connection = self.setup_connection()?;
250
251 let start_time = time::Instant::now();
252 let timeout = event.timeout();
253 let target_uuid_string = Uuid::from_u128(event.uuid()).to_string();
254
255 let mut current_time = start_time;
256 let mut response_key: Option<String> = None;
257 let mut last_response_id: String = self.get_last_response_id(&mut connection)?;
258
259 self.enqueue(event)?;
260
261 while start_time + time::Duration::new(timeout.into(), 0) >= current_time {
262 let new_responses: Vec<StreamMap> = match connection.xread(
264 &[&self.response_stream_name],
265 &[&last_response_id]
266 ) {
267 Err(error) => return Err(EventQueueError::DequeueError(error.to_string())),
268 Ok(response_vec) => response_vec
269 };
270
271 if new_responses.is_empty() {
273 current_time = time::Instant::now();
274 continue;
275 }
276
277 let response_map = &new_responses[0];
279
280 let new_responses = match response_map.get(&self.response_stream_name) {
282 None => return Err(EventQueueError::DequeueError(String::from("invalid stream name in response map"))),
283 Some(response_vec) => response_vec
284 };
285
286 for response in new_responses {
287 let response_id = match response.keys().next() {
289 None => return Err(EventQueueError::DequeueError(String::from("no response ID in response map"))),
290 Some(id) => id.clone()
291 };
292
293 let response_metadata = match response.get(&response_id) {
295 None => return Err(EventQueueError::DequeueError(std::format!("no metadata stored for response ID {}", response_id))),
296 Some(data) => data
297 };
298
299 let found_uuid_string = match response_metadata.keys().next() {
301 None => return Err(EventQueueError::DequeueError(std::format!("UUID string not found in metadata {:#?}", response_metadata))),
302 Some(uuid) => uuid.clone()
303 };
304
305 if found_uuid_string != target_uuid_string {
307 last_response_id = response_id;
308 continue;
309 }
310
311 response_key = match response_metadata.get(&target_uuid_string) {
313 None => return Err(EventQueueError::DequeueError(std::format!("failed to get response key from metadata {:#?}", response_metadata))),
314 Some(key) => Some(key.clone())
315 };
316
317 break;
320 }
321
322 if response_key.is_some() {
324 break;
325 }
326
327 current_time = time::Instant::now();
329 }
330
331 let response_key = match response_key {
333 None => return Err(EventQueueError::TimeoutExpired),
334 Some(response) => response
335 };
336
337 let response = self.get_service_event_by_key(&mut connection, &response_key, "response")?;
339 let timestamp = Self::extract_timestamp_from_event_key(&response_key);
340
341 Ok(TimestampedEvent(timestamp, response))
342 }
343}
344
345#[cfg(test)]
346mod tests {
347 use std::time::Duration;
348 use std::thread;
349 use super::*;
350
351 #[test]
352 fn create_ok() {
353 let _interface = EventQueue::new(
354 "test_queue",
355 "redis://127.0.0.1"
356 );
357 }
358
359 #[test]
360 fn enqueue_dequeue_ok() {
361 let mut interface = EventQueue::new(
362 "test_event_enqueue_dequeue",
363 "redis://127.0.0.1"
364 );
365
366 let event = ServiceEvent::new(
367 10,
368 "test_enqueue",
369 None
370 );
371
372 interface.enqueue(&event).unwrap();
373
374 let result = interface.dequeue().unwrap();
375
376 assert_eq!(&event, result.event());
377 }
378
379 #[test]
380 fn dequeue_blocking_ok() {
381 let mut interface = EventQueue::new(
382 "test_event_dequeue_blocking",
383 "redis://127.0.0.1"
384 );
385
386 let event = ServiceEvent::new(
387 10,
388 "test_enqueue",
389 Some(String::from("Payload!"))
390 );
391
392 let event_uuid = event.uuid();
393
394 let handle = thread::spawn(move || {
395 thread::sleep(Duration::from_secs(2));
396
397 let mut local_interface = EventQueue::new(
398 "test_event_dequeue_blocking",
399 "redis://127.0.0.1"
400 );
401
402 local_interface.enqueue(&event).unwrap();
403 });
404
405 let result = interface.dequeue_blocking(10).unwrap();
406
407 handle.join().unwrap();
408
409 assert_eq!(event_uuid, result.event().uuid());
410 assert_eq!(result.event().payload(), Some(String::from("Payload!")));
411 }
412
413 #[test]
414 #[should_panic(expected="called `Result::unwrap()` on an `Err` value: EmptyQueue")]
415 fn dequeue_blocking_timeout() {
416 let mut interface = EventQueue::new(
417 "test_event_dequeue_blocking_timeout",
418 "redis://127.0.0.1"
419 );
420
421 interface.dequeue_blocking(1).unwrap();
422 }
423
424 #[test]
425 fn await_ok() {
426 let mut interface = EventQueue::new(
427 "test_event_await",
428 "redis://127.0.0.1"
429 );
430
431 let event = ServiceEvent::new(
432 10,
433 "await_test",
434 Some(String::from("ping"))
435 );
436
437 let join_handle = thread::spawn(|| {
438 let mut thread_interface = EventQueue::new(
439 "test_event_await",
440 "redis://127.0.0.1"
441 );
442
443 let event = thread_interface.dequeue_blocking(10).unwrap();
444 let event = event.event();
445
446 println!("{:#?}", event);
447
448 assert_eq!(event.payload(), Some(String::from("ping")));
449
450 let response = ServiceEvent::new_response(event, "await_response", Some(String::from("pong")));
451 thread_interface.enqueue_response(&response).unwrap();
452 });
453
454 let response = interface.await_response(&event).unwrap();
455 let response = response.event();
456
457 join_handle.join().unwrap();
458
459 assert_eq!(response.action(), "await_response");
460 assert_eq!(response.payload(), Some(String::from("pong")));
461 assert_eq!(response.uuid(), event.uuid());
462 }
463
464 #[test]
465 fn simultaneous_await_ok() {
466 let mut interface = EventQueue::new(
467 "test_event_await_sim",
468 "redis://127.0.0.1"
469 );
470
471 let answer_thread = thread::spawn(|| {
472 let mut thread_interface = EventQueue::new(
473 "test_event_await_sim",
474 "redis://127.0.0.1"
475 );
476
477 for _ in 0..2 {
478 let event = thread_interface.dequeue_blocking(10).unwrap();
479 let event = event.event();
480
481 assert_eq!(event.payload(), Some(String::from("ping")));
482
483 let response = ServiceEvent::new_response(event, "await_response", Some(String::from("pong")));
484 thread_interface.enqueue_response(&response).unwrap();
485 }
486 });
487
488 let event_thread = thread::spawn(|| {
489 let mut thread_interface = EventQueue::new(
490 "test_event_await_sim",
491 "redis://127.0.0.1"
492 );
493
494 let event = ServiceEvent::new(
495 1,
496 "await_test",
497 Some(String::from("ping"))
498 );
499
500 let response = thread_interface.await_response(&event).unwrap();
501 let response = response.event();
502
503 assert_eq!(response.action(), "await_response");
504 assert_eq!(response.payload(), Some(String::from("pong")));
505 assert_eq!(response.uuid(), event.uuid());
506 });
507
508 let event = ServiceEvent::new(
509 1,
510 "await_test",
511 Some(String::from("ping"))
512 );
513
514 let response = interface.await_response(&event).unwrap();
515 let response = response.event();
516
517 assert_eq!(response.action(), "await_response");
518 assert_eq!(response.payload(), Some(String::from("pong")));
519 assert_eq!(response.uuid(), event.uuid());
520
521 answer_thread.join().unwrap();
522 event_thread.join().unwrap();
523 }
524
525 #[test]
526 #[should_panic(expected="called `Result::unwrap()` on an `Err` value: TimeoutExpired")]
527 fn await_timeout() {
528 let mut interface = EventQueue::new(
529 "test_event_await_timeout",
530 "redis://127.0.0.1"
531 );
532
533 let event = ServiceEvent::new(
534 1,
535 "await_test",
536 Some(String::from("ping"))
537 );
538
539 interface.await_response(&event).unwrap();
540 }
541}