use diesel::deserialize::FromSql;
use diesel::pg::{Pg, PgValue};
use diesel::prelude::*;
use diesel::sql_types::{BigInt, Integer, Jsonb, Nullable, Text, Timestamptz};
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;
}