diff --git a/Cargo.lock b/Cargo.lock index ffddf69..db18679 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -70,6 +70,17 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "async-trait" +version = "0.1.91" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae36dc4177970ef04fde5178d3e2429882def40e57a451f919c098f72baa6cec" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "atomic-waker" version = "1.1.2" @@ -140,13 +151,13 @@ dependencies = [ name = "audquotes" version = "0.1.0" dependencies = [ + "async-trait", "bsky-sdk", "chrono", "cron-lite", "futures", "glob", "grep", - "kameo", "opentelemetry", "opentelemetry_sdk", "rand", @@ -445,7 +456,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn", + "syn 2.0.117", ] [[package]] @@ -456,7 +467,7 @@ checksum = "fc34b93ccb385b40dc71c6fceac4b2ad23662c7eeb248cf10d529b7e055b6ead" dependencies = [ "darling_core", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -496,7 +507,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7ab67060fc6b8ef687992d439ca0fa36e7ed17e9a0b16b25b601e8757df720de" dependencies = [ "data-encoding", - "syn", + "syn 2.0.117", ] [[package]] @@ -517,7 +528,7 @@ dependencies = [ "darling", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -527,7 +538,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ab63b0e2bf4d5928aff72e83a7dace85d7bba5fe12dcc3c5a572d78caffd3f3c" dependencies = [ "derive_builder_core", - "syn", + "syn 2.0.117", ] [[package]] @@ -538,21 +549,9 @@ checksum = "97369cbbc041bc366949bc74d34658d6cda5621039731c6310521892a3a20ae0" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] -[[package]] -name = "downcast-rs" -version = "2.0.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "117240f60069e65410b3ae1bb213295bd828f707b5bec6596a1afc8793ce0cbc" - -[[package]] -name = "dyn-clone" -version = "1.0.20" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555" - [[package]] name = "encoding_rs" version = "0.8.35" @@ -722,7 +721,7 @@ checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1198,35 +1197,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "kameo" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "41a73be96f616ca2784f597b5b6635582f5a7b3ba73b1dbe7afa5d9667955d39" -dependencies = [ - "downcast-rs", - "dyn-clone", - "futures", - "kameo_macros", - "once_cell", - "serde", - "tokio", - "tracing", -] - -[[package]] -name = "kameo_macros" -version = "0.17.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b3f384b32bf6426ae93a8b37da62c85073b676a31a82a86d608ad86453878de0" -dependencies = [ - "heck", - "proc-macro2", - "quote", - "syn", - "uuid", -] - [[package]] name = "langtag" version = "0.3.4" @@ -1298,7 +1268,7 @@ checksum = "757aee279b8bdbb9f9e676796fd459e4207a1f986e87886700abf589f5abf771" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1424,7 +1394,7 @@ checksum = "ed3955f1a9c7c0c15e092f9c887db08b1fc683305fdf6eb6684f22555355e202" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1474,7 +1444,7 @@ checksum = "a948666b637a0f465e8564c73e89d4dde00d72d4d473cc972f390fc3dcee7d9c" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1594,7 +1564,7 @@ checksum = "d9b20ed30f105399776b9c883e68e536ef602a16ae6f596d2c473591d6ad64c6" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -1646,7 +1616,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "479ca8adacdd7ce8f1fb39ce9ecccbfe93a3f1344b3d0d97f20bc0196208f62b" dependencies = [ "proc-macro2", - "syn", + "syn 2.0.117", ] [[package]] @@ -1936,7 +1906,7 @@ checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2082,7 +2052,7 @@ dependencies = [ "heck", "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2096,6 +2066,17 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "syn" +version = "3.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + [[package]] name = "sync_wrapper" version = "1.0.2" @@ -2113,7 +2094,7 @@ checksum = "728a70f3dbaf5bab7f0c4b1ac8d7ae5ea60a4b5549c8a5914361c99147a709d2" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2170,7 +2151,7 @@ checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2181,7 +2162,7 @@ checksum = "ebc4ee7f67670e9b64d05fa4253e753e016c6c95ff35b89b7941d6b856dec1d5" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2217,7 +2198,6 @@ dependencies = [ "signal-hook-registry", "socket2 0.6.3", "tokio-macros", - "tracing", "windows-sys 0.61.2", ] @@ -2245,7 +2225,7 @@ checksum = "5c55a2eff8b69ce66c84f85e1da1c233edc36ceb85a2058d11b0d6a3c7e7569c" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2340,7 +2320,7 @@ checksum = "7490cfa5ec963746568740651ac6781f701c9c5ea257c58e057f3ba8cf69e8da" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2402,7 +2382,7 @@ checksum = "70977707304198400eb4835a78f6a9f928bf41bba420deb8fdb175cd965d77a7" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2555,7 +2535,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn", + "syn 2.0.117", "wasm-bindgen-shared", ] @@ -2652,7 +2632,7 @@ checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2663,7 +2643,7 @@ checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2802,7 +2782,7 @@ dependencies = [ "heck", "indexmap", "prettyplease", - "syn", + "syn 2.0.117", "wasm-metadata", "wit-bindgen-core", "wit-component", @@ -2818,7 +2798,7 @@ dependencies = [ "prettyplease", "proc-macro2", "quote", - "syn", + "syn 2.0.117", "wit-bindgen-core", "wit-bindgen-rust", ] @@ -2885,7 +2865,7 @@ checksum = "b659052874eb698efe5b9e8cf382204678a0086ebf46982b79d6ca3182927e5d" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", "synstructure", ] @@ -2906,7 +2886,7 @@ checksum = "0e8bc7269b54418e7aeeef514aa68f8690b8c0489a06b0136e5f57c4c5ccab89" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -2926,7 +2906,7 @@ checksum = "d71e5d6e06ab090c67b5e44993ec16b72dcbaabc526db883a360057678b48502" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", "synstructure", ] @@ -2966,7 +2946,7 @@ checksum = "eadce39539ca5cb3985590102671f2567e659fca9666581ad3411d59207951f3" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index e9e851a..4f649bb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,13 +5,13 @@ edition = "2024" rust-version = "1.93" [dependencies] +async-trait = "0.1.91" bsky-sdk = "0.1.16" chrono = "0.4.42" cron-lite = { version = "0.3.0", features = ["async"] } futures = "0.3.31" glob = "0.3.2" grep = "0.3.2" -kameo = "0.17.2" opentelemetry = "0.31.0" opentelemetry_sdk = "0.31.0" rand = "0.9.0" diff --git a/src/data.rs b/src/data.rs index e4d3dcf..40620d6 100644 --- a/src/data.rs +++ b/src/data.rs @@ -22,3 +22,12 @@ impl Quote { &self.0 } } + +pub trait Message { + type Reply; +} + +#[async_trait::async_trait(?Send)] +pub trait Handle { + async fn handle(&'_ mut self, msg: M) -> M::Reply; +} diff --git a/src/lib.rs b/src/lib.rs index 6e0a82b..045f3d8 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,6 +5,7 @@ pub mod storage; pub mod run { use std::time::Duration; + use crate::data::Handle; use crate::sink::{BskySink, PostQuote, SinkManager, StdoutSink}; use crate::storage::queue::QueueManager; use crate::storage::source::{QuoteFilter, SourceManager}; @@ -14,21 +15,22 @@ pub mod run { }; use cron_lite::CronEvent; use futures::StreamExt; - use kameo::prelude::*; use tokio::time::timeout; pub async fn entrypoint() -> Result<(), Box> { // TODO: Clean up this function's internals. // The current structure is alright, but it was stitched together // quickly just to confirm that everything is functioning as it should. + let mut sinks = SinkManager::new(); + let use_bsky = std::env::var("USE_BLUESKY").unwrap_or("0".to_string()) == "1"; - let bsky = if use_bsky { - Some(BskySink::spawn(BskySink::new_from_env().await.expect( - "Could not connect to Bluesky with supplied credentials", - ))) - } else { - None - }; + if use_bsky { + sinks = sinks.add_sink( + BskySink::new_from_env() + .await + .expect("Could not connect to Bluesky with supplied credentials"), + ); + } // If we're in "debug" mode, posts become more frequent. let is_debug = std::env::var("DEBUG").unwrap_or("0".to_string()) == "1"; @@ -38,38 +40,34 @@ pub mod run { // Debug schedule "*/10 * * * * * *" }; - let source = FsFilterSourceManager::spawn(if !is_debug { + let source = if !is_debug { FsFilterSourceManager::new(QuoteFilter::new_from_glob("quotes/**/*.txt".to_string())) } else { // Debug quotes fall back to the default configuration FsFilterSourceManager::default() - }); - let stdout = if !is_debug { - None - } else { - Some(StdoutSink::spawn(StdoutSink)) + }; + + if is_debug { + sinks = sinks.add_sink(StdoutSink); }; let use_redis = std::env::var("USE_REDIS").unwrap_or("0".to_string()) == "1"; if use_redis { - let queue = RedisQueueStorage::spawn(RedisQueueStorage::new_from_env().unwrap()); - run_cycle(schedule, source, queue, stdout, bsky).await + let queue = RedisQueueStorage::new_from_env().unwrap(); + run_cycle(schedule, source, queue, sinks).await } else { - let queue = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); - run_cycle(schedule, source, queue, stdout, bsky).await + let queue = MemoryQueueStorage::new(); + run_cycle(schedule, source, queue, sinks).await } } async fn run_cycle( schedule: &str, - source: ActorRef, - queue: ActorRef, - stdout: Option>, - bsky: Option>, + source: S, + queue: Q, + mut sink: SinkManager, ) -> Result<(), Box> { - let sink = { SinkManager::spawn(SinkManager::new(stdout, bsky)) }; - - let cycle = { QuoteCycle::spawn(QuoteCycle::with_thread_rng(source, queue)) }; + let mut cycle = { QuoteCycle::with_thread_rng(source, queue) }; use cron_lite::Schedule; const POSTING_TIMEOUT: Duration = Duration::from_secs(60); @@ -94,10 +92,10 @@ pub mod run { // We store the code to perform the next posting iteration as one atomic future which we wrap with a timeout. // This means that, if we miss a posting window due to the timeout, we will not get multiple consecutive or late posts. - let next_post_iteration = async || -> Result<(), Box> { + let mut next_post_iteration = async || -> Result<(), Box> { tracing::debug!("Fetching next quote"); let next_quote = cycle - .ask(FetchQuote) + .handle(FetchQuote) .await .map_err(|_| "fetch quote should always succeed")?; @@ -108,7 +106,7 @@ pub mod run { quote = next_quote.get().to_string(), "Forwarding quote to sink" ); - sink.tell(PostQuote(next_quote)).await?; + sink.handle(PostQuote(next_quote)).await?; println!(); Ok(()) diff --git a/src/sink.rs b/src/sink.rs index 22bb3a6..616e5ac 100644 --- a/src/sink.rs +++ b/src/sink.rs @@ -1,6 +1,6 @@ use crate::data::Quote; +use crate::data::{Handle, Message}; use bsky_sdk::{BskyAgent, api::types::Object}; -use kameo::prelude::*; /// A newtype over [Quote] used to prompt the [SinkManager] to /// submit a new quote to all its configured sinks. @@ -29,26 +29,35 @@ pub enum PostFailure { Unrecoverable, } +impl std::fmt::Display for PostFailure { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + PostFailure::Retry { .. } => f.write_str("PostFailure::Retry"), + PostFailure::Unsupported => f.write_str("PostFailure::Unsupported"), + PostFailure::Unrecoverable => f.write_str("PostFailure::Unrecoverable"), + } + } +} + +impl std::error::Error for PostFailure {} + pub type PostResult = Result<(), PostFailure>; +impl Message for PostQuote { + type Reply = PostResult; +} /// Represents internal implementation details of the interactions between /// the [SinkManager] and its sinks. -pub trait QuoteSink: Actor + Message {} +pub trait QuoteSink: Handle {} /// A [QuoteSink] which will output the contents of each quote /// over Stdout. Is primarily meant for testing and observing sink behavior. -#[derive(Actor)] pub struct StdoutSink; -impl Message for StdoutSink { - type Reply = PostResult; - - #[tracing::instrument(name = "stdout::post", skip(self, _ctx))] - async fn handle( - &mut self, - PostQuote(quote): PostQuote, - _ctx: &mut Context, - ) -> Self::Reply { +#[async_trait::async_trait(?Send)] +impl Handle for StdoutSink { + #[tracing::instrument(name = "stdout::post", skip(self))] + async fn handle(&mut self, PostQuote(quote): PostQuote) -> PostResult { tracing::debug!(quote = quote.get()); println!("{}", quote.get()); Ok(()) @@ -56,7 +65,6 @@ impl Message for StdoutSink { } /// A [QuoteSink] which will post the contents of each quote to Bluesky. -#[derive(Actor)] pub struct BskySink { bsky_agent: BskyAgent, bsky_session: Object, @@ -146,14 +154,9 @@ impl BskySink { } } -impl Message for BskySink { - type Reply = PostResult; - - async fn handle( - &mut self, - PostQuote(quote): PostQuote, - _ctx: &mut Context, - ) -> Self::Reply { +#[async_trait::async_trait(?Send)] +impl Handle for BskySink { + async fn handle(&mut self, PostQuote(quote): PostQuote) -> PostResult { match self.submit_post(quote).await { Ok(_) => Ok(()), Err(_) => Err(PostFailure::Unrecoverable), @@ -165,56 +168,35 @@ impl Message for BskySink { /// [PostQuote] messages to them as they are received. /// The SinkManager will attempt to reinitialize failed sinks upon /// encountering recoverable errors. -#[derive(Actor)] pub struct SinkManager { - // Uh oh. As the [Actor] trait is *not* dyn-compatible, - // and I do not own its definition, I'm fairly certain that I cannot - // do asynchronous dynamic dispatch for it here. - // I've decided I'll limit this to one sink per implementation right now. - stdout_sink: Option>, - bsky_sink: Option>, - // ... + sinks: Vec>>, } impl SinkManager { - pub fn new( - stdout_sink: Option>, - bsky_sink: Option>, - ) -> Self { - Self { - stdout_sink, - bsky_sink, - } + pub fn new() -> Self { + Self { sinks: vec![] } + } + + pub fn add_sink + 'static>(mut self, sink: S) -> Self { + self.sinks.push(Box::new(sink)); + self } } pub type SinkReplies = Vec>; -impl Message for SinkManager { - type Reply = SinkReplies; - - #[tracing::instrument(name = "sink::post", skip(self, _ctx))] - async fn handle( - &mut self, - msg: PostQuote, - _ctx: &mut Context, - ) -> Self::Reply { +#[async_trait::async_trait(?Send)] +impl Handle for SinkManager { + #[tracing::instrument(name = "sink::post", skip(self))] + async fn handle(&mut self, msg: PostQuote) -> PostResult { use futures::future::join_all; - let stdout_result = self - .stdout_sink - .as_ref() - .map(|s| s.ask(msg.clone()).into_future()); - - let bsky_result = self - .bsky_sink - .as_ref() - .map(|s| s.ask(msg.clone()).into_future()); - - let futures = [stdout_result, bsky_result].into_iter().flatten(); - let results = join_all(futures).await; + let results = join_all(self.sinks.iter_mut().map(|s| s.handle(msg.clone()))).await; - results.iter().map(|r| r.clone().or(Err(()))).collect() + results + .iter() + .map(|r| r.clone().or(Err(PostFailure::Unrecoverable))) + .collect() } } @@ -223,12 +205,12 @@ mod test { async fn stdout_sink() { use super::*; - let stdout = StdoutSink::spawn(StdoutSink); - let manager = SinkManager::spawn(SinkManager::new(Some(stdout), None)); + let stdout = StdoutSink; + let mut manager = SinkManager::new().add_sink(stdout); let messages = ["First test!", "Second test.", "Third..."]; for msg in messages { - manager.tell(PostQuote(msg.into())).await.unwrap(); + manager.handle(PostQuote(msg.into())).await.unwrap(); } // Hopefully we don't crash...! diff --git a/src/storage.rs b/src/storage.rs index 60116a3..0389313 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -1,11 +1,10 @@ -use kameo::prelude::*; - use crate::storage::{ queue::{DequeueQuote, EnqueueQuotes}, source::SourceQuotes, }; use crate::data::Quote; +use crate::data::{Handle, Message}; mod rng { use rand::SeedableRng; @@ -55,17 +54,19 @@ pub mod source { /// Message to request that a SourceManager source its quotes once again. pub struct SourceQuotes; pub type SourceReply = Result, ()>; + impl Message for SourceQuotes { + type Reply = SourceReply; + } /// Subtrait of Actor which specifically /// denotes actors that can handle all relevant source messages. - pub trait SourceManager: Actor + Message {} + pub trait SourceManager: Handle {} - impl SourceManager for T where T: Message {} + impl SourceManager for T where T: Handle {} /// Implementation of [SourceManager] which sources quotes from a Vec /// that it holds in memory, without accessing external services. /// Its main purpose is to be used for testing. - #[derive(Actor)] pub struct MemorySourceManager { quotes: Vec, } @@ -78,14 +79,9 @@ pub mod source { } } - impl Message for MemorySourceManager { - type Reply = SourceReply; - - async fn handle( - &mut self, - _msg: SourceQuotes, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for MemorySourceManager { + async fn handle(&mut self, _msg: SourceQuotes) -> Result, ()> { // We just clone the quotes we've been holding onto since startup Ok(self.quotes.clone()) } @@ -93,7 +89,6 @@ pub mod source { /// Uses a [QuoteFilter] to source quotes from the local filesystem /// at the beginning of each cycle. - #[derive(Actor)] pub struct FsFilterSourceManager { filter: QuoteFilter, } @@ -117,14 +112,9 @@ pub mod source { } } - impl Message for FsFilterSourceManager { - type Reply = SourceReply; - - async fn handle( - &mut self, - _msg: SourceQuotes, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for FsFilterSourceManager { + async fn handle(&mut self, _msg: SourceQuotes) -> Result, ()> { self.filter.read_files() } } @@ -199,6 +189,9 @@ pub mod queue { #[derive(Debug)] pub struct DequeueQuote; pub type DequeueReply = Result, ()>; + impl Message for DequeueQuote { + type Reply = DequeueReply; + } #[derive(Debug)] pub struct EnqueueQuotes(pub Vec); @@ -209,26 +202,20 @@ pub mod queue { /// A [usize] is used as that is the maximum amount of items a [Vec] /// of quotes within [EnqueueQuotes] can contain. pub type EnqueueReply = Result<(), usize>; + impl Message for EnqueueQuotes { + type Reply = EnqueueReply; + } /// Subtrait of Actor which specifically /// denotes actors that can handle all relevant queue messages. - pub trait QueueManager: - Actor - + Message - + Message - { - } + pub trait QueueManager: Handle + Handle {} - impl QueueManager for T where - T: Message - + Message - { - } + impl QueueManager for T where T: Handle + Handle {} /// 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, Default)] + #[derive(Default)] pub struct MemoryQueueStorage { quotes: VecDeque, } @@ -241,14 +228,9 @@ pub mod queue { } } - impl Message for MemoryQueueStorage { - type Reply = EnqueueReply; - - async fn handle( - &mut self, - msg: EnqueueQuotes, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for MemoryQueueStorage { + async fn handle(&mut self, msg: EnqueueQuotes) -> Result<(), usize> { // Note: as `.push_back` into a `VecDeque` is an infallible operation, // we never need to worry about failure cases in this implementation. @@ -260,21 +242,15 @@ pub mod queue { } } - impl Message for MemoryQueueStorage { - type Reply = DequeueReply; - - async fn handle( - &mut self, - _msg: DequeueQuote, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for MemoryQueueStorage { + async fn handle(&mut self, _msg: DequeueQuote) -> Result, ()> { // Note: this can never fail, since the quotes are stored in memory Ok(self.quotes.pop_front()) } } /// Implementation of persistent quote queue storage through a Redis server. - #[derive(Actor)] pub struct RedisQueueStorage { /// Key that points to the Redis queue to be used for fetching and storing quotes. queue_key: String, @@ -313,15 +289,10 @@ pub mod queue { } } - impl Message for RedisQueueStorage { - type Reply = EnqueueReply; - - #[tracing::instrument(name = "redis::enqueue_many", skip(self, _ctx, msg))] - async fn handle( - &mut self, - msg: EnqueueQuotes, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for RedisQueueStorage { + #[tracing::instrument(name = "redis::enqueue_many", skip(self, msg))] + async fn handle(&mut self, msg: EnqueueQuotes) -> Result<(), usize> { tracing::debug!( to_enqueue = msg.0.len(), "Enqueuing {} quotes...", @@ -354,15 +325,10 @@ pub mod queue { } } - impl Message for RedisQueueStorage { - type Reply = DequeueReply; - - #[tracing::instrument(name = "redis::dequeue_one", skip(self, _ctx, _msg))] - async fn handle( - &mut self, - _msg: DequeueQuote, - _ctx: &mut Context, - ) -> Self::Reply { + #[async_trait::async_trait(?Send)] + impl Handle for RedisQueueStorage { + #[tracing::instrument(name = "redis::dequeue_one", skip(self, _msg))] + async fn handle(&mut self, _msg: DequeueQuote) -> Result, ()> { tracing::debug!("Dequeuing quote..."); let quote: Result, ()> = self .con() @@ -387,19 +353,14 @@ pub mod queue { } } -#[derive(Actor)] pub struct QuoteCycle { rng: rng::PrngState, - source_manager: ActorRef, - queue_manager: ActorRef, + source_manager: S, + queue_manager: Q, } impl QuoteCycle { - pub fn new( - rng: rng::PrngState, - source_manager: ActorRef, - queue_manager: ActorRef, - ) -> Self { + pub fn new(rng: rng::PrngState, source_manager: S, queue_manager: Q) -> Self { Self { rng, source_manager, @@ -407,7 +368,7 @@ impl QuoteCycle { } } - pub fn with_thread_rng(source_manager: ActorRef, queue_manager: ActorRef) -> Self { + pub fn with_thread_rng(source_manager: S, queue_manager: Q) -> Self { Self { rng: rng::PrngState::from_thread_rng(), source_manager, @@ -419,23 +380,26 @@ impl QuoteCycle { /// A message to [QuoteCycle] to fetch one more quote from its storage. #[derive(Debug)] pub struct FetchQuote; +impl Message for FetchQuote { + type Reply = Result; +} -impl Message for QuoteCycle +#[async_trait::async_trait(?Send)] +impl Handle for QuoteCycle where S: source::SourceManager, Q: queue::QueueManager, { - type Reply = Result; - - #[tracing::instrument(name = "quote_cycle::fetch", skip(self, _ctx))] - async fn handle( - &mut self, - _msg: FetchQuote, - _ctx: &mut Context, - ) -> Self::Reply { + #[tracing::instrument(name = "quote_cycle::fetch", skip(self))] + async fn handle(&mut self, _msg: FetchQuote) -> Result { // 1. We query our queue storage for the next quote tracing::debug!("Fetching new quotes..."); - if let Some(next_quote) = self.queue_manager.ask(DequeueQuote).await.map_err(|_| ())? { + if let Some(next_quote) = self + .queue_manager + .handle(DequeueQuote) + .await + .map_err(|_| ())? + { // if there is a quote, we simply return it and move on tracing::debug!( quote = next_quote.get().to_string(), @@ -445,10 +409,14 @@ where } // 2. Otherwise, we must repopulate the queue through our source - let mut refreshed_quotes = self.source_manager.ask(SourceQuotes).await.map_err(|_| { - tracing::error!("Could not retrieve new quotes from source"); - () - })?; + let mut refreshed_quotes = + self.source_manager + .handle(SourceQuotes) + .await + .map_err(|_| { + tracing::error!("Could not retrieve new quotes from source"); + () + })?; tracing::debug!("Retrieved {} quotes.", refreshed_quotes.len()); // 3. We shuffle the newly-sourced quotes @@ -460,7 +428,7 @@ where // 4. We enqueue the newly-sourced quotes... tracing::debug!("Enqueuing quotes..."); self.queue_manager - .ask(EnqueueQuotes(refreshed_quotes)) + .handle(EnqueueQuotes(refreshed_quotes)) .await .map_err(|_| { tracing::error!("Could not enqueue quotes to storage"); @@ -468,7 +436,7 @@ where })?; // 5. and, finally, we return the first among them. - match self.queue_manager.ask(DequeueQuote).await { + match self.queue_manager.handle(DequeueQuote).await { Ok(Some(q)) => Ok(q), Ok(None) => panic!("Newly-enqueued quotes should never be empty"), Err(_) => Err(()), @@ -481,14 +449,14 @@ mod test { async fn memory_queue() { use super::Quote; use super::queue::*; - use kameo::prelude::*; + use super::*; - let queue_manager = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); + let mut queue_manager = MemoryQueueStorage::new(); let sample_quotes = ["Test no.1", "Test no.2", "Test no.3"]; queue_manager - .ask(EnqueueQuotes( + .handle(EnqueueQuotes( sample_quotes.iter().cloned().map(Quote::from).collect(), )) .await @@ -498,7 +466,7 @@ mod test { assert_eq!( *text, queue_manager - .ask(DequeueQuote) + .handle(DequeueQuote) .await .expect("In-memory queue storage should never panic on dequeue") .expect("In-memory queue storage should never be initialized as empty") @@ -510,14 +478,14 @@ mod test { #[tokio::test] async fn memory_source() { use super::source::*; - use kameo::prelude::*; + use super::*; let sample_quotes = ["Minie", "Miney", "Moe", "and", "some", "more"]; - let source_manager = MemorySourceManager::spawn(MemorySourceManager::new(sample_quotes)); + let mut source_manager = MemorySourceManager::new(sample_quotes); let quotes = source_manager - .ask(SourceQuotes) + .handle(SourceQuotes) .await .expect("In-memory quote queue storage should be valid for insertion"); @@ -541,14 +509,14 @@ mod test { use super::QuoteCycle; use super::queue::*; use super::source::*; - use kameo::prelude::*; + use super::*; let sample_quotes = ["Minie", "Miney", "Moe"]; - let cycle = { - let source = MemorySourceManager::spawn(MemorySourceManager::new(sample_quotes)); - let queue = MemoryQueueStorage::spawn(MemoryQueueStorage::new()); + let mut cycle = { + let source = MemorySourceManager::new(sample_quotes); + let queue = MemoryQueueStorage::new(); - QuoteCycle::spawn(QuoteCycle::with_thread_rng(source, queue)) + QuoteCycle::with_thread_rng(source, queue) }; // We loop over `sample_quotes` twice to simulate the queue being exhausted fully, then re-sourced @@ -557,7 +525,7 @@ mod test { const LOOPS: usize = 3; let mut quote_counts = HashMap::new(); for _ in 0..(sample_quotes.len() * LOOPS) { - let next_quote = cycle.ask(FetchQuote).await.unwrap(); + let next_quote = cycle.handle(FetchQuote).await.unwrap(); quote_counts .entry(next_quote.get().to_owned()) .or_insert(0)