worked on a thread safe message queue
This commit is contained in:
parent
19d5911f1e
commit
7934c4f71c
@ -10,3 +10,5 @@ lfcm = { git = "https://gitea.linvogel.ch/linus/lfcm.git", tag = "v0.1.2" }
|
|||||||
minijinja = "2.24.0"
|
minijinja = "2.24.0"
|
||||||
serde-saphyr = "1.1.0"
|
serde-saphyr = "1.1.0"
|
||||||
serde_json = "1.0.151"
|
serde_json = "1.0.151"
|
||||||
|
serde = { version = "1.0.229", features = ["derive"] }
|
||||||
|
uuid = { version = "1.26.1", features = ["v4"] }
|
||||||
|
|||||||
13
src/com/message.rs
Normal file
13
src/com/message.rs
Normal file
@ -0,0 +1,13 @@
|
|||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
|
#[derive(Serialize, Deserialize)]
|
||||||
|
pub enum Message {
|
||||||
|
EchoRequest,
|
||||||
|
EchoResponse,
|
||||||
|
GrainRequest,
|
||||||
|
GrainResponse,
|
||||||
|
PillarRequest,
|
||||||
|
PillarResponse,
|
||||||
|
StateApplyRequest,
|
||||||
|
StateApplyResponse,
|
||||||
|
}
|
||||||
104
src/com/message_queue.rs
Normal file
104
src/com/message_queue.rs
Normal file
@ -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<Mutex<HashMap<Uuid, Mutex<VecDeque<(Uuid, Message)>>>>>,
|
||||||
|
cardinality: Arc<AtomicUsize>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Option<(Uuid, Message)>, 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(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
2
src/com/mod.rs
Normal file
2
src/com/mod.rs
Normal file
@ -0,0 +1,2 @@
|
|||||||
|
pub mod message;
|
||||||
|
pub mod message_queue;
|
||||||
@ -1,10 +1,15 @@
|
|||||||
|
#![feature(mpmc_channel)]
|
||||||
|
|
||||||
|
|
||||||
#[macro_use]
|
#[macro_use]
|
||||||
extern crate lfcm;
|
extern crate lfcm;
|
||||||
|
|
||||||
|
|
||||||
use crate::config::LFCMError;
|
use crate::config::LFCMError;
|
||||||
|
|
||||||
pub mod config;
|
pub mod config;
|
||||||
pub mod sls;
|
pub mod sls;
|
||||||
|
pub mod com;
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
// TODO: handle some parameters
|
// TODO: handle some parameters
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user