use std::marker::PhantomData;
use apalis::prelude::Storage;
use apalis_redis::{Config, ConnectionManager, RedisStorage};
use async_trait::async_trait;
use nest_rs_queue::{Job, JobProducer, QueueError, WIRE_FORMAT_VERSION};
use serde_json::json;
use crate::error::ConnectionError;
#[derive(Clone)]
pub struct QueueConnection {
conn: ConnectionManager,
}
impl QueueConnection {
pub async fn connect(redis_url: &str) -> Result<Self, ConnectionError> {
let conn = apalis_redis::connect(redis_url).await?;
Ok(Self { conn })
}
pub fn manager(&self) -> ConnectionManager {
self.conn.clone()
}
pub fn of<J: Job>(&self, queue: &str) -> Queue<J> {
Queue {
storage: self.value_storage(queue),
_phantom: PhantomData,
}
}
pub(crate) fn value_storage(&self, queue: &str) -> RedisStorage<serde_json::Value> {
RedisStorage::new_with_config(self.conn.clone(), Config::default().set_namespace(queue))
}
pub(crate) fn consumer_storage(
&self,
queue: &str,
concurrency: usize,
) -> RedisStorage<serde_json::Value> {
RedisStorage::new_with_config(
self.conn.clone(),
Config::default()
.set_namespace(queue)
.set_buffer_size(concurrency.max(1)),
)
}
}
pub struct Queue<J: Job> {
storage: RedisStorage<serde_json::Value>,
_phantom: PhantomData<fn(J)>,
}
impl<J: Job> Queue<J> {
pub async fn push(&self, job: J) -> Result<(), QueueError> {
let payload = serde_json::to_value(&job)?;
let mut storage = self.storage.clone();
storage
.push(envelope(payload))
.await
.map_err(QueueError::backend)?;
Ok(())
}
}
#[async_trait]
impl JobProducer for QueueConnection {
async fn push_json(&self, queue: &str, payload: serde_json::Value) -> Result<(), QueueError> {
let mut storage = self.value_storage(queue);
storage
.push(envelope(payload))
.await
.map_err(QueueError::backend)?;
Ok(())
}
}
fn envelope(payload: serde_json::Value) -> serde_json::Value {
json!({
"v": WIRE_FORMAT_VERSION,
"payload": payload,
})
}