#[cfg(feature = "diesel-async")]
pub mod diesel_async;
#[cfg(feature = "diesel-sync")]
pub mod diesel_sync;
mod query;
pub mod sql;
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;