diff --git a/src/lib.rs b/src/lib.rs index ececcca..5776215 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,8 +6,10 @@ pub mod run { use std::time::Duration; use crate::sink::{BskySink, PostQuote, SinkManager, StdoutSink}; + use crate::storage::queue::QueueManager; use crate::storage::{ - FetchQuote, QuoteCycle, queue::MemoryQueueStorage, source::FsFilterSourceManager, + FetchQuote, QuoteCycle, queue::MemoryQueueStorage, queue::RedisQueueStorage, + source::FsFilterSourceManager, }; use cron_lite::CronEvent; use futures::StreamExt; @@ -33,6 +35,30 @@ pub mod run { None }; + let use_redis = std::env::var("USE_REDIS").unwrap_or("0".to_string()) == "1"; + if use_redis { + let redis = redis::Client::open( + std::env::var("REDIS_URL").unwrap_or("redis://localhost".to_string()), + )?; + let con = redis.get_multiplexed_async_connection().await?; + + /// The default queue key we use to fetch and store Redis quote. + // TODO(redis): Parameterize the quote queue key. + const DEFAULT_KEY: &str = "queue:default"; + + let queue = + RedisQueueStorage::spawn(RedisQueueStorage::new(con, DEFAULT_KEY.to_string())); + run_cycle(queue, bsky).await + } else { + let queue = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); + run_cycle(queue, bsky).await + } + } + + async fn run_cycle( + queue: ActorRef, + bsky: Option>, + ) -> Result<(), Box> { let sink = { let stdout = StdoutSink::spawn(StdoutSink); SinkManager::spawn(SinkManager::new(Some(stdout), bsky)) @@ -40,8 +66,6 @@ pub mod run { let cycle = { let source = FsFilterSourceManager::spawn(FsFilterSourceManager::default()); - let queue = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); - QuoteCycle::spawn(QuoteCycle::with_thread_rng(source, queue)) }; diff --git a/src/storage.rs b/src/storage.rs index 6f555c5..87db09c 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -245,6 +245,82 @@ pub mod queue { Ok(self.quotes.pop_front()) } } + + /// Implementation of persistent quote queue storage through a Redis server. + #[derive(Actor, Default)] + pub struct RedisQueueStorage + where + R: redis::aio::ConnectionLike + redis::AsyncCommands + 'static, + { + /// Key that points to the Redis queue to be used for fetching and storing quotes. + queue_key: String, + + /// Underlying connection representing the Redis client; must be boxed to derive [Actor]. + con: Box, + } + + impl RedisQueueStorage + where + R: redis::aio::ConnectionLike + redis::AsyncCommands + 'static, + { + pub fn new(con: R, queue_key: String) -> Self { + Self { + queue_key, + con: Box::new(con), + } + } + } + + impl Message for RedisQueueStorage + where + R: redis::aio::ConnectionLike + redis::AsyncCommands + 'static, + { + /// We only need to signal success or failure in this instance, + /// with no added metadata in either case. + type Reply = EnqueueReply; + + async fn handle( + &mut self, + msg: EnqueueQuotes, + _ctx: &mut Context, + ) -> Self::Reply { + for q in msg.0 { + let _: () = self + .con + .rpush(&self.queue_key, q.get()) + .await + .map_err(|_| ())?; + } + + Ok(()) + } + } + + impl Message for RedisQueueStorage + where + R: redis::aio::ConnectionLike + redis::AsyncCommands + 'static, + { + type Reply = DequeueReply; + + async fn handle( + &mut self, + _msg: DequeueQuote, + _ctx: &mut Context, + ) -> Self::Reply { + let quote: Result, ()> = + self.con.lpop(&self.queue_key, None).await.map_err(|_| ()); + + match quote { + Ok(Some(quote)) => Ok(Some(Quote::from(quote))), + + // We need to separate these next two cases, + // as the type of quote: `Result>` + // does not match the return type: `Result>` + Ok(None) => Ok(None), + Err(()) => Err(()), + } + } + } } #[derive(Actor)]