diff --git a/Cargo.lock b/Cargo.lock index 713e120..e4da3dd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -53,6 +53,12 @@ version = "1.0.97" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dcfed56ad506cb2c684a14971b8861fdc3baaaae314b9e5f9bb532cbe3ba7a4f" +[[package]] +name = "arc-swap" +version = "1.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69f7f8c3906b62b754cd5326047894316021dcfe5a194c8ea52bdd94934a3457" + [[package]] name = "async-compression" version = "0.4.20" @@ -145,6 +151,7 @@ dependencies = [ "glob", "grep", "rand", + "redis", "tokio", "tokio-cron-scheduler", ] @@ -280,6 +287,20 @@ dependencies = [ "unsigned-varint", ] +[[package]] +name = "combine" +version = "4.6.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba5a308b75df32fe02788e748662718f03fde005016435c444eea572398219fd" +dependencies = [ + "bytes", + "futures-core", + "memchr", + "pin-project-lite", + "tokio", + "tokio-util", +] + [[package]] name = "concurrent-queue" version = "2.5.0" @@ -560,6 +581,7 @@ checksum = "9fa08315bb612088cc391249efdc3bc77536f16c91f6cf495e6fbe85b20a4a81" dependencies = [ "futures-core", "futures-macro", + "futures-sink", "futures-task", "pin-project-lite", "pin-utils", @@ -1203,6 +1225,16 @@ dependencies = [ "winapi", ] +[[package]] +name = "num-bigint" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a5e44f723f1133c9deac646763579fdb3ac745e418f2a7af9cd0c431da1f20b9" +dependencies = [ + "num-integer", + "num-traits", +] + [[package]] name = "num-derive" version = "0.4.2" @@ -1214,6 +1246,15 @@ dependencies = [ "syn", ] +[[package]] +name = "num-integer" +version = "0.1.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7969661fd2958a5cb096e56c8e1ad0444ac2bbcd0061bd28660485a44879858f" +dependencies = [ + "num-traits", +] + [[package]] name = "num-traits" version = "0.2.19" @@ -1419,6 +1460,28 @@ dependencies = [ "getrandom", ] +[[package]] +name = "redis" +version = "0.29.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8034fb926579ff49d3fe58d288d5dcb580bf11e9bccd33224b45adebf0fd0c23" +dependencies = [ + "arc-swap", + "bytes", + "combine", + "futures-util", + "itoa", + "num-bigint", + "percent-encoding", + "pin-project-lite", + "ryu", + "sha1_smol", + "socket2", + "tokio", + "tokio-util", + "url", +] + [[package]] name = "redox_syscall" version = "0.5.10" @@ -1685,6 +1748,12 @@ dependencies = [ "serde", ] +[[package]] +name = "sha1_smol" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbfa15b3dddfee50a0fff136974b3e1bde555604ba463834a7eb7deb6417705d" + [[package]] name = "sharded-slab" version = "0.1.7" diff --git a/Cargo.toml b/Cargo.toml index f0ced60..5312bfd 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,5 +11,6 @@ bsky-sdk = "0.1.16" glob = "0.3.2" grep = "0.3.2" rand = "0.9.0" +redis = { version = "0.29.1", features = ["aio", "tokio-comp"] } tokio = { version = "1.44.0", features = ["full"] } tokio-cron-scheduler = "0.13.0" diff --git a/src/main.rs b/src/main.rs index 108d868..6a48632 100644 --- a/src/main.rs +++ b/src/main.rs @@ -5,11 +5,14 @@ use bsky_sdk::BskyAgent; use glob::glob; use grep::{matcher::Matcher, regex, searcher::sinks}; use rand::random_range; +use rand::seq::SliceRandom; use std::{sync::Arc, time::Duration}; use tokio::sync::Mutex; use tokio_cron_scheduler::{Job, JobScheduler, JobSchedulerError}; +use redis::{aio::{self, MultiplexedConnection}, AsyncCommands, Client}; + fn prepare_post>(text: I) -> post::RecordData { post::RecordData { text: text.into(), @@ -30,7 +33,7 @@ struct QuoteFilter { dates: Vec, } -fn read_files(filter: QuoteFilter) -> Vec { +fn read_files(filter: &QuoteFilter) -> Vec { let matcher = regex::RegexMatcher::new(&filter.content).unwrap(); let mut searcher = grep::searcher::Searcher::new(); let mut results = Vec::new(); @@ -59,13 +62,33 @@ fn read_files(filter: QuoteFilter) -> Vec { results } +async fn reshuffle_quotes(filter: &QuoteFilter, mut con: impl redis::aio::ConnectionLike + AsyncCommands, output_queue: &str) -> Result<(), ()> { + let len: u64 = con.llen(output_queue).await.unwrap(); + // NOTE: The following assumes the queue hasn't been repopulated by any other client + // in-between the call to llen and the execution of the pipeline. + // Hopefully won't be a problem :) + if len == 0 { + let mut file_contents = read_files(filter); + + { + let mut rand = rand::rng(); + file_contents.shuffle(&mut rand); + } + + let mut pipeline = redis::pipe(); + for file_contents in file_contents.into_iter() { + pipeline.lpush(output_queue,file_contents.as_str()); + } + let _: () = pipeline.query_async(&mut con).await.unwrap(); + } + + Ok(()) +} + #[tokio::main] async fn main() -> Result<(), Box> { - let file_contents = Arc::new(read_files(QuoteFilter { - content: r"\b(?i:mother|mommy|mama|mom)\b".to_string(), - path: "quotes/**/*.txt".to_string(), - dates: vec![], - })); + let redis = redis::Client::open(std::env::var("REDIS_URL").unwrap_or("redis://localhost".to_string()))?; + let con = redis.get_multiplexed_async_connection().await?; let agent = BskyAgent::builder().build().await?; let _session = agent @@ -77,16 +100,23 @@ async fn main() -> Result<(), Box> { let sched = JobScheduler::new().await?; let agent = Arc::new(Mutex::new(agent)); + let filter = Arc::new(QuoteFilter { + content: r"\b(?i:mother|mommy|mama|mom)\b".to_string(), + path: "quotes/**/*.txt".to_string(), + dates: vec![], + }); // Add async job sched .add(Job::new_async("0/10,5/10 * * * * *", move |_uuid, _| { - let file_contents = file_contents.clone(); + let filter = filter.clone(); + let mut con = con.clone(); let agent = agent.clone(); Box::pin(async move { - let text = file_contents[random_range(..file_contents.len())].as_str(); - let post = prepare_post(text); + let _ = reshuffle_quotes(&filter, con.clone(), "test:queue").await.unwrap(); + let text: String = con.lpop("test:queue", None).await.unwrap(); + let post = prepare_post(text.as_str()); let agent = agent.lock().await; if let Err(_) = agent.create_record(post).await { println!("{}\n", text)