#![allow(clippy::unwrap_used)]
use std::{
collections::HashMap,
sync::{Arc, Barrier},
};
use fraiseql_observers::{ActionConfig, DeadLetterQueue, DlqItem, EntityEvent, EventKind};
use uuid::Uuid;
use super::InMemoryDlq;
fn test_event() -> EntityEvent {
EntityEvent::new(
EventKind::Created,
"TestEntity".to_string(),
Uuid::new_v4(),
serde_json::json!({}),
)
}
fn test_action() -> ActionConfig {
ActionConfig::Webhook {
url: Some("http://localhost/hook".to_string()),
url_env: None,
method: None,
headers: HashMap::new(),
body_template: None,
signing_secret: None,
signing_secret_env: None,
}
}
async fn push(dlq: &InMemoryDlq) -> Uuid {
dlq.push(test_event(), test_action(), "boom".to_string()).await.unwrap()
}
#[tokio::test]
async fn unbounded_dlq_grows_without_limit() {
let dlq = InMemoryDlq::new_with_max(None);
for _ in 0..5 {
push(&dlq).await;
}
assert_eq!(dlq.count(), 5);
assert_eq!(dlq.overflow_count(), 0);
}
#[tokio::test]
async fn capped_dlq_drops_newest_at_capacity() {
let dlq = InMemoryDlq::new_with_max(Some(2));
let first = push(&dlq).await;
let second = push(&dlq).await;
push(&dlq).await;
assert_eq!(dlq.count(), 2, "cap should hold the queue at 2 entries");
assert_eq!(dlq.overflow_count(), 1, "the dropped entry should bump the overflow counter");
let ids: Vec<Uuid> = dlq.list_all().into_iter().map(|i| i.id).collect();
assert!(
ids.contains(&first) && ids.contains(&second),
"the first two entries are retained"
);
}
#[tokio::test]
async fn mark_retry_failed_keeps_item_and_records_failure() {
let dlq = InMemoryDlq::new_with_max(None);
let id = push(&dlq).await;
dlq.mark_retry_failed(id, "second failure").await.unwrap();
let item = dlq.get(id).expect("item must still be present after a failed retry");
assert_eq!(item.attempts, 1, "attempts should be incremented");
assert_eq!(item.error_message, "second failure", "error_message should be updated");
assert_eq!(dlq.count(), 1);
}
#[tokio::test]
async fn mark_success_removes_item() {
let dlq = InMemoryDlq::new_with_max(None);
let id = push(&dlq).await;
dlq.mark_success(id).await.unwrap();
assert!(dlq.get(id).is_none(), "a succeeded item should be removed");
assert_eq!(dlq.count(), 0);
}
#[tokio::test]
async fn try_claim_removes_and_is_idempotent() {
let dlq = InMemoryDlq::new_with_max(None);
let id = push(&dlq).await;
assert!(dlq.try_claim(id).is_some(), "first claim returns the item");
assert_eq!(dlq.count(), 0, "claim removes the item");
assert!(dlq.try_claim(id).is_none(), "second claim finds nothing");
}
#[tokio::test]
async fn try_claim_is_atomic_under_concurrency() {
let dlq = Arc::new(InMemoryDlq::new_with_max(None));
let id = push(&dlq).await;
let n = 8;
let barrier = Arc::new(Barrier::new(n));
#[allow(clippy::needless_collect)] let handles: Vec<_> = (0..n)
.map(|_| {
let dlq = Arc::clone(&dlq);
let barrier = Arc::clone(&barrier);
std::thread::spawn(move || {
barrier.wait();
dlq.try_claim(id).is_some()
})
})
.collect();
let winners = handles.into_iter().map(|h| h.join().unwrap()).filter(|&won| won).count();
assert_eq!(winners, 1, "exactly one of {n} concurrent claimers should win");
assert_eq!(dlq.count(), 0);
}
#[tokio::test]
async fn reinsert_bypasses_the_cap() {
let dlq = InMemoryDlq::new_with_max(Some(1));
push(&dlq).await;
let claimed = DlqItem {
id: Uuid::new_v4(),
event: test_event(),
action: test_action(),
error_message: "retry failed".to_string(),
attempts: 1,
};
dlq.reinsert(claimed);
assert_eq!(dlq.count(), 2, "reinsert must bypass the cap");
assert_eq!(dlq.overflow_count(), 0, "reinsert is not an overflow");
}
mod function_dlq {
use fraiseql_observers::{DeadLetterQueue, DispatchSource, FunctionDispatchRecord};
use super::InMemoryDlq;
fn record(error: &str) -> FunctionDispatchRecord {
FunctionDispatchRecord::new(
DispatchSource::AfterMutation,
"onUserCreated",
"after:mutation:onUserCreated",
"0123456789abcdef0123456789abcdef",
serde_json::json!({ "event_kind": "insert", "new": { "id": "u1" } }),
error,
3,
)
}
#[tokio::test]
async fn exhausted_dispatch_lands_one_row() {
let dlq = InMemoryDlq::new_with_max(None);
let id = dlq.push_function(record("upstream 503")).await.unwrap();
assert_eq!(dlq.function_count(), 1, "one function DLQ row after exhaustion");
let pending = dlq.get_pending_functions(10).await.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, id);
assert_eq!(pending[0].function_name, "onUserCreated");
assert_eq!(pending[0].attempts, 3);
assert_eq!(pending[0].error_message, "upstream 503");
}
#[tokio::test]
async fn function_and_observer_entries_are_separate() {
let dlq = InMemoryDlq::new_with_max(None);
super::push(&dlq).await; dlq.push_function(record("boom")).await.unwrap();
assert_eq!(dlq.count(), 1, "observer entries counted separately");
assert_eq!(dlq.function_count(), 1, "function entries counted separately");
}
#[tokio::test]
async fn in_memory_dlq_loses_function_entries_on_restart() {
let dlq = InMemoryDlq::new_with_max(None);
dlq.push_function(record("upstream 503")).await.unwrap();
assert_eq!(dlq.function_count(), 1, "entry present before restart");
drop(dlq);
let after_restart = InMemoryDlq::new_with_max(None);
assert_eq!(
after_restart.function_count(),
0,
"M-598: the in-memory DLQ loses dead-lettered function dispatches on restart — \
phase 07's Postgres-backed store must make this survive."
);
}
#[tokio::test]
async fn capped_function_dlq_drops_newest() {
let dlq = InMemoryDlq::new_with_max(Some(2));
dlq.push_function(record("e1")).await.unwrap();
dlq.push_function(record("e2")).await.unwrap();
dlq.push_function(record("e3")).await.unwrap();
assert_eq!(dlq.function_count(), 2, "cap holds the function queue at 2");
assert_eq!(dlq.overflow_count(), 1, "the dropped entry bumps the overflow counter");
}
}
mod listener_selection {
use fraiseql_observers::config::TransportKind;
use super::super::{ListenerSelection, listener_selection};
#[test]
fn postgres_uses_the_change_log_listener() {
assert_eq!(
listener_selection(TransportKind::Postgres),
ListenerSelection::PostgresChangeLog,
);
}
#[test]
fn nats_uses_the_transport_stream_not_the_pg_listener() {
assert_eq!(listener_selection(TransportKind::Nats), ListenerSelection::TransportStream,);
}
#[test]
fn in_memory_uses_the_transport_stream() {
assert_eq!(listener_selection(TransportKind::InMemory), ListenerSelection::TransportStream,);
}
}
mod transport_requires_broker {
use fraiseql_observers::config::TransportKind;
use super::super::{ObserverRuntime, ObserverRuntimeConfig};
fn runtime_with(kind: TransportKind) -> ObserverRuntime {
let pool = sqlx::PgPool::connect_lazy("postgres://u:u@127.0.0.1:1/db")
.expect("lazy pool construction does not connect");
let mut config = ObserverRuntimeConfig::new(pool);
config.transport.transport = kind;
ObserverRuntime::new(config)
}
#[tokio::test]
async fn postgres_start_failure_is_not_boot_fatal() {
assert!(!runtime_with(TransportKind::Postgres).transport_requires_broker());
}
#[tokio::test]
async fn nats_start_failure_is_boot_fatal() {
assert!(runtime_with(TransportKind::Nats).transport_requires_broker());
}
}
mod nats_unrunnable_gate {
use fraiseql_observers::config::TransportKind;
use super::super::{ObserverRuntime, ObserverRuntimeConfig};
#[tokio::test]
async fn nats_that_cannot_run_fails_start_with_no_pg_fallback() {
let pool = sqlx::PgPool::connect_lazy("postgres://u:u@127.0.0.1:1/db")
.expect("lazy pool construction does not connect");
let mut config = ObserverRuntimeConfig::new(pool);
config.transport.transport = TransportKind::Nats;
config.transport.nats.url = "nats://127.0.0.1:4222".to_string();
let mut runtime = ObserverRuntime::new(config);
let result = runtime.start().await;
assert!(
result.is_err(),
"NATS that cannot run must fail start(), not fall back to PostgreSQL"
);
assert!(
!runtime.is_running(),
"runtime must not report running after a failed NATS start"
);
assert!(
runtime.transport_requires_broker(),
"a NATS transport is broker-backed, so the start failure is boot-fatal in production"
);
}
}
mod log_payload_truncation {
use super::super::{MAX_LOG_PAYLOAD_BYTES, truncate_log_payload};
#[test]
fn small_payload_is_passed_through_unchanged() {
let data = serde_json::json!({"id": "abc", "status": "new"});
assert_eq!(truncate_log_payload(&data), data);
}
#[test]
fn oversized_payload_is_replaced_with_a_size_marker() {
let big = "x".repeat(MAX_LOG_PAYLOAD_BYTES + 1_024);
let data = serde_json::json!({ "blob": big });
let out = truncate_log_payload(&data);
assert_eq!(out["_truncated"], serde_json::Value::Bool(true));
let recorded = out["_original_size_bytes"].as_u64().unwrap();
assert!(
recorded > u64::try_from(MAX_LOG_PAYLOAD_BYTES).unwrap(),
"marker must record the original (oversized) byte length"
);
assert!(out.get("blob").is_none());
}
}
mod subscription_forwarding {
use fraiseql_core::runtime::subscription::SubscriptionOperation;
use fraiseql_observers::EventKind;
use super::super::subscription_operation_for;
#[test]
fn custom_events_are_not_forwarded_to_subscribers() {
assert_eq!(
subscription_operation_for(EventKind::Custom),
None,
"a snapshot/read/no-op row must never become a subscriber-visible event (#773)"
);
}
#[test]
fn real_changes_map_one_to_one() {
assert_eq!(
subscription_operation_for(EventKind::Created),
Some(SubscriptionOperation::Create)
);
assert_eq!(
subscription_operation_for(EventKind::Updated),
Some(SubscriptionOperation::Update)
);
assert_eq!(
subscription_operation_for(EventKind::Deleted),
Some(SubscriptionOperation::Delete)
);
}
}
mod bridge_backpressure {
use std::sync::Arc;
use fraiseql_core::{
runtime::subscription::{SubscriptionManager, SubscriptionOperation},
schema::{CompiledSchema, SubscriptionDefinition},
};
use super::super::forward_to_bridge;
use crate::subscriptions::{EntityEvent, EventBridge, EventBridgeConfig};
#[tokio::test]
async fn burst_beyond_bridge_capacity_loses_no_events() {
const BURST: usize = 200;
let mut schema = CompiledSchema::new();
schema.subscriptions.push(SubscriptionDefinition::new("orderChanged", "Order"));
let manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
manager
.subscribe("orderChanged", serde_json::json!({}), serde_json::json!({}), "conn-1")
.expect("subscribe");
let mut rx = manager.receiver();
let bridge = EventBridge::new(
Arc::clone(&manager),
EventBridgeConfig::new().with_channel_capacity(2),
);
let sender = bridge.sender();
let handle = bridge.spawn();
for i in 0..BURST {
let event = EntityEvent::new(
"Order",
format!("order_{i}"),
SubscriptionOperation::Create,
serde_json::json!({"id": format!("order_{i}")}),
);
forward_to_bridge(&sender, event, &format!("evt-{i}")).await;
}
let mut delivered = 0_usize;
while delivered < BURST {
match tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()).await {
Ok(result) => {
let _payload = result.expect("broadcast receiver must stay healthy");
delivered += 1;
},
Err(_) => break, }
}
assert_eq!(
delivered, BURST,
"every event in the burst must reach the subscriber — a full bridge channel \
must apply backpressure, never silently drop (#772)"
);
handle.abort();
}
}
#[test]
fn observer_log_statuses_are_accepted_by_the_shipped_check() {
use fraiseql_observers::migrations::OBSERVER_LOG_STATUSES;
for status in [
super::OBSERVER_LOG_STATUS_SUCCESS,
super::OBSERVER_LOG_STATUS_FAILED,
] {
assert!(
OBSERVER_LOG_STATUSES.contains(&status),
"the runtime writes tb_observer_log.status = {status:?}, which migration 06's \
ck_observer_log_status rejects (accepts: {OBSERVER_LOG_STATUSES:?}) — the row \
is silently dropped"
);
}
}