udb 0.3.6

Universal Data Broker — a Rust gRPC broker over multiple databases (Postgres, MySQL, SQLite, MongoDB, ClickHouse, Cassandra, MSSQL, Redis, Qdrant, S3, Neo4j, …) with per-tenant RLS, 2PC, sagas, and CDC.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
//! Phase 13 (Production Readiness) multi-node / cluster-convergence live tests.
//!
//! Extends [`super::ha_convergence_live`] (HA-1 session revocation, HA-2 logout-all)
//! with the remaining scenarios from `docs/auth-ha-test-plan.md` plus the Phase-9
//! control-plane distribution (ACK/NACK/rollback) and CDC idempotency convergence.
//!
//! Topology, exactly as the convergence suite: two independent service instances
//! (`node1`/`node2`) share ONE live Postgres pool to simulate two cluster nodes,
//! each with its own in-process caches over the shared durable store. A mutation
//! performed on `node1` must converge on `node2` (the node that did NOT perform it)
//! within the configured window — the property single-instance tests cannot prove.
//!
//! Every test is env-gated like the convergence suite (`#[ignore]`; run with
//! `UDB_LIVE_AUTH_TESTS=1 cargo test --lib <name> -- --ignored --nocapture`) and
//! self-cleans via `cleanup_native_auth_db`. No hardcoded schema/table names: the
//! control-plane store resolves every identifier through the proto-driven
//! `native_model`, and the auth services through the native catalog DDL.

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;

/// Snapshot TTL (seconds) used for the cross-node authz convergence tests so a
/// non-mutating node's per-instance snapshot cache expires within the test
/// window. `authz_snapshot_ttl()` reads `UDB_AUTHZ_SNAPSHOT_TTL_SECS` at
/// construction and floors it at 1s, so we set it before building `node2`.
const TEST_SNAPSHOT_TTL_SECS: u64 = 1;

fn set_short_snapshot_ttl() {
    // SAFETY: tests in this module run serially under `live_auth_db_lock`; the env
    // is read only at `AuthzServiceImpl` construction below.
    unsafe {
        std::env::set_var(
            "UDB_AUTHZ_SNAPSHOT_TTL_SECS",
            TEST_SNAPSHOT_TTL_SECS.to_string(),
        );
    }
}

/// Build a CheckAccess request for a role-bound principal on (tenant, project).
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
}

/// Seed a user + a role + a role binding + an allow policy so a principal is
/// authorized for `(object, action)` on both nodes. Returns the user id and the
/// role code the policy is keyed on.
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)
}

/// HA-3: authz revision invalidation propagates across nodes. `node1` mutates a
/// policy (adds an explicit DENY that overrides the seeded ALLOW); `node2` — which
/// performed no mutation, so its per-instance snapshot cache was NOT invalidated —
/// returns the NEW decision once its snapshot TTL elapses, and `GetAuthzRevision`
/// (a durable read) reports a changed content-hash `policy_version`/revision.
#[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;

    // Baseline: both nodes ALLOW. node2 reads through to the shared store on its
    // first call (its cache was seeded empty + invalidated by `with_postgres`).
    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)"
    );

    // Capture node2's durable revision/content-hash BEFORE the mutation.
    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 adds an explicit DENY for the same (role, action, resource). An
    // explicit deny overrides the role-bound allow.
    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");

    // The durable revision must have advanced; node2 sees it immediately because
    // `GetAuthzRevision` is a durable read (not the cached snapshot).
    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
    );
    // TODO(host): assert the revision-row content_hash CHANGES once put_authz_policy
    // feeds a payload-derived hash into bump_authz_revision. Today it passes the
    // constant "policy-put" (see authz/mod.rs ~L1307), so the AuthzRevision.content_hash
    // column is identical across policy puts and cannot distinguish them. The
    // content-addressed value that DOES change on every policy edit is the snapshot
    // version (`policy_content_version`), asserted indirectly via node2's decision
    // flip below; the durable monotonic `policy_revision` counter (asserted above)
    // is the reliable cross-node invalidation signal in the meantime.
    assert!(
        !after.content_hash.is_empty(),
        "node2's revision content_hash must be populated"
    );

    // node2's DECISION must flip to DENY within its snapshot TTL window. The cache
    // was last loaded before the mutation, so wait out the TTL then re-check.
    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;
}

/// Policy-bundle revocation: `node1` invalidates the cached bundles for a tenant
/// via the REAL `InvalidatePolicyBundles` governance RPC (driven with an
/// `authz:admin` governance actor). The revision bump is durable, so `node2`'s
/// `GetAuthzRevision` reflects the new revision — every SDK bundle pinned below it
/// is now stale. We also assert the bundle a node would re-issue carries the new
/// policy_version (bundle signing requires `UDB_POLICY_BUNDLE_SECRET`; without it
/// the signing precondition is asserted instead so the test never fakes a pass).
#[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;

    // A policy must exist so a revision row exists to bump.
    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();

    // node1 revokes/invalidates every cached bundle for the tenant via the real
    // governance RPC. An `authz:admin`-scoped actor satisfies the standing-scope
    // gate (and the snapshot governance check permits when no governance policy
    // is configured), so this exercises `invalidate_policy_bundles_impl` end to
    // end — the revision bump + the durable bundle-revocation record.
    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
    );

    // node2 reflects the new revision durably: any bundle pinned to `before` is
    // now stale (its policy_revision < node2's reported 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
    );

    // A bundle node2 would re-issue must carry the new revision (or, if bundle
    // signing isn't configured, the signing precondition must be reported — we
    // never fake a signed-bundle pass).
    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) => {
            // Bundle signing not configured in this environment: assert the exact
            // precondition rather than passing silently.
            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;
}

/// Policy distribution ACK/NACK + rollback against the REAL Postgres-backed
/// control-plane store (`control_plane::store`), not the in-memory mirror in
/// `control_plane/tests.rs`. Exercises the full acceptance property "a node can
/// reject a bad policy without silently diverging":
///   1. push v1 (a routing policy) and ACK it → `accepted_version` advances;
///   2. push a bad v2 and NACK it → `accepted_version` stays at v1 (no divergence),
///      `last_good_version` is preserved, the NACK error is recorded;
///   3. rollback the registry to the v1 content → a fresh push+ACK restores the
///      node onto the good world version.
/// Two nodes share the registry; node B's NACK must not advance node A.
#[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());

    // --- v1: a good routing policy is published to the registry. ---
    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));

    // The version a node ACKs is the aggregate world version of the type it is
    // subscribed to (matches the DiscoveryResponse the stream would carry).
    let world_v1 = store::world_version(&pool, rt, None, &[])
        .await
        .expect("world version v1");

    // node A subscribes and ACKs v1 under the nonce the control plane stamped.
    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());

    // A stale-nonce ACK must be IGNORED (no silent divergence on a stale echo).
    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"
    );

    // --- v2: a BAD routing policy is published; node A NACKs it. ---
    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"
    );
    // INVARIANT: node A keeps running v1 — accepted/last_good are preserved, only
    // the NACK error is set. The bad v2 did NOT silently take effect.
    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"
    );

    // A second node B NACKing the same bad v2 must not advance node A either
    // (per-node ledgers; one node's reject cannot move another's accepted state).
    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"
    );

    // --- rollback: restore the v1 content; a fresh push+ACK returns to good. ---
    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"
    );

    // Cross-check the aggregate-version helper agrees with the store's world view.
    let resources = store::list_resources(&pool, rt, None, &[name.clone()])
        .await
        .expect("list resources");
    assert_eq!(aggregate_version(&resources), world_rb);

    // Self-clean: the control-plane registry + node-state rows live in the native
    // schemas dropped by cleanup_native_auth_db.
    cleanup_native_auth_db(&pool).await;
}

/// API-key revocation propagation: `node1` revokes an API key; `node2`'s
/// `ValidateApiKey` (which did not create or revoke it) denies it immediately,
/// because revocation is a durable store read on the shared pool — the API-key
/// analogue of HA-1's session revocation.
#[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 node1 = api_key_service(pool.clone());
    let node2 = api_key_service(pool.clone());
    let principal_id = Uuid::new_v4().to_string();

    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;

    // node2 validates the key created on node1 (shared durable store).
    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");

    // node1 revokes it — as the authenticated tenant-`acme` admin (the validated
    // claim the transport layer installs over the wire); the guard stays strict.
    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");

    // node2 must now deny it — durable revocation, no per-node cache to lag.
    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;
}

/// Refresh-token replay race across two instances (HA-2 in the plan): `node1`
/// rotates a refresh token; replaying the ORIGINAL (rotated-away) token on
/// `node2` is reuse — `node2` detects it (the token-family ledger is durable),
/// revokes the family, and the legitimately rotated token then fails on BOTH
/// nodes. Single-instance reuse is covered by `authn_jwt_live`; this proves the
/// detection holds when the rotate and the replay land on different nodes.
#[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"
    );

    // node1 rotates the refresh token (the original is now rotated-away).
    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);

    // node2 (which did NOT rotate) sees the ORIGINAL token replayed: reuse. The
    // family ledger is durable, so node2 detects it and rejects the replay.
    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);

    // Reuse detection revoked the whole family: the legitimately rotated token now
    // fails on BOTH nodes (durable family revocation converged).
    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;
}

/// CDC idempotency: double-processing the same outbox event publishes it ONCE.
/// Drives the REAL CDC idempotency path (`process_outbox_event` →
/// `was_durably_published` consulted via the durable CDC journal): the first pass
/// publishes + journals + acks-and-deletes the outbox row; the second pass, with
/// the same event_id re-enqueued and a Redis idempotency claim already present,
/// is recognized as durably published and acked WITHOUT a second publish.
/// Needs Postgres + Kafka + Redis (the idempotency `SET NX` fast-path only fires
/// with a Redis connection); gated behind the `kafka` + `redis` features.
#[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;

    // First pass: publishes to Kafka, records the durable journal entry, and
    // acks-and-deletes the outbox row.
    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"
    );

    // Re-enqueue the SAME event_id (simulates a redelivery / replay), then process
    // again with the Redis idempotency claim already set. The engine must consult
    // the durable journal (`was_durably_published`), recognize it as published,
    // and ack WITHOUT a second publish — the outbox row is removed, not republished.
    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")
}

/// Native service enable/disable/migration matrix. Single-instance, config-only
/// (no DB): for every descriptor-derived native service, drive the enabled /
/// disabled / serving-only / migration-only toggles through the real
/// `resolved_native_service_statuses` resolver and assert the mount/health/
/// migration status stays internally coherent (the property a node's health
/// report and its actually-mounted routes must agree on). No hardcoded service
/// id list: the matrix iterates `native_service_ids()` (descriptor-sourced).
#[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"
    );

    // Helper: resolve the status for one service under a given per-service override.
    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 {
        // 1. Explicitly DISABLED: never mounted, never migrating, never healthy.
        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"
        );

        // 2. ENABLED (serve + migrate). Coherence invariants regardless of whether
        //    a backend dependency happens to be missing in this config.
        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"
        );
        // mounted iff enabled AND its listener is on AND no missing dependency.
        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 {
            // Not mounted while enabled => degraded (missing dep) or listener off,
            // and it must say why.
            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"
        );

        // 3. MIGRATION-ONLY (serve disabled, migrate enabled) — distinct states:
        //    a service can own its schema migration while refusing RPC traffic.
        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'"
        );
    }

    // 4. GLOBAL kill switch overrides everything: nothing serves or migrates.
    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
        );
    }
}