From 1e330210530cd0fc0b26f4c9dfbba3e5b602784e Mon Sep 17 00:00:00 2001 From: "crashkeys.dev" Date: Wed, 29 Jul 2026 21:17:52 +0200 Subject: [PATCH] actor: refactor to remove kameo usage; enables dynamic dispatch. Kameo's definition of its `Actor` and `Message` traits prevented me from creating trait objects for my actors, as I did not own the trait definitions and could not use `async_trait` on them. For my use case in this project, most of Kameo's functionality was superfluous and went unused; in the future, introducing an abstraction over the usage of channels for message passing and sharing of actors may be worth considering, though there is no necessity as of yet. Leveraging dynamic dispatch for quote "sinks" has the following benefits: - The logic to submit quotes to any number of sinks has been greatly simplified. - New sink implementations may be supplied transparently even by 3rd-party crates. - Multiple instances of the same sink implementation may be used at once. The previous implementation did not support this feature, as it would have required even more complex logic spanning the whole lifetime of the program. --- Cargo.lock | 124 ++++++++++++++------------------- Cargo.toml | 2 +- src/data.rs | 9 +++ src/lib.rs | 54 +++++++-------- src/sink.rs | 106 ++++++++++++---------------- src/storage.rs | 182 ++++++++++++++++++++----------------------------- 6 files changed, 207 insertions(+), 270 deletions(-) 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) -- 2.51.2