use fraiseql_observers::{
ActionConfig, DeadLetterQueue, DispatchSource, DlqItem, EntityEvent, FunctionDispatchRecord,
ObserverError, Result,
};
use sqlx::PgPool;
use uuid::Uuid;
pub struct PgFunctionDlq {
pool: PgPool,
max_size: Option<usize>,
}
fn dlq_err(op: &str, error: &sqlx::Error) -> ObserverError {
ObserverError::DlqError {
reason: format!("function DLQ {op}: {error}"),
}
}
impl PgFunctionDlq {
#[must_use]
pub const fn new(pool: PgPool, max_size: Option<usize>) -> Self {
Self { pool, max_size }
}
pub async fn init(&self) -> std::result::Result<(), sqlx::Error> {
sqlx::raw_sql(fraiseql_functions::migrations::dlq_migration_sql())
.execute(&self.pool)
.await
.map(|_| ())
}
async fn count(&self) -> Result<i64> {
sqlx::query_scalar::<_, i64>("SELECT count(*) FROM _fraiseql_function_dlq")
.fetch_one(&self.pool)
.await
.map_err(|error| dlq_err("count", &error))
}
}
#[async_trait::async_trait]
impl DeadLetterQueue for PgFunctionDlq {
async fn push_function(&self, record: FunctionDispatchRecord) -> Result<Uuid> {
let id = record.id;
if let Some(max) = self.max_size {
let cap = i64::try_from(max).unwrap_or(i64::MAX);
let current = self.count().await?;
if current >= cap {
tracing::warn!(
max_dlq_size = max,
function = %record.function_name,
trigger = %record.trigger_type,
"function DLQ full; dropping failed function dispatch entry"
);
crate::function_metrics::record_dlq_eviction();
crate::function_metrics::set_dlq_size(
usize::try_from(current).unwrap_or(usize::MAX),
);
return Ok(id);
}
}
let payload = record.payload.to_string();
sqlx::query(
"INSERT INTO _fraiseql_function_dlq \
(id, source, function_name, trigger_type, idempotency_token, \
payload, error_message, attempts) \
VALUES ($1, $2, $3, $4, $5, $6::jsonb, $7, $8)",
)
.bind(id)
.bind(record.source.label())
.bind(&record.function_name)
.bind(&record.trigger_type)
.bind(&record.idempotency_token)
.bind(payload)
.bind(&record.error_message)
.bind(i64::from(record.attempts))
.execute(&self.pool)
.await
.map_err(|error| dlq_err("push", &error))?;
let current = self.count().await?;
crate::function_metrics::set_dlq_size(usize::try_from(current).unwrap_or(usize::MAX));
Ok(id)
}
async fn get_pending_functions(&self, limit: i64) -> Result<Vec<FunctionDispatchRecord>> {
let rows =
sqlx::query_as::<_, (Uuid, String, String, String, String, String, String, i64)>(
"SELECT id, source, function_name, trigger_type, idempotency_token, \
payload::text, error_message, attempts \
FROM _fraiseql_function_dlq \
ORDER BY created_at ASC, pk_function_dlq ASC \
LIMIT $1",
)
.bind(limit)
.fetch_all(&self.pool)
.await
.map_err(|error| dlq_err("list", &error))?;
Ok(rows
.into_iter()
.map(
|(
id,
source,
function_name,
trigger_type,
idempotency_token,
payload_text,
error_message,
attempts,
)| {
let source = DispatchSource::from_label(&source).unwrap_or_else(|| {
tracing::warn!(%source, "function DLQ row has an unknown source label");
DispatchSource::Source
});
let payload = serde_json::from_str(&payload_text).unwrap_or_else(|error| {
tracing::warn!(%error, "function DLQ row payload is not valid JSON");
serde_json::Value::Null
});
FunctionDispatchRecord {
id,
source,
function_name,
trigger_type,
idempotency_token,
payload,
error_message,
attempts: u32::try_from(attempts).unwrap_or(u32::MAX),
}
},
)
.collect())
}
async fn push(&self, event: EntityEvent, action: ActionConfig, error: String) -> Result<Uuid> {
let id = Uuid::new_v4();
tracing::warn!(
action_type = action.action_type(),
event_id = %event.id,
%error,
"PgFunctionDlq is function-dispatch only — observer-action failure not persisted"
);
Ok(id)
}
async fn get_pending(&self, _limit: i64) -> Result<Vec<DlqItem>> {
Ok(Vec::new())
}
async fn mark_success(&self, _id: Uuid) -> Result<()> {
Ok(())
}
async fn mark_retry_failed(&self, _id: Uuid, _error: &str) -> Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests;