diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 205517393..c5f46ccd4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -699,7 +699,7 @@ jobs: -c "CREATE DATABASE buzz_identity_tests" cargo nextest run \ --archive-file target/ci/backend-integration-tests.tar.zst \ - -E '(package(buzz-db) and test(/identity_binding::tests/)) or (package(buzz-relay) and test(/corporate_identity::tests/))' \ + -E '(package(buzz-db) and test(/identity_binding::tests/)) or (package(buzz-db) and test(/identity_lifecycle::tests|identity_lifecycle::deterministic_tests/)) or (package(buzz-db) and test(/migration::tests::identity_0029_|migration::deterministic_tests::identity_0029_/)) or (package(buzz-relay) and test(/corporate_identity::tests/))' \ --test-threads 1 \ --run-ignored ignored-only env: diff --git a/crates/buzz-db/src/identity_binding.rs b/crates/buzz-db/src/identity_binding.rs index d33a417e5..7e2cfa9b2 100644 --- a/crates/buzz-db/src/identity_binding.rs +++ b/crates/buzz-db/src/identity_binding.rs @@ -16,6 +16,312 @@ use uuid::Uuid; use crate::error::{DbError, Result}; use buzz_core::CommunityId; +#[cfg(test)] +pub(crate) mod test_lock_schedule { + use std::cell::Cell; + use std::future::Future; + use std::sync::{Mutex, OnceLock}; + + use sqlx::{Postgres, Transaction}; + use tokio::sync::{mpsc, oneshot}; + + tokio::task_local! { + static ACTOR: &'static str; + static ROW_REQUEST_REPORTED: Cell; + static ROW_ACQUIRED_REPORTED: Cell; + } + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub(crate) enum LockPhase { + Request, + Acquired, + } + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub(crate) enum RowLockPhase { + Request, + Acquired, + } + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub(crate) struct AdvisoryLockKey { + class_id: u32, + object_id: u32, + } + + impl AdvisoryLockKey { + pub(crate) const fn class_id(self) -> u32 { + self.class_id + } + + pub(crate) const fn object_id(self) -> u32 { + self.object_id + } + } + + pub(crate) struct LockEvent { + actor: &'static str, + phase: LockPhase, + isolation: Option, + transaction_id: Option, + backend_pid: i32, + database_oid: u32, + lock_keys: Vec, + resume: oneshot::Sender<()>, + } + + impl LockEvent { + pub(crate) const fn actor(&self) -> &'static str { + self.actor + } + + pub(crate) const fn phase(&self) -> LockPhase { + self.phase + } + + pub(crate) fn isolation(&self) -> Option<&str> { + self.isolation.as_deref() + } + + pub(crate) const fn transaction_id(&self) -> Option { + self.transaction_id + } + + pub(crate) const fn backend_pid(&self) -> i32 { + self.backend_pid + } + + pub(crate) const fn database_oid(&self) -> u32 { + self.database_oid + } + + pub(crate) fn lock_keys(&self) -> &[AdvisoryLockKey] { + &self.lock_keys + } + + pub(crate) const fn coordinate_count(&self) -> usize { + self.lock_keys.len() + } + + pub(crate) fn resume(self) { + let _ = self.resume.send(()); + } + } + + fn controller() -> &'static Mutex>> { + static CONTROLLER: OnceLock>>> = + OnceLock::new(); + CONTROLLER.get_or_init(|| Mutex::new(None)) + } + + pub(crate) struct ControllerGuard; + + impl Drop for ControllerGuard { + fn drop(&mut self) { + *controller().lock().expect("lock test controller") = None; + } + } + + pub(crate) fn install() -> (mpsc::UnboundedReceiver, ControllerGuard) { + let (sender, receiver) = mpsc::unbounded_channel(); + let mut current = controller().lock().expect("lock test controller"); + assert!( + current.is_none(), + "only one deterministic lock controller may be active" + ); + *current = Some(sender); + (receiver, ControllerGuard) + } + + pub(crate) struct RowLockEvent { + actor: &'static str, + phase: RowLockPhase, + transaction_id: i64, + backend_pid: i32, + database_oid: u32, + resume: oneshot::Sender<()>, + } + + impl RowLockEvent { + pub(crate) const fn actor(&self) -> &'static str { + self.actor + } + + pub(crate) const fn phase(&self) -> RowLockPhase { + self.phase + } + + pub(crate) const fn transaction_id(&self) -> i64 { + self.transaction_id + } + + pub(crate) const fn backend_pid(&self) -> i32 { + self.backend_pid + } + + pub(crate) const fn database_oid(&self) -> u32 { + self.database_oid + } + + pub(crate) fn resume(self) { + let _ = self.resume.send(()); + } + } + + fn row_controller() -> &'static Mutex>> { + static CONTROLLER: OnceLock>>> = + OnceLock::new(); + CONTROLLER.get_or_init(|| Mutex::new(None)) + } + + pub(crate) struct RowControllerGuard; + + impl Drop for RowControllerGuard { + fn drop(&mut self) { + *row_controller().lock().expect("lock row test controller") = None; + } + } + + pub(crate) fn install_row() -> (mpsc::UnboundedReceiver, RowControllerGuard) { + let (sender, receiver) = mpsc::unbounded_channel(); + let mut current = row_controller().lock().expect("lock row test controller"); + assert!( + current.is_none(), + "only one deterministic row-lock controller may be active" + ); + *current = Some(sender); + (receiver, RowControllerGuard) + } + + pub(crate) async fn actor_scope(actor: &'static str, future: F) -> F::Output + where + F: Future, + { + ACTOR + .scope( + actor, + ROW_REQUEST_REPORTED.scope( + Cell::new(false), + ROW_ACQUIRED_REPORTED.scope(Cell::new(false), future), + ), + ) + .await + } + + pub(super) async fn checkpoint( + tx: &mut Transaction<'_, Postgres>, + phase: LockPhase, + coordinates: &[Vec], + ) { + let Ok(actor) = ACTOR.try_with(|actor| *actor) else { + return; + }; + let sender = controller().lock().expect("lock test controller").clone(); + let Some(sender) = sender else { + return; + }; + let (backend_pid, database_oid): (i32, i64) = sqlx::query_as( + "SELECT pg_backend_pid(), oid::BIGINT \ + FROM pg_database WHERE datname=current_database()", + ) + .fetch_one(&mut **tx) + .await + .expect("read test lock backend identity"); + let class_id: i32 = sqlx::query_scalar("SELECT hashtext('buzz_nip_fi_v1')") + .fetch_one(&mut **tx) + .await + .expect("hash test lock namespace"); + let mut lock_keys = Vec::with_capacity(coordinates.len()); + for coordinate in coordinates { + let object_id: i32 = sqlx::query_scalar("SELECT hashtext(encode($1, 'hex'))") + .bind(coordinate.as_slice()) + .fetch_one(&mut **tx) + .await + .expect("hash test lock coordinate"); + lock_keys.push(AdvisoryLockKey { + class_id: class_id as u32, + object_id: object_id as u32, + }); + } + let (isolation, transaction_id) = if phase == LockPhase::Acquired { + let isolation = sqlx::query_scalar("SHOW transaction_isolation") + .fetch_one(&mut **tx) + .await + .ok(); + let transaction_id = sqlx::query_scalar("SELECT txid_current()::BIGINT") + .fetch_one(&mut **tx) + .await + .ok(); + (isolation, transaction_id) + } else { + (None, None) + }; + let (resume, resumed) = oneshot::channel(); + if sender + .send(LockEvent { + actor, + phase, + isolation, + transaction_id, + backend_pid, + database_oid: u32::try_from(database_oid) + .expect("database OID fits the PostgreSQL OID type"), + lock_keys, + resume, + }) + .is_ok() + { + let _ = resumed.await; + } + } + + pub(crate) async fn row_checkpoint(tx: &mut Transaction<'_, Postgres>, phase: RowLockPhase) { + let Ok(actor) = ACTOR.try_with(|actor| *actor) else { + return; + }; + let should_report = match phase { + RowLockPhase::Request => ROW_REQUEST_REPORTED + .try_with(|reported| !reported.replace(true)) + .unwrap_or(false), + RowLockPhase::Acquired => ROW_ACQUIRED_REPORTED + .try_with(|reported| !reported.replace(true)) + .unwrap_or(false), + }; + if !should_report { + return; + } + let sender = row_controller() + .lock() + .expect("lock row test controller") + .clone(); + let Some(sender) = sender else { + return; + }; + let (backend_pid, database_oid, transaction_id): (i32, i64, i64) = sqlx::query_as( + "SELECT pg_backend_pid(), oid::BIGINT, txid_current()::BIGINT \ + FROM pg_database WHERE datname=current_database()", + ) + .fetch_one(&mut **tx) + .await + .expect("read test row-lock backend identity"); + let (resume, resumed) = oneshot::channel(); + if sender + .send(RowLockEvent { + actor, + phase, + transaction_id, + backend_pid, + database_oid: u32::try_from(database_oid) + .expect("database OID fits the PostgreSQL OID type"), + resume, + }) + .is_ok() + { + let _ = resumed.await; + } + } +} + /// Binding source when the IdP JWT carries the pubkey claim. pub const SOURCE_JWT_NPUB: &str = "jwt_npub"; /// Binding source when the relay falls back to the stored uid/pubkey binding. @@ -615,16 +921,20 @@ pub(crate) async fn lock_identity_coordinates_tx( ) -> Result<()> { coordinates.sort(); coordinates.dedup(); - for coordinate in coordinates { + #[cfg(test)] + test_lock_schedule::checkpoint(tx, test_lock_schedule::LockPhase::Request, &coordinates).await; + for coordinate in &coordinates { sqlx::query( "SELECT pg_advisory_xact_lock(\ hashtext('buzz_nip_fi_v1'), hashtext(encode($1, 'hex'))\ )", ) - .bind(coordinate) + .bind(coordinate.as_slice()) .execute(&mut **tx) .await?; } + #[cfg(test)] + test_lock_schedule::checkpoint(tx, test_lock_schedule::LockPhase::Acquired, &coordinates).await; Ok(()) } diff --git a/crates/buzz-db/src/identity_lifecycle.rs b/crates/buzz-db/src/identity_lifecycle.rs index 4adc66cd8..f750ccd70 100644 --- a/crates/buzz-db/src/identity_lifecycle.rs +++ b/crates/buzz-db/src/identity_lifecycle.rs @@ -390,6 +390,12 @@ async fn active_principal_tx( community_id: CommunityId, principal: IdentityPrincipal<'_>, ) -> Result> { + #[cfg(test)] + crate::identity_binding::test_lock_schedule::row_checkpoint( + tx, + crate::identity_binding::test_lock_schedule::RowLockPhase::Request, + ) + .await; let row = sqlx::query( "SELECT binding_id, issuer, uid, pubkey, binding_version, binding_provenance \ FROM identity_bindings \ @@ -401,6 +407,12 @@ async fn active_principal_tx( .bind(principal.subject) .fetch_optional(&mut **tx) .await?; + #[cfg(test)] + crate::identity_binding::test_lock_schedule::row_checkpoint( + tx, + crate::identity_binding::test_lock_schedule::RowLockPhase::Acquired, + ) + .await; active_binding_from_row(row.as_ref()) } @@ -409,6 +421,12 @@ async fn active_key_tx( community_id: CommunityId, pubkey: &[u8], ) -> Result> { + #[cfg(test)] + crate::identity_binding::test_lock_schedule::row_checkpoint( + tx, + crate::identity_binding::test_lock_schedule::RowLockPhase::Request, + ) + .await; let row = sqlx::query( "SELECT binding_id, issuer, uid, pubkey, binding_version, binding_provenance \ FROM identity_bindings \ @@ -419,6 +437,12 @@ async fn active_key_tx( .bind(pubkey) .fetch_optional(&mut **tx) .await?; + #[cfg(test)] + crate::identity_binding::test_lock_schedule::row_checkpoint( + tx, + crate::identity_binding::test_lock_schedule::RowLockPhase::Acquired, + ) + .await; active_binding_from_row(row.as_ref()) } @@ -1845,6 +1869,10 @@ pub async fn enable_identity_principal( .await } +#[cfg(test)] +#[path = "identity_lifecycle_deterministic_tests.rs"] +mod deterministic_tests; + #[cfg(test)] mod tests { use std::sync::Arc; diff --git a/crates/buzz-db/src/identity_lifecycle_deterministic_tests.rs b/crates/buzz-db/src/identity_lifecycle_deterministic_tests.rs new file mode 100644 index 000000000..fad602b84 --- /dev/null +++ b/crates/buzz-db/src/identity_lifecycle_deterministic_tests.rs @@ -0,0 +1,735 @@ +use super::*; + +use crate::identity_binding::{ + get_active_identity_binding_by_pubkey, resolve_identity_binding, test_lock_schedule, + BindingDenial, ResolveBindingInput, ResolveBindingResult, +}; +use serde_json::Value; +use sqlx::PgPool; +use uuid::Uuid; + +const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; +const ISSUER: &str = "https://idp.example"; +const SUBJECT: &str = "deterministic-subject"; +const OLD_KEY: [u8; 32] = [31; 32]; +const ENROLL_KEY: [u8; 32] = [32; 32]; +const ROTATE_KEY: [u8; 32] = [33; 32]; +const RECOVER_KEY: [u8; 32] = [34; 32]; +const ENABLE_KEY: [u8; 32] = [35; 32]; +const DOMAIN_B_KEY: [u8; 32] = [201; 32]; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Action { + Enroll, + Retire, + Rotate, + Recover, + Disable, + Revoke, + Enable, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Outcome { + Applied, + Existing, + Denied, + Error, +} + +#[derive(Clone)] +struct Fixture { + community_id: CommunityId, + expected_pending: Option, + enrollment_key: [u8; 32], +} + +async fn setup_pool() -> PgPool { + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| TEST_DB_URL.to_owned()); + let pool = PgPool::connect(&database_url) + .await + .expect("connect deterministic test DB"); + crate::migration::run_migrations(&pool) + .await + .expect("run migrations"); + pool +} + +async fn make_community(pool: &PgPool, label: &str) -> CommunityId { + let id = Uuid::new_v4(); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1,$2)") + .bind(id) + .bind(format!("{label}-{}.example", id.simple())) + .execute(pool) + .await + .expect("insert deterministic community"); + CommunityId::from_uuid(id) +} + +fn principal() -> IdentityPrincipal<'static> { + IdentityPrincipal { + issuer: ISSUER, + subject: SUBJECT, + } +} + +fn context(operation_id: Uuid, reason: &'static str) -> LifecycleContext<'static> { + LifecycleContext { + operation_id, + actor: None, + reason, + } +} + +fn replacement(pubkey: &'static [u8; 32]) -> VerifiedReplacementKey<'static> { + VerifiedReplacementKey::after_verified_proof( + pubkey, + None, + BindingProvenance::AttestedKey, + Some("deterministic-policy-v1"), + ) + .expect("construct deterministic replacement") +} + +async fn enroll_key(pool: &PgPool, community_id: CommunityId, pubkey: &[u8]) { + let result = resolve_identity_binding( + pool, + community_id, + &ResolveBindingInput { + issuer: ISSUER, + subject: SUBJECT, + pubkey, + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("seed enrollment"); + assert!(matches!(result, ResolveBindingResult::Enrolled(_))); +} + +fn pair_contains(pair: (Action, Action), action: Action) -> bool { + pair.0 == action || pair.1 == action +} + +async fn setup_fixture(pool: &PgPool, pair: (Action, Action), label: &str) -> Fixture { + let community_id = make_community(pool, label).await; + let mut expected_pending = None; + let mut enrollment_key = ENROLL_KEY; + + if pair_contains(pair, Action::Enroll) && pair_contains(pair, Action::Rotate) { + enrollment_key = OLD_KEY; + } else if pair_contains(pair, Action::Enroll) && pair_contains(pair, Action::Recover) { + enroll_key(pool, community_id, &OLD_KEY).await; + retire_identity_pair( + pool, + community_id, + context(Uuid::from_u128(10), "prepare enrollment recovery"), + principal(), + &OLD_KEY, + ) + .await + .expect("prepare pending recovery"); + expected_pending = get_pending_lineage(pool, community_id, principal()) + .await + .expect("read recovery selector"); + enrollment_key = RECOVER_KEY; + } else { + enroll_key(pool, community_id, &OLD_KEY).await; + let needs_enabled_pending = + pair_contains(pair, Action::Recover) && !pair_contains(pair, Action::Enable); + let needs_disabled_pending = pair_contains(pair, Action::Enable); + if needs_enabled_pending { + retire_identity_pair( + pool, + community_id, + context(Uuid::from_u128(11), "prepare enabled pending"), + principal(), + &OLD_KEY, + ) + .await + .expect("prepare enabled pending"); + } else if needs_disabled_pending { + disable_identity_principal( + pool, + community_id, + context(Uuid::from_u128(12), "prepare disabled pending"), + principal(), + ) + .await + .expect("prepare disabled pending"); + } + expected_pending = get_pending_lineage(pool, community_id, principal()) + .await + .expect("read prepared selector"); + } + + Fixture { + community_id, + expected_pending, + enrollment_key, + } +} + +async fn run_action( + pool: &PgPool, + fixture: &Fixture, + action: Action, + operation_id: Uuid, +) -> Outcome { + match action { + Action::Enroll => match resolve_identity_binding( + pool, + fixture.community_id, + &ResolveBindingInput { + issuer: ISSUER, + subject: SUBJECT, + pubkey: &fixture.enrollment_key, + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + { + Ok(ResolveBindingResult::Enrolled(_)) => Outcome::Applied, + Ok(ResolveBindingResult::Existing(_)) => Outcome::Existing, + Ok(ResolveBindingResult::Denied( + BindingDenial::Conflict + | BindingDenial::Revoked + | BindingDenial::BindingRequired + | BindingDenial::KeyAttestationRequired, + )) => Outcome::Denied, + Err(_) => Outcome::Error, + }, + Action::Rotate => rotate_identity_binding( + pool, + fixture.community_id, + context(operation_id, "deterministic rotate"), + principal(), + &OLD_KEY, + replacement(&ROTATE_KEY), + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied), + Action::Retire => retire_identity_pair( + pool, + fixture.community_id, + context(operation_id, "deterministic retire"), + principal(), + &OLD_KEY, + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied), + Action::Recover => { + let Some(expected) = fixture.expected_pending.as_ref() else { + return Outcome::Error; + }; + recover_identity_binding( + pool, + fixture.community_id, + context(operation_id, "deterministic recover"), + principal(), + expected, + replacement(&RECOVER_KEY), + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied) + } + Action::Disable => disable_identity_principal( + pool, + fixture.community_id, + context(operation_id, "deterministic disable"), + principal(), + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied), + Action::Revoke => revoke_identity_key( + pool, + fixture.community_id, + context(operation_id, "deterministic revoke"), + &OLD_KEY, + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied), + Action::Enable => enable_identity_principal( + pool, + fixture.community_id, + context(operation_id, "deterministic enable"), + principal(), + fixture.expected_pending.as_ref(), + replacement(&ENABLE_KEY), + ) + .await + .map_or(Outcome::Error, |_| Outcome::Applied), + } +} + +async fn normalized_rows( + pool: &PgPool, + community_id: CommunityId, + query: &'static str, +) -> Vec { + let mut rows = sqlx::query_scalar::<_, Value>(query) + .bind(community_id.as_uuid()) + .fetch_all(pool) + .await + .expect("read normalized identity projection") + .into_iter() + .map(|value| value.to_string()) + .collect::>(); + rows.sort(); + rows +} + +async fn logical_snapshot(pool: &PgPool, community_id: CommunityId) -> Vec { + let queries = [ + "SELECT jsonb_build_object('t','binding','v',to_jsonb(b)-ARRAY['community_id','binding_id','replacement_binding_id','created_at','updated_at','last_seen_at','revoked_at','revoked_by','rotation_completed_at','rotation_by']::text[]) FROM identity_bindings b WHERE community_id=$1", + "SELECT jsonb_build_object('t','principal','v',jsonb_build_object('issuer',issuer,'subject',uid,'disabled',disabled_at IS NOT NULL,'reason',disabled_reason)) FROM identity_principals WHERE community_id=$1", + "SELECT jsonb_build_object('t','revoked_key','v',jsonb_build_object('pubkey',encode(pubkey,'hex'),'reason',reason)) FROM identity_revoked_keys WHERE community_id=$1", + "SELECT jsonb_build_object('t','retired','v',to_jsonb(r)-ARRAY['community_id','retired_binding_id','retired_at','retired_by']::text[]) FROM identity_retired_pairs r WHERE community_id=$1", + "SELECT jsonb_build_object('t','pending','v',to_jsonb(p)-ARRAY['community_id','retired_binding_id','created_at','created_operation_id','cleared_at','cleared_operation_id']::text[] || jsonb_build_object('cleared',cleared_at IS NOT NULL)) FROM identity_pending_replacements p WHERE community_id=$1", + "SELECT jsonb_build_object('t','history','v',to_jsonb(h)-ARRAY['community_id','history_id','binding_id','replacement_binding_id','operation_id','recorded_at']::text[]) FROM identity_binding_history h WHERE community_id=$1", + "SELECT jsonb_build_object('t','operation','v',to_jsonb(o)-ARRAY['community_id','operation_id','request_fingerprint','binding_id','replacement_binding_id','created_at']::text[]) FROM identity_lifecycle_operations o WHERE community_id=$1", + "SELECT jsonb_build_object('t','lineage','v',jsonb_build_object('predecessor',encode(p.pubkey,'hex'),'successor',encode(s.pubkey,'hex'))) FROM identity_binding_lineage l JOIN identity_bindings p ON p.community_id=l.community_id AND p.binding_id=l.predecessor_binding_id JOIN identity_bindings s ON s.community_id=l.community_id AND s.binding_id=l.successor_binding_id WHERE l.community_id=$1", + ]; + let mut snapshot = Vec::new(); + for query in queries { + snapshot.extend(normalized_rows(pool, community_id, query).await); + } + snapshot.sort(); + snapshot +} + +async fn raw_domain_snapshot(pool: &PgPool, community_id: CommunityId) -> Vec { + let community: String = sqlx::query_scalar( + "SELECT to_jsonb(community)::text FROM communities community WHERE id=$1", + ) + .bind(community_id.as_uuid()) + .fetch_one(pool) + .await + .expect("read exact domain community sentinel"); + let tables = [ + "identity_bindings", + "identity_principals", + "identity_revoked_keys", + "identity_migration_denials", + "identity_migration_denied_keys", + "identity_binding_lineage", + "identity_retired_pairs", + "identity_pending_replacements", + "identity_binding_history", + "identity_lifecycle_operations", + "audit_log", + ]; + let mut snapshot = vec![format!("communities:{community}")]; + for table in tables { + let query = format!( + "SELECT to_jsonb(row_value)::text FROM {table} row_value WHERE community_id=$1 ORDER BY 1" + ); + let rows = sqlx::query_scalar::<_, String>(sqlx::AssertSqlSafe(query)) + .bind(community_id.as_uuid()) + .fetch_all(pool) + .await + .expect("read exact domain sentinel"); + snapshot.extend(rows.into_iter().map(|row| format!("{table}:{row}"))); + } + snapshot +} + +async fn authorization_sentinel(pool: &PgPool, community_id: CommunityId) -> (Uuid, u64) { + let binding = get_active_identity_binding_by_pubkey(pool, community_id, &DOMAIN_B_KEY) + .await + .expect("read domain-B authorization") + .expect("domain-B binding remains active"); + (binding.binding_id, binding.binding_version) +} + +async fn wait_for_advisory_waiter( + pool: &PgPool, + database_oid: u32, + holder_pid: i32, + waiter_pid: i32, + lock_keys: &[test_lock_schedule::AdvisoryLockKey], +) { + assert_ne!(holder_pid, waiter_pid); + assert!(!lock_keys.is_empty(), "actors must share an identity lock"); + for _ in 0..10_000 { + for lock_key in lock_keys { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS(\ + SELECT 1 FROM pg_locks held JOIN pg_locks waiting \ + ON waiting.locktype=held.locktype \ + AND waiting.database=held.database \ + AND waiting.classid=held.classid \ + AND waiting.objid=held.objid \ + AND waiting.objsubid=held.objsubid \ + WHERE held.locktype='advisory' \ + AND held.database::BIGINT=$1 \ + AND held.classid::BIGINT=$2 \ + AND held.objid::BIGINT=$3 \ + AND held.pid=$4 AND held.granted \ + AND waiting.pid=$5 AND NOT waiting.granted\ + )", + ) + .bind(i64::from(database_oid)) + .bind(i64::from(lock_key.class_id())) + .bind(i64::from(lock_key.object_id())) + .bind(holder_pid) + .bind(waiter_pid) + .fetch_one(pool) + .await + .expect("inspect exact identity advisory waiter"); + if waiting { + return; + } + } + tokio::task::yield_now().await; + } + panic!("second actor never blocked on the shared identity advisory lock"); +} + +async fn wait_for_transaction_waiter( + pool: &PgPool, + database_oid: u32, + holder_pid: i32, + holder_transaction_id: i64, + waiter_pid: i32, +) { + assert_ne!(holder_pid, waiter_pid); + for _ in 0..10_000 { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS(\ + SELECT 1 FROM pg_locks held JOIN pg_locks waiting \ + ON waiting.locktype=held.locktype \ + AND waiting.transactionid=held.transactionid \ + WHERE held.locktype='transactionid' \ + AND held.database IS NULL AND waiting.database IS NULL \ + AND held.transactionid::text::BIGINT=$1 \ + AND held.pid=$2 AND held.mode='ExclusiveLock' AND held.granted \ + AND waiting.pid=$3 AND waiting.mode='ShareLock' AND NOT waiting.granted \ + AND EXISTS(SELECT 1 FROM pg_stat_activity activity \ + WHERE activity.pid=waiting.pid \ + AND activity.datid::BIGINT=$4 \ + AND activity.wait_event_type='Lock' \ + AND activity.wait_event='transactionid')\ + )", + ) + .bind(holder_transaction_id) + .bind(holder_pid) + .bind(waiter_pid) + .bind(i64::from(database_oid)) + .fetch_one(pool) + .await + .expect("inspect exact identity row-lock waiter"); + if waiting { + return; + } + tokio::task::yield_now().await; + } + panic!("second actor never blocked on the first actor's exact row transaction lock"); +} + +fn serializes_on_active_binding_row(first: Action, second: Action) -> bool { + matches!( + (first, second), + (Action::Disable, Action::Revoke) | (Action::Revoke, Action::Disable) + ) +} + +async fn force_order( + pool: &PgPool, + fixture: &Fixture, + first: Action, + second: Action, + domain_b: CommunityId, + domain_b_auth: (Uuid, u64), +) -> (Outcome, Outcome) { + let (mut events, _controller) = test_lock_schedule::install(); + let row_schedule = serializes_on_active_binding_row(first, second); + let (mut row_events, _row_controller) = row_schedule + .then(test_lock_schedule::install_row) + .map_or((None, None), |(events, guard)| (Some(events), Some(guard))); + let first_pool = pool.clone(); + let first_fixture = fixture.clone(); + let first_task = tokio::spawn(test_lock_schedule::actor_scope("first", async move { + run_action(&first_pool, &first_fixture, first, Uuid::from_u128(100)).await + })); + let event = events.recv().await.expect("first lock request trace"); + assert_eq!( + (event.actor(), event.phase()), + ("first", test_lock_schedule::LockPhase::Request) + ); + let first_pid = event.backend_pid(); + let database_oid = event.database_oid(); + let first_lock_keys = event.lock_keys().to_vec(); + event.resume(); + let first_acquired = events.recv().await.expect("first lock acquired trace"); + assert_eq!( + (first_acquired.actor(), first_acquired.phase()), + ("first", test_lock_schedule::LockPhase::Acquired) + ); + assert_eq!(first_acquired.isolation(), Some("read committed")); + assert!(first_acquired.transaction_id().is_some()); + assert!(first_acquired.coordinate_count() >= 2); + assert_eq!(first_acquired.backend_pid(), first_pid); + assert_eq!(first_acquired.database_oid(), database_oid); + assert_eq!(first_acquired.lock_keys(), first_lock_keys); + + let second_pool = pool.clone(); + let second_fixture = fixture.clone(); + let second_task = tokio::spawn(test_lock_schedule::actor_scope("second", async move { + run_action(&second_pool, &second_fixture, second, Uuid::from_u128(101)).await + })); + let event = events.recv().await.expect("second lock request trace"); + assert_eq!( + (event.actor(), event.phase()), + ("second", test_lock_schedule::LockPhase::Request) + ); + let second_pid = event.backend_pid(); + assert_eq!(event.database_oid(), database_oid); + let shared_lock_keys = first_lock_keys + .iter() + .copied() + .filter(|lock_key| event.lock_keys().contains(lock_key)) + .collect::>(); + event.resume(); + if row_schedule { + assert!( + shared_lock_keys.is_empty(), + "disable/revoke must prove its actual row-lock serialization point" + ); + let second_acquired = events + .recv() + .await + .expect("second independent advisory lock acquired trace"); + assert_eq!( + (second_acquired.actor(), second_acquired.phase()), + ("second", test_lock_schedule::LockPhase::Acquired) + ); + assert_eq!(second_acquired.isolation(), Some("read committed")); + assert_eq!(second_acquired.backend_pid(), second_pid); + assert_eq!(second_acquired.database_oid(), database_oid); + let first_advisory_transaction_id = first_acquired.transaction_id(); + let second_advisory_transaction_id = second_acquired.transaction_id(); + + first_acquired.resume(); + let row_events = row_events.as_mut().expect("installed row-lock schedule"); + let first_row_request = row_events + .recv() + .await + .expect("first row-lock request trace"); + assert_eq!( + (first_row_request.actor(), first_row_request.phase()), + ("first", test_lock_schedule::RowLockPhase::Request) + ); + assert_eq!(first_row_request.backend_pid(), first_pid); + assert_eq!(first_row_request.database_oid(), database_oid); + assert_eq!( + Some(first_row_request.transaction_id()), + first_advisory_transaction_id + ); + first_row_request.resume(); + let first_row_acquired = row_events + .recv() + .await + .expect("first row-lock acquired trace"); + assert_eq!( + (first_row_acquired.actor(), first_row_acquired.phase()), + ("first", test_lock_schedule::RowLockPhase::Acquired) + ); + assert_eq!(first_row_acquired.backend_pid(), first_pid); + assert_eq!(first_row_acquired.database_oid(), database_oid); + let first_transaction_id = first_row_acquired.transaction_id(); + assert_eq!(Some(first_transaction_id), first_advisory_transaction_id); + + second_acquired.resume(); + let second_row_request = row_events + .recv() + .await + .expect("second row-lock request trace"); + assert_eq!( + (second_row_request.actor(), second_row_request.phase()), + ("second", test_lock_schedule::RowLockPhase::Request) + ); + assert_eq!(second_row_request.backend_pid(), second_pid); + assert_eq!(second_row_request.database_oid(), database_oid); + assert_eq!( + Some(second_row_request.transaction_id()), + second_advisory_transaction_id + ); + second_row_request.resume(); + wait_for_transaction_waiter( + pool, + database_oid, + first_pid, + first_transaction_id, + second_pid, + ) + .await; + assert_eq!(authorization_sentinel(pool, domain_b).await, domain_b_auth); + + first_row_acquired.resume(); + let second_row_acquired = row_events + .recv() + .await + .expect("second row-lock acquired trace"); + assert_eq!( + (second_row_acquired.actor(), second_row_acquired.phase()), + ("second", test_lock_schedule::RowLockPhase::Acquired) + ); + assert_eq!(second_row_acquired.backend_pid(), second_pid); + assert_eq!(second_row_acquired.database_oid(), database_oid); + assert_eq!( + Some(second_row_acquired.transaction_id()), + second_advisory_transaction_id + ); + second_row_acquired.resume(); + } else { + wait_for_advisory_waiter(pool, database_oid, first_pid, second_pid, &shared_lock_keys) + .await; + assert_eq!(authorization_sentinel(pool, domain_b).await, domain_b_auth); + + first_acquired.resume(); + let second_acquired = events.recv().await.expect("second lock acquired trace"); + assert_eq!( + (second_acquired.actor(), second_acquired.phase()), + ("second", test_lock_schedule::LockPhase::Acquired) + ); + assert_eq!(second_acquired.isolation(), Some("read committed")); + assert_eq!(second_acquired.backend_pid(), second_pid); + assert_eq!(second_acquired.database_oid(), database_oid); + second_acquired.resume(); + } + + ( + first_task.await.expect("join first scheduled actor"), + second_task.await.expect("join second scheduled actor"), + ) +} + +async fn exercise_ordered_case(pool: &PgPool, first: Action, second: Action) { + let pair = (first, second); + let reference = setup_fixture(pool, pair, "identity-reference").await; + let scheduled = setup_fixture(pool, pair, "identity-scheduled").await; + let domain_b = make_community(pool, "identity-domain-b").await; + enroll_key(pool, domain_b, &DOMAIN_B_KEY).await; + let domain_b_bytes = raw_domain_snapshot(pool, domain_b).await; + let domain_b_auth = authorization_sentinel(pool, domain_b).await; + + let expected_first = run_action(pool, &reference, first, Uuid::from_u128(100)).await; + let expected_second = run_action(pool, &reference, second, Uuid::from_u128(101)).await; + let actual = force_order(pool, &scheduled, first, second, domain_b, domain_b_auth).await; + + assert_eq!(actual, (expected_first, expected_second)); + assert_eq!( + logical_snapshot(pool, scheduled.community_id).await, + logical_snapshot(pool, reference.community_id).await, + "forced schedule must match its independently executed sequential reference: {first:?} then {second:?}" + ); + assert_eq!(raw_domain_snapshot(pool, domain_b).await, domain_b_bytes); + assert_eq!(authorization_sentinel(pool, domain_b).await, domain_b_auth); +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn forced_lock_schedules_match_every_ordered_lifecycle_reference() { + let pool = setup_pool().await; + let actions = [ + Action::Retire, + Action::Disable, + Action::Revoke, + Action::Rotate, + Action::Recover, + Action::Enable, + ]; + for left in 0..actions.len() { + for right in (left + 1)..actions.len() { + exercise_ordered_case(&pool, actions[left], actions[right]).await; + exercise_ordered_case(&pool, actions[right], actions[left]).await; + } + } +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn enrollment_rotate_and_recover_have_forced_orders_and_references() { + let pool = setup_pool().await; + for lifecycle in [Action::Rotate, Action::Recover] { + exercise_ordered_case(&pool, Action::Enroll, lifecycle).await; + exercise_ordered_case(&pool, lifecycle, Action::Enroll).await; + } +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn compare_and_clear_rejects_aba_recreation_and_preserves_domain_b() { + let pool = setup_pool().await; + let fixture = setup_fixture(&pool, (Action::Recover, Action::Disable), "identity-aba").await; + let expected = fixture.expected_pending.expect("pending selector"); + let domain_b = make_community(&pool, "identity-aba-domain-b").await; + enroll_key(&pool, domain_b, &DOMAIN_B_KEY).await; + let domain_b_bytes = raw_domain_snapshot(&pool, domain_b).await; + let domain_b_auth = authorization_sentinel(&pool, domain_b).await; + + let mut tx = pool.begin().await.expect("begin ABA transaction"); + sqlx::query( + "UPDATE identity_pending_replacements SET cleared_at=NOW(), cleared_operation_id=$4 \ + WHERE community_id=$1 AND issuer=$2 AND subject=$3 AND cleared_at IS NULL", + ) + .bind(fixture.community_id.as_uuid()) + .bind(ISSUER) + .bind(SUBJECT) + .bind(Uuid::from_u128(200)) + .execute(&mut *tx) + .await + .expect("clear G1"); + sqlx::query( + "INSERT INTO identity_pending_replacements \ + (community_id,issuer,subject,selector_version,retired_pubkey,retired_binding_id,retired_binding_version,created_operation_id) \ + VALUES ($1,$2,$3,$4,$5,$6,$7,$8)", + ) + .bind(fixture.community_id.as_uuid()) + .bind(ISSUER) + .bind(SUBJECT) + .bind(i64::try_from(expected.selector_version + 1).expect("selector fits i64")) + .bind(&expected.retired_pubkey) + .bind(expected.retired_binding_id) + .bind(i64::try_from(expected.retired_binding_version).expect("version fits i64")) + .bind(Uuid::from_u128(201)) + .execute(&mut *tx) + .await + .expect("recreate G2"); + assert!(compare_and_clear_pending_tx( + &mut tx, + fixture.community_id, + context(Uuid::from_u128(202), "stale ABA compare"), + principal(), + &expected, + ) + .await + .is_err()); + let active_selector: i64 = sqlx::query_scalar( + "SELECT selector_version FROM identity_pending_replacements \ + WHERE community_id=$1 AND issuer=$2 AND subject=$3 AND cleared_at IS NULL", + ) + .bind(fixture.community_id.as_uuid()) + .bind(ISSUER) + .bind(SUBJECT) + .fetch_one(&mut *tx) + .await + .expect("G2 remains active"); + assert_eq!( + active_selector, + i64::try_from(expected.selector_version + 1).unwrap() + ); + tx.rollback() + .await + .expect("rollback synthetic ABA recreation"); + + assert_eq!(raw_domain_snapshot(&pool, domain_b).await, domain_b_bytes); + assert_eq!(authorization_sentinel(&pool, domain_b).await, domain_b_auth); +} diff --git a/crates/buzz-db/src/migration.rs b/crates/buzz-db/src/migration.rs index 0bf579359..65b37327b 100644 --- a/crates/buzz-db/src/migration.rs +++ b/crates/buzz-db/src/migration.rs @@ -10,6 +10,10 @@ use crate::Result; static MIGRATOR: sqlx::migrate::Migrator = sqlx::migrate!("../../migrations"); +#[cfg(test)] +#[path = "migration_deterministic_tests.rs"] +mod deterministic_tests; + /// Run all pending Buzz database migrations. pub async fn run_migrations(pool: &PgPool) -> Result<()> { reject_legacy_nip_rs_cardinality_ambiguity(pool).await?; diff --git a/crates/buzz-db/src/migration_deterministic_tests.rs b/crates/buzz-db/src/migration_deterministic_tests.rs new file mode 100644 index 000000000..e92078ed6 --- /dev/null +++ b/crates/buzz-db/src/migration_deterministic_tests.rs @@ -0,0 +1,846 @@ +use super::MIGRATOR; + +use crate::identity_binding::{ + get_active_identity_binding_by_pubkey, resolve_identity_binding, BindingDenial, + BindingProvenance, EnrollmentMode, ResolveBindingInput, ResolveBindingResult, +}; +use crate::identity_lifecycle::{ + disable_identity_principal, enable_identity_principal, provision_identity_binding, + recover_identity_binding, retire_identity_pair, revoke_identity_key, rotate_identity_binding, + IdentityPrincipal, LifecycleContext, PendingLineage, VerifiedReplacementKey, +}; +use buzz_core::CommunityId; +use sqlx::{Acquire, PgPool}; +use uuid::Uuid; + +const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; + +async fn connect_pool() -> PgPool { + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| TEST_DB_URL.to_owned()); + PgPool::connect(&database_url) + .await + .expect("connect deterministic migration DB") +} + +async fn reset_to_0028(pool: &PgPool) -> (Uuid, Uuid) { + sqlx::query("DROP SCHEMA IF EXISTS public CASCADE") + .execute(pool) + .await + .expect("drop public schema"); + sqlx::query("CREATE SCHEMA public") + .execute(pool) + .await + .expect("create public schema"); + MIGRATOR + .run_to(28, pool) + .await + .expect("migrate through 0028"); + let domain_a = Uuid::new_v4(); + let domain_b = Uuid::new_v4(); + for (domain, label, key) in [ + (domain_a, "migration-fault-a", vec![71_u8; 32]), + (domain_b, "migration-fault-b", vec![72_u8; 32]), + ] { + sqlx::query("INSERT INTO communities (id,host) VALUES ($1,$2)") + .bind(domain) + .bind(format!("{label}-{}.example", domain.simple())) + .execute(pool) + .await + .expect("insert migration fault domain"); + sqlx::query( + "INSERT INTO identity_bindings (community_id,issuer,uid,pubkey,source) \ + VALUES ($1,'https://idp.example',$2,$3,'db_binding')", + ) + .bind(domain) + .bind(format!("{label}-subject")) + .bind(key) + .execute(pool) + .await + .expect("insert migration fault binding"); + } + (domain_a, domain_b) +} + +fn split_statements(sql: &str) -> Vec { + let bytes = sql.as_bytes(); + let mut statements = Vec::new(); + let mut start = 0; + let mut index = 0; + let mut single_quote = false; + let mut line_comment = false; + while index < bytes.len() { + if line_comment { + if bytes[index] == b'\n' { + line_comment = false; + } + index += 1; + continue; + } + if !single_quote && index + 1 < bytes.len() && &bytes[index..index + 2] == b"--" { + line_comment = true; + index += 2; + continue; + } + match bytes[index] { + b'\'' => { + if single_quote && index + 1 < bytes.len() && bytes[index + 1] == b'\'' { + index += 2; + continue; + } + single_quote = !single_quote; + index += 1; + } + b';' if !single_quote => { + let statement = sql[start..index].trim(); + if !statement.is_empty() { + statements.push(statement.to_owned()); + } + start = index + 1; + index += 1; + } + _ => index += 1, + } + } + let tail = sql[start..].trim(); + if !tail.is_empty() { + statements.push(tail.to_owned()); + } + statements +} + +fn migration_0029() -> &'static sqlx::migrate::Migration { + MIGRATOR + .iter() + .find(|migration| migration.version == 29) + .expect("embedded migration 0029") +} + +async fn legacy_snapshot(pool: &PgPool) -> Vec { + let tables = [ + "communities", + "identity_bindings", + "identity_principals", + "identity_revoked_keys", + "audit_log", + ]; + let mut snapshot = Vec::new(); + for table in tables { + let query = format!("SELECT to_jsonb(t)::text FROM {table} t ORDER BY 1"); + let rows = sqlx::query_scalar::<_, String>(sqlx::AssertSqlSafe(query)) + .fetch_all(pool) + .await + .expect("snapshot legacy rows"); + snapshot.extend(rows.into_iter().map(|row| format!("row:{table}:{row}"))); + } + let catalog = sqlx::query_scalar::<_, String>( + "SELECT value FROM (\ + SELECT 'column:'||table_name||':'||column_name||':'||data_type||':'||is_nullable||':'||COALESCE(column_default,'') AS value \ + FROM information_schema.columns WHERE table_schema='public' AND table_name LIKE 'identity_%' \ + UNION ALL \ + SELECT 'constraint:'||conrelid::regclass::text||':'||conname||':'||pg_get_constraintdef(oid) \ + FROM pg_constraint WHERE conrelid::regclass::text LIKE 'identity_%' \ + UNION ALL \ + SELECT 'index:'||tablename||':'||indexname||':'||indexdef \ + FROM pg_indexes WHERE schemaname='public' AND tablename LIKE 'identity_%'\ + ) catalog ORDER BY value", + ) + .fetch_all(pool) + .await + .expect("snapshot legacy catalog"); + snapshot.extend(catalog); + let versions = sqlx::query_scalar::<_, i64>( + "SELECT version FROM _sqlx_migrations WHERE success ORDER BY version", + ) + .fetch_all(pool) + .await + .expect("snapshot migration versions"); + snapshot.extend( + versions + .into_iter() + .map(|version| format!("version:{version}")), + ); + snapshot.sort(); + snapshot +} + +async fn full_identity_snapshot(pool: &PgPool) -> Vec { + let tables = [ + "identity_bindings", + "identity_principals", + "identity_revoked_keys", + "identity_migration_denials", + "identity_migration_denied_keys", + "identity_binding_lineage", + "identity_retired_pairs", + "identity_pending_replacements", + "identity_binding_history", + "identity_lifecycle_operations", + "audit_log", + ]; + let mut snapshot = Vec::new(); + for table in tables { + let query = format!("SELECT to_jsonb(t)::text FROM {table} t ORDER BY 1"); + let rows = sqlx::query_scalar::<_, String>(sqlx::AssertSqlSafe(query)) + .fetch_all(pool) + .await + .expect("snapshot migrated identity rows"); + snapshot.extend(rows.into_iter().map(|row| format!("{table}:{row}"))); + } + let marker: Vec = + sqlx::query_scalar("SELECT to_jsonb(m)::text FROM _sqlx_migrations m WHERE version=29") + .fetch_all(pool) + .await + .expect("snapshot 0029 marker"); + snapshot.extend(marker.into_iter().map(|row| format!("marker:{row}"))); + snapshot.sort(); + snapshot +} + +async fn raw_domain_authorized(pool: &PgPool, domain: Uuid, key: &[u8]) -> bool { + sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM identity_bindings \ + WHERE community_id=$1 AND pubkey=$2 AND revoked_at IS NULL)", + ) + .bind(domain) + .bind(key) + .fetch_one(pool) + .await + .expect("read raw legacy authorization sentinel") +} + +async fn domain_audit_snapshot(pool: &PgPool, domain: Uuid) -> Vec { + sqlx::query_scalar::<_, String>( + "SELECT to_jsonb(row_value)::text FROM audit_log row_value \ + WHERE community_id=$1 ORDER BY 1", + ) + .bind(domain) + .fetch_all(pool) + .await + .expect("read domain audit sentinel") +} + +async fn legacy_identity_facts(pool: &PgPool, domain: Uuid) -> Vec { + let mut facts = sqlx::query_scalar::<_, serde_json::Value>( + "SELECT to_jsonb(binding)-ARRAY[\ + 'binding_id','binding_version','binding_state','binding_provenance',\ + 'replacement_binding_id','created_by','created_policy_version']::text[] \ + FROM identity_bindings binding WHERE community_id=$1", + ) + .bind(domain) + .fetch_all(pool) + .await + .expect("read legacy binding facts") + .into_iter() + .map(|value| format!("binding:{value}")) + .collect::>(); + for (table, query) in [ + ( + "principal", + "SELECT to_jsonb(row_value) FROM identity_principals row_value WHERE community_id=$1", + ), + ( + "revoked_key", + "SELECT to_jsonb(row_value) FROM identity_revoked_keys row_value WHERE community_id=$1", + ), + ] { + let rows = sqlx::query_scalar::<_, serde_json::Value>(query) + .bind(domain) + .fetch_all(pool) + .await + .expect("read legacy selector facts"); + facts.extend(rows.into_iter().map(|value| format!("{table}:{value}"))); + } + facts.sort(); + facts +} + +async fn domain_identity_snapshot(pool: &PgPool, domain: Uuid) -> Vec { + let tables = [ + "identity_bindings", + "identity_principals", + "identity_revoked_keys", + "identity_migration_denials", + "identity_migration_denied_keys", + "identity_binding_lineage", + "identity_retired_pairs", + "identity_pending_replacements", + "identity_binding_history", + "identity_lifecycle_operations", + "audit_log", + ]; + let mut snapshot = Vec::new(); + for table in tables { + let query = format!( + "SELECT to_jsonb(row_value)::text FROM {table} row_value WHERE community_id=$1 ORDER BY 1" + ); + let rows = sqlx::query_scalar::<_, String>(sqlx::AssertSqlSafe(query)) + .bind(domain) + .fetch_all(pool) + .await + .expect("read domain identity snapshot"); + snapshot.extend(rows.into_iter().map(|row| format!("{table}:{row}"))); + } + snapshot +} + +async fn insert_rotated_legacy_row( + pool: &PgPool, + domain: Uuid, + subject: &str, + key: &[u8], + target: &[u8], +) { + sqlx::query( + "INSERT INTO identity_bindings \ + (community_id,issuer,uid,pubkey,source,revoked_at,revoked_reason,revocation_scope,\ + rotation_completed_at,rotated_to_pubkey,rotation_reason) \ + VALUES ($1,'https://idp.example',$2,$3,'db_binding',NOW(),'legacy rotation',\ + 'rotation',NOW(),$4,'legacy rotation')", + ) + .bind(domain) + .bind(subject) + .bind(key) + .bind(target) + .execute(pool) + .await + .expect("insert rotated legacy row"); +} + +async fn insert_revoked_legacy_row(pool: &PgPool, domain: Uuid, subject: &str, key: &[u8]) { + sqlx::query( + "INSERT INTO identity_bindings \ + (community_id,issuer,uid,pubkey,source,revoked_at,revoked_reason,revocation_scope) \ + VALUES ($1,'https://idp.example',$2,$3,'db_binding',NOW(),'legacy revoke','key')", + ) + .bind(domain) + .bind(subject) + .bind(key) + .execute(pool) + .await + .expect("insert revoked legacy row"); +} + +fn lifecycle_context(id: u128, reason: &'static str) -> LifecycleContext<'static> { + LifecycleContext { + operation_id: Uuid::from_u128(id), + actor: None, + reason, + } +} + +fn replacement( + key: &'static [u8; 32], + provenance: BindingProvenance, +) -> VerifiedReplacementKey<'static> { + VerifiedReplacementKey::after_verified_proof( + key, + None, + provenance, + Some("migration-denial-policy-v1"), + ) + .expect("construct migrated denial replacement") +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn identity_0029_fault_at_every_statement_boundary_restarts_and_retries() { + let statements = split_statements(migration_0029().sql.as_ref()); + assert_eq!( + statements.len(), + 28, + "0029 boundary count is part of the oracle adapter" + ); + + for boundary in 0..=statements.len() { + let pool = connect_pool().await; + let (_domain_a, domain_b) = reset_to_0028(&pool).await; + let before = legacy_snapshot(&pool).await; + let domain_b_facts = legacy_identity_facts(&pool, domain_b).await; + let domain_b_audit = domain_audit_snapshot(&pool, domain_b).await; + assert!(raw_domain_authorized(&pool, domain_b, &[72_u8; 32]).await); + + let mut connection = pool.acquire().await.expect("acquire crashable backend"); + let backend_pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut *connection) + .await + .expect("read crashable backend pid"); + let mut tx = connection + .begin() + .await + .expect("begin crashable boundary migration"); + for statement in &statements[..boundary] { + sqlx::raw_sql(sqlx::AssertSqlSafe(statement.as_str())) + .execute(&mut *tx) + .await + .unwrap_or_else(|error| { + panic!("0029 statement before boundary {boundary} failed: {error}") + }); + } + let terminated: bool = sqlx::query_scalar("SELECT pg_terminate_backend($1)") + .bind(backend_pid) + .fetch_one(&pool) + .await + .unwrap_or_else(|error| { + panic!("terminate boundary {boundary} migration backend: {error}") + }); + assert!(terminated, "boundary {boundary} backend must terminate"); + assert!( + sqlx::query("SELECT 1").execute(&mut *tx).await.is_err(), + "boundary {boundary} transaction must observe backend loss" + ); + drop(tx); + drop(connection); + pool.close().await; + + let pool = connect_pool().await; + assert_eq!(legacy_snapshot(&pool).await, before, "boundary {boundary}"); + assert_eq!( + legacy_identity_facts(&pool, domain_b).await, + domain_b_facts, + "boundary {boundary} domain-B facts before retry" + ); + assert_eq!( + domain_audit_snapshot(&pool, domain_b).await, + domain_b_audit, + "boundary {boundary} domain-B audit before retry" + ); + assert!(raw_domain_authorized(&pool, domain_b, &[72_u8; 32]).await); + + MIGRATOR.run_to(29, &pool).await.unwrap_or_else(|error| { + panic!("boundary {boundary} retry after restart failed: {error}") + }); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM _sqlx_migrations WHERE version=29 AND success" + ) + .fetch_one(&pool) + .await + .expect("count boundary retry marker"), + 1, + "boundary {boundary} must produce one durable migration marker" + ); + assert_eq!( + legacy_identity_facts(&pool, domain_b).await, + domain_b_facts, + "boundary {boundary} domain-B facts after retry" + ); + assert_eq!( + domain_audit_snapshot(&pool, domain_b).await, + domain_b_audit, + "boundary {boundary} domain-B audit after retry" + ); + assert!(raw_domain_authorized(&pool, domain_b, &[72_u8; 32]).await); + let complete_post = full_identity_snapshot(&pool).await; + pool.close().await; + + let pool = connect_pool().await; + assert_eq!( + full_identity_snapshot(&pool).await, + complete_post, + "boundary {boundary} complete post-state must survive restart" + ); + MIGRATOR.run_to(29, &pool).await.unwrap_or_else(|error| { + panic!("boundary {boundary} idempotent post-restart retry failed: {error}") + }); + assert_eq!( + full_identity_snapshot(&pool).await, + complete_post, + "boundary {boundary} retry must remain an exact no-op" + ); + pool.close().await; + } +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn identity_0029_commit_failure_rolls_back_projection_and_success_marker() { + let pool = connect_pool().await; + let (_domain_a, domain_b) = reset_to_0028(&pool).await; + let before = legacy_snapshot(&pool).await; + let migration = migration_0029(); + let statements = split_statements(migration.sql.as_ref()); + let mut tx = pool.begin().await.expect("begin commit-failure migration"); + for statement in statements { + sqlx::raw_sql(sqlx::AssertSqlSafe(statement)) + .execute(&mut *tx) + .await + .expect("execute 0029 before commit failure"); + } + sqlx::query( + "INSERT INTO _sqlx_migrations \ + (version,description,installed_on,success,checksum,execution_time) \ + VALUES ($1,$2,NOW(),TRUE,$3,0)", + ) + .bind(migration.version) + .bind(migration.description.as_ref()) + .bind(migration.checksum.as_ref()) + .execute(&mut *tx) + .await + .expect("insert transactional 0029 marker"); + sqlx::query( + "UPDATE identity_bindings SET replacement_binding_id=$1 \ + WHERE community_id=$2", + ) + .bind(Uuid::from_u128(u128::MAX)) + .bind(domain_b) + .execute(&mut *tx) + .await + .expect("deferred FK accepts invalid replacement before commit"); + assert!( + tx.commit().await.is_err(), + "deferred FK must fail at commit" + ); + assert_eq!(legacy_snapshot(&pool).await, before); + assert!(raw_domain_authorized(&pool, domain_b, &[72_u8; 32]).await); + + MIGRATOR + .run_to(29, &pool) + .await + .expect("retry after commit failure"); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM _sqlx_migrations WHERE version=29 AND success" + ) + .fetch_one(&pool) + .await + .expect("count successful 0029 marker"), + 1 + ); +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn identity_0029_backend_loss_restarts_cleanly_and_retry_converges() { + let pool = connect_pool().await; + let (_domain_a, domain_b) = reset_to_0028(&pool).await; + let before = legacy_snapshot(&pool).await; + let statements = split_statements(migration_0029().sql.as_ref()); + let mut connection = pool.acquire().await.expect("acquire migration backend"); + let backend_pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut *connection) + .await + .expect("read migration backend pid"); + let mut tx = connection.begin().await.expect("begin crashable migration"); + for statement in &statements[..statements.len() / 2] { + sqlx::raw_sql(sqlx::AssertSqlSafe(statement.as_str())) + .execute(&mut *tx) + .await + .expect("execute pre-crash migration prefix"); + } + let terminated: bool = sqlx::query_scalar("SELECT pg_terminate_backend($1)") + .bind(backend_pid) + .fetch_one(&pool) + .await + .expect("terminate migration backend"); + assert!(terminated); + assert!(sqlx::query("SELECT 1").execute(&mut *tx).await.is_err()); + drop(tx); + drop(connection); + pool.close().await; + + let pool = connect_pool().await; + assert_eq!(legacy_snapshot(&pool).await, before); + assert!(raw_domain_authorized(&pool, domain_b, &[72_u8; 32]).await); + MIGRATOR + .run_to(29, &pool) + .await + .expect("retry after backend restart"); +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn identity_0029_post_commit_response_loss_is_idempotent_after_restart() { + let pool = connect_pool().await; + reset_to_0028(&pool).await; + MIGRATOR + .run_to(29, &pool) + .await + .expect("commit 0029 before response loss"); + let committed = full_identity_snapshot(&pool).await; + let synthetic_response: Result<(), &str> = Err("response lost after confirmed commit"); + assert!(synthetic_response.is_err()); + pool.close().await; + + let pool = connect_pool().await; + assert_eq!(full_identity_snapshot(&pool).await, committed); + MIGRATOR + .run_to(29, &pool) + .await + .expect("idempotent retry after response loss"); + assert_eq!(full_identity_snapshot(&pool).await, committed); +} + +#[tokio::test] +#[ignore = "requires a dedicated disposable Postgres database"] +async fn identity_0029_readable_ambiguities_preserve_facts_and_never_create_authority() { + static PROVISION_KEY: [u8; 32] = [91; 32]; + static ROTATE_KEY: [u8; 32] = [92; 32]; + static RECOVER_KEY: [u8; 32] = [93; 32]; + static ENABLE_KEY: [u8; 32] = [94; 32]; + + let pool = connect_pool().await; + let (domain_a, domain_b) = reset_to_0028(&pool).await; + let active_key = vec![71_u8; 32]; + let missing_predecessor = vec![73_u8; 32]; + let missing_target = vec![74_u8; 32]; + insert_rotated_legacy_row( + &pool, + domain_a, + "migration-fault-a-subject", + &missing_predecessor, + &missing_target, + ) + .await; + + let duplicate_key = vec![80_u8; 32]; + insert_revoked_legacy_row(&pool, domain_a, "duplicate", &duplicate_key).await; + insert_revoked_legacy_row(&pool, domain_a, "duplicate", &duplicate_key).await; + + let cycle_a = vec![81_u8; 32]; + let cycle_b = vec![82_u8; 32]; + insert_rotated_legacy_row(&pool, domain_a, "cycle", &cycle_a, &cycle_b).await; + insert_rotated_legacy_row(&pool, domain_a, "cycle", &cycle_b, &cycle_a).await; + + let fork_source = vec![83_u8; 32]; + let fork_target = vec![84_u8; 32]; + insert_rotated_legacy_row(&pool, domain_a, "fork", &fork_source, &fork_target).await; + insert_revoked_legacy_row(&pool, domain_a, "fork", &fork_target).await; + insert_revoked_legacy_row(&pool, domain_a, "fork", &fork_target).await; + + let conflicting_key = vec![85_u8; 32]; + sqlx::query( + "INSERT INTO identity_bindings (community_id,issuer,uid,pubkey,source) \ + VALUES ($1,'https://idp.example','active-with-tombstone',$2,'db_binding')", + ) + .bind(domain_a) + .bind(&conflicting_key) + .execute(&pool) + .await + .expect("insert active binding selected by legacy tombstone"); + sqlx::query( + "INSERT INTO identity_revoked_keys (community_id,pubkey,reason) \ + VALUES ($1,$2,'legacy key tombstone')", + ) + .bind(domain_a) + .bind(&conflicting_key) + .execute(&pool) + .await + .expect("insert conflicting legacy tombstone"); + + let legacy_a = legacy_identity_facts(&pool, domain_a).await; + let legacy_b = legacy_identity_facts(&pool, domain_b).await; + MIGRATOR + .run_to(29, &pool) + .await + .expect("readable ambiguity migrates into fail-closed quarantine"); + assert_eq!(legacy_identity_facts(&pool, domain_a).await, legacy_a); + assert_eq!(legacy_identity_facts(&pool, domain_b).await, legacy_b); + + let quarantined: Vec = sqlx::query_scalar( + "SELECT subject FROM identity_migration_denials WHERE community_id=$1 ORDER BY subject", + ) + .bind(domain_a) + .fetch_all(&pool) + .await + .expect("read principal quarantines"); + assert_eq!( + quarantined, + vec![ + "cycle".to_owned(), + "duplicate".to_owned(), + "fork".to_owned(), + "migration-fault-a-subject".to_owned(), + ] + ); + for key in [ + active_key.as_slice(), + missing_predecessor.as_slice(), + missing_target.as_slice(), + duplicate_key.as_slice(), + cycle_a.as_slice(), + cycle_b.as_slice(), + fork_source.as_slice(), + fork_target.as_slice(), + ] { + let denied: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM identity_migration_denied_keys \ + WHERE community_id=$1 AND pubkey=$2)", + ) + .bind(domain_a) + .bind(key) + .fetch_one(&pool) + .await + .expect("read implicated key quarantine"); + assert!( + denied, + "every stored or referenced implicated key is denied" + ); + } + let imported_ambiguous_edges: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM identity_binding_lineage lineage \ + JOIN identity_bindings binding \ + ON binding.community_id=lineage.community_id \ + AND binding.binding_id=lineage.predecessor_binding_id \ + WHERE binding.community_id=$1 AND binding.uid IN \ + ('cycle','duplicate','fork','migration-fault-a-subject')", + ) + .bind(domain_a) + .fetch_one(&pool) + .await + .expect("count ambiguous imported lineage"); + assert_eq!(imported_ambiguous_edges, 0); + + let domain_a_id = CommunityId::from_uuid(domain_a); + let domain_b_id = CommunityId::from_uuid(domain_b); + let main_principal = IdentityPrincipal { + issuer: "https://idp.example", + subject: "migration-fault-a-subject", + }; + assert!( + get_active_identity_binding_by_pubkey(&pool, domain_a_id, &active_key) + .await + .is_err() + ); + assert!( + get_active_identity_binding_by_pubkey(&pool, domain_a_id, &conflicting_key) + .await + .is_err() + ); + for (subject, key) in [ + (main_principal.subject, active_key.as_slice()), + ("different-principal", missing_target.as_slice()), + ("active-with-tombstone", conflicting_key.as_slice()), + ] { + let result = resolve_identity_binding( + &pool, + domain_a_id, + &ResolveBindingInput { + issuer: "https://idp.example", + subject, + pubkey: key, + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("migrated denial resolves without authority"); + assert_eq!(result, ResolveBindingResult::Denied(BindingDenial::Revoked)); + } + + let domain_b_cross = resolve_identity_binding( + &pool, + domain_b_id, + &ResolveBindingInput { + issuer: "https://idp.example", + subject: "independent-domain-subject", + pubkey: &missing_target, + display_name: None, + enrollment_mode: EnrollmentMode::AttestedKey, + key_attested: true, + }, + ) + .await + .expect("same bytes remain independent in domain B"); + assert!(matches!(domain_b_cross, ResolveBindingResult::Enrolled(_))); + let domain_b_sentinel = domain_identity_snapshot(&pool, domain_b).await; + + let (retired_binding_id, retired_binding_version): (Uuid, i64) = sqlx::query_as( + "SELECT binding_id,binding_version FROM identity_bindings \ + WHERE community_id=$1 AND issuer='https://idp.example' \ + AND uid='migration-fault-a-subject' AND pubkey=$2", + ) + .bind(domain_a) + .bind(&missing_predecessor) + .fetch_one(&pool) + .await + .expect("read quarantined retired coordinate"); + let fabricated_pending = PendingLineage { + retired_pubkey: missing_predecessor.clone(), + retired_binding_id, + retired_binding_version: u64::try_from(retired_binding_version).expect("positive version"), + selector_version: 1, + }; + let before_denials = domain_identity_snapshot(&pool, domain_a).await; + assert!(provision_identity_binding( + &pool, + domain_a_id, + lifecycle_context(301, "quarantine provision denial"), + main_principal, + EnrollmentMode::Provisioned, + replacement(&PROVISION_KEY, BindingProvenance::Provisioned), + ) + .await + .is_err()); + assert!(retire_identity_pair( + &pool, + domain_a_id, + lifecycle_context(302, "quarantine retire denial"), + main_principal, + &active_key, + ) + .await + .is_err()); + assert!(disable_identity_principal( + &pool, + domain_a_id, + lifecycle_context(303, "quarantine disable denial"), + main_principal, + ) + .await + .is_err()); + assert!(revoke_identity_key( + &pool, + domain_a_id, + lifecycle_context(304, "quarantine revoke denial"), + &active_key, + ) + .await + .is_err()); + assert!(rotate_identity_binding( + &pool, + domain_a_id, + lifecycle_context(305, "quarantine rotate denial"), + main_principal, + &active_key, + replacement(&ROTATE_KEY, BindingProvenance::AttestedKey), + ) + .await + .is_err()); + assert!(recover_identity_binding( + &pool, + domain_a_id, + lifecycle_context(306, "quarantine recover denial"), + main_principal, + &fabricated_pending, + replacement(&RECOVER_KEY, BindingProvenance::AttestedKey), + ) + .await + .is_err()); + assert!(enable_identity_principal( + &pool, + domain_a_id, + lifecycle_context(307, "quarantine enable denial"), + main_principal, + Some(&fabricated_pending), + replacement(&ENABLE_KEY, BindingProvenance::AttestedKey), + ) + .await + .is_err()); + assert_eq!( + domain_identity_snapshot(&pool, domain_a).await, + before_denials + ); + assert_eq!( + domain_identity_snapshot(&pool, domain_b).await, + domain_b_sentinel + ); + assert!( + get_active_identity_binding_by_pubkey(&pool, domain_b_id, &missing_target) + .await + .expect("domain-B authorization sentinel") + .is_some() + ); +}