pgmq 0.34.0-alpha.5

A distributed message queue for Rust applications, on Postgres.
Documentation
#[cfg(feature = "diesel-async")]
pub mod diesel_async;
#[cfg(feature = "diesel-sync")]
pub mod diesel_sync;
mod query;
pub mod sql;

/// This macro defines all the functions required to implement [`crate::queue::Queue`] for both
/// sync and async diesel. This is possible because the sync and async diesel implementations are
/// nearly identical. The only differences are which `RunQueryDsl` trait is used and that one
/// requires `await`-ing DB results.
///
/// The functions defined by this macro are designed to be used by the [`crate::queue::Queue`]
/// definition generated by the [`crate::queue::macros::impl_queue`] macro.
macro_rules! diesel_functions {
    (
        /// The common connection/executor trait to use to implement all the functions
        $executor_trait:path,
        /// Path to the `RunQueryDsl` trait to use
        $run_dsl_trait:path,
        /// How to transform the result (e.g., either simply use it directly, or perform an `await`)
        $transform_result:tt
    ) => {
        use $run_dsl_trait;

        async fn create<C>(
            executor: &mut C,
            queue_name: crate::types::queue_name::QueueName<'_>,
        ) -> Result<(), crate::errors::PgmqError>
        where
            C: $executor_trait,
        {
            let result =
                crate::queue::diesel::query::create_queue_query(queue_name).execute(executor);
            $transform_result!(result)?;
            Ok(())
        }

        async fn send<C>(
            executor: &mut C,
            queue_name: crate::types::queue_name::QueueName<'_>,
            message: serde_json::Value,
            headers: serde_json::Value,
            delay: crate::types::visibility_timeout_offset::VisibilityTimeoutOffset,
        ) -> Result<i64, crate::errors::PgmqError>
        where
            C: $executor_trait,
        {
            let message_id =
                crate::queue::diesel::query::send_query(queue_name, message, headers, delay)
                    .get_result(executor);
            let message_id = $transform_result!(message_id)?;
            Ok(message_id)
        }

        async fn read<C, T, H>(
            executor: &mut C,
            queue_name: crate::types::queue_name::QueueName<'_>,
            visibility_timeout: crate::types::visibility_timeout_offset::VisibilityTimeoutOffset,
            quantity: i32,
        ) -> Result<Vec<crate::types::Message<T, H>>, crate::errors::PgmqError>
        where
            C: $executor_trait,
            T: 'static + Send + for<'de> serde::Deserialize<'de>,
            H: 'static + Send + for<'de> serde::Deserialize<'de>,
        {
            let messages =
                crate::queue::diesel::query::read_query(queue_name, visibility_timeout, quantity)
                    .get_results::<crate::types::Message<T, H>>(executor);
            let messages = $transform_result!(messages)?;
            Ok(messages)
        }
    };
}
pub(crate) use diesel_functions;