Skip to main content

elk_mq/
event_queue.rs

1//  Copyright 2022 Tijmen Menno Verhoef
2
3//  Licensed under the Apache License, Version 2.0 (the "License");
4//  you may not use this file except in compliance with the License.
5//  You may obtain a copy of the License at
6
7//      http://www.apache.org/licenses/LICENSE-2.0
8
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14
15mod 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            // read new response entries from last seen ID onward
263            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 no new responses are found, we continue with polling
272            if new_responses.is_empty() {
273                current_time = time::Instant::now();
274                continue;
275            }
276
277            // only 1 stream is read, convert [ hashmap ] -> hashmap
278            let response_map = &new_responses[0];
279
280            // extract the stream name and verify it actually matches read stream
281            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                // extract response id for this entry, we know only 1 exists because of structure (id, (key, data))
288                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                // extract metadata
294                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                // extract uuid string from metadata
300                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                // check if we are looking for this string
306                if found_uuid_string != target_uuid_string {
307                    last_response_id = response_id;
308                    continue;
309                }
310
311                // fetch the key we are looking for
312                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                // after extracting the key we are done with the loop, so early break
318                // UUID is guaranteed unique with low collisions, so looking further will provide no benefit
319                break;
320            }
321
322            // break polling if a response key is found
323            if response_key.is_some() {
324                break;
325            }
326
327            // update our current time to detect when timeout is done
328            current_time = time::Instant::now();
329        }
330
331        // check if we found a response key
332        let response_key = match response_key {
333            None => return Err(EventQueueError::TimeoutExpired),
334            Some(response) => response
335        };
336
337        // create a timestamped event from found data
338        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}