use super::super::control_plane::resources::{aggregate_version, content_version};
use super::super::control_plane::store;
use super::super::{AuthnServiceImpl, AuthzServiceImpl};
use super::support::*;
use crate::proto::udb::core::apikey::services::v1 as apikey_pb;
use crate::proto::udb::core::apikey::services::v1::api_key_service_server::ApiKeyService;
use crate::proto::udb::core::authn::services::v1 as authn_pb;
use crate::proto::udb::core::authn::services::v1::authn_service_server::AuthnService;
use crate::proto::udb::core::authz::services::v1 as authz_pb;
use crate::proto::udb::core::authz::services::v1::authz_service_server::AuthzService;
use crate::proto::udb::core::common::v1 as common_pb;
use crate::proto::udb::core::control::entity::v1::ResourceType;
use crate::runtime::service::method_security::{scope_claim_context_for_test, test_claim_context};
use tonic::Request;
use uuid::Uuid;
const TEST_SNAPSHOT_TTL_SECS: u64 = 1;
fn set_short_snapshot_ttl() {
unsafe {
std::env::set_var(
"UDB_AUTHZ_SNAPSHOT_TTL_SECS",
TEST_SNAPSHOT_TTL_SECS.to_string(),
);
}
}
fn check_access(user_id: &str, object: &str, action: &str) -> authz_pb::CheckAccessRequest {
authz_pb::CheckAccessRequest {
user_id: user_id.to_string(),
domain: "acme".to_string(),
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
object: object.to_string(),
action: action.to_string(),
..Default::default()
}
}
async fn allowed_on(node: &AuthzServiceImpl, user_id: &str, object: &str, action: &str) -> bool {
node.check_access(Request::new(check_access(user_id, object, action)))
.await
.expect("check_access")
.into_inner()
.allowed
}
async fn seed_authorized_principal(
authn: &AuthnServiceImpl,
authz: &AuthzServiceImpl,
prefix: &str,
object: &str,
action: &str,
) -> (String, String) {
let user = create_verified_user(authn, prefix, "CorrectHorse1!").await;
let suffix = Uuid::new_v4().simple().to_string();
let role_code = format!("{prefix}_{suffix}");
authz
.put_role_binding(Request::new(authz_pb::PutRoleBindingRequest {
binding: Some(authz_pb::RoleBinding {
subject: user.user_id.clone(),
role: role_code.clone(),
tenant: "acme".to_string(),
project: "billing".to_string(),
expires_at_unix: 0,
source: "ha_multinode_test".to_string(),
}),
}))
.await
.expect("put_role_binding");
authz
.put_authz_policy(Request::new(authz_pb::PutAuthzPolicyRequest {
policy: Some(authz_pb::AuthzPolicyRecord {
id: Uuid::new_v4().to_string(),
enabled: true,
effect: "allow".to_string(),
tenant: "acme".to_string(),
project: "billing".to_string(),
role: role_code.clone(),
action: action.to_string(),
resource: object.to_string(),
..Default::default()
}),
}))
.await
.expect("put allow policy");
(user.user_id, role_code)
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_ha_authz_revision_invalidation -- --ignored --nocapture"]
async fn live_postgres_ha_authz_revision_invalidation() {
let _guard = live_auth_db_lock().lock().await;
set_short_snapshot_ttl();
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let authn = authn_service(pool.clone());
let node1 = authz_service(pool.clone()).await;
let node2 = authz_service(pool.clone()).await;
let (user_id, role_code) =
seed_authorized_principal(&authn, &node1, "ha-authz", "invoice", "data.update").await;
assert!(
allowed_on(&node1, &user_id, "invoice", "data.update").await,
"node1 must ALLOW the seeded principal"
);
assert!(
allowed_on(&node2, &user_id, "invoice", "data.update").await,
"node2 must ALLOW the seeded principal (shared durable snapshot)"
);
let before = node2
.get_authz_revision(Request::new(authz_pb::GetAuthzRevisionRequest {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
}))
.await
.expect("node2 revision before")
.into_inner();
node1
.put_authz_policy(Request::new(authz_pb::PutAuthzPolicyRequest {
policy: Some(authz_pb::AuthzPolicyRecord {
id: Uuid::new_v4().to_string(),
enabled: true,
effect: "deny".to_string(),
tenant: "acme".to_string(),
project: "billing".to_string(),
role: role_code,
action: "data.update".to_string(),
resource: "invoice".to_string(),
..Default::default()
}),
}))
.await
.expect("node1 put explicit DENY");
let after = node2
.get_authz_revision(Request::new(authz_pb::GetAuthzRevisionRequest {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
}))
.await
.expect("node2 revision after")
.into_inner();
assert!(
after.policy_revision > before.policy_revision,
"node2 must observe an advanced policy_revision ({} -> {}) after node1's mutation",
before.policy_revision,
after.policy_revision
);
assert!(
!after.content_hash.is_empty(),
"node2's revision content_hash must be populated"
);
tokio::time::sleep(std::time::Duration::from_secs(TEST_SNAPSHOT_TTL_SECS + 2)).await;
assert!(
!allowed_on(&node2, &user_id, "invoice", "data.update").await,
"node2 must DENY within the snapshot TTL after node1 added an explicit DENY"
);
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_ha_policy_bundle_revocation -- --ignored --nocapture"]
async fn live_postgres_ha_policy_bundle_revocation() {
let _guard = live_auth_db_lock().lock().await;
set_short_snapshot_ttl();
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let authn = authn_service(pool.clone());
let node1 = authz_service(pool.clone()).await;
let node2 = authz_service(pool.clone()).await;
seed_authorized_principal(&authn, &node1, "ha-bundle", "report", "data.export").await;
let before = node2
.get_authz_revision(Request::new(authz_pb::GetAuthzRevisionRequest {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
}))
.await
.expect("node2 revision before invalidation")
.into_inner();
let invalidated = node1
.invalidate_policy_bundles(Request::new(authz_pb::InvalidatePolicyBundlesRequest {
actor: Some(authz_pb::GovernanceActor {
subject: "ha-admin".to_string(),
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
scopes: vec!["authz:admin".to_string()],
..Default::default()
}),
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
reason: "ha_multinode_bundle_revocation".to_string(),
}))
.await
.expect("node1 invalidate_policy_bundles")
.into_inner();
assert!(invalidated.ok);
assert!(
invalidated.policy_revision > before.policy_revision,
"invalidation must advance the policy revision ({} -> {})",
before.policy_revision,
invalidated.policy_revision
);
let after = node2
.get_authz_revision(Request::new(authz_pb::GetAuthzRevisionRequest {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
}))
.await
.expect("node2 revision after invalidation")
.into_inner();
assert!(
after.policy_revision >= invalidated.policy_revision,
"node2 must observe the post-invalidation revision (>= {}), saw {}",
invalidated.policy_revision,
after.policy_revision
);
match node2
.get_policy_bundle(Request::new(authz_pb::PolicyBundleRequest {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
domain: "acme".to_string(),
}))
.await
{
Ok(resp) => {
let bundle = resp.into_inner().bundle.expect("signed bundle");
assert!(
!bundle.policy_version.is_empty(),
"a freshly issued bundle must carry a non-empty policy_version"
);
}
Err(status) => {
assert_eq!(
status.code(),
tonic::Code::FailedPrecondition,
"without UDB_POLICY_BUNDLE_SECRET, bundle issuance must fail the \
signing precondition (got {status:?})"
);
}
}
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_ha_policy_distribution_ack_nack_rollback -- --ignored --nocapture"]
async fn live_postgres_ha_policy_distribution_ack_nack_rollback() {
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let rt = ResourceType::RoutingPolicy;
let name = format!("ha-routing-{}", Uuid::new_v4().simple());
let node_a = format!("node-a-{}", Uuid::new_v4().simple());
let node_b = format!("node-b-{}", Uuid::new_v4().simple());
let v1_payload = r#"{"route":"primary","weight":100}"#;
let v1 = store::upsert_resource(&pool, rt, &name, "", "billing", v1_payload, "ha-test")
.await
.expect("publish v1 resource");
assert_eq!(v1.content_hash, content_version(v1_payload));
let world_v1 = store::world_version(&pool, rt, None, &[])
.await
.expect("world version v1");
store::ensure_node_state(&pool, &node_a, rt, &[name.clone()])
.await
.expect("ensure node A state");
let nonce_a1 = store::next_response_nonce(&pool, &node_a, rt)
.await
.expect("nonce A v1");
assert!(
store::record_ack(&pool, &node_a, rt, &world_v1, &nonce_a1)
.await
.expect("record ACK A v1"),
"ACK with the matching nonce must be applied"
);
let after_ack = store::get_node_state(&pool, &node_a, rt)
.await
.expect("get node A state")
.expect("node A row exists");
assert_eq!(after_ack.accepted_version, world_v1);
assert_eq!(after_ack.last_good_version, world_v1);
assert!(after_ack.nack_error_detail.is_empty());
assert!(
!store::record_ack(&pool, &node_a, rt, "bogus-world", "stale-nonce")
.await
.expect("stale ACK A"),
"an ACK whose nonce does not match the last response nonce must be ignored"
);
let still_v1 = store::get_node_state(&pool, &node_a, rt)
.await
.expect("get node A after stale ack")
.expect("row");
assert_eq!(
still_v1.accepted_version, world_v1,
"a stale-nonce ACK must not advance accepted_version"
);
let v2_payload = r#"{"route":"","weight":-1}"#;
let v2 = store::upsert_resource(&pool, rt, &name, "", "billing", v2_payload, "ha-test")
.await
.expect("publish v2 resource");
assert_ne!(
v2.content_hash, v1.content_hash,
"a content change must bump the version"
);
let world_v2 = store::world_version(&pool, rt, None, &[])
.await
.expect("world version v2");
assert_ne!(world_v1, world_v2, "v2 must change the world version");
let nonce_a2 = store::next_response_nonce(&pool, &node_a, rt)
.await
.expect("nonce A v2");
assert!(
store::record_nack(
&pool,
&node_a,
rt,
&nonce_a2,
r#"{"code":3,"message":"invalid routing policy: empty route"}"#,
)
.await
.expect("record NACK A v2"),
"NACK with the matching nonce must be applied"
);
let after_nack = store::get_node_state(&pool, &node_a, rt)
.await
.expect("get node A after NACK")
.expect("row");
assert_eq!(
after_nack.accepted_version, world_v1,
"NACK must NOT advance accepted_version onto the bad v2"
);
assert_eq!(
after_nack.last_good_version, world_v1,
"NACK must preserve last_good_version"
);
assert!(
!after_nack.nack_error_detail.is_empty(),
"NACK must record the structured error detail"
);
store::ensure_node_state(&pool, &node_b, rt, &[name.clone()])
.await
.expect("ensure node B state");
let nonce_b2 = store::next_response_nonce(&pool, &node_b, rt)
.await
.expect("nonce B v2");
assert!(
store::record_nack(&pool, &node_b, rt, &nonce_b2, "node B rejects v2")
.await
.expect("record NACK B v2")
);
let node_a_unchanged = store::get_node_state(&pool, &node_a, rt)
.await
.expect("node A after B's NACK")
.expect("row");
assert_eq!(
node_a_unchanged.accepted_version, world_v1,
"node B's NACK must not affect node A's accepted_version"
);
let rolled_back =
store::upsert_resource(&pool, rt, &name, "", "billing", v1_payload, "ha-test")
.await
.expect("rollback to v1 content");
assert_eq!(
rolled_back.content_hash,
content_version(v1_payload),
"rollback restores the v1 content hash"
);
let world_rb = store::world_version(&pool, rt, None, &[])
.await
.expect("world version after rollback");
assert_eq!(
world_rb, world_v1,
"rolling the content back to v1 restores the v1 world version"
);
let nonce_a3 = store::next_response_nonce(&pool, &node_a, rt)
.await
.expect("nonce A rollback");
assert!(
store::record_ack(&pool, &node_a, rt, &world_rb, &nonce_a3)
.await
.expect("record ACK A rollback")
);
let after_rollback = store::get_node_state(&pool, &node_a, rt)
.await
.expect("node A after rollback")
.expect("row");
assert_eq!(after_rollback.accepted_version, world_v1);
assert!(
after_rollback.nack_error_detail.is_empty(),
"a successful ACK after rollback must clear the prior NACK error"
);
let resources = store::list_resources(&pool, rt, None, &[name.clone()])
.await
.expect("list resources");
assert_eq!(aggregate_version(&resources), world_rb);
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_control_reload_metric_fires_on_served_path -- --ignored --nocapture"]
async fn live_postgres_control_reload_metric_fires_on_served_path() {
use super::super::control_plane::subscriber::{self, LastSeen, SubscriberHandle};
use crate::runtime::config::UdbConfig;
use crate::runtime::metrics::{MetricsRecorder, PrometheusMetrics};
use std::sync::Arc;
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let metrics = Arc::new(PrometheusMetrics::new().expect("build prometheus metrics"));
let metrics_sink: Arc<dyn MetricsRecorder> = metrics.clone();
let handle = SubscriberHandle::new(pool.clone(), Arc::new(UdbConfig::default()))
.with_metrics(metrics_sink);
let baseline = subscriber::run_once(&handle, &LastSeen::default(), false)
.await
.expect("seed subscriber baseline");
let seeded_text = metrics.gather_text("");
assert!(
!seeded_text.contains("udb_control_reload_applied_total"),
"the seeding pass must NOT count a reload:\n{seeded_text}"
);
let rt = ResourceType::RoutingPolicy;
let name = format!("served-reload-{}", Uuid::new_v4().simple());
store::upsert_resource(
&pool,
rt,
&name,
"",
"billing",
r#"{"route":"primary"}"#,
"metric-test",
)
.await
.expect("publish routing policy");
let _next = subscriber::run_once(&handle, &baseline, true)
.await
.expect("served reload tick");
let text = metrics.gather_text("");
assert!(
text.lines().any(|line| {
line.starts_with("udb_control_reload_applied_total{")
&& line.contains("RESOURCE_TYPE_ROUTING_POLICY")
}),
"the served reload path must increment udb_control_reload_applied_total \
for the changed type:\n{text}"
);
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_control_rollback_resources_restores_retained_snapshot -- --ignored --nocapture"]
async fn live_postgres_control_rollback_resources_restores_retained_snapshot() {
use super::super::control_plane::ControlPlaneServiceImpl;
use crate::proto::udb::core::control::services::v1 as control_pb;
use crate::proto::udb::core::control::services::v1::control_plane_service_server::ControlPlaneService;
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let rt = ResourceType::RoutingPolicy;
let node = format!("rollback-node-{}", Uuid::new_v4().simple());
let name = format!("rollback-routing-{}", Uuid::new_v4().simple());
let v1_payload = r#"{"route":"primary","weight":100}"#;
store::upsert_resource(&pool, rt, &name, "", "billing", v1_payload, "rb-test")
.await
.expect("publish v1");
store::ensure_node_state(&pool, &node, rt, &[name.clone()])
.await
.expect("seed node state");
let v1_resources = store::list_resources(&pool, rt, None, &[name.clone()])
.await
.expect("list v1 resources");
let world_v1 = store::world_version(&pool, rt, None, &[name.clone()])
.await
.expect("world v1");
store::retain_served_snapshot(&pool, &node, rt, &world_v1, &v1_resources)
.await
.expect("retain v1 snapshot");
let v2_payload = r#"{"route":"","weight":-1}"#;
store::upsert_resource(&pool, rt, &name, "", "billing", v2_payload, "rb-test")
.await
.expect("publish v2");
let v2_resources = store::list_resources(&pool, rt, None, &[name.clone()])
.await
.expect("list v2 resources");
let world_v2 = store::world_version(&pool, rt, None, &[name.clone()])
.await
.expect("world v2");
store::retain_served_snapshot(&pool, &node, rt, &world_v2, &v2_resources)
.await
.expect("retain v2 snapshot");
assert_ne!(world_v1, world_v2, "v2 must change the world version");
let svc = ControlPlaneServiceImpl::new().with_postgres(Some(pool.clone()));
let resp = svc
.rollback_resources(Request::new(control_pb::RollbackResourcesRequest {
node_id: node.clone(),
resource_type: rt as i32,
target_version: world_v1.clone(),
..Default::default()
}))
.await
.expect("rollback to retained v1")
.into_inner();
assert_eq!(resp.rolled_back_to_version, world_v1);
assert_eq!(resp.resources_restored, v1_resources.len() as i32);
let world_after = store::world_version(&pool, rt, None, &[name.clone()])
.await
.expect("world after rollback");
assert_eq!(
world_after, world_v1,
"rollback must restore the retained v1 world version"
);
let err = svc
.rollback_resources(Request::new(control_pb::RollbackResourcesRequest {
node_id: node.clone(),
resource_type: rt as i32,
target_version: "no-such-retained-version".to_string(),
..Default::default()
}))
.await
.expect_err("an unretained target must fail closed");
assert_eq!(err.code(), tonic::Code::FailedPrecondition);
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_ha_apikey_revocation_propagation -- --ignored --nocapture"]
async fn live_postgres_ha_apikey_revocation_propagation() {
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let authn = authn_service(pool.clone());
let node1 = api_key_service(pool.clone());
let node2 = api_key_service(pool.clone());
let (owner, _) =
create_service_account_with_grant(&authn, "ha_apikey", "CorrectHorse1!", &["data:read"])
.await;
let principal_id = owner.user_id;
let created = node1
.create_api_key(Request::new(apikey_pb::CreateApiKeyRequest {
name: "ha-key".to_string(),
owner_id: principal_id.clone(),
scopes: vec!["data:read".to_string()],
context: Some(common_pb::RequestContext {
principal_id: principal_id.clone(),
tenant: Some(common_pb::TenantContext {
tenant_id: "acme".to_string(),
project_id: "billing".to_string(),
..Default::default()
}),
..Default::default()
}),
..Default::default()
}))
.await
.expect("node1 create API key")
.into_inner();
let key_id = created.key.expect("created key").key_id;
let valid = node2
.validate_api_key(Request::new(apikey_pb::ValidateApiKeyRequest {
plain_key: created.plain_key.clone(),
required_scope: "data:read".to_string(),
..Default::default()
}))
.await
.expect("node2 validate live key")
.into_inner();
assert!(valid.valid, "key created on node1 must validate on node2");
scope_claim_context_for_test(
test_claim_context("live-admin", "acme", "billing", &[], &[]),
node1.revoke_api_key(Request::new(apikey_pb::RevokeApiKeyRequest {
key_id,
revoke_reason: "ha_multinode_test".to_string(),
..Default::default()
})),
)
.await
.expect("node1 revoke API key");
let after = node2
.validate_api_key(Request::new(apikey_pb::ValidateApiKeyRequest {
plain_key: created.plain_key,
required_scope: "data:read".to_string(),
..Default::default()
}))
.await
.expect("node2 validate after revoke")
.into_inner();
assert!(
!after.valid,
"API-key revocation on node1 must propagate to node2 (durable read)"
);
cleanup_native_auth_db(&pool).await;
}
#[tokio::test]
#[ignore = "requires live Postgres; run with UDB_LIVE_AUTH_TESTS=1 cargo test --lib live_postgres_ha_refresh_token_replay_race -- --ignored --nocapture"]
async fn live_postgres_ha_refresh_token_replay_race() {
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
let node1 = authn_service_with_jwt(pool.clone());
let node2 = authn_service_with_jwt(pool.clone());
let user = create_verified_user(&node1, "ha-rtrace", "CorrectHorse1!").await;
let login = node1
.login(Request::new(authn_pb::LoginRequest {
username: user.email.clone(),
password: "CorrectHorse1!".to_string(),
device_name: "ha-rtf".to_string(),
..Default::default()
}))
.await
.expect("node1 login")
.into_inner();
assert!(
login.refresh_token.starts_with("rt_"),
"login must issue a token-family refresh token"
);
let rotated = node1
.refresh_token(Request::new(authn_pb::RefreshTokenRequest {
refresh_token: login.refresh_token.clone(),
..Default::default()
}))
.await
.expect("node1 rotate")
.into_inner();
assert!(rotated.refresh_token.starts_with("rt_"));
assert_ne!(rotated.refresh_token, login.refresh_token);
let reuse = node2
.refresh_token(Request::new(authn_pb::RefreshTokenRequest {
refresh_token: login.refresh_token.clone(),
..Default::default()
}))
.await
.expect_err("replaying the rotated-away token on node2 must be rejected as reuse");
assert_eq!(reuse.code(), tonic::Code::Unauthenticated);
let on_node1 = node1
.refresh_token(Request::new(authn_pb::RefreshTokenRequest {
refresh_token: rotated.refresh_token.clone(),
..Default::default()
}))
.await
.expect_err("rotated token must be dead on node1 once the family is revoked");
assert_eq!(on_node1.code(), tonic::Code::Unauthenticated);
let on_node2 = node2
.refresh_token(Request::new(authn_pb::RefreshTokenRequest {
refresh_token: rotated.refresh_token,
..Default::default()
}))
.await
.expect_err("rotated token must be dead on node2 once the family is revoked");
assert_eq!(on_node2.code(), tonic::Code::Unauthenticated);
cleanup_native_auth_db(&pool).await;
}
#[cfg(all(feature = "kafka", feature = "redis"))]
#[tokio::test]
#[ignore = "requires live Postgres + Kafka + Redis. UDB_LIVE_AUTH_TESTS=1 \
UDB_INTEGRATION_KAFKA_BROKERS=localhost:59192 \
UDB_INTEGRATION_REDIS_URL=redis://127.0.0.1:56379 \
cargo test --lib live_postgres_ha_cdc_idempotent_double_process -- --ignored --nocapture"]
async fn live_postgres_ha_cdc_idempotent_double_process() {
use crate::runtime::cdc::{CdcConfig, CdcEngine};
use crate::runtime::metrics::NoopMetrics;
use std::sync::Arc;
let _guard = live_auth_db_lock().lock().await;
let pool = live_pg_pool().await;
migrate_native_auth_db(&pool).await;
crate::runtime::system::ensure_system_catalog(&pool)
.await
.expect("ensure live UDB system catalog");
ensure_outbox_table(&pool).await;
let brokers = std::env::var("UDB_INTEGRATION_KAFKA_BROKERS")
.or_else(|_| std::env::var("UDB_KAFKA_BROKERS"))
.unwrap_or_else(|_| "localhost:59192".to_string());
let redis_url = std::env::var("UDB_LIVE_AUTH_REDIS_URL")
.or_else(|_| std::env::var("UDB_INTEGRATION_REDIS_URL"))
.unwrap_or_else(|_| "redis://127.0.0.1:56379".to_string());
let topic = "udb.authn.user.registered.v1";
let dlq_topic = format!("udb.cdc.dlq.{}.v1", Uuid::new_v4().simple());
ensure_kafka_topic(&brokers, topic).await;
ensure_kafka_topic(&brokers, &dlq_topic).await;
let config = CdcConfig {
outbox_table: "outbox_events".to_string(),
dlq_topic,
..CdcConfig::default()
};
let engine = CdcEngine::new(
pool.clone(),
None,
&brokers,
live_pg_dsn(),
Arc::new(NoopMetrics),
config,
)
.expect("build CDC engine");
let redis_client = redis::Client::open(redis_url.as_str()).expect("redis client");
let mut redis_conn = redis_client
.get_multiplexed_async_connection()
.await
.expect("connect redis");
let event_id = Uuid::new_v4();
let partition_key = Uuid::new_v4().to_string();
let payload = serde_json::json!({
"event_id": event_id.to_string(),
"event_type": topic,
"correlation_id": format!("cdc-idem:{event_id}"),
"document_id": partition_key,
"tenant_id": "acme",
"timestamp": chrono::Utc::now().to_rfc3339(),
"payload": {"user_id": partition_key}
});
insert_outbox_for_idempotency(&pool, event_id, topic, &partition_key, payload.clone()).await;
engine
.process_outbox_event(
event_id,
topic.to_string(),
partition_key.clone(),
payload.clone(),
chrono::Utc::now(),
101,
Some(&mut redis_conn),
)
.await;
assert_eq!(
outbox_event_count(&pool, event_id).await,
0,
"the first pass must ack-and-delete the outbox row after publishing"
);
insert_outbox_for_idempotency(&pool, event_id, topic, &partition_key, payload.clone()).await;
engine
.process_outbox_event(
event_id,
topic.to_string(),
partition_key,
payload,
chrono::Utc::now(),
102,
Some(&mut redis_conn),
)
.await;
assert_eq!(
outbox_event_count(&pool, event_id).await,
0,
"the second (duplicate) pass must be a durable-publish-skip + ack, not a republish"
);
cleanup_native_auth_db(&pool).await;
}
#[cfg(all(feature = "kafka", feature = "redis"))]
async fn insert_outbox_for_idempotency(
pool: &sqlx::PgPool,
event_id: Uuid,
topic: &str,
partition_key: &str,
payload: serde_json::Value,
) {
sqlx::query(
"INSERT INTO udb_system.outbox_events (event_id, topic, partition_key, payload, created_at) \
VALUES ($1, $2, $3, $4::JSONB, NOW()) ON CONFLICT (event_id) DO NOTHING",
)
.bind(event_id)
.bind(topic)
.bind(partition_key)
.bind(payload)
.execute(pool)
.await
.expect("insert idempotency outbox event");
}
#[cfg(all(feature = "kafka", feature = "redis"))]
async fn outbox_event_count(pool: &sqlx::PgPool, event_id: Uuid) -> i64 {
sqlx::query_scalar("SELECT COUNT(*) FROM udb_system.outbox_events WHERE event_id = $1")
.bind(event_id)
.fetch_one(pool)
.await
.expect("count outbox event")
}
#[test]
fn native_service_enable_disable_migration_matrix_is_coherent() {
use crate::runtime::config::{NativeServiceConfig, NativeServicesSettings, UdbConfig};
use crate::runtime::service::native_registry::{
native_service_ids, resolved_native_service_statuses,
};
let service_ids = native_service_ids();
assert!(
!service_ids.is_empty(),
"the descriptor must yield at least one native service id"
);
let resolve_one = |service_id: &str, override_cfg: NativeServiceConfig| {
let config = UdbConfig {
native_services: NativeServicesSettings {
enabled: true,
default_enabled: true,
migrate_enabled: true,
control_plane_enabled: true,
webrtc_peer_enabled: true,
services: [(service_id.to_string(), override_cfg)]
.into_iter()
.collect(),
..NativeServicesSettings::default()
},
..UdbConfig::default()
};
resolved_native_service_statuses(&config)
.into_iter()
.find(|status| status.service_id == service_id)
.unwrap_or_else(|| panic!("missing status for {service_id}"))
};
for service_id in &service_ids {
let disabled = resolve_one(
service_id,
NativeServiceConfig {
enabled: Some(false),
migrate: Some(false),
..NativeServiceConfig::default()
},
);
assert!(
!disabled.enabled,
"{service_id}: disabled must report enabled=false"
);
assert!(
!disabled.mounted,
"{service_id}: disabled must not be mounted"
);
assert!(
!disabled.migration_enabled,
"{service_id}: disabled must not migrate"
);
assert_eq!(
disabled.migration_status, "disabled",
"{service_id}: disabled migration_status must be 'disabled'"
);
assert!(
!disabled.disabled_reason.is_empty(),
"{service_id}: a disabled service must report WHY"
);
let enabled = resolve_one(
service_id,
NativeServiceConfig {
enabled: Some(true),
migrate: Some(true),
..NativeServiceConfig::default()
},
);
assert!(
enabled.enabled,
"{service_id}: enabled must report enabled=true"
);
if enabled.mounted {
assert!(
enabled.missing_dependencies.is_empty(),
"{service_id}: a mounted service must have no missing dependencies"
);
assert!(
!enabled.surface.is_empty() && !enabled.listener_kind.is_empty(),
"{service_id}: a mounted service must report surface + listener_kind"
);
assert!(
enabled.healthy,
"{service_id}: a mounted, non-degraded service must be healthy"
);
} else {
assert!(
!enabled.disabled_reason.is_empty() || !enabled.missing_dependencies.is_empty(),
"{service_id}: an enabled-but-unmounted service must report a reason"
);
}
assert_eq!(
enabled.background_worker_enabled,
enabled.mounted && enabled.owns_background_workers,
"{service_id}: worker enablement must follow mounted worker-owner status"
);
let migrate_only = resolve_one(
service_id,
NativeServiceConfig {
enabled: Some(false),
migrate: Some(true),
..NativeServiceConfig::default()
},
);
assert!(
!migrate_only.enabled,
"{service_id}: migrate-only must not serve"
);
assert!(
!migrate_only.mounted,
"{service_id}: migrate-only must not mount"
);
assert!(
migrate_only.migration_enabled,
"{service_id}: migrate-only must report migration_enabled"
);
assert_eq!(
migrate_only.migration_status, "enabled",
"{service_id}: migrate-only migration_status must be 'enabled'"
);
}
let killed = UdbConfig {
native_services: NativeServicesSettings {
enabled: false,
..NativeServicesSettings::default()
},
..UdbConfig::default()
};
for status in resolved_native_service_statuses(&killed) {
assert!(
!status.enabled,
"global disable must disable {}",
status.service_id
);
assert!(
!status.mounted,
"global disable must unmount {}",
status.service_id
);
}
}