use std::time::Duration;
use crate::redis::connection::RedisPool;
use crate::storage::DedupeStore;
#[derive(Clone)]
pub struct RedisDedupeStore {
pool: RedisPool,
}
impl RedisDedupeStore {
pub fn new(pool: RedisPool) -> Self {
Self { pool }
}
fn set_key(&self, topic: &str) -> String {
self.pool.topic_key("dedupe", topic)
}
}
impl DedupeStore for RedisDedupeStore {
fn check_and_record(&self, topic: &str, key: &str, window: Duration) -> bool {
let set_key = self.set_key(topic);
let msg_id = key.to_string();
let window_secs = window.as_secs().max(1) as usize;
let (set_key2, msg_id2) = (set_key.clone(), msg_id.clone());
let added: i32 = self.pool.sync_cmd(move |c| {
redis::cmd("SADD")
.arg(&set_key2)
.arg(&msg_id2)
.query::<i32>(c)
});
let (set_key3, window_secs2) = (set_key.clone(), window_secs);
let _: () = self.pool.sync_cmd(move |c| {
redis::cmd("EXPIRE")
.arg(&set_key3)
.arg(window_secs2)
.query::<()>(c)
});
added > 0
}
fn sweep(&self) -> usize {
0
}
}