diff --git a/Cargo.lock b/Cargo.lock index 2ad6637..1c06c3d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -143,6 +143,7 @@ version = "0.1.0" dependencies = [ "bsky-sdk", "cron-lite", + "futures", "glob", "grep", "kameo", diff --git a/Cargo.toml b/Cargo.toml index 162f5d4..2993fbc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -9,6 +9,7 @@ rust-version = "1.85" [dependencies] bsky-sdk = "0.1.16" cron-lite = { version = "0.3.0", features = ["async"] } +futures = "0.3.31" glob = "0.3.2" grep = "0.3.2" kameo = "0.17.2" diff --git a/src/sink.rs b/src/sink.rs index 020569f..d0898de 100644 --- a/src/sink.rs +++ b/src/sink.rs @@ -3,6 +3,7 @@ use kameo::prelude::*; /// A newtype over [Quote] used to prompt the [SinkManager] to /// submit a new quote to all its configured sinks. +#[derive(Debug, Clone)] pub struct PostQuote(pub Quote); /// Error type for internal communication between @@ -56,10 +57,61 @@ impl Message for StdoutSink { /// The SinkManager will attempt to reinitialize failed sinks upon /// encountering recoverable errors. #[derive(Actor)] -pub struct SinkManager { +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. - // We will have to use multiple vectors with static dispatch, instead... - sinks: Vec>, + // I've decided I'll limit this to one sink per implementation right now. + stdout_sink: Option>, + // ... +} + +impl SinkManager { + pub fn new(stdout_sink: Option>) -> Self { + Self { stdout_sink } + } +} + +pub type SinkReplies = Vec>; + +impl Message for SinkManager { + type Reply = SinkReplies; + + async fn handle( + &mut self, + msg: PostQuote, + _ctx: &mut Context, + ) -> Self::Reply { + use futures::future::join_all; + + // We'll see if this monstrosity actually works + let sinks = [self.stdout_sink.clone()]; + let futures = sinks + .iter() + .flatten() + .map(|s| s.ask(msg.clone()).into_future()); + let results = join_all(futures).await; + + results.iter().map(|r| r.clone().or(Err(()))).collect() + } +} + +mod test { + #[tokio::test] + async fn stdout_sink() { + use super::*; + + let stdout = StdoutSink::spawn(StdoutSink); + let manager = SinkManager::spawn(SinkManager::new(Some(stdout))); + + let messages = ["First test!", "Second test.", "Third..."]; + for msg in messages { + manager.tell(PostQuote(msg.into())).await.unwrap(); + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + } + + // Hopefully we don't crash...! + // TODO: Sink that actually stores every quote it "posts"? + // Could help in verifying everything was sent correctly. + } }