use fraiseql_functions::InboundMessage;
use sqlx::{PgPool, Postgres, Transaction};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Emitted {
New(uuid::Uuid),
Duplicate,
}
impl Emitted {
#[must_use]
pub const fn is_new(&self) -> bool {
matches!(self, Emitted::New(_))
}
}
fn db_err(context: &str, error: &sqlx::Error) -> fraiseql_error::FraiseQLError {
fraiseql_error::FraiseQLError::database(format!("inbound spine: {context}: {error}"))
}
pub async fn emit_in_tx(
tx: &mut Transaction<'_, Postgres>,
message: &InboundMessage,
) -> fraiseql_error::Result<Emitted> {
let payload = serde_json::to_string(message).map_err(|error| {
fraiseql_error::FraiseQLError::database(format!(
"inbound spine: serialize message: {error}"
))
})?;
let id = sqlx::query_scalar::<_, uuid::Uuid>(
"INSERT INTO _fraiseql_inbound_message \
(source, idempotency_key, thread_key, payload, received_at) \
VALUES ($1, $2, $3, $4::jsonb, $5) \
ON CONFLICT (source, idempotency_key) DO NOTHING \
RETURNING id",
)
.bind(message.source.as_key())
.bind(&message.idempotency_key)
.bind(message.thread_key.as_deref())
.bind(payload)
.bind(message.received_at)
.fetch_optional(&mut **tx)
.await
.map_err(|error| db_err("claim", &error))?;
Ok(id.map_or(Emitted::Duplicate, Emitted::New))
}
pub struct PostgresInboundSpine {
pool: PgPool,
}
impl PostgresInboundSpine {
#[must_use]
pub const fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn init(&self) -> fraiseql_error::Result<()> {
sqlx::raw_sql(fraiseql_functions::migrations::inbound_migration_sql())
.execute(&self.pool)
.await
.map_err(|error| db_err("init", &error))?;
Ok(())
}
pub async fn emit(&self, message: &InboundMessage) -> fraiseql_error::Result<Emitted> {
let mut tx = self.pool.begin().await.map_err(|error| db_err("begin", &error))?;
let emitted = emit_in_tx(&mut tx, message).await?;
tx.commit().await.map_err(|error| db_err("commit", &error))?;
Ok(emitted)
}
}
#[cfg(test)]
mod tests;