#![cfg(feature = "observers")]
#![allow(
clippy::unwrap_used,
clippy::panic,
clippy::print_stdout,
clippy::print_stderr
)] mod observer_test_helpers;
use std::time::Duration;
use fraiseql_server::observers::runtime::{ObserverRuntime, ObserverRuntimeConfig};
use observer_test_helpers::*;
use uuid::Uuid;
fn init_test_tracing() {
use std::sync::Once;
static INIT: Once = Once::new();
INIT.call_once(|| {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "info".into()),
)
.with_test_writer()
.init();
});
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_runtime_start_stop_lifecycle() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let entity_type = format!("Order_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-lifecycle-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let probe_url = mock_server.webhook_url();
let probe = reqwest::Client::new().get(&probe_url).send().await;
eprintln!("[DIAG] mock server probe at {probe_url}: {probe:?}");
let order_id = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "status": "new"}),
None,
)
.await
.expect("Failed to insert change log entry");
let cl_count: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM core.tb_entity_change_log WHERE object_type = $1")
.bind(&entity_type)
.fetch_one(&pool)
.await
.expect("Failed to query change log");
eprintln!("[DIAG] change log entries for {entity_type}: {}", cl_count.0);
assert!(cl_count.0 > 0, "change log entry must exist in DB");
tokio::time::sleep(Duration::from_secs(3)).await;
let health = runtime.health();
eprintln!(
"[DIAG] runtime health: running={}, observer_count={}, events_processed={}, errors={}",
health.running, health.observer_count, health.events_processed, health.errors
);
assert!(health.errors == 0, "runtime should have zero errors, got {}", health.errors);
wait_for_webhook(&mock_server, 1, Duration::from_secs(15)).await;
let requests = mock_server.received_requests().await;
assert_eq!(
requests.len(),
1,
"Expected 1 webhook call during lifecycle, got {}",
requests.len()
);
let log_count = get_observer_log_count(&pool, "success")
.await
.expect("Failed to query observer logs");
assert!(log_count > 0, "Expected at least 1 successful observer log");
let checkpoint_exists = check_checkpoint_exists(&pool, "change_log")
.await
.expect("Failed to check checkpoint");
assert!(checkpoint_exists, "Expected checkpoint to be saved after processing");
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_observer_log_populates_audit_columns() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let entity_type = format!("AuditOrder_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-audit-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone())
.with_poll_interval(50)
.with_log_payloads(true);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let order_id = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "status": "new"}),
None,
)
.await
.expect("Failed to insert change log entry");
wait_for_webhook(&mock_server, 1, Duration::from_secs(15)).await;
let mut audit = None;
for _ in 0..50 {
if let Some(a) = get_observer_log_audit(&pool, &order_id.to_string())
.await
.expect("Failed to query observer log audit columns")
{
audit = Some(a);
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
let audit = audit.expect("Expected an observer_log row for the processed event");
assert_eq!(audit.status, "success", "delivery should be logged as success");
assert_eq!(
audit.action_type.as_deref(),
Some("webhook"),
"action_type must be populated (#468)"
);
assert_eq!(audit.action_index, Some(0), "action_index must be populated (#468)");
assert_eq!(
audit.response_status_code,
Some(200),
"response_status_code must record the HTTP status (#468)"
);
assert!(
audit.response_payload.is_some(),
"response_payload summary must be populated (#468)"
);
assert!(
audit.duration_ms.is_some_and(|d| d >= 0),
"duration_ms must be recorded from the action detail"
);
let request_payload = audit
.request_payload
.expect("request_payload must be stored when log_payloads=true");
assert_eq!(
request_payload["id"].as_str(),
Some(order_id.to_string().as_str()),
"request_payload must contain the triggering event data"
);
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_checkpoint_recovery_after_restart() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
sqlx::query("DELETE FROM observer_checkpoints")
.execute(&pool)
.await
.expect("Failed to clean checkpoints");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let entity_type = format!("Order_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-checkpoint-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
for i in 0..5 {
let order_id = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "sequence": i}),
None,
)
.await
.expect("Failed to insert change log entry");
}
wait_for_webhook(&mock_server, 5, Duration::from_secs(20)).await;
assert_eq!(mock_server.request_count().await, 5, "Expected 5 webhooks before restart");
runtime.stop().await.expect("Failed to stop first runtime");
let config2 = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime2 = ObserverRuntime::new(config2);
runtime2.start().await.expect("Failed to start second runtime");
let sentinel_id = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&sentinel_id.to_string(),
serde_json::json!({"id": sentinel_id.to_string(), "sequence": 5}),
None,
)
.await
.expect("Failed to insert sentinel change log entry");
wait_for_webhook(&mock_server, 6, Duration::from_secs(20)).await;
tokio::time::sleep(Duration::from_millis(500)).await;
let requests = mock_server.received_requests().await;
let ids: Vec<String> = requests
.iter()
.filter_map(|r| r["id"].as_str().map(|s| s.to_string()))
.collect();
let unique_ids: std::collections::HashSet<_> = ids.iter().cloned().collect();
assert_eq!(
ids.len(),
unique_ids.len(),
"restart re-dispatched already-processed events: {} webhooks for {} unique entities",
ids.len(),
unique_ids.len()
);
assert_eq!(
requests.len(),
6,
"expected exactly 6 webhooks total (5 before restart + 1 sentinel), got {}",
requests.len()
);
let checkpoint = get_checkpoint_value(&pool, "change_log")
.await
.expect("Failed to get checkpoint");
assert!(
checkpoint > 0,
"checkpoint must be persisted under the stable listener id 'change_log'"
);
let by_entity_type = get_checkpoint_value(&pool, &entity_type)
.await
.expect("Failed to query checkpoint by entity type");
assert_eq!(
by_entity_type, 0,
"checkpoint must not be keyed on object_type (found a row under '{entity_type}')"
);
runtime2.stop().await.expect("Failed to stop second runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_hot_reload_observers() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server_1 = MockWebhookServer::start().await;
let mock_server_2 = MockWebhookServer::start().await;
mock_server_1.mock_success().await;
mock_server_2.mock_success().await;
let entity_type = format!("Order_{}", test_id);
let _observer_id_1 = create_test_observer(
&pool,
&format!("test-reload-1-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server_1.webhook_url(),
)
.await
.expect("Failed to create observer 1");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let order_id_1 = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id_1.to_string(),
serde_json::json!({"id": order_id_1.to_string(), "status": "created"}),
None,
)
.await
.expect("Failed to insert change log entry 1");
wait_for_webhook(&mock_server_1, 1, Duration::from_secs(15)).await;
assert_eq!(mock_server_1.request_count().await, 1);
let _observer_id_2 = create_test_observer(
&pool,
&format!("test-reload-2-{}", test_id),
Some(&entity_type),
Some("UPDATE"),
None,
&mock_server_2.webhook_url(),
)
.await
.expect("Failed to create observer 2");
let observer_count = runtime.reload_observers().await.expect("Failed to reload observers");
assert_eq!(observer_count, 2, "Should have 2 observers after reload");
let order_id_2 = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"UPDATE",
&entity_type,
&order_id_2.to_string(),
serde_json::json!({"id": order_id_2.to_string(), "status": "updated"}),
Some(serde_json::json!({"id": order_id_2.to_string(), "status": "created"})),
)
.await
.expect("Failed to insert change log entry 2");
wait_for_webhook(&mock_server_2, 1, Duration::from_secs(15)).await;
assert_eq!(mock_server_1.request_count().await, 1, "Observer 1 should have 1 event");
assert_eq!(
mock_server_2.request_count().await,
1,
"Observer 2 should have 1 event after reload"
);
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_graceful_shutdown_mid_processing() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_delayed_response(Duration::from_secs(2)).await;
let entity_type = format!("Order_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-shutdown-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let order_ids: Vec<_> = (0..5)
.map(|i| {
let order_id = Uuid::new_v4();
(order_id, i)
})
.collect();
for (order_id, i) in &order_ids {
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "sequence": i}),
None,
)
.await
.expect("Failed to insert change log entry");
}
tokio::time::sleep(Duration::from_secs(11)).await;
let checkpoint_exists = check_checkpoint_exists(&pool, "change_log")
.await
.expect("Failed to check checkpoint");
assert!(checkpoint_exists, "Expected checkpoint to exist");
let initial_count = mock_server.request_count().await;
assert!(initial_count > 0, "Expected at least one event to start processing");
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_runtime_continues_after_errors() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_transient_failure(2).await;
let entity_type = format!("Order_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-error-resilience-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let order_id_1 = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id_1.to_string(),
serde_json::json!({"id": order_id_1.to_string(), "sequence": 1}),
None,
)
.await
.expect("Failed to insert change log entry 1");
wait_for_webhook(&mock_server, 1, Duration::from_secs(20)).await;
let requests = mock_server.received_requests().await;
assert_eq!(requests.len(), 1, "Expected 1 successful webhook after retries");
let logs =
wait_for_observer_logs(&pool, &order_id_1.to_string(), 1, Duration::from_secs(10)).await;
assert!(!logs.is_empty(), "Expected observer logs for event with retries");
mock_server.reset().await;
mock_server.mock_success().await;
let order_id_2 = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id_2.to_string(),
serde_json::json!({"id": order_id_2.to_string(), "sequence": 2}),
None,
)
.await
.expect("Failed to insert change log entry 2");
wait_for_webhook(&mock_server, 1, Duration::from_secs(15)).await;
let second_count = mock_server.request_count().await;
assert_eq!(second_count, 1, "Expected runtime to continue processing after errors");
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_high_throughput_processing() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let entity_type = format!("Order_{}", test_id);
let _observer_id = create_test_observer(
&pool,
&format!("test-throughput-{}", test_id),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
let event_count = 100;
for i in 0..event_count {
let order_id = Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
&entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "sequence": i, "batch": "throughput"}),
None,
)
.await
.expect("Failed to insert change log entry");
}
wait_for_webhook(&mock_server, event_count, Duration::from_mins(1)).await;
let request_count = mock_server.request_count().await;
assert_eq!(
request_count, event_count,
"Expected {} webhooks for high throughput test, got {}",
event_count, request_count
);
let requests = mock_server.received_requests().await;
let ids: Vec<String> = requests
.iter()
.filter_map(|r| r["id"].as_str().map(|s| s.to_string()))
.collect();
let unique_ids: std::collections::HashSet<_> = ids.iter().cloned().collect();
assert_eq!(ids.len(), unique_ids.len(), "Expected no duplicates in high throughput test");
let success_count = get_observer_log_count(&pool, "success")
.await
.expect("Failed to query observer logs");
assert!(
usize::try_from(success_count).unwrap_or(0) >= event_count * 90 / 100,
"Expected at least 90% of events logged as success, got {}",
success_count
);
runtime.stop().await.expect("Failed to stop runtime");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_runtime_basic_lifecycle() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(50);
let mut runtime = ObserverRuntime::new(config);
let start_result = runtime.start().await;
assert!(start_result.is_ok(), "Failed to start runtime: {:?}", start_result);
let stop_result = runtime.stop().await;
assert!(stop_result.is_ok(), "Failed to stop runtime: {:?}", stop_result);
let mut runtime2 = ObserverRuntime::new(ObserverRuntimeConfig::new(pool));
let start_result2 = runtime2.start().await;
assert!(start_result2.is_ok(), "Failed to start runtime second time");
runtime2.stop().await.ok();
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_debug_event_processing() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let _observer_id = create_test_observer(
&pool,
"debug-observer",
Some("TestOrder"),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let observer_count: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM tb_observer WHERE enabled = true")
.fetch_one(&pool)
.await
.expect("Failed to count observers");
println!("✓ Created observer. Count in DB: {}", observer_count.0);
assert_eq!(observer_count.0, 1, "Observer not in database");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(10);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
println!("✓ Runtime started");
let order_id = uuid::Uuid::new_v4();
let _change_log_id = insert_change_log_entry(
&pool,
"INSERT",
"TestOrder",
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "amount": 100}),
None,
)
.await
.expect("Failed to insert change log entry");
println!("✓ Inserted change log entry");
let entry_count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM core.tb_entity_change_log")
.fetch_one(&pool)
.await
.expect("Failed to count entries");
println!("✓ Change log entries in DB: {}", entry_count.0);
tokio::time::sleep(Duration::from_secs(2)).await;
let requests = mock_server.received_requests().await;
println!("✓ Webhook calls received: {}", requests.len());
let log_count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM tb_observer_log")
.fetch_one(&pool)
.await
.ok()
.unwrap_or((0,));
println!("✓ Observer log entries: {}", log_count.0);
runtime.stop().await.expect("Failed to stop runtime");
println!("\nDebug Results:");
println!(" Observers in DB: {}", observer_count.0);
println!(" Change log entries: {}", entry_count.0);
println!(" Webhook calls: {}", requests.len());
println!(" Observer logs: {}", log_count.0);
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_observer_loading() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let _observer_id = create_test_observer(
&pool,
"load-test-observer",
Some("Product"),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let observer: Option<(String, Option<String>, Option<String>)> = sqlx::query_as(
"SELECT name, entity_type, event_type FROM tb_observer WHERE name = 'load-test-observer'",
)
.fetch_optional(&pool)
.await
.expect("Failed to query observer");
let (name, entity_type, event_type) = observer.expect("Observer not found");
println!("✓ Observer in DB:");
println!(" name: {}", name);
println!(" entity_type: {:?}", entity_type);
println!(" event_type: {:?}", event_type);
let actions: Option<(serde_json::Value,)> =
sqlx::query_as("SELECT actions FROM tb_observer WHERE name = 'load-test-observer'")
.fetch_optional(&pool)
.await
.expect("Failed to query actions");
if let Some((actions_json,)) = actions {
println!("✓ Actions:");
println!(" {}", serde_json::to_string_pretty(&actions_json).unwrap());
}
assert_eq!(name, "load-test-observer");
assert_eq!(entity_type.as_deref(), Some("Product"));
assert_eq!(event_type.as_deref(), Some("INSERT"));
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_runtime_loads_observers() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
let _observer_id = create_test_observer(
&pool,
"runtime-load-test",
Some("User"),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone());
let mut runtime = ObserverRuntime::new(config);
let start_result = runtime.start().await;
println!("✓ Runtime.start() result: {:?}", start_result);
assert!(start_result.is_ok(), "Failed to start: {:?}", start_result);
tokio::time::sleep(Duration::from_millis(100)).await;
runtime.stop().await.ok();
println!("✓ Runtime started and stopped successfully");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_debug_debezium_envelope() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let order_id = uuid::Uuid::new_v4();
let _ = insert_change_log_entry(
&pool,
"INSERT",
"Order",
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string(), "total": 50}),
None,
)
.await
.expect("Failed to insert");
let entry: Option<(i64, Option<i64>, String, uuid::Uuid, String, Option<String>, chrono::DateTime<chrono::Utc>, serde_json::Value)> = sqlx::query_as(
"SELECT pk_entity_change_log, fk_customer_org, object_type, object_id, modification_type, change_status, created_at, object_data FROM core.tb_entity_change_log LIMIT 1"
)
.fetch_optional(&pool)
.await
.expect("Query failed");
if let Some((pk, _fk_cust, obj_type, obj_id, mod_type, change_status, _created_at, obj_data)) =
entry
{
println!("✓ Change log entry found:");
println!(" pk: {}", pk);
println!(" object_type: {}", obj_type);
println!(" object_id: {}", obj_id);
println!(" modification_type: {}", mod_type);
println!(" change_status: {:?}", change_status);
println!(" object_data (Debezium envelope):");
println!(" {}", serde_json::to_string_pretty(&obj_data).unwrap());
if let Some(op_val) = obj_data.get("op") {
println!(" ✓ op field: {:?}", op_val);
if let Some(op_char) = op_val.as_str().and_then(|s| s.chars().next()) {
println!(" ✓ op first char: '{}'", op_char);
match op_char {
'c' => println!(" → Recognized as CREATE"),
'u' => println!(" → Recognized as UPDATE"),
'd' => println!(" → Recognized as DELETE"),
x => println!(" → UNRECOGNIZED: '{}'", x),
}
}
}
} else {
println!("✗ No change log entry found!");
}
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_action_parsing() {
let actions_json = serde_json::json!([
{
"type": "webhook",
"url": "http://127.0.0.1:8080/webhook",
"method": "POST",
"headers": {
"Content-Type": "application/json"
}
}
]);
println!("Input JSON: {}", serde_json::to_string_pretty(&actions_json).unwrap());
match serde_json::from_value::<Vec<fraiseql_observers::config::ActionConfig>>(actions_json) {
Ok(actions) => {
println!("✓ Successfully parsed {} actions", actions.len());
for (i, action) in actions.iter().enumerate() {
println!(" Action {}: {:?}", i, action);
}
},
Err(e) => {
println!("✗ Failed to parse actions: {}", e);
panic!("Action parsing failed");
},
}
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_with_longer_polling() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let mock_server = MockWebhookServer::start().await;
mock_server.mock_success().await;
let _observer_id = create_test_observer(
&pool,
"long-poll-test",
Some("Widget"),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
let config = ObserverRuntimeConfig::new(pool.clone()).with_poll_interval(5);
let mut runtime = ObserverRuntime::new(config);
runtime.start().await.expect("Failed to start runtime");
println!("Runtime initialized");
let widget_id = uuid::Uuid::new_v4();
println!("Inserting change log entry...");
let _ = insert_change_log_entry(
&pool,
"INSERT",
"Widget",
&widget_id.to_string(),
serde_json::json!({"id": widget_id.to_string(), "name": "Test Widget"}),
None,
)
.await
.expect("Failed to insert");
println!("Change log entry inserted");
println!("Waiting for event processing...");
for i in 0..50 {
tokio::time::sleep(Duration::from_millis(10)).await;
let requests = mock_server.received_requests().await;
if !requests.is_empty() {
println!("✓ Webhook called after {} ms", (i + 1) * 10);
break;
}
if i % 10 == 0 {
println!(" Still waiting... ({} ms elapsed)", (i + 1) * 10);
}
}
let requests = mock_server.received_requests().await;
println!("Final webhook calls: {}", requests.len());
println!("Expected: 1");
runtime.stop().await.ok();
if requests.is_empty() {
println!("\nDEBUG: Checking database state...");
let cl_count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM core.tb_entity_change_log")
.fetch_one(&pool)
.await
.ok()
.unwrap_or((0,));
println!(" Change log entries: {}", cl_count.0);
let obs_count: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM tb_observer WHERE enabled = true")
.fetch_one(&pool)
.await
.ok()
.unwrap_or((0,));
println!(" Observers enabled: {}", obs_count.0);
let logs: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM tb_observer_log")
.fetch_one(&pool)
.await
.ok()
.unwrap_or((0,));
println!(" Observer logs: {}", logs.0);
}
assert!(!requests.is_empty(), "No webhook calls received after 500ms");
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_listener_direct() {
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("Failed to clean observer logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("Failed to clean observers");
sqlx::query("DELETE FROM core.tb_entity_change_log")
.execute(&pool)
.await
.expect("Failed to clean change log");
let product_id = uuid::Uuid::new_v4();
insert_change_log_entry(
&pool,
"INSERT",
"Product",
&product_id.to_string(),
serde_json::json!({"id": product_id.to_string(), "name": "Test"}),
None,
)
.await
.expect("Failed to insert");
let config =
fraiseql_observers::listener::change_log::ChangeLogListenerConfig::new(pool.clone())
.with_poll_interval(10);
let mut listener = fraiseql_observers::listener::change_log::ChangeLogListener::new(config);
println!("Calling listener.next_batch()...");
let result = listener.next_batch().await;
match result {
Ok(entries) => {
println!("✓ Got {} entries from listener", entries.len());
assert!(!entries.is_empty(), "Listener should have found entries");
for entry in entries {
println!(
" Entry: pk={}, object_type={}, op={:?}",
entry.id,
entry.object_type,
entry.object_data.get("op")
);
match entry.to_entity_event() {
Ok(event) => {
println!(" ✓ Converted to EntityEvent: {:?}", event.event_type);
},
Err(e) => {
println!(" ✗ Failed to convert: {}", e);
panic!("Failed to convert: {}", e);
},
}
}
},
Err(e) => {
panic!("Listener failed: {}", e);
},
}
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_admin_create_reloads_matcher() {
use std::sync::Arc;
use axum::{Json, extract::State, http::StatusCode, response::IntoResponse};
use fraiseql_server::{
extractors::OptionalSecurityContext,
observers::{
CreateObserverRequest, ObserverRepository, ObserverState, handlers::create_observer,
},
};
use tokio::sync::RwLock;
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
sqlx::query("DELETE FROM tb_observer_log")
.execute(&pool)
.await
.expect("clean logs");
sqlx::query("DELETE FROM tb_observer")
.execute(&pool)
.await
.expect("clean observers");
let runtime =
Arc::new(RwLock::new(ObserverRuntime::new(ObserverRuntimeConfig::new(pool.clone()))));
let baseline = runtime.read().await.reload_observers().await.expect("baseline reload");
assert_eq!(baseline, 0, "no observers should be live before the create");
let state = ObserverState {
repository: ObserverRepository::new(pool.clone()),
runtime: Some(Arc::clone(&runtime)),
};
let entity_type = format!("Order_{test_id}");
let request: CreateObserverRequest = serde_json::from_value(serde_json::json!({
"name": format!("reload-on-create-{test_id}"),
"entity_type": entity_type,
"event_type": "INSERT",
"actions": [{ "type": "webhook", "url": "http://127.0.0.1:9/never-called" }],
}))
.expect("valid create request");
let response = create_observer(State(state), OptionalSecurityContext(None), Json(request))
.await
.into_response();
assert_eq!(response.status(), StatusCode::CREATED, "create should succeed");
let live = runtime.read().await.health().observer_count;
assert_eq!(live, 1, "create handler should have refreshed the in-process matcher (#466)");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
}
async fn seed_replica_test(
pool: &sqlx::PgPool,
test_id: &str,
webhook_delay: Duration,
) -> (MockWebhookServer, String) {
setup_observer_schema(pool).await.expect("Failed to setup schema");
for table in [
"tb_observer_log",
"tb_observer",
"core.tb_entity_change_log",
] {
sqlx::query(&format!("DELETE FROM {table}"))
.execute(pool)
.await
.expect("Failed to clean table");
}
let mock_server = MockWebhookServer::start().await;
mock_server.mock_delayed_response(webhook_delay).await;
let entity_type = format!("Order_{test_id}");
create_test_observer(
pool,
&format!("test-replicas-{test_id}"),
Some(&entity_type),
Some("INSERT"),
None,
&mock_server.webhook_url(),
)
.await
.expect("Failed to create observer");
(mock_server, entity_type)
}
async fn insert_order(pool: &sqlx::PgPool, entity_type: &str) -> Uuid {
let order_id = Uuid::new_v4();
insert_change_log_entry(
pool,
"INSERT",
entity_type,
&order_id.to_string(),
serde_json::json!({"id": order_id.to_string()}),
None,
)
.await
.expect("Failed to insert change log entry");
order_id
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_two_replicas_dispatch_a_row_once() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let listener_id = format!("replicas-{test_id}");
let pool_a = create_test_pool().await;
let pool_b = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool_a, &test_id, Duration::from_millis(1500)).await;
let mut replica_a = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool_a.clone())
.with_poll_interval(50)
.with_listener_id(listener_id.clone()),
);
let mut replica_b = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool_b.clone())
.with_poll_interval(50)
.with_listener_id(listener_id.clone()),
);
replica_a.start().await.expect("Failed to start replica A");
replica_b.start().await.expect("Failed to start replica B");
let order_id = insert_order(&pool_a, &entity_type).await;
wait_for_webhook(&mock_server, 1, Duration::from_secs(20)).await;
tokio::time::sleep(Duration::from_secs(4)).await;
let deliveries = mock_server.request_count().await;
replica_a.stop().await.expect("Failed to stop replica A");
replica_b.stop().await.expect("Failed to stop replica B");
cleanup_test_data(&pool_a, &test_id).await.expect("Failed to cleanup");
assert_eq!(
deliveries, 1,
"one change-log row ({order_id}) was dispatched {deliveries} times across two replicas"
);
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_standby_replica_takes_over_when_the_poller_stops() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let listener_id = format!("takeover-{test_id}");
let pool_a = create_test_pool().await;
let pool_b = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool_a, &test_id, Duration::from_millis(0)).await;
let mut replica_a = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool_a.clone())
.with_poll_interval(50)
.with_listener_id(listener_id.clone()),
);
replica_a.start().await.expect("Failed to start replica A");
insert_order(&pool_a, &entity_type).await;
wait_for_webhook(&mock_server, 1, Duration::from_secs(20)).await;
let mut replica_b = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool_b.clone())
.with_poll_interval(50)
.with_listener_id(listener_id.clone()),
);
replica_b.start().await.expect("Failed to start replica B");
insert_order(&pool_a, &entity_type).await;
wait_for_webhook(&mock_server, 2, Duration::from_secs(20)).await;
replica_a.stop().await.expect("Failed to stop replica A");
let third = insert_order(&pool_b, &entity_type).await;
wait_for_webhook(&mock_server, 3, Duration::from_secs(20)).await;
tokio::time::sleep(Duration::from_millis(500)).await;
let ids: Vec<String> = mock_server
.received_requests()
.await
.iter()
.filter_map(|r| r["id"].as_str().map(str::to_string))
.collect();
let (event_count,): (i32,) =
sqlx::query_as("SELECT event_count FROM observer_checkpoints WHERE listener_id = $1")
.bind(&listener_id)
.fetch_one(&pool_a)
.await
.expect("the checkpoint is stored");
replica_b.stop().await.expect("Failed to stop replica B");
cleanup_test_data(&pool_a, &test_id).await.expect("Failed to cleanup");
assert_eq!(ids.len(), 3, "expected each of three rows once, got {ids:?}");
assert_eq!(
ids[2],
third.to_string(),
"the standby must dispatch the row inserted after takeover"
);
assert_eq!(
event_count, 3,
"the standby must continue the stored checkpoint, not the one it read at start"
);
}
#[cfg(feature = "observers-nats")]
#[tokio::test]
#[ignore = "requires PostgreSQL and NATS"]
async fn test_nats_runtime_creates_its_stream_and_consumer_with_the_configured_limits() {
use fraiseql_observers::config::{TransportConfig, TransportKind};
init_test_tracing();
let pool = create_test_pool().await;
setup_observer_schema(&pool).await.expect("Failed to setup schema");
let test_id = Uuid::new_v4().simple().to_string();
let url =
std::env::var("NATS_URL").expect("NATS_URL must be set (the observers leg binds NATS)");
let mut transport = TransportConfig {
transport: TransportKind::Nats,
..TransportConfig::default()
};
transport.nats.url.clone_from(&url);
transport.nats.stream_name = format!("limits-{test_id}");
transport.nats.consumer_name = format!("limits-consumer-{test_id}");
transport.nats.subject_prefix = format!("limits.{test_id}");
transport.nats.jetstream.dedup_window_minutes = 3;
transport.nats.jetstream.max_age_days = 2;
transport.nats.jetstream.max_deliver = 9;
transport.nats.jetstream.max_bytes = 1024 * 1024;
let mut runtime =
ObserverRuntime::new(ObserverRuntimeConfig::new(pool.clone()).with_transport(transport));
runtime.start().await.expect("the runtime starts on the NATS transport");
let jetstream = async_nats::jetstream::new(async_nats::connect(&url).await.expect("connect"));
let stream_name = format!("limits-{test_id}");
let mut stream = jetstream
.get_stream(&stream_name)
.await
.expect("the runtime created the stream");
let stream_config = stream.info().await.expect("stream info").config.clone();
let mut consumer: async_nats::jetstream::consumer::PullConsumer = stream
.get_consumer(&format!("limits-consumer-{test_id}"))
.await
.expect("the runtime created the durable consumer");
let max_deliver = consumer.info().await.expect("consumer info").config.max_deliver;
runtime.stop().await.expect("Failed to stop runtime");
jetstream.delete_stream(&stream_name).await.expect("clean up the stream");
assert_eq!(stream_config.duplicate_window, Duration::from_mins(3));
assert_eq!(stream_config.max_age, Duration::from_hours(48));
assert_eq!(max_deliver, 9);
}
async fn collect_bridge_ids(
rx: &mut tokio::sync::mpsc::Receiver<fraiseql_server::subscriptions::EntityEvent>,
want: usize,
timeout: Duration,
) -> Vec<String> {
let mut ids = Vec::new();
let deadline = tokio::time::Instant::now() + timeout;
loop {
let distinct: std::collections::HashSet<&String> = ids.iter().collect();
if distinct.len() >= want {
while let Ok(Some(e)) =
tokio::time::timeout(Duration::from_millis(300), rx.recv()).await
{
ids.push(e.entity_id);
}
return ids;
}
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(e)) => ids.push(e.entity_id),
Ok(None) | Err(_) => return ids,
}
}
}
#[tokio::test]
#[ignore = "requires PostgreSQL"]
async fn test_every_replica_delivers_every_change_to_its_subscribers() {
init_test_tracing();
let test_id = Uuid::new_v4().to_string();
let listener_id = format!("fanout-{test_id}");
let pool_a = create_test_pool().await;
let pool_b = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool_a, &test_id, Duration::from_millis(0)).await;
let mut replicas = Vec::new();
let mut receivers = Vec::new();
for pool in [&pool_a, &pool_b] {
let (tx, rx) = tokio::sync::mpsc::channel(64);
let mut replica = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool.clone())
.with_poll_interval(50)
.with_listener_id(listener_id.clone()),
);
replica.set_event_bridge_sender(tx);
replica.start().await.expect("Failed to start replica");
replicas.push(replica);
receivers.push(rx);
}
let mut orders = Vec::new();
for _ in 0..3 {
orders.push(insert_order(&pool_a, &entity_type).await.to_string());
}
wait_for_webhook(&mock_server, 3, Duration::from_secs(20)).await;
let mut seen = Vec::new();
for rx in &mut receivers {
seen.push(collect_bridge_ids(rx, 3, Duration::from_secs(20)).await);
}
tokio::time::sleep(Duration::from_millis(500)).await;
let deliveries = mock_server.request_count().await;
for replica in &mut replicas {
replica.stop().await.expect("Failed to stop replica");
}
cleanup_test_data(&pool_a, &test_id).await.expect("Failed to cleanup");
for (name, ids) in ["A", "B"].iter().zip(&seen) {
let mut sorted = ids.clone();
sorted.sort();
let mut want = orders.clone();
want.sort();
assert_eq!(sorted, want, "replica {name}'s subscribers must see each change once");
}
assert_eq!(deliveries, 3, "each change must still be dispatched once, not once per replica");
}
#[cfg(feature = "observers-nats")]
#[tokio::test]
#[ignore = "requires PostgreSQL and NATS"]
async fn test_every_nats_replica_delivers_every_change_to_its_subscribers() {
use fraiseql_observers::{
EntityEvent, EventKind,
config::{TransportConfig, TransportKind},
};
init_test_tracing();
let test_id = Uuid::new_v4().simple().to_string();
let pool = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool, &test_id, Duration::from_millis(0)).await;
let url =
std::env::var("NATS_URL").expect("NATS_URL must be set (the observers leg binds NATS)");
let mut transport = TransportConfig {
transport: TransportKind::Nats,
..TransportConfig::default()
};
transport.nats.url.clone_from(&url);
transport.nats.stream_name = format!("fanout-{test_id}");
transport.nats.consumer_name = format!("fanout-consumer-{test_id}");
transport.nats.subject_prefix = format!("fanout.{test_id}");
transport.nats.jetstream.max_bytes = 1024 * 1024;
let mut replicas = Vec::new();
let mut receivers = Vec::new();
for _ in 0..2 {
let (tx, rx) = tokio::sync::mpsc::channel(64);
let mut replica = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool.clone()).with_transport(transport.clone()),
);
replica.set_event_bridge_sender(tx);
replica.start().await.expect("the replica starts on the NATS transport");
replicas.push(replica);
receivers.push(rx);
}
let jetstream = async_nats::jetstream::new(async_nats::connect(&url).await.expect("connect"));
let mut orders = Vec::new();
for _ in 0..4 {
let id = Uuid::new_v4();
let event = EntityEvent::new(
EventKind::Created,
entity_type.clone(),
id,
serde_json::json!({"id": id.to_string()}),
);
jetstream
.publish(
format!("fanout.{test_id}.{entity_type}.INSERT"),
serde_json::to_vec(&event).expect("serialize").into(),
)
.await
.expect("publish")
.await
.expect("the stream stores the event");
orders.push(id.to_string());
}
wait_for_webhook(&mock_server, 4, Duration::from_secs(20)).await;
let mut seen = Vec::new();
for rx in &mut receivers {
seen.push(collect_bridge_ids(rx, 4, Duration::from_secs(20)).await);
}
tokio::time::sleep(Duration::from_millis(500)).await;
let deliveries = mock_server.request_count().await;
for replica in &mut replicas {
replica.stop().await.expect("Failed to stop replica");
}
jetstream
.delete_stream(&format!("fanout-{test_id}"))
.await
.expect("clean up the stream");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
for (name, ids) in ["A", "B"].iter().zip(&seen) {
let mut sorted = ids.clone();
sorted.sort();
let mut want = orders.clone();
want.sort();
assert_eq!(sorted, want, "replica {name}'s subscribers must see each change once");
}
assert_eq!(deliveries, 4, "each change must still be dispatched once, not once per replica");
}
#[cfg(feature = "observers-nats")]
#[tokio::test]
#[ignore = "requires PostgreSQL and NATS"]
async fn test_nats_runtime_acknowledges_an_event_after_its_actions_ran() {
use fraiseql_observers::{
EntityEvent, EventKind,
config::{TransportConfig, TransportKind},
};
init_test_tracing();
let test_id = Uuid::new_v4().simple().to_string();
let pool = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool, &test_id, Duration::from_secs(3)).await;
let url =
std::env::var("NATS_URL").expect("NATS_URL must be set (the observers leg binds NATS)");
let mut transport = TransportConfig {
transport: TransportKind::Nats,
..TransportConfig::default()
};
transport.nats.url.clone_from(&url);
transport.nats.stream_name = format!("ack-{test_id}");
transport.nats.consumer_name = format!("ack-consumer-{test_id}");
transport.nats.subject_prefix = format!("ack.{test_id}");
transport.nats.jetstream.max_bytes = 1024 * 1024;
let mut runtime =
ObserverRuntime::new(ObserverRuntimeConfig::new(pool.clone()).with_transport(transport));
runtime.start().await.expect("the runtime starts on the NATS transport");
let jetstream = async_nats::jetstream::new(async_nats::connect(&url).await.expect("connect"));
let id = Uuid::new_v4();
let event = EntityEvent::new(
EventKind::Created,
entity_type.clone(),
id,
serde_json::json!({"id": id.to_string()}),
);
jetstream
.publish(
format!("ack.{test_id}.{entity_type}.INSERT"),
serde_json::to_vec(&event).expect("serialize").into(),
)
.await
.expect("publish")
.await
.expect("the stream stores the event");
let pending = |jetstream: async_nats::jetstream::Context| {
let (stream, consumer) = (format!("ack-{test_id}"), format!("ack-consumer-{test_id}"));
async move {
let mut consumer: async_nats::jetstream::consumer::PullConsumer = jetstream
.get_stream(&stream)
.await
.expect("stream")
.get_consumer(&consumer)
.await
.expect("consumer");
consumer.info().await.expect("consumer info").num_ack_pending
}
};
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while mock_server.request_count().await == 0 && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(50)).await;
}
let while_running = pending(jetstream.clone()).await;
tokio::time::sleep(Duration::from_secs(5)).await;
let after = pending(jetstream.clone()).await;
runtime.stop().await.expect("Failed to stop runtime");
jetstream
.delete_stream(&format!("ack-{test_id}"))
.await
.expect("clean up the stream");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
assert_eq!(while_running, 1, "the event must stay unacknowledged while its webhook runs");
assert_eq!(after, 0, "the event must be acknowledged once its actions ran");
}
#[cfg(feature = "observers-nats")]
async fn publish_order_event(
jetstream: &async_nats::jetstream::Context,
prefix: &str,
entity_type: &str,
) {
use fraiseql_observers::{EntityEvent, EventKind};
let id = Uuid::new_v4();
let event = EntityEvent::new(
EventKind::Created,
entity_type.to_string(),
id,
serde_json::json!({"id": id.to_string()}),
);
jetstream
.publish(
format!("{prefix}.{entity_type}.INSERT"),
serde_json::to_vec(&event).expect("serialize").into(),
)
.await
.expect("publish")
.await
.expect("the stream stores the event");
}
#[cfg(feature = "observers-nats")]
fn nats_test_transport(
url: &str,
name: &str,
ack_wait_secs: u64,
) -> fraiseql_observers::config::TransportConfig {
use fraiseql_observers::config::{TransportConfig, TransportKind};
let mut transport = TransportConfig {
transport: TransportKind::Nats,
..TransportConfig::default()
};
transport.nats.url = url.to_string();
transport.nats.stream_name = name.to_string();
transport.nats.consumer_name = format!("{name}-consumer");
transport.nats.subject_prefix = name.to_string();
transport.nats.jetstream.max_bytes = 1024 * 1024;
transport.nats.jetstream.ack_wait_secs = ack_wait_secs;
transport.nats.jetstream.max_deliver = 3;
transport
}
#[cfg(feature = "observers-nats")]
async fn nats_consumer_info(
jetstream: &async_nats::jetstream::Context,
name: &str,
) -> async_nats::jetstream::consumer::Info {
let mut consumer: async_nats::jetstream::consumer::PullConsumer = jetstream
.get_stream(name)
.await
.expect("stream")
.get_consumer(&format!("{name}-consumer"))
.await
.expect("consumer");
consumer.info().await.expect("consumer info").clone()
}
#[cfg(feature = "observers-nats")]
#[tokio::test]
#[ignore = "requires PostgreSQL and NATS"]
async fn test_nats_runtime_does_not_redeliver_a_dispatch_longer_than_ack_wait() {
init_test_tracing();
let test_id = Uuid::new_v4().simple().to_string();
let pool = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool, &test_id, Duration::from_secs(5)).await;
let url =
std::env::var("NATS_URL").expect("NATS_URL must be set (the observers leg binds NATS)");
let name = format!("progress-{test_id}");
let mut runtime = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool.clone())
.with_transport(nats_test_transport(&url, &name, 2)),
);
runtime.start().await.expect("the runtime starts on the NATS transport");
let jetstream = async_nats::jetstream::new(async_nats::connect(&url).await.expect("connect"));
publish_order_event(&jetstream, &name, &entity_type).await;
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while mock_server.request_count().await == 0 && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(50)).await;
}
tokio::time::sleep(Duration::from_secs(12)).await;
let info = nats_consumer_info(&jetstream, &name).await;
let requests = mock_server.request_count().await;
runtime.stop().await.expect("Failed to stop runtime");
jetstream.delete_stream(&name).await.expect("clean up the stream");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
assert_eq!(
info.delivered.consumer_sequence, 1,
"the broker must not redeliver a running dispatch"
);
assert_eq!(requests, 1, "the webhook must be called once");
assert_eq!(info.num_ack_pending, 0, "the event must be acknowledged once its actions ran");
}
#[cfg(feature = "observers-nats")]
#[tokio::test]
#[ignore = "requires PostgreSQL and NATS"]
async fn test_nats_event_of_a_runtime_that_died_mid_dispatch_is_redelivered() {
init_test_tracing();
let test_id = Uuid::new_v4().simple().to_string();
let pool = create_test_pool().await;
let (mock_server, entity_type) =
seed_replica_test(&pool, &test_id, Duration::from_secs(6)).await;
let url =
std::env::var("NATS_URL").expect("NATS_URL must be set (the observers leg binds NATS)");
let name = format!("crash-{test_id}");
let (started_tx, started_rx) = std::sync::mpsc::channel::<()>();
let (crash_tx, crash_rx) = std::sync::mpsc::channel::<()>();
let transport = nats_test_transport(&url, &name, 2);
let doomed = std::thread::spawn(move || {
let rt = tokio::runtime::Builder::new_multi_thread().enable_all().build().unwrap();
rt.block_on(async {
let pool = create_test_pool().await;
let mut runtime =
ObserverRuntime::new(ObserverRuntimeConfig::new(pool).with_transport(transport));
runtime.start().await.expect("the first runtime starts");
Box::leak(Box::new(runtime));
});
started_tx.send(()).unwrap();
crash_rx.recv().unwrap();
rt.shutdown_background();
});
started_rx.recv().expect("the first runtime started");
let jetstream = async_nats::jetstream::new(async_nats::connect(&url).await.expect("connect"));
publish_order_event(&jetstream, &name, &entity_type).await;
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
while mock_server.request_count().await == 0 && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(50)).await;
}
tokio::time::sleep(Duration::from_millis(2500)).await;
assert_eq!(mock_server.request_count().await, 1, "the first runtime must be dispatching");
crash_tx.send(()).unwrap();
doomed.join().unwrap();
let mut survivor = ObserverRuntime::new(
ObserverRuntimeConfig::new(pool.clone())
.with_transport(nats_test_transport(&url, &name, 2)),
);
survivor.start().await.expect("the second runtime starts");
let deadline = tokio::time::Instant::now() + Duration::from_secs(45);
while mock_server.request_count().await < 2 && tokio::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(100)).await;
}
let requests = mock_server.request_count().await;
survivor.stop().await.expect("Failed to stop runtime");
jetstream.delete_stream(&name).await.expect("clean up the stream");
cleanup_test_data(&pool, &test_id).await.expect("Failed to cleanup");
assert_eq!(requests, 2, "the event must be redelivered to the surviving runtime");
}