pgmq 0.34.0-alpha.5

A distributed message queue for Rust applications, on Postgres.
Documentation
//! Diesel SQL function definitions for the `pgmq` extension.

use diesel::deserialize::FromSql;
use diesel::pg::{Pg, PgValue};
use diesel::prelude::*;
use diesel::sql_types::{BigInt, Integer, Jsonb, Nullable, Text, Timestamptz};

/// Diesel SQL type that can be used as the return type of diesel SQL functions generated by
/// the [`declare_sql_function`] macro.
// Relevant example: https://github.com/diesel-rs/diesel/blob/main/examples/postgres/composite_types/examples/composite2rust_colors.rs
type PgMessage = diesel::sql_types::Record<(
    // msg_id
    BigInt,
    // read_ct
    Integer,
    // enqueued_at
    Timestamptz,
    // last_read_at
    Nullable<Timestamptz>,
    // vt
    Timestamptz,
    // message
    Jsonb,
    // headers
    Nullable<Jsonb>,
)>;

impl<T, H> FromSql<PgMessage, Pg> for crate::types::Message<T, H>
where
    T: for<'de> serde::Deserialize<'de>,
    H: for<'de> serde::Deserialize<'de>,
{
    fn from_sql(bytes: PgValue) -> diesel::deserialize::Result<Self> {
        let (msg_id, read_ct, enqueued_at, last_read_at, vt, message, headers): (
            _,
            _,
            _,
            _,
            _,
            serde_json::Value,
            Option<serde_json::Value>,
        ) = FromSql::<PgMessage, Pg>::from_sql(bytes)?;

        let headers = if let Some(headers) = headers {
            Some(H::deserialize(headers)?)
        } else {
            None
        };

        Ok(Self {
            msg_id,
            read_ct,
            enqueued_at,
            last_read_at,
            vt,
            message: T::deserialize(message)?,
            headers,
        })
    }
}

#[declare_sql_function]
extern "SQL" {
    #[sql_name = "pgmq.create"]
    fn pgmq_create(queue_name: Text);

    #[sql_name = "pgmq.send"]
    fn pgmq_send(queue_name: Text, msg: Jsonb, headers: Jsonb, delay: Integer) -> BigInt;

    #[sql_name = "pgmq.read"]
    fn pgmq_read(queue_name: Text, vt: Integer, qty: Integer) -> PgMessage;
}