#![allow(clippy::unwrap_used, clippy::expect_used)] #![allow(clippy::print_stderr)]
use fraiseql_observers::{DeadLetterQueue, DispatchSource, FunctionDispatchRecord};
use sqlx::{
PgPool,
postgres::{PgConnectOptions, PgPoolOptions},
};
use super::PgFunctionDlq;
async fn isolated_pool(url: &str, schema: &str) -> PgPool {
let options: PgConnectOptions = url.parse().unwrap();
let options = options.options([("search_path", format!("{schema},public"))]);
PgPoolOptions::new().connect_with(options).await.unwrap()
}
async fn connect_pool(schema: &str) -> Option<(PgPool, fraiseql_test_support::Service)> {
let svc = fraiseql_test_support::postgres().await?;
let admin = PgPool::connect(svc.url()).await.unwrap();
sqlx::query(&format!("CREATE SCHEMA IF NOT EXISTS {schema}"))
.execute(&admin)
.await
.unwrap();
admin.close().await;
Some((isolated_pool(svc.url(), schema).await, svc))
}
fn record(error: &str) -> FunctionDispatchRecord {
FunctionDispatchRecord::new(
DispatchSource::AfterMutation,
"notify_ops",
"after:mutation:createOrder",
"idem-tok-abc123",
serde_json::json!({ "order_id": 42, "total": "19.99" }),
error,
3,
)
}
async fn truncate(pool: &PgPool) {
sqlx::query("TRUNCATE _fraiseql_function_dlq").execute(pool).await.unwrap();
}
#[tokio::test]
async fn dead_lettered_dispatch_survives_a_restart() {
const SCHEMA: &str = "dlq_test_restart";
let Some((pool, svc)) = connect_pool(SCHEMA).await else {
eprintln!("SKIP dead_lettered_dispatch_survives_a_restart: no postgres");
return;
};
let store_a = PgFunctionDlq::new(pool.clone(), None);
store_a.init().await.unwrap();
truncate(&pool).await;
let rec = record("upstream 503");
let want_id = rec.id;
store_a.push_function(rec).await.unwrap();
assert_eq!(store_a.get_pending_functions(10).await.unwrap().len(), 1);
drop(store_a);
drop(pool);
let pool_b = isolated_pool(svc.url(), SCHEMA).await;
let store_b = PgFunctionDlq::new(pool_b, None);
store_b.init().await.unwrap();
let pending = store_b.get_pending_functions(10).await.unwrap();
assert_eq!(
pending.len(),
1,
"M-598: the Postgres-backed DLQ must survive a restart (the in-memory store loses it)"
);
let reloaded = &pending[0];
assert_eq!(reloaded.id, want_id, "the stable dead-letter id round-trips");
assert_eq!(reloaded.function_name, "notify_ops");
assert_eq!(reloaded.trigger_type, "after:mutation:createOrder");
assert_eq!(reloaded.idempotency_token, "idem-tok-abc123");
assert_eq!(reloaded.source, DispatchSource::AfterMutation, "the source label round-trips");
assert_eq!(reloaded.payload, serde_json::json!({ "order_id": 42, "total": "19.99" }));
assert_eq!(reloaded.attempts, 3);
}
#[tokio::test]
async fn get_pending_returns_records_oldest_first() {
const SCHEMA: &str = "dlq_test_oldest_first";
let Some((pool, _svc)) = connect_pool(SCHEMA).await else {
eprintln!("SKIP get_pending_returns_records_oldest_first: no postgres");
return;
};
let store = PgFunctionDlq::new(pool.clone(), None);
store.init().await.unwrap();
truncate(&pool).await;
store.push_function(record("first")).await.unwrap();
store.push_function(record("second")).await.unwrap();
store.push_function(record("third")).await.unwrap();
let pending = store.get_pending_functions(10).await.unwrap();
assert_eq!(pending.len(), 3);
assert_eq!(pending[0].error_message, "first", "oldest drains first");
assert_eq!(pending[2].error_message, "third");
assert_eq!(store.get_pending_functions(2).await.unwrap().len(), 2);
}
#[tokio::test]
async fn capacity_drops_newest_and_holds_the_cap() {
const SCHEMA: &str = "dlq_test_capacity";
let Some((pool, _svc)) = connect_pool(SCHEMA).await else {
eprintln!("SKIP capacity_drops_newest_and_holds_the_cap: no postgres");
return;
};
let store = PgFunctionDlq::new(pool.clone(), Some(2));
store.init().await.unwrap();
truncate(&pool).await;
store.push_function(record("e1")).await.unwrap();
store.push_function(record("e2")).await.unwrap();
store.push_function(record("e3")).await.unwrap();
let pending = store.get_pending_functions(10).await.unwrap();
assert_eq!(pending.len(), 2, "the cap holds the durable queue at 2");
assert_eq!(pending[0].error_message, "e1", "the two oldest survive; the newest is dropped");
assert_eq!(pending[1].error_message, "e2");
}