From 7934c4f71cc12fc282bf714b9ee77d4295a69a10 Mon Sep 17 00:00:00 2001 From: Linus Vogel Date: Mon, 14 Sep 2026 22:28:31 +0200 Subject: [PATCH] worked on a thread safe message queue --- Cargo.toml | 2 + src/com/message.rs | 13 +++++ src/com/message_queue.rs | 104 +++++++++++++++++++++++++++++++++++++++ src/com/mod.rs | 2 + src/main.rs | 5 ++ 5 files changed, 126 insertions(+) create mode 100644 src/com/message.rs create mode 100644 src/com/message_queue.rs create mode 100644 src/com/mod.rs diff --git a/Cargo.toml b/Cargo.toml index 759c8a0..c2ed533 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,3 +10,5 @@ lfcm = { git = "https://gitea.linvogel.ch/linus/lfcm.git", tag = "v0.1.2" } minijinja = "2.24.0" serde-saphyr = "1.1.0" serde_json = "1.0.151" +serde = { version = "1.0.229", features = ["derive"] } +uuid = { version = "1.26.1", features = ["v4"] } diff --git a/src/com/message.rs b/src/com/message.rs new file mode 100644 index 0000000..5a06617 --- /dev/null +++ b/src/com/message.rs @@ -0,0 +1,13 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Serialize, Deserialize)] +pub enum Message { + EchoRequest, + EchoResponse, + GrainRequest, + GrainResponse, + PillarRequest, + PillarResponse, + StateApplyRequest, + StateApplyResponse, +} \ No newline at end of file diff --git a/src/com/message_queue.rs b/src/com/message_queue.rs new file mode 100644 index 0000000..0015189 --- /dev/null +++ b/src/com/message_queue.rs @@ -0,0 +1,104 @@ +use std::collections::vec_deque::VecDeque; +use std::collections::HashMap; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use uuid::Uuid; +use crate::com::message::Message; + +pub struct AddressableMessageQueue { + queues: Arc>>>>, + cardinality: Arc, +} + +pub enum MessageQueueError { + LockingPoisonError (String), + UnableToCreateRecipientQueue, +} + +impl AddressableMessageQueue { + pub fn new() -> Self { + Self { + queues: Arc::new(Mutex::new(HashMap::new())), + cardinality: Arc::new(AtomicUsize::new(0)), + } + } + + pub fn push(&mut self, sender: Uuid, recipient: Uuid, msg: Message) -> Result<(), MessageQueueError> { + // lock the data structure + let mut queues = match self.queues.lock() { + Ok(x) => x, + Err(e) => return Err(MessageQueueError::LockingPoisonError(e.to_string())), + }; + + // ensure that the recipient is in the data structure + let mut addressed_mutex = { + match queues.get_mut(&recipient) { + Some(addressed_mutex) => match addressed_mutex.lock() { + Ok(queue) => queue, + Err(e) => return Err(MessageQueueError::LockingPoisonError(e.to_string())), + }, + None => { + let addressed_mutex = match queues.insert(recipient, Mutex::new(VecDeque::new())) { + Some(x) => x, + None => return Err(MessageQueueError::UnableToCreateRecipientQueue) + }; + + let mut queue = match addressed_mutex.lock() { + Ok(y) => y, + Err(e) => return Err(MessageQueueError::LockingPoisonError(e.to_string())) + }; + + // insert the message in question + queue.push_back((sender, msg)); + self.cardinality.fetch_add(1, Ordering::SeqCst); + + return Ok (()) + } + } + }; + + // insert the message in question + addressed_mutex.push_back((sender, msg)); + self.cardinality.fetch_add(1, Ordering::SeqCst); + + // return success + Ok (()) + } + + pub fn pop(&mut self, recipient: Uuid) -> Result, MessageQueueError> { + // lock the data structure + let mut queues = match self.queues.lock() { + Ok(x) => x, + Err(e) => return Err(MessageQueueError::LockingPoisonError(e.to_string())), + }; + + // check if the recipient is present in the data structure + let mut addressed_mutex = { + match queues.get_mut(&recipient) { + Some(addressed_mutex) => match addressed_mutex.lock() { + Ok(queue) => queue, + Err(e) => return Err(MessageQueueError::LockingPoisonError(e.to_string())), + }, + None => return Ok (None), + } + }; + + // if recipient is present in the data structure + match addressed_mutex.pop_front() { + Some(x) => { + self.cardinality.fetch_sub(1, Ordering::SeqCst); + Ok(Some(x)) + }, + None => return Ok (None) + } + } +} + +impl Clone for AddressableMessageQueue { + fn clone(&self) -> Self { + Self { + queues: self.queues.clone(), + cardinality: self.cardinality.clone(), + } + } +} \ No newline at end of file diff --git a/src/com/mod.rs b/src/com/mod.rs new file mode 100644 index 0000000..da645bb --- /dev/null +++ b/src/com/mod.rs @@ -0,0 +1,2 @@ +pub mod message; +pub mod message_queue; \ No newline at end of file diff --git a/src/main.rs b/src/main.rs index ec9026b..1f4f778 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,10 +1,15 @@ +#![feature(mpmc_channel)] + + #[macro_use] extern crate lfcm; + use crate::config::LFCMError; pub mod config; pub mod sls; +pub mod com; fn main() { // TODO: handle some parameters