diff --git a/src/storage.rs b/src/storage.rs index 05d89d6..a1ab246 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -27,10 +27,10 @@ mod rng { } mod test { - use super::*; - #[tokio::test] async fn shuffle_slice() { + use super::*; + let mut data = vec![1, 2, 3, 4]; let mut rng = PrngState::from_thread_rng(); @@ -40,13 +40,22 @@ mod rng { } } +#[derive(Debug, Clone)] pub struct Quote(String); -impl Quote { - pub fn new(quote: String) -> Self { - Self(quote) +impl> From for Quote { + fn from(value: S) -> Self { + Self(value.as_ref().to_owned()) } +} +impl From for String { + fn from(value: Quote) -> Self { + value.0 + } +} + +impl Quote { pub fn get(&self) -> &str { &self.0 } @@ -57,20 +66,120 @@ pub mod source { } pub mod queue { + use std::collections::VecDeque; + use super::*; // Messages to interact with the quote queue pub struct DequeueQuote; pub struct EnqueueQuotes(pub Vec); + /// Subtrait of Actor which specifically + /// denotes actors that can handle all relevant queue messages. pub trait QueueManager: Actor {} impl QueueManager for T where T: Message + Message {} + + /// A basic implementation of an in-memory queue of quotes. + /// Its contents are *not* persisted across application restarts, so it + /// is only suited for testing purposes. + #[derive(Actor)] + pub struct MemoryQueueStorage { + quotes: VecDeque, + } + + impl MemoryQueueStorage { + pub fn new() -> Self { + Self { + quotes: VecDeque::new(), + } + } + } + + impl Message for MemoryQueueStorage { + /// We only need to signal success or failure in this instance, + /// with no added metadata in either case. + type Reply = Result<(), ()>; + + async fn handle( + &mut self, + msg: EnqueueQuotes, + _ctx: &mut Context, + ) -> Self::Reply { + for q in msg.0 { + self.quotes.push_back(q); + } + + Ok(()) + } + } + + impl Message for MemoryQueueStorage { + type Reply = Result, ()>; + + async fn handle( + &mut self, + _msg: DequeueQuote, + _ctx: &mut Context, + ) -> Self::Reply { + // Note: this can never fail, since the quotes are stored in memory + Ok(self.quotes.pop_front()) + } + } } -use self::queue::*; struct QuoteCycle { rng: rng::PrngState, source_manager: (), queue_manager: ActorRef, } + +impl QuoteCycle { + pub fn new(rng: rng::PrngState, source_manager: (), queue_manager: ActorRef) -> Self { + Self { + rng, + source_manager, + queue_manager, + } + } + + pub fn with_thread_rng(source_manager: (), queue_manager: ActorRef) -> Self { + Self { + rng: rng::PrngState::from_thread_rng(), + source_manager, + queue_manager, + } + } +} + +mod test { + #[tokio::test] + async fn memory_queue() { + use super::Quote; + use super::queue::*; + use kameo::prelude::*; + + let queue_manager = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); + + let sample_quotes = ["Test no.1", "Test no.2", "Test no.3"]; + + let _ = queue_manager + .ask(EnqueueQuotes( + sample_quotes.iter().cloned().map(Quote::from).collect(), + )) + .await + .expect("In-memory quote queue storage should be valid for insertion"); + + for text in sample_quotes.iter() { + assert_eq!( + *text, + queue_manager + .ask(DequeueQuote) + .await + .expect("In-memory queue storage should never panic on dequeue") + .expect("In-memory queue storage should never be initialized as empty") + .get() + ); + } + } +}