bullrs 0.0.3

A BullMQ compatible Job Queue based on Redis
Documentation
use chrono::Utc;
use redis::{ErrorKind, RedisError, Value};
use serde::Serialize;

use crate::{
    JobOptions,
    error::AddJobErr,
    luacommands::{ADD_STANDARD_JOB, InvokeLuaScript},
    queue::QueueName,
};

pub struct AddStandardJob<'a, D> {
    // Name of the queue we want to add the job to
    pub queue: &'a QueueName,
    pub job_name: &'a str,
    pub data: &'a D,
    pub job_options: &'a JobOptions,
}

impl<'a, D> InvokeLuaScript for AddStandardJob<'a, D>
where
    D: Serialize,
{
    type RedisOutput = Value;
    type DomainOk = String;
    type DomainErr = AddJobErr;

    fn generate_invocation(&self) -> Result<redis::ScriptInvocation<'static>, Self::DomainErr> {
        let key_prefix = self.queue.prefix();
        let custom_id: &str = self.job_options.job_id.as_deref().unwrap_or("");
        let parent_key: Option<String> = None;
        let _wait_children_key = "";
        let parent_dependencies_key = "";
        let parent: Option<String> = None;
        let repeat_job_key = "";
        let deduplication_key = "";
        let job_name = self.job_name;
        let timestamp = self
            .job_options
            .timestamp
            .unwrap_or_else(Utc::now)
            .timestamp_millis();
        let arguments_tuple = (
            key_prefix,
            custom_id,
            job_name,
            timestamp,
            parent_key,
            parent_dependencies_key,
            parent,
            repeat_job_key,
            deduplication_key,
        );
        let payload_serialized = serde_json::to_string(self.data)?;
        let mut invocation = ADD_STANDARD_JOB.prepare_invoke();
        invocation
            .key(self.queue.wait())
            .key(self.queue.paused())
            .key(self.queue.meta())
            .key(self.queue.id())
            .key(self.queue.completed())
            .key(self.queue.delayed())
            .key(self.queue.active())
            .key(self.queue.events())
            .key(self.queue.marker())
            .arg(rmp_serde::to_vec(&arguments_tuple).expect("should never fails"))
            .arg(payload_serialized)
            .arg(rmp_serde::to_vec_named(self.job_options).expect("serializing never fails"));
        Ok(invocation)
    }

    fn map_value(&self, value: Self::RedisOutput) -> Result<Self::DomainOk, Self::DomainErr> {
        match value {
            Value::Int(-5) => Err(AddJobErr::MissingParentKey),
            Value::BulkString(s) => Ok(String::from_utf8_lossy(&s).into()),
            Value::SimpleString(s) => Ok(s),
            x => Err(RedisError::from((
                ErrorKind::ResponseError,
                "Unexpected response from AddStandardJob lua script",
                format!("Response was {x:?}"),
            )))?,
        }
    }
}