From a84a32f4b62d1fe9a7ded6a9592bfff160b06982 Mon Sep 17 00:00:00 2001 From: tornquist Date: Tue, 22 Sep 2026 21:53:16 +0000 Subject: [PATCH 01/10] Add owner deletion admission control plane Signed-off-by: tornquist Co-authored-by: Codex --- ARCHITECTURE.md | 12 + crates/buzz-db/src/lib.rs | 2 +- crates/buzz-db/src/runtime/migration.rs | 45 +- .../tests/thread_window_postgres_tests.rs | 2 +- crates/buzz-db/src/store/community.rs | 119 ++++-- crates/buzz-db/src/store/deletion.rs | 403 +++++++++++++++++- crates/buzz-db/src/store/relay_members.rs | 43 +- crates/buzz-relay/src/api/operator.rs | 274 +++++++++++- crates/buzz-relay/src/router.rs | 6 +- ...050_owner_community_deletion_admission.sql | 66 +++ schema/schema.sql | 27 +- 11 files changed, 943 insertions(+), 56 deletions(-) create mode 100644 migrations/0050_owner_community_deletion_admission.sql diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index ae82c131ec7..7d7da313d7b 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -14,6 +14,18 @@ EVENT, REQ, REST, media, git, search, workflow, or pub/sub handling. Unknown hosts fail closed, and NIP-98/API-token stamps must agree with the host-derived community rather than overriding it. +Deployment-root community management uses operator-signed NIP-98 HTTP requests. +`POST /operator/communities/delete` accepts only an exact normalized, archived +community whose asserted pubkey is still its owner. The caller supplies the +request UUID as the stable correlation/idempotency identity; the durable row +records owner intent, mediating operator, and acknowledgement version. Admission +returns `202` at the `submitted` stage and performs no inventory, approval, +quiescing, object-store access, or deletion execution synchronously. While that +non-aborted request exists, unarchive and ownership transfer conflict and owner +management lists suppress the archived row. Replaying the same UUID converges +to its current stage; a different UUID conflicts with the existing one-active- +request invariant until that request is aborted. + Buzz is a Rust monorepo, licensed Apache 2.0 under Block, Inc. --- diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index e70c5dd85b1..2f5d3e8aacc 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -72,7 +72,7 @@ pub use allowlist::AllowlistEntry; pub use api_token::{ApiTokenRecord, TokenSummary}; pub use community::{ ArchivedCommunityRecord, CommunityRecord, CreateCommunityWithOwnerResult, - CreatedCommunityRecord, EnsuredCommunityRecord, OwnedCommunityRecord, + CreatedCommunityRecord, EnsuredCommunityRecord, OwnedCommunityRecord, UnarchiveCommunityResult, UnarchivedCommunityRecord, }; pub use error::{DbError, Result}; diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 20496ac1cfc..a6727c4ac72 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -703,12 +703,17 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 49); + assert_eq!(migrations.len(), 50); assert_eq!(migrations[48].version, 49); + assert_eq!(migrations[49].version, 50); assert!(migrations[48] .sql .as_str() .contains("idx_thread_metadata_window")); + assert!(migrations[49] + .sql + .as_str() + .contains("community_deletion_owner_provenance")); assert_eq!(migrations[0].version, 1); assert_eq!(&*migrations[0].description, "initial schema"); assert!(migrations[0] @@ -1765,6 +1770,12 @@ mod postgres_tests { .expect("embedded migration 0029") .sql .as_ref(); + let migration_0050: &str = MIGRATOR + .iter() + .find(|migration| migration.version == 50) + .expect("embedded migration 0050") + .sql + .as_ref(); let workspace_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) .parent() .and_then(std::path::Path::parent) @@ -1773,6 +1784,7 @@ mod postgres_tests { .expect("read schema/schema.sql"); let migration = surface(migration_0029); + let owner_admission_migration = surface(migration_0050); let schema = surface(&schema_sql); assert_eq!( @@ -1801,13 +1813,42 @@ mod postgres_tests { .functions .get(function) .unwrap_or_else(|| panic!("schema.sql is missing deletion function {function}")); - if function != "community_write_fence_excluded_table" { + if function != "community_write_fence_excluded_table" + && function != "prevent_community_deletion_request_retargeting" + { assert_eq!( in_schema, definition, "schema.sql definition of {function}() drifted from migration 0029" ); } } + assert_eq!( + schema + .functions + .get("prevent_community_deletion_request_retargeting") + .expect("schema.sql deletion retargeting guard"), + owner_admission_migration + .functions + .get("prevent_community_deletion_request_retargeting") + .expect("0050 deletion retargeting guard"), + "schema.sql must carry the latest immutable owner-provenance guard" + ); + let request_table = schema + .tables + .get("community_deletion_requests") + .expect("schema.sql deletion request table"); + for owner_provenance_fragment in [ + "request_origin text not null default 'operator'", + "owner_pubkey text", + "mediating_operator_pubkey text", + "acknowledgement_version integer", + "constraint community_deletion_owner_provenance check", + ] { + assert!( + request_table.contains(owner_provenance_fragment), + "schema.sql deletion requests are missing {owner_provenance_fragment}" + ); + } for (trigger, definition) in &migration.triggers { let in_schema = schema .triggers diff --git a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs index 6062528d550..5ca0185af56 100644 --- a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs +++ b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs @@ -311,7 +311,7 @@ async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer() production_result.is_ok(), "production migrator must preserve ingestion progress: {production_result:?}" ); - assert_eq!(version, 49); + assert_eq!(version, 50); assert_eq!(final_oid, oid, "prebuild must not be replaced"); assert_eq!(count, 4, "all writer witnesses must persist"); } diff --git a/crates/buzz-db/src/store/community.rs b/crates/buzz-db/src/store/community.rs index bcd38f4e5ce..461348a5cd3 100644 --- a/crates/buzz-db/src/store/community.rs +++ b/crates/buzz-db/src/store/community.rs @@ -81,6 +81,17 @@ pub struct UnarchivedCommunityRecord { pub host: String, } +/// Result of an owner-authorized unarchive attempt. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum UnarchiveCommunityResult { + /// The community is active, with archive state cleared idempotently. + Unarchived(UnarchivedCommunityRecord), + /// Durable deletion intent exists and wins over restoration. + DeletionPending, + /// The host is absent, unavailable, or not owned by the asserted pubkey. + NotFound, +} + impl Db { /// Returns the community mapped to a normalized request host, if one exists. /// @@ -209,6 +220,10 @@ impl Db { JOIN relay_members rm ON rm.community_id = c.id WHERE rm.pubkey = $1 AND rm.role = 'owner' + AND NOT EXISTS ( + SELECT 1 FROM community_deletion_requests request + WHERE request.community_id = c.id AND request.stage <> 'aborted' + ) ORDER BY c.created_at ASC, c.host ASC "#, ) @@ -531,40 +546,69 @@ impl Db { } /// Idempotently restores a community when the asserted pubkey is its current owner. + /// + /// Locks the community row so owner-deletion admission and restoration have + /// one serial order. A non-aborted deletion request returns + /// [`UnarchiveCommunityResult::DeletionPending`] without clearing archive state. #[datastore_span(name = "unarchive_community_owned_by", system = "postgresql")] pub async fn unarchive_community_owned_by( &self, normalized_host: &str, owner_pubkey: &str, - ) -> Result> { - let mut connection = crate::observability::acquire_writer( + ) -> Result { + let connection = crate::observability::acquire_writer( &self.pool, crate::observability::WriterOperation::Authorization, ) .await?; - let row = sqlx::query( - r#"UPDATE communities c - SET archived_at = NULL - FROM relay_members rm - WHERE lower(c.host) = lower($1) - AND rm.community_id = c.id - AND lower(rm.pubkey) = lower($2) - AND rm.role = 'owner' - AND c.deletion_state = 'active' - AND c.deleted_at IS NULL - RETURNING c.id, c.host"#, + let mut tx = sqlx::Transaction::begin(connection, None).await?; + let target = sqlx::query( + "SELECT id, host FROM communities \ + WHERE lower(host) = lower($1) AND deletion_state = 'active' \ + AND deleted_at IS NULL FOR UPDATE", ) .bind(normalized_host) + .fetch_optional(&mut *tx) + .await?; + let Some(target) = target else { + tx.rollback().await?; + return Ok(UnarchiveCommunityResult::NotFound); + }; + let community_id: Uuid = target.try_get("id")?; + let is_owner: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM relay_members \ + WHERE community_id = $1 AND lower(pubkey) = lower($2) AND role = 'owner')", + ) + .bind(community_id) .bind(owner_pubkey) - .fetch_optional(&mut *connection) + .fetch_one(&mut *tx) .await?; - row.map(|row| { - Ok(UnarchivedCommunityRecord { - id: CommunityId::from_uuid(row.try_get("id")?), - host: row.try_get("host")?, - }) - }) - .transpose() + if !is_owner { + tx.rollback().await?; + return Ok(UnarchiveCommunityResult::NotFound); + } + let deletion_pending: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ + WHERE community_id = $1 AND stage <> 'aborted')", + ) + .bind(community_id) + .fetch_one(&mut *tx) + .await?; + if deletion_pending { + tx.rollback().await?; + return Ok(UnarchiveCommunityResult::DeletionPending); + } + sqlx::query("UPDATE communities SET archived_at = NULL WHERE id = $1") + .bind(community_id) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(UnarchiveCommunityResult::Unarchived( + UnarchivedCommunityRecord { + id: CommunityId::from_uuid(community_id), + host: target.try_get("host")?, + }, + )) } /// Returns the community that owns a channel, if the channel exists. @@ -900,22 +944,26 @@ mod postgres_tests { .is_none(), "archived communities must fail admission" ); - assert!(db - .unarchive_community_owned_by(&host, &outsider) - .await - .expect("wrong-owner unarchive") - .is_none()); - assert!(db - .unarchive_community_owned_by("missing.example", &owner) - .await - .expect("unknown-host unarchive") - .is_none()); + assert_eq!( + db.unarchive_community_owned_by(&host, &outsider) + .await + .expect("wrong-owner unarchive"), + UnarchiveCommunityResult::NotFound + ); + assert_eq!( + db.unarchive_community_owned_by("missing.example", &owner) + .await + .expect("unknown-host unarchive"), + UnarchiveCommunityResult::NotFound + ); let restored = db .unarchive_community_owned_by(&host.to_ascii_uppercase(), &owner) .await - .expect("unarchive community") - .expect("owned community"); + .expect("unarchive community"); + let UnarchiveCommunityResult::Unarchived(restored) = restored else { + panic!("expected owned community") + }; assert_eq!(restored.id, created.id); assert_eq!(restored.host, host); assert_eq!( @@ -938,9 +986,8 @@ mod postgres_tests { let retry = db .unarchive_community_owned_by(&host, &owner) .await - .expect("idempotent retry") - .expect("owned community"); - assert_eq!(retry, restored); + .expect("idempotent retry"); + assert_eq!(retry, UnarchiveCommunityResult::Unarchived(restored)); } #[tokio::test] diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 0e184e00d88..090aa973223 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -27,6 +27,8 @@ pub const POSTGRES_STORE_NAME: &str = "postgres"; pub const OBJECT_STORE_NAME: &str = "object_store"; /// Durable name of the Redis/cache manifest component. pub const REDIS_STORE_NAME: &str = "redis"; +/// Owner acknowledgement contract accepted by the first self-serve deletion API. +pub const OWNER_DELETION_ACKNOWLEDGEMENT_VERSION: i32 = 1; /// Deployment-global advisory-lock key serializing schema migration with /// destructive deletion. @@ -220,7 +222,7 @@ impl FromStr for DeletionStage { } /// Durable community deletion request. -#[derive(Debug, Clone, Serialize)] +#[derive(Debug, Clone, PartialEq, Serialize)] pub struct DeletionRequest { /// Request identifier. pub id: Uuid, @@ -233,8 +235,16 @@ pub struct DeletionRequest { pub stage: DeletionStage, /// Stage at which the current consecutive retry streak started. pub retry_stage: Option, - /// Operator identity that submitted the request. + /// Legacy display identity that submitted the request. pub requested_by: String, + /// Whether the request originated from an operator or authenticated owner intent. + pub request_origin: DeletionRequestOrigin, + /// Current owner identity authenticated at owner-request admission. + pub owner_pubkey: Option, + /// Deployment operator that mediated the authenticated owner intent. + pub mediating_operator_pubkey: Option, + /// Owner-facing destructive-action acknowledgement contract version. + pub acknowledgement_version: Option, /// Optional request reason. pub reason: Option, /// Frozen catalog manifest. @@ -283,6 +293,47 @@ pub struct DeletionRequest { pub completed_at: Option>, } +/// Durable provenance class for a community deletion request. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum DeletionRequestOrigin { + /// Request was submitted directly by a deployment operator. + Operator, + /// Request records authenticated owner intent mediated by an operator. + Owner, +} + +impl FromStr for DeletionRequestOrigin { + type Err = DbError; + + fn from_str(value: &str) -> std::result::Result { + match value { + "operator" => Ok(Self::Operator), + "owner" => Ok(Self::Owner), + other => Err(DbError::InvalidData(format!( + "unknown community deletion request origin: {other}" + ))), + } + } +} + +/// Result of atomically admitting authenticated owner deletion intent. +#[derive(Debug, Clone, PartialEq)] +pub enum OwnerDeletionAdmission { + /// A new request was created or the stable request UUID converged to its row. + Accepted(Box), + /// The exact host is absent or the asserted owner is no longer current. + NotFoundOrNotOwner, + /// The community exists and is current-owner controlled, but is not archived. + NotArchived, + /// The community is already quiescing, fenced, or deleted. + LifecycleConflict, + /// The request UUID targets different intent, or another active request exists. + RequestConflict, + /// The owner acknowledgement contract is not supported. + UnsupportedAcknowledgementVersion, +} + /// Frozen PostgreSQL catalog inventory. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct SchemaManifest { @@ -742,6 +793,131 @@ impl DeletionStore { } } + /// Atomically admit an archived current owner's deletion intent. + /// + /// `request_id` is both the durable request id and the caller's stable + /// correlation/idempotency identity. Replays return the existing request at + /// its current stage. This operation only persists intent; it never + /// inventories, approves, quiesces, or executes deletion. + pub async fn admit_owner_request( + &self, + normalized_community_host: &str, + owner_pubkey: &str, + mediating_operator_pubkey: &str, + acknowledgement_version: i32, + request_id: Uuid, + ) -> Result { + if acknowledgement_version != OWNER_DELETION_ACKNOWLEDGEMENT_VERSION { + return Ok(OwnerDeletionAdmission::UnsupportedAcknowledgementVersion); + } + let owner_pubkey = owner_pubkey.to_ascii_lowercase(); + let mediating_operator_pubkey = mediating_operator_pubkey.to_ascii_lowercase(); + let mut tx = self.pool.begin().await?; + + sqlx::query( + "SELECT pg_advisory_xact_lock(hashtextextended('buzz-owner-deletion-intent:' || $1::text, 0))", + ) + .bind(request_id) + .execute(&mut *tx) + .await?; + + if let Some(row) = sqlx::query("SELECT * FROM community_deletion_requests WHERE id = $1") + .bind(request_id) + .fetch_optional(&mut *tx) + .await? + { + let existing = row_to_request(row)?; + let converges = existing.community_host == normalized_community_host + && existing.request_origin == DeletionRequestOrigin::Owner + && existing.owner_pubkey.as_deref() == Some(owner_pubkey.as_str()) + && existing.acknowledgement_version == Some(acknowledgement_version); + tx.rollback().await?; + return Ok(if converges { + OwnerDeletionAdmission::Accepted(Box::new(existing)) + } else { + OwnerDeletionAdmission::RequestConflict + }); + } + + let target = sqlx::query( + "SELECT id, host, archived_at, deletion_state, deleted_at \ + FROM communities WHERE host = $1 FOR UPDATE", + ) + .bind(normalized_community_host) + .fetch_optional(&mut *tx) + .await?; + let Some(target) = target else { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::NotFoundOrNotOwner); + }; + let community_id: Uuid = target.try_get("id")?; + let canonical_host: String = target.try_get("host")?; + let deletion_state: String = target.try_get("deletion_state")?; + let deleted_at: Option> = target.try_get("deleted_at")?; + if deletion_state != "active" || deleted_at.is_some() { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::LifecycleConflict); + } + + let owner_exists = sqlx::query_scalar::<_, String>( + "SELECT pubkey FROM relay_members \ + WHERE community_id = $1 AND pubkey = $2 AND role = 'owner' FOR UPDATE", + ) + .bind(community_id) + .bind(&owner_pubkey) + .fetch_optional(&mut *tx) + .await? + .is_some(); + if !owner_exists { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::NotFoundOrNotOwner); + } + if target + .try_get::>, _>("archived_at")? + .is_none() + { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::NotArchived); + } + let active_request_exists: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ + WHERE community_id = $1 AND stage <> 'aborted')", + ) + .bind(community_id) + .fetch_one(&mut *tx) + .await?; + if active_request_exists { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::RequestConflict); + } + + let row = sqlx::query( + r#" + INSERT INTO community_deletion_requests ( + id, community_id, community_host, requested_by, request_origin, + owner_pubkey, mediating_operator_pubkey, acknowledgement_version + ) VALUES ($1, $2, $3, $4, 'owner', $4, $5, $6) + ON CONFLICT (community_id) WHERE stage <> 'aborted' DO NOTHING + RETURNING * + "#, + ) + .bind(request_id) + .bind(community_id) + .bind(canonical_host) + .bind(owner_pubkey) + .bind(mediating_operator_pubkey) + .bind(acknowledgement_version) + .fetch_optional(&mut *tx) + .await?; + let Some(row) = row else { + tx.rollback().await?; + return Ok(OwnerDeletionAdmission::RequestConflict); + }; + let request = row_to_request(row)?; + tx.commit().await?; + Ok(OwnerDeletionAdmission::Accepted(Box::new(request))) + } + /// List requests newest first with a hard bound. pub async fn list(&self, limit: i64) -> Result> { let rows = sqlx::query( @@ -3112,6 +3288,10 @@ fn row_to_request(row: sqlx::postgres::PgRow) -> Result { .map(|stage| stage.parse()) .transpose()?, requested_by: row.try_get("requested_by")?, + request_origin: row.try_get::("request_origin")?.parse()?, + owner_pubkey: row.try_get("owner_pubkey")?, + mediating_operator_pubkey: row.try_get("mediating_operator_pubkey")?, + acknowledgement_version: row.try_get("acknowledgement_version")?, reason: row.try_get("reason")?, schema_manifest: row.try_get("schema_manifest")?, storage_manifest: row.try_get("storage_manifest")?, @@ -3416,7 +3596,10 @@ mod tests { #[cfg(test)] mod postgres_tests { use super::*; - use crate::{CreateCommunityWithOwnerResult, Db, DbConfig}; + use crate::{ + relay_members::TransferResult, CreateCommunityWithOwnerResult, Db, DbConfig, + UnarchiveCommunityResult, + }; async fn store() -> (Db, DeletionStore) { let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") @@ -3485,6 +3668,220 @@ mod postgres_tests { (request, inventory) } + async fn archived_owned_community(db: &Db) -> (String, String, CommunityId) { + let host = format!("owner-delete-{}.example", Uuid::new_v4().simple()); + let owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let CreateCommunityWithOwnerResult::Created(created) = db + .create_community_with_owner(&host, &owner) + .await + .expect("create owned community") + else { + panic!("expected a fresh community") + }; + db.archive_community_owned_by(&host, &owner, "protected.example") + .await + .expect("archive community") + .expect("owned community"); + (host, owner, created.id) + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_admission_requires_current_owner_and_archived_active_target() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + + let admitted = store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit archived owner request"); + let OwnerDeletionAdmission::Accepted(request) = admitted else { + panic!("expected accepted owner request") + }; + assert_eq!(request.id, request_id); + assert_eq!(request.community_id, community); + assert_eq!(request.community_host, host); + assert_eq!(request.stage, DeletionStage::Submitted); + assert_eq!(request.request_origin, DeletionRequestOrigin::Owner); + assert_eq!(request.owner_pubkey.as_deref(), Some(owner.as_str())); + assert_eq!(request.mediating_operator_pubkey.as_deref(), Some(operator)); + assert_eq!(request.acknowledgement_version, Some(1)); + assert_eq!( + sqlx::query_scalar::<_, String>("SELECT deletion_state FROM communities WHERE id = $1") + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("community lifecycle"), + "active", + "admission must not prematurely quiesce the community" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_admission_rejects_non_archived_non_owner_and_stale_owner() { + let (db, store) = store().await; + let host = format!("owner-delete-active-{}.example", Uuid::new_v4().simple()); + let owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let replacement = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let outsider = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let CreateCommunityWithOwnerResult::Created(created) = db + .create_community_with_owner(&host, &owner) + .await + .expect("create community") + else { + panic!("expected fresh community") + }; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + + assert_eq!( + store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("non-archived admission result"), + OwnerDeletionAdmission::NotArchived + ); + db.archive_community_owned_by(&host, &owner, "protected.example") + .await + .expect("archive") + .expect("owned community"); + assert_eq!( + store + .admit_owner_request(&host, &outsider, operator, 1, Uuid::new_v4()) + .await + .expect("non-owner admission result"), + OwnerDeletionAdmission::NotFoundOrNotOwner + ); + assert!(matches!( + db.transfer_ownership(created.id, &replacement, &owner) + .await + .expect("rotate owner"), + TransferResult::Transferred { .. } + )); + assert_eq!( + store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("stale-owner admission result"), + OwnerDeletionAdmission::NotFoundOrNotOwner + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_admission_replay_converges_after_stage_advancement_and_rejects_retargeting() { + let (db, store) = store().await; + let (host, owner, _) = archived_owned_community(&db).await; + let (other_host, other_owner, _) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + let OwnerDeletionAdmission::Accepted(first) = store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("first admission") + else { + panic!("expected accepted request") + }; + let inventory = FrozenInventory { + schema: store + .inventory_schema(first.community_id) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(first.community_id), + }; + store + .freeze_inventory(first.id, &inventory) + .await + .expect("advance request"); + + let OwnerDeletionAdmission::Accepted(replayed) = store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("replay admission") + else { + panic!("expected converged replay") + }; + assert_eq!(replayed.id, first.id); + assert_eq!(replayed.stage, DeletionStage::Inventoried); + assert_eq!( + store + .admit_owner_request(&other_host, &other_owner, operator, 1, request_id) + .await + .expect("retargeting result"), + OwnerDeletionAdmission::RequestConflict + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn concurrent_owner_admission_duplicates_return_one_stable_request() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + let (first, second) = tokio::join!( + store.admit_owner_request(&host, &owner, operator, 1, request_id), + store.admit_owner_request(&host, &owner, operator, 1, request_id), + ); + let accepted_id = |result: Result| { + let OwnerDeletionAdmission::Accepted(request) = result.expect("admission") else { + panic!("expected accepted request") + }; + request.id + }; + assert_eq!(accepted_id(first), request_id); + assert_eq!(accepted_id(second), request_id); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM community_deletion_requests WHERE community_id = $1 AND stage <> 'aborted'" + ) + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("active request count"), + 1 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn accepted_owner_request_blocks_unarchive_and_transfer_and_suppresses_owner_list_row() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let new_owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let OwnerDeletionAdmission::Accepted(_) = store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("admit owner request") + else { + panic!("expected accepted request") + }; + + assert_eq!( + db.unarchive_community_owned_by(&host, &owner) + .await + .expect("unarchive result"), + UnarchiveCommunityResult::DeletionPending + ); + assert_eq!( + db.transfer_ownership(community, &new_owner, &owner) + .await + .expect("transfer result"), + TransferResult::DeletionPending + ); + assert!( + db.list_communities_owned_by(&owner) + .await + .expect("owner list") + .iter() + .all(|row| row.id != community), + "accepted deletion requests must not remain actionable archived rows" + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn approval_boundary_blocks_claim_until_exact_inventory_is_approved() { diff --git a/crates/buzz-db/src/store/relay_members.rs b/crates/buzz-db/src/store/relay_members.rs index ecde3924fef..a54adc139ee 100644 --- a/crates/buzz-db/src/store/relay_members.rs +++ b/crates/buzz-db/src/store/relay_members.rs @@ -444,6 +444,8 @@ pub enum TransferResult { /// concurrent transfer or owner rotation has already changed ownership. /// The caller must NOT retry blindly — re-read ownership and re-evaluate. OwnerConflict, + /// Durable deletion intent exists and wins over ownership mutation. + DeletionPending, /// The transferee already owns the maximum number of communities. /// Enforced atomically inside the transfer transaction so concurrent /// transfers to the same recipient cannot both pass the limit. @@ -500,13 +502,16 @@ pub fn owner_count_advisory_lock_key(pubkey_hex: &str) -> i64 { /// so that concurrent transfers to the same recipient serialize. The same /// lock key is also used by `Db::create_community_with_owner` to prevent /// transfer-vs-create races. -/// 2. Locks the current owner row `FOR UPDATE` and verifies +/// 2. Locks the community row and rejects any non-aborted deletion request. +/// Owner-deletion admission takes the same row lock, so whichever operation +/// commits first makes the other re-evaluate and conflict. +/// 3. Locks the current owner row `FOR UPDATE` and verifies /// `expected_owner_pubkey` matches. This prevents a stale-owner race where /// a delayed/retried request overwrites a completed transfer. -/// 3. Enforces the [`MAX_COMMUNITIES_PER_OWNER`] limit on the transferee by +/// 4. Enforces the [`MAX_COMMUNITIES_PER_OWNER`] limit on the transferee by /// counting owned communities inside the same transaction. -/// 4. Upserts `new_owner_pubkey` as `owner` (insert or promote). -/// 5. Demotes every other owner in this community to `member` — **not** +/// 5. Upserts `new_owner_pubkey` as `owner` (insert or promote). +/// 6. Demotes every other owner in this community to `member` — **not** /// `admin`, per product decision: the former owner retains no management /// capabilities. /// @@ -533,7 +538,29 @@ pub async fn transfer_ownership( ) .await?; - // 2. Lock the current owner row FOR UPDATE and verify the expected owner. + let community_exists = + sqlx::query_scalar::<_, Uuid>("SELECT id FROM communities WHERE id = $1 FOR UPDATE") + .bind(community.as_uuid()) + .fetch_optional(&mut *tx) + .await? + .is_some(); + if !community_exists { + tx.rollback().await?; + return Ok(TransferResult::NoOwner); + } + let deletion_pending: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ + WHERE community_id = $1 AND stage <> 'aborted')", + ) + .bind(community.as_uuid()) + .fetch_one(&mut *tx) + .await?; + if deletion_pending { + tx.rollback().await?; + return Ok(TransferResult::DeletionPending); + } + + // 3. Lock the current owner row FOR UPDATE and verify the expected owner. // FOR UPDATE prevents the stale-owner race: a concurrent transfer that // already changed the owner will block on this lock until our txn // completes (or vice versa), and the expected_owner check will fail. @@ -570,7 +597,7 @@ pub async fn transfer_ownership( existing_owners.iter().find(|p| **p != pubkey).cloned() }; - // 3. Enforce the transferee's community ownership limit inside the same + // 4. Enforce the transferee's community ownership limit inside the same // transaction that holds the advisory lock. This is the authoritative // check — kgoose's preflight count is advisory only. let owned_count: i64 = sqlx::query_scalar( @@ -585,7 +612,7 @@ pub async fn transfer_ownership( return Ok(TransferResult::LimitReached); } - // 4. Upsert the new owner. + // 5. Upsert the new owner. sqlx::query( "INSERT INTO relay_members (community_id, pubkey, role, added_by) \ VALUES ($1, $2, 'owner', NULL) \ @@ -596,7 +623,7 @@ pub async fn transfer_ownership( .execute(&mut *tx) .await?; - // 5. Demote all other owners to member (not admin). + // 6. Demote all other owners to member (not admin). sqlx::query( "UPDATE relay_members SET role = 'member', updated_at = now() \ WHERE community_id = $1 AND role = 'owner' AND pubkey <> $2", diff --git a/crates/buzz-relay/src/api/operator.rs b/crates/buzz-relay/src/api/operator.rs index 2c49ca6a5c3..c672f837227 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -203,6 +203,15 @@ pub struct ArchiveCommunityRequest { owner_pubkey: String, } +/// Authenticated owner intent mediated by a trusted deployment operator. +#[derive(Debug, Deserialize)] +pub struct DeleteCommunityRequest { + host: String, + owner_pubkey: String, + request_id: Uuid, + acknowledgement_version: i32, +} + /// Idempotently archive a community owned by the asserted end-user identity. pub async fn archive_community( State(state): State>, @@ -287,12 +296,23 @@ pub async fn unarchive_community( "invalid owner_pubkey: expected 64-char hex pubkey", ) })?; - let record = state + let result = state .db .unarchive_community_owned_by(&normalized_host, &owner) .await - .map_err(|e| internal_error(&format!("unarchive community: {e}")))? - .ok_or_else(|| api_error(StatusCode::NOT_FOUND, "community not found"))?; + .map_err(|e| internal_error(&format!("unarchive community: {e}")))?; + let record = match result { + buzz_db::UnarchiveCommunityResult::Unarchived(record) => record, + buzz_db::UnarchiveCommunityResult::DeletionPending => { + return Err(api_error( + StatusCode::CONFLICT, + "community deletion is pending", + )); + } + buzz_db::UnarchiveCommunityResult::NotFound => { + return Err(api_error(StatusCode::NOT_FOUND, "community not found")); + } + }; tracing::info!(community = %record.id, host = %record.host, "community unarchived"); Ok(Json(serde_json::json!({ "community_id": record.id.to_string(), @@ -302,6 +322,106 @@ pub async fn unarchive_community( }))) } +/// Persist authenticated owner deletion intent without executing deletion work. +/// +/// `POST /operator/communities/delete`, NIP-98 signed by a pubkey in +/// `RELAY_OPERATOR_PUBKEYS`, body: +/// +/// ```json +/// { +/// "host": "archived.communities.example", +/// "owner_pubkey": "<64-char hex>", +/// "request_id": "", +/// "acknowledgement_version": 1 +/// } +/// ``` +/// +/// The request UUID is the correlation/idempotency key. Acceptance is a fast +/// PostgreSQL-only transaction and returns `202`; inventory, approval, +/// quiescing, object-store access, and executor work remain asynchronous. +pub async fn delete_community( + State(state): State>, + headers: HeaderMap, + body: axum::body::Bytes, +) -> Result<(StatusCode, Json), (StatusCode, Json)> { + const PATH: &str = "/operator/communities/delete"; + let operator = + authorize_operator_request(&state, &headers, "POST", PATH, None, Some(&body)).await?; + let request: DeleteCommunityRequest = serde_json::from_slice(&body).map_err(|e| { + api_error( + StatusCode::BAD_REQUEST, + &format!("invalid delete-community JSON: {e}"), + ) + })?; + let normalized_host = normalize_candidate_host(&request.host) + .map_err(|msg| api_error(StatusCode::BAD_REQUEST, &msg))?; + let deployment_host = buzz_core::tenant::relay_url_authority(&state.config.relay_url); + if normalized_host == deployment_host { + return Err(api_error( + StatusCode::CONFLICT, + "the deployment community cannot be deleted", + )); + } + let owner = validate_pubkey_hex(&request.owner_pubkey).ok_or_else(|| { + api_error( + StatusCode::BAD_REQUEST, + "invalid owner_pubkey: expected 64-char hex pubkey", + ) + })?; + let operator = operator.to_hex(); + let admission = state + .db + .deletion_store() + .admit_owner_request( + &normalized_host, + &owner, + &operator, + request.acknowledgement_version, + request.request_id, + ) + .await + .map_err(|error| internal_error(&format!("admit owner deletion request: {error}")))?; + let accepted = match admission { + buzz_db::deletion::OwnerDeletionAdmission::Accepted(request) => request, + buzz_db::deletion::OwnerDeletionAdmission::NotFoundOrNotOwner => { + return Err(api_error(StatusCode::NOT_FOUND, "community not found")); + } + buzz_db::deletion::OwnerDeletionAdmission::NotArchived => { + return Err(api_error( + StatusCode::CONFLICT, + "community must be archived before deletion", + )); + } + buzz_db::deletion::OwnerDeletionAdmission::LifecycleConflict => { + return Err(api_error( + StatusCode::CONFLICT, + "community deletion lifecycle is already active", + )); + } + buzz_db::deletion::OwnerDeletionAdmission::RequestConflict => { + return Err(api_error( + StatusCode::CONFLICT, + "deletion request conflicts with existing intent", + )); + } + buzz_db::deletion::OwnerDeletionAdmission::UnsupportedAcknowledgementVersion => { + return Err(api_error( + StatusCode::BAD_REQUEST, + "unsupported acknowledgement_version", + )); + } + }; + Ok(( + StatusCode::ACCEPTED, + Json(serde_json::json!({ + "request_id": accepted.id, + "community_id": accepted.community_id.to_string(), + "host": accepted.community_host, + "status": accepted.stage.to_string(), + })), + )) +} + /// List communities where a pubkey currently holds the `owner` role. pub async fn list_owned_communities( State(state): State>, @@ -424,6 +544,12 @@ pub async fn transfer_community( "owner_conflict: the current owner no longer matches expected_owner_pubkey", )); } + buzz_db::relay_members::TransferResult::DeletionPending => { + return Err(api_error( + StatusCode::CONFLICT, + "community deletion is pending", + )); + } buzz_db::relay_members::TransferResult::LimitReached => { return Err(api_error( StatusCode::CONFLICT, @@ -669,6 +795,29 @@ mod postgres_tests { signed_operator_request(state, operator, "POST", "/operator/communities", Some(body)).await } + async fn archive_for_owner_deletion(state: &AppState, host: &str, owner: &Keys) { + state + .db + .archive_community_owned_by( + host, + &owner.public_key().to_hex(), + &buzz_core::tenant::relay_url_authority(&state.config.relay_url), + ) + .await + .expect("archive community") + .expect("owned community"); + } + + fn owner_delete_body(host: &str, owner: &Keys, request_id: Uuid) -> String { + serde_json::json!({ + "host": host, + "owner_pubkey": owner.public_key().to_hex(), + "request_id": request_id, + "acknowledgement_version": 1, + }) + .to_string() + } + fn is_member_tag(tag: &Tag, pubkey: &str, role: &str) -> bool { let values = tag.as_slice(); values.first().is_some_and(|value| value == "member") @@ -738,6 +887,125 @@ mod postgres_tests { assert_eq!(response.status(), StatusCode::FORBIDDEN); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_delete_endpoint_accepts_archived_owner_intent_without_running_inventory() { + let operator = Keys::generate(); + let owner = Keys::generate(); + let Some(state) = operator_test_state(std::slice::from_ref(&operator)).await else { + return; + }; + let host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &host, &owner).await; + let request_id = Uuid::new_v4(); + let response = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&host, &owner, request_id)), + ) + .await; + + assert_eq!(response.status(), StatusCode::ACCEPTED); + let json = read_json(response).await; + assert_eq!(json["request_id"], request_id.to_string()); + assert_eq!(json["status"], "submitted"); + let request = state + .db + .deletion_store() + .get(request_id) + .await + .expect("persisted deletion request"); + assert_eq!(request.stage, buzz_db::deletion::DeletionStage::Submitted); + assert!(request.inventory_manifest.is_none()); + assert!(request.inventory_digest.is_none()); + assert_eq!( + request.mediating_operator_pubkey.as_deref(), + Some(operator.public_key().to_hex().as_str()) + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_delete_endpoint_rejects_protected_host_and_malformed_or_mismatched_target() { + let operator = Keys::generate(); + let owner = Keys::generate(); + let Some(state) = operator_test_state(std::slice::from_ref(&operator)).await else { + return; + }; + let protected_host = buzz_core::tenant::relay_url_authority(&state.config.relay_url); + let protected = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&protected_host, &owner, Uuid::new_v4())), + ) + .await; + assert_eq!(protected.status(), StatusCode::CONFLICT); + + let malformed = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some( + serde_json::json!({ + "host": "https://not-an-authority.example/path", + "owner_pubkey": "not-a-pubkey", + "request_id": "not-a-uuid", + "acknowledgement_version": 1, + }) + .to_string(), + ), + ) + .await; + assert_eq!(malformed.status(), StatusCode::BAD_REQUEST); + + let host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &host, &owner).await; + let request_id = Uuid::new_v4(); + let first = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&host, &owner, request_id)), + ) + .await; + assert_eq!(first.status(), StatusCode::ACCEPTED); + let other_host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &other_host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &other_host, &owner).await; + let mismatched = signed_operator_request( + state, + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&other_host, &owner, request_id)), + ) + .await; + assert_eq!(mismatched.status(), StatusCode::CONFLICT); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn post_operator_body_requires_payload_tag() { diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 61aedf70be0..5c9bf04db6a 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -96,6 +96,10 @@ pub fn build_router(state: Arc) -> Router { "/operator/communities/unarchive", post(api::operator::unarchive_community), ) + .route( + "/operator/communities/delete", + post(api::operator::delete_community), + ) .route( "/operator/communities/availability", get(api::operator::community_availability), @@ -710,7 +714,7 @@ mod tests { async fn readiness_state(evaluator: Arc) -> Arc { let mut config = crate::config::Config::from_env().expect("default config loads"); config.require_relay_membership = false; - config.database_url = "postgres://buzz:buzz_dev@127.0.0.1:1/buzz".to_string(); + config.database_url = "postgres://buzz:buzz_dev@127.0.0.1:1/buzz".to_string(); // sadscan:disable np.postgres.1 -- local test-only credentials config.redis_url = "redis://127.0.0.1:1".to_string(); let pool = sqlx::PgPool::connect_lazy(&config.database_url).expect("lazy pg pool"); let db = buzz_db::Db::from_pool(pool.clone()); diff --git a/migrations/0050_owner_community_deletion_admission.sql b/migrations/0050_owner_community_deletion_admission.sql new file mode 100644 index 00000000000..3fc78a9d92c --- /dev/null +++ b/migrations/0050_owner_community_deletion_admission.sql @@ -0,0 +1,66 @@ +-- Structured provenance for owner-origin whole-community deletion requests. +-- +-- The existing request UUID is the stable correlation/idempotency identity. +-- Owner admission supplies that UUID instead of creating a second identity +-- column, while these bounded columns distinguish owner intent from the +-- deployment operator that mediated it. +SET LOCAL lock_timeout = '5s'; + +ALTER TABLE community_deletion_requests + ADD COLUMN request_origin TEXT NOT NULL DEFAULT 'operator' + CHECK (request_origin IN ('operator', 'owner')), + ADD COLUMN owner_pubkey TEXT, + ADD COLUMN mediating_operator_pubkey TEXT, + ADD COLUMN acknowledgement_version INTEGER, + ADD CONSTRAINT community_deletion_owner_provenance CHECK ( + (request_origin = 'operator' + AND owner_pubkey IS NULL + AND mediating_operator_pubkey IS NULL + AND acknowledgement_version IS NULL) + OR + (request_origin = 'owner' + AND owner_pubkey ~ '^[0-9a-f]{64}$' + AND mediating_operator_pubkey ~ '^[0-9a-f]{64}$' + AND acknowledgement_version BETWEEN 1 AND 32767 + AND requested_by = owner_pubkey) + ); + +CREATE OR REPLACE FUNCTION prevent_community_deletion_request_retargeting() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF NEW.community_id IS DISTINCT FROM OLD.community_id + OR NEW.community_host IS DISTINCT FROM OLD.community_host + THEN + RAISE EXCEPTION 'community deletion target identity is immutable' + USING ERRCODE = 'integrity_constraint_violation'; + END IF; + IF NEW.request_origin IS DISTINCT FROM OLD.request_origin + OR NEW.owner_pubkey IS DISTINCT FROM OLD.owner_pubkey + OR NEW.mediating_operator_pubkey IS DISTINCT FROM OLD.mediating_operator_pubkey + OR NEW.acknowledgement_version IS DISTINCT FROM OLD.acknowledgement_version + THEN + RAISE EXCEPTION 'community deletion request provenance is immutable' + USING ERRCODE = 'integrity_constraint_violation'; + END IF; + IF OLD.inventory_frozen_at IS NOT NULL AND ( + NEW.schema_manifest IS DISTINCT FROM OLD.schema_manifest + OR NEW.storage_manifest IS DISTINCT FROM OLD.storage_manifest + OR NEW.inventory_manifest IS DISTINCT FROM OLD.inventory_manifest + OR NEW.inventory_digest IS DISTINCT FROM OLD.inventory_digest + OR NEW.inventory_frozen_at IS DISTINCT FROM OLD.inventory_frozen_at + ) THEN + RAISE EXCEPTION 'frozen community deletion inventory is immutable' + USING ERRCODE = 'integrity_constraint_violation'; + END IF; + IF OLD.destructive_storage_frozen_at IS NOT NULL AND ( + NEW.destructive_storage_manifest IS DISTINCT FROM OLD.destructive_storage_manifest + OR NEW.destructive_storage_frozen_at IS DISTINCT FROM OLD.destructive_storage_frozen_at + ) THEN + RAISE EXCEPTION 'frozen destructive storage manifest is immutable' + USING ERRCODE = 'integrity_constraint_violation'; + END IF; + RETURN NEW; +END; +$$; diff --git a/schema/schema.sql b/schema/schema.sql index 1a0a849f7ac..d18e2eec513 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1232,6 +1232,11 @@ CREATE TABLE community_deletion_requests ( 'logically_verified', 'retention_pending', 'aborted' )), requested_by TEXT NOT NULL, + request_origin TEXT NOT NULL DEFAULT 'operator' + CHECK (request_origin IN ('operator', 'owner')), + owner_pubkey TEXT, + mediating_operator_pubkey TEXT, + acknowledgement_version INTEGER, reason TEXT, schema_manifest JSONB, storage_manifest JSONB, @@ -1268,6 +1273,18 @@ CREATE TABLE community_deletion_requests ( CHECK ((aborted_at IS NULL) = (aborted_by IS NULL)), CHECK ((aborted_at IS NULL) = (abort_reason IS NULL)), CHECK ((inventory_frozen_at IS NULL) = (inventory_digest IS NULL)), + CONSTRAINT community_deletion_owner_provenance CHECK ( + (request_origin = 'operator' + AND owner_pubkey IS NULL + AND mediating_operator_pubkey IS NULL + AND acknowledgement_version IS NULL) + OR + (request_origin = 'owner' + AND owner_pubkey ~ '^[0-9a-f]{64}$' + AND mediating_operator_pubkey ~ '^[0-9a-f]{64}$' + AND acknowledgement_version BETWEEN 1 AND 32767 + AND requested_by = owner_pubkey) + ), UNIQUE (id, community_id, inventory_digest) ); CREATE UNIQUE INDEX community_deletion_requests_active_community @@ -1293,7 +1310,7 @@ CREATE TABLE community_deletion_approvals ( ON DELETE RESTRICT ); -CREATE FUNCTION prevent_community_deletion_request_retargeting() +CREATE OR REPLACE FUNCTION prevent_community_deletion_request_retargeting() RETURNS trigger LANGUAGE plpgsql AS $$ @@ -1304,6 +1321,14 @@ BEGIN RAISE EXCEPTION 'community deletion target identity is immutable' USING ERRCODE = 'integrity_constraint_violation'; END IF; + IF NEW.request_origin IS DISTINCT FROM OLD.request_origin + OR NEW.owner_pubkey IS DISTINCT FROM OLD.owner_pubkey + OR NEW.mediating_operator_pubkey IS DISTINCT FROM OLD.mediating_operator_pubkey + OR NEW.acknowledgement_version IS DISTINCT FROM OLD.acknowledgement_version + THEN + RAISE EXCEPTION 'community deletion request provenance is immutable' + USING ERRCODE = 'integrity_constraint_violation'; + END IF; IF OLD.inventory_frozen_at IS NOT NULL AND ( NEW.schema_manifest IS DISTINCT FROM OLD.schema_manifest OR NEW.storage_manifest IS DISTINCT FROM OLD.storage_manifest From acf561219f4223bc6090ce5c82555ee294d11ea4 Mon Sep 17 00:00:00 2001 From: tornquist Date: Wed, 23 Sep 2026 00:05:22 +0000 Subject: [PATCH 02/10] Fence archived community owner rotation Co-authored-by: Codex Signed-off-by: tornquist --- ARCHITECTURE.md | 6 + crates/buzz-db/src/store/deletion.rs | 377 +++++++++++++++++- crates/buzz-db/src/store/relay_members.rs | 227 +++++++++-- crates/buzz-relay/src/api/operator.rs | 74 +++- .../src/handlers/community_provisioning.rs | 32 +- 5 files changed, 653 insertions(+), 63 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 7d7da313d7b..ede53a4b838 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -26,6 +26,12 @@ management lists suppress the archived row. Replaying the same UUID converges to its current stage; a different UUID conflicts with the existing one-active- request invariant until that request is aborted. +Ownership is mutable only while a community is active. Archiving freezes the +current owner. Normal transfer and deployment-root legacy convergence take the +same community-row lock as owner-deletion admission, then reject archived, +quiescing, deleted, or deletion-pending rotation without changing membership. +Initial owner bootstrap for a newly created community remains supported. + Buzz is a Rust monorepo, licensed Apache 2.0 under Block, Inc. --- diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 090aa973223..01e9a09c8b1 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -3597,9 +3597,11 @@ mod tests { mod postgres_tests { use super::*; use crate::{ - relay_members::TransferResult, CreateCommunityWithOwnerResult, Db, DbConfig, - UnarchiveCommunityResult, + relay_members::{ProvisionOwnerResult, TransferResult}, + CreateCommunityWithOwnerResult, Db, DbConfig, UnarchiveCommunityResult, }; + use sqlx::postgres::{PgConnectOptions, PgPoolOptions}; + use std::str::FromStr; async fn store() -> (Db, DeletionStore) { let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") @@ -3685,6 +3687,188 @@ mod postgres_tests { (host, owner, created.id) } + struct AdmissionInsertGate { + connection: sqlx::pool::PoolConnection, + trigger_name: String, + function_name: String, + first_key: i32, + second_key: i32, + } + + async fn install_admission_insert_gate(db: &Db, request_id: Uuid) -> AdmissionInsertGate { + let suffix = Uuid::new_v4().simple().to_string(); + let trigger_name = format!("owner_admission_gate_{suffix}"); + let function_name = format!("owner_admission_gate_fn_{suffix}"); + let first_key = (request_id.as_u128() as u32 & 0x7fff_ffff) as i32; + let second_key = ((request_id.as_u128() >> 32) as u32 & 0x7fff_ffff) as i32; + let mut connection = db.pool.acquire().await.expect("acquire gate connection"); + sqlx::query("SELECT pg_advisory_lock(712345, 193847)") + .execute(&mut *connection) + .await + .expect("serialize admission gate fixtures"); + sqlx::query("SELECT pg_advisory_lock($1, $2)") + .bind(first_key) + .bind(second_key) + .execute(&mut *connection) + .await + .expect("hold admission gate"); + sqlx::query(AssertSqlSafe(format!( + "CREATE FUNCTION {function_name}() RETURNS trigger LANGUAGE plpgsql AS $$ \ + BEGIN \ + IF NEW.id = '{request_id}'::uuid THEN \ + PERFORM pg_advisory_xact_lock({first_key}, {second_key}); \ + END IF; \ + RETURN NEW; \ + END $$" + ))) + .execute(&db.pool) + .await + .expect("install admission gate function"); + sqlx::query(AssertSqlSafe(format!( + "CREATE TRIGGER {trigger_name} BEFORE INSERT ON community_deletion_requests \ + FOR EACH ROW EXECUTE FUNCTION {function_name}()" + ))) + .execute(&db.pool) + .await + .expect("install admission gate trigger"); + AdmissionInsertGate { + connection, + trigger_name, + function_name, + first_key, + second_key, + } + } + + async fn wait_for_admission_gate(db: &Db, gate: &AdmissionInsertGate) { + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM pg_locks \ + WHERE locktype = 'advisory' AND classid = $1::oid AND objid = $2::oid \ + AND objsubid = 2 AND NOT granted)", + ) + .bind(gate.first_key) + .bind(gate.second_key) + .fetch_one(&db.pool) + .await + .expect("inspect admission gate"); + if waiting { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("owner admission reached gated insert"); + } + + async fn release_admission_insert_gate(gate: &mut AdmissionInsertGate) { + sqlx::query("SELECT pg_advisory_unlock($1, $2)") + .bind(gate.first_key) + .bind(gate.second_key) + .execute(&mut *gate.connection) + .await + .expect("release admission gate"); + } + + async fn remove_admission_insert_gate(db: &Db, mut gate: AdmissionInsertGate) { + sqlx::query(AssertSqlSafe(format!( + "DROP TRIGGER {} ON community_deletion_requests", + gate.trigger_name + ))) + .execute(&db.pool) + .await + .expect("remove admission gate trigger"); + sqlx::query(AssertSqlSafe(format!( + "DROP FUNCTION {}()", + gate.function_name + ))) + .execute(&db.pool) + .await + .expect("remove admission gate function"); + sqlx::query("SELECT pg_advisory_unlock(712345, 193847)") + .execute(&mut *gate.connection) + .await + .expect("release admission fixture serialization"); + } + + async fn contender_db(application_name: &str) -> Db { + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()); + let options = PgConnectOptions::from_str(&database_url) + .expect("parse test database URL") + .application_name(application_name); + let pool = PgPoolOptions::new() + .max_connections(1) + .connect_with(options) + .await + .expect("connect contender DB"); + Db::from_pool(pool) + } + + async fn wait_for_contender_lock(db: &Db, application_name: &str) { + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM pg_stat_activity \ + WHERE datname = current_database() AND application_name = $1 \ + AND wait_event_type = 'Lock')", + ) + .bind(application_name) + .fetch_one(&db.pool) + .await + .expect("inspect contender lock"); + if waiting { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("contender reached community row lock"); + } + + async fn membership_roles(db: &Db, community: CommunityId) -> Vec<(String, String)> { + sqlx::query_as( + "SELECT pubkey, role FROM relay_members WHERE community_id = $1 ORDER BY pubkey", + ) + .bind(community.as_uuid()) + .fetch_all(&db.pool) + .await + .expect("read membership roles") + } + + async fn assert_owner_admission_is_only_committed_mutation( + db: &Db, + community: CommunityId, + request_id: Uuid, + ) { + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM community_deletion_requests \ + WHERE id = $1 AND community_id = $2 AND stage = 'submitted'", + ) + .bind(request_id) + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("count admitted request"), + 1 + ); + assert!( + sqlx::query_scalar::<_, bool>( + "SELECT archived_at IS NOT NULL FROM communities WHERE id = $1", + ) + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("read archive state"), + "the losing mutation must not clear archive state" + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn owner_admission_requires_current_owner_and_archived_active_target() { @@ -3754,19 +3938,196 @@ mod postgres_tests { .expect("non-owner admission result"), OwnerDeletionAdmission::NotFoundOrNotOwner ); - assert!(matches!( + assert_eq!( db.transfer_ownership(created.id, &replacement, &owner) .await - .expect("rotate owner"), - TransferResult::Transferred { .. } + .expect("reject archived owner rotation"), + TransferResult::LifecycleConflict + ); + assert!(matches!( + store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("current-owner admission result"), + OwnerDeletionAdmission::Accepted(_) )); - assert_eq!( + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn accepted_owner_deletion_blocks_legacy_owner_convergence_without_membership_change() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let replacement = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let before = membership_roles(&db, community).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + assert!(matches!( store .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) .await - .expect("stale-owner admission result"), - OwnerDeletionAdmission::NotFoundOrNotOwner + .expect("admit owner request"), + OwnerDeletionAdmission::Accepted(_) + )); + + assert_eq!( + db.provision_owner(community, &replacement) + .await + .expect("legacy convergence result"), + ProvisionOwnerResult::DeletionPending ); + assert_eq!(membership_roles(&db, community).await, before); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn concurrent_owner_admission_serializes_before_normal_transfer() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let replacement = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let before = membership_roles(&db, community).await; + let request_id = Uuid::new_v4(); + let mut gate = install_admission_insert_gate(&db, request_id).await; + let admission = tokio::spawn({ + let store = store.clone(); + let host = host.clone(); + let owner = owner.clone(); + async move { + store + .admit_owner_request( + &host, + &owner, + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + 1, + request_id, + ) + .await + } + }); + wait_for_admission_gate(&db, &gate).await; + + let application_name = format!("owner-transfer-{}", Uuid::new_v4().simple()); + let contender = contender_db(&application_name).await; + let transfer = tokio::spawn({ + let owner = owner.clone(); + let replacement = replacement.clone(); + async move { + contender + .transfer_ownership(community, &replacement, &owner) + .await + } + }); + wait_for_contender_lock(&db, &application_name).await; + release_admission_insert_gate(&mut gate).await; + + let admission_result = admission.await.expect("join admission").expect("admission"); + let transfer_result = transfer.await.expect("join transfer").expect("transfer"); + remove_admission_insert_gate(&db, gate).await; + assert!(matches!( + admission_result, + OwnerDeletionAdmission::Accepted(_) + )); + assert_eq!(transfer_result, TransferResult::DeletionPending); + assert_eq!(membership_roles(&db, community).await, before); + assert_owner_admission_is_only_committed_mutation(&db, community, request_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn concurrent_owner_admission_serializes_before_unarchive() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let before = membership_roles(&db, community).await; + let request_id = Uuid::new_v4(); + let mut gate = install_admission_insert_gate(&db, request_id).await; + let admission = tokio::spawn({ + let store = store.clone(); + let host = host.clone(); + let owner = owner.clone(); + async move { + store + .admit_owner_request( + &host, + &owner, + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + 1, + request_id, + ) + .await + } + }); + wait_for_admission_gate(&db, &gate).await; + + let application_name = format!("owner-unarchive-{}", Uuid::new_v4().simple()); + let contender = contender_db(&application_name).await; + let unarchive = tokio::spawn({ + let host = host.clone(); + let owner = owner.clone(); + async move { contender.unarchive_community_owned_by(&host, &owner).await } + }); + wait_for_contender_lock(&db, &application_name).await; + release_admission_insert_gate(&mut gate).await; + + let admission_result = admission.await.expect("join admission").expect("admission"); + let unarchive_result = unarchive.await.expect("join unarchive").expect("unarchive"); + remove_admission_insert_gate(&db, gate).await; + assert!(matches!( + admission_result, + OwnerDeletionAdmission::Accepted(_) + )); + assert_eq!(unarchive_result, UnarchiveCommunityResult::DeletionPending); + assert_eq!(membership_roles(&db, community).await, before); + assert_owner_admission_is_only_committed_mutation(&db, community, request_id).await; + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn concurrent_owner_admission_serializes_before_legacy_owner_rotation() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let replacement = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let before = membership_roles(&db, community).await; + let request_id = Uuid::new_v4(); + let mut gate = install_admission_insert_gate(&db, request_id).await; + let admission = tokio::spawn({ + let store = store.clone(); + let host = host.clone(); + let owner = owner.clone(); + async move { + store + .admit_owner_request( + &host, + &owner, + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + 1, + request_id, + ) + .await + } + }); + wait_for_admission_gate(&db, &gate).await; + + let application_name = format!("owner-legacy-{}", Uuid::new_v4().simple()); + let contender = contender_db(&application_name).await; + let provision = tokio::spawn({ + let replacement = replacement.clone(); + async move { contender.provision_owner(community, &replacement).await } + }); + wait_for_contender_lock(&db, &application_name).await; + release_admission_insert_gate(&mut gate).await; + + let admission_result = admission.await.expect("join admission").expect("admission"); + let provision_result = provision + .await + .expect("join legacy convergence") + .expect("legacy convergence"); + remove_admission_insert_gate(&db, gate).await; + assert!(matches!( + admission_result, + OwnerDeletionAdmission::Accepted(_) + )); + assert_eq!(provision_result, ProvisionOwnerResult::DeletionPending); + assert_eq!(membership_roles(&db, community).await, before); + assert_owner_admission_is_only_committed_mutation(&db, community, request_id).await; } #[tokio::test] diff --git a/crates/buzz-db/src/store/relay_members.rs b/crates/buzz-db/src/store/relay_members.rs index a54adc139ee..3bd09c625fe 100644 --- a/crates/buzz-db/src/store/relay_members.rs +++ b/crates/buzz-db/src/store/relay_members.rs @@ -366,8 +366,10 @@ pub async fn update_relay_member_role( /// Ensures the configured owner pubkey holds the `"owner"` role *in /// `community`*, and demotes any other owners in that community to `"admin"`. /// This handles owner rotation: if `RELAY_OWNER_PUBKEY` changes, the old owner -/// is automatically demoted. Scoped to one community — an owner of community A -/// is never bootstrapped into community B. +/// is automatically demoted only while the community is active and has no +/// durable deletion intent. Scoped to one community — an owner of community A +/// is never bootstrapped into community B. Initial insertion and exact-owner +/// convergence remain allowed because neither rotates existing ownership. /// /// Runs in a single transaction. Safe to call at every startup — idempotent. /// @@ -383,13 +385,92 @@ pub async fn bootstrap_owner( community: CommunityId, owner_pubkey: &str, ) -> Result<()> { - bootstrap_owner_with_operation( + match bootstrap_owner_with_operation( pool, community, owner_pubkey, observability::WriterOperation::Bootstrap, ) - .await + .await? + { + ProvisionOwnerResult::Applied => Ok(()), + ProvisionOwnerResult::LifecycleConflict => Err(DbError::AccessDenied( + "community lifecycle freezes owner rotation".to_string(), + )), + ProvisionOwnerResult::DeletionPending => Err(DbError::AccessDenied( + "community deletion is pending".to_string(), + )), + } +} + +#[derive(Clone, Copy)] +enum OwnerMutationMode { + Transfer, + Converge, +} + +enum OwnerMutationAdmission { + Allowed(Vec), + NotFound, + LifecycleConflict, + DeletionPending, +} + +async fn lock_owner_mutation_admission( + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + community: CommunityId, + proposed_owner: &str, + mode: OwnerMutationMode, +) -> Result { + let target = sqlx::query( + "SELECT archived_at, deletion_state, deleted_at FROM communities \ + WHERE id = $1 FOR UPDATE", + ) + .bind(community.as_uuid()) + .fetch_optional(&mut **tx) + .await?; + let Some(target) = target else { + return Ok(OwnerMutationAdmission::NotFound); + }; + let existing_owners: Vec = sqlx::query_scalar( + "SELECT pubkey FROM relay_members \ + WHERE community_id = $1 AND role = 'owner' \ + ORDER BY pubkey FOR UPDATE", + ) + .bind(community.as_uuid()) + .fetch_all(&mut **tx) + .await?; + + let changes_existing_owner = match mode { + OwnerMutationMode::Transfer => true, + OwnerMutationMode::Converge => { + !(existing_owners.is_empty() + || existing_owners.len() == 1 && existing_owners[0] == proposed_owner) + } + }; + if !changes_existing_owner { + return Ok(OwnerMutationAdmission::Allowed(existing_owners)); + } + + let deletion_pending: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ + WHERE community_id = $1 AND stage <> 'aborted')", + ) + .bind(community.as_uuid()) + .fetch_one(&mut **tx) + .await?; + if deletion_pending { + return Ok(OwnerMutationAdmission::DeletionPending); + } + + let archived_at: Option> = target.try_get("archived_at")?; + let deletion_state: String = target.try_get("deletion_state")?; + let deleted_at: Option> = target.try_get("deleted_at")?; + if archived_at.is_some() || deletion_state != "active" || deleted_at.is_some() { + return Ok(OwnerMutationAdmission::LifecycleConflict); + } + + Ok(OwnerMutationAdmission::Allowed(existing_owners)) } async fn bootstrap_owner_with_operation( @@ -397,11 +478,25 @@ async fn bootstrap_owner_with_operation( community: CommunityId, owner_pubkey: &str, operation: observability::WriterOperation, -) -> Result<()> { +) -> Result { let pubkey = owner_pubkey.to_ascii_lowercase(); let connection = observability::acquire_writer(pool, operation).await?; let mut tx = sqlx::Transaction::begin(connection, None).await?; + match lock_owner_mutation_admission(&mut tx, community, &pubkey, OwnerMutationMode::Converge) + .await? + { + OwnerMutationAdmission::Allowed(_) => {} + OwnerMutationAdmission::NotFound | OwnerMutationAdmission::LifecycleConflict => { + tx.rollback().await?; + return Ok(ProvisionOwnerResult::LifecycleConflict); + } + OwnerMutationAdmission::DeletionPending => { + tx.rollback().await?; + return Ok(ProvisionOwnerResult::DeletionPending); + } + } + // 1. Upsert the configured owner for this community. sqlx::query( "INSERT INTO relay_members (community_id, pubkey, role, added_by) \ @@ -424,7 +519,7 @@ async fn bootstrap_owner_with_operation( .await?; tx.commit().await?; - Ok(()) + Ok(ProvisionOwnerResult::Applied) } /// The result of a transfer-ownership attempt. @@ -444,6 +539,8 @@ pub enum TransferResult { /// concurrent transfer or owner rotation has already changed ownership. /// The caller must NOT retry blindly — re-read ownership and re-evaluate. OwnerConflict, + /// The community is archived, quiescing, or deleted; ownership is frozen. + LifecycleConflict, /// Durable deletion intent exists and wins over ownership mutation. DeletionPending, /// The transferee already owns the maximum number of communities. @@ -452,6 +549,17 @@ pub enum TransferResult { LimitReached, } +/// Result of converging an owner through deployment-root provisioning. +#[derive(Debug, PartialEq, Eq)] +pub enum ProvisionOwnerResult { + /// The initial owner was inserted or existing ownership converged. + Applied, + /// The community lifecycle freezes ownership mutation. + LifecycleConflict, + /// Durable deletion intent exists and wins over owner convergence. + DeletionPending, +} + /// Default maximum number of communities a single pubkey can own. Enforced at /// the relay layer — the authoritative layer — so that concurrent transfers or /// transfer-vs-create races cannot both pass a preflight count. @@ -502,7 +610,8 @@ pub fn owner_count_advisory_lock_key(pubkey_hex: &str) -> i64 { /// so that concurrent transfers to the same recipient serialize. The same /// lock key is also used by `Db::create_community_with_owner` to prevent /// transfer-vs-create races. -/// 2. Locks the community row and rejects any non-aborted deletion request. +/// 2. Locks the community row and rejects archive/non-active lifecycle or any +/// non-aborted deletion request. /// Owner-deletion admission takes the same row lock, so whichever operation /// commits first makes the other re-evaluate and conflict. /// 3. Locks the current owner row `FOR UPDATE` and verifies @@ -538,40 +647,28 @@ pub async fn transfer_ownership( ) .await?; - let community_exists = - sqlx::query_scalar::<_, Uuid>("SELECT id FROM communities WHERE id = $1 FOR UPDATE") - .bind(community.as_uuid()) - .fetch_optional(&mut *tx) - .await? - .is_some(); - if !community_exists { - tx.rollback().await?; - return Ok(TransferResult::NoOwner); - } - let deletion_pending: bool = sqlx::query_scalar( - "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ - WHERE community_id = $1 AND stage <> 'aborted')", - ) - .bind(community.as_uuid()) - .fetch_one(&mut *tx) - .await?; - if deletion_pending { - tx.rollback().await?; - return Ok(TransferResult::DeletionPending); - } - - // 3. Lock the current owner row FOR UPDATE and verify the expected owner. - // FOR UPDATE prevents the stale-owner race: a concurrent transfer that - // already changed the owner will block on this lock until our txn - // completes (or vice versa), and the expected_owner check will fail. - let existing_owners: Vec = sqlx::query_scalar( - "SELECT pubkey FROM relay_members \ - WHERE community_id = $1 AND role = 'owner' \ - FOR UPDATE", + let existing_owners = match lock_owner_mutation_admission( + &mut tx, + community, + &pubkey, + OwnerMutationMode::Transfer, ) - .bind(community.as_uuid()) - .fetch_all(&mut *tx) - .await?; + .await? + { + OwnerMutationAdmission::Allowed(existing_owners) => existing_owners, + OwnerMutationAdmission::NotFound => { + tx.rollback().await?; + return Ok(TransferResult::NoOwner); + } + OwnerMutationAdmission::LifecycleConflict => { + tx.rollback().await?; + return Ok(TransferResult::LifecycleConflict); + } + OwnerMutationAdmission::DeletionPending => { + tx.rollback().await?; + return Ok(TransferResult::DeletionPending); + } + }; if existing_owners.is_empty() { tx.rollback().await?; @@ -825,7 +922,11 @@ impl Db { /// Ensure an owner during operator-driven community provisioning. #[datastore_span(name = "provision_owner", system = "postgresql")] - pub async fn provision_owner(&self, community: CommunityId, owner_pubkey: &str) -> Result<()> { + pub async fn provision_owner( + &self, + community: CommunityId, + owner_pubkey: &str, + ) -> Result { bootstrap_owner_with_operation( &self.pool, community, @@ -1368,6 +1469,50 @@ mod postgres_tests { assert_role(&pool, community, &owner, "owner").await; } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn transfer_ownership_rejects_archived_community_without_membership_change() { + let pool = setup_pool().await; + let (community, owner) = owned_community(&pool).await; + let replacement = test_pubkey(); + sqlx::query("UPDATE communities SET archived_at = now() WHERE id = $1") + .bind(community.as_uuid()) + .execute(&pool) + .await + .expect("archive community"); + + let result = transfer_ownership(&pool, community, &replacement, &owner) + .await + .expect("transfer archived community"); + + assert_eq!(result, TransferResult::LifecycleConflict); + assert_role(&pool, community, &owner, "owner").await; + assert!( + get_relay_member(&pool, community, &replacement) + .await + .expect("get replacement") + .is_none(), + "archived transfer must not add the replacement owner" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn provision_owner_still_bootstraps_initial_owner() { + let pool = setup_pool().await; + let community = make_test_community(&pool).await; + let owner = test_pubkey(); + let db = Db::from_pool(pool.clone()); + + assert_eq!( + db.provision_owner(community, &owner) + .await + .expect("provision initial owner"), + ProvisionOwnerResult::Applied + ); + assert_role(&pool, community, &owner, "owner").await; + } + /// Transferring a community with no owner row returns `NoOwner`. #[tokio::test] #[ignore = "requires Postgres"] diff --git a/crates/buzz-relay/src/api/operator.rs b/crates/buzz-relay/src/api/operator.rs index c672f837227..20c80efd3d9 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -182,7 +182,11 @@ pub async fn provision_community( Err(msg) if msg.starts_with("actor not authorized") => { Err(api_error(StatusCode::FORBIDDEN, &msg)) } - Err(msg) if msg == "community already exists" || msg.starts_with("limit_reached:") => { + Err(msg) + if msg == "community already exists" + || msg.starts_with("limit_reached:") + || msg.starts_with("owner_conflict:") => + { Err(api_error(StatusCode::CONFLICT, &msg)) } Err(msg) @@ -544,6 +548,12 @@ pub async fn transfer_community( "owner_conflict: the current owner no longer matches expected_owner_pubkey", )); } + buzz_db::relay_members::TransferResult::LifecycleConflict => { + return Err(api_error( + StatusCode::CONFLICT, + "community must be active to transfer ownership", + )); + } buzz_db::relay_members::TransferResult::DeletionPending => { return Err(api_error( StatusCode::CONFLICT, @@ -1346,6 +1356,68 @@ mod postgres_tests { assert_snapshot_roles(&state, community.id, &[(&owner_hex, "owner")]).await; } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn archived_legacy_owner_convergence_returns_conflict_without_membership_change() { + let operator = Keys::generate(); + let owner = Keys::generate(); + let replacement = Keys::generate(); + let Some(state) = operator_test_state(std::slice::from_ref(&operator)).await else { + return; + }; + let host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(state.clone(), &operator, &host, &owner) + .await + .status(), + StatusCode::OK + ); + let community = state + .db + .lookup_community_by_host(&host) + .await + .expect("lookup community") + .expect("community exists"); + archive_for_owner_deletion(&state, &host, &owner).await; + + let body = serde_json::json!({ + "host": host, + "initial_owner_pubkey": replacement.public_key().to_hex(), + "create_only": false, + }) + .to_string(); + let response = signed_operator_request( + state.clone(), + &operator, + "POST", + "/operator/communities", + Some(body), + ) + .await; + + assert_eq!(response.status(), StatusCode::CONFLICT); + assert_eq!( + read_json(response).await["error"], + "owner_conflict: community must be active to rotate ownership" + ); + assert_eq!( + state + .db + .get_relay_member(community.id, &owner.public_key().to_hex()) + .await + .expect("get owner") + .expect("owner exists") + .role, + "owner" + ); + assert!(state + .db + .get_relay_member(community.id, &replacement.public_key().to_hex()) + .await + .expect("get replacement") + .is_none()); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn fresh_host_at_owner_limit_returns_limit_reached_conflict() { diff --git a/crates/buzz-relay/src/handlers/community_provisioning.rs b/crates/buzz-relay/src/handlers/community_provisioning.rs index 229b5f37161..1b13e6d921a 100644 --- a/crates/buzz-relay/src/handlers/community_provisioning.rs +++ b/crates/buzz-relay/src/handlers/community_provisioning.rs @@ -21,9 +21,9 @@ //! { "host": "acme.communities.buzz.xyz", "initial_owner_pubkey": "" } //! ``` //! -//! `initial_owner_pubkey` is optional. When present for an existing community, -//! it rotates that community owner through the same bootstrap path used by -//! `RELAY_OWNER_PUBKEY`; relay operators are deployment-root authorities. +//! `initial_owner_pubkey` is optional. Existing-community convergence may +//! rotate ownership only while the community is active and has no durable +//! deletion intent. Archive freezes ownership. use std::sync::Arc; @@ -238,14 +238,9 @@ async fn publish_membership_snapshot_if_required( /// signer here for the deployment-level `RELAY_OPERATOR_PUBKEYS` allowlist. /// /// Idempotency and owner semantics: the request is idempotent on the host row -/// (re-sending it never duplicates a community). When `initial_owner_pubkey` is -/// present, the owner is (re)bootstrapped via [`buzz_db::Db::bootstrap_owner`] -/// even if the community already existed — any previous owner is demoted to -/// admin, exactly like rotating `RELAY_OWNER_PUBKEY` for the deployment -/// community. This makes a retry after a partial failure (row created, owner -/// bootstrap crashed) converge, at the cost that an operator-signed request can -/// rotate an existing community's owner. The operator allowlist is therefore -/// documented as deployment-root authority, not create-only authority. +/// (re-sending it never duplicates a community). Initial owner bootstrap still +/// converges after a partial create, but rotating an existing owner requires an +/// active, non-archived community with no non-aborted deletion request. pub async fn provision_community( state: &Arc, operator_pubkey: &nostr::PublicKey, @@ -325,11 +320,22 @@ pub async fn provision_community( .map_err(|e| format!("failed to create community: {e}"))?; if let Some(owner_hex) = &initial_owner { - state + match state .db .provision_owner(record.id, owner_hex) .await - .map_err(|e| format!("community provisioned but owner bootstrap failed: {e}"))?; + .map_err(|e| format!("community provisioned but owner bootstrap failed: {e}"))? + { + buzz_db::relay_members::ProvisionOwnerResult::Applied => {} + buzz_db::relay_members::ProvisionOwnerResult::LifecycleConflict => { + return Err( + "owner_conflict: community must be active to rotate ownership".to_string(), + ); + } + buzz_db::relay_members::ProvisionOwnerResult::DeletionPending => { + return Err("owner_conflict: community deletion is pending".to_string()); + } + } publish_membership_snapshot_if_required(state, record.id, &record.host).await; } From 4432b5700067fb80692cfea19baee6cd362191e7 Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 18:03:38 +0000 Subject: [PATCH 03/10] Allow privileged abort at the reversible deletion boundary Owner admission records durable intent with no owner-facing cancellation, so an operator needs a recovery path when preparation cannot continue. Abort already reversed approved and fenced requests; extend it to the submitted and inventoried stages, which have destroyed nothing. Aborting releases the durable request fence over owner listing, unarchive, and owner rotation. It reverses deletion intent only: the community stays archived and the owner restores it explicitly. Stages from drained onward stay closed. Lock ordering and the active-serving-write-lease guard are unchanged. Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- crates/buzz-db/src/store/deletion.rs | 140 ++++++++++++++++++++++++++- 1 file changed, 138 insertions(+), 2 deletions(-) diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 01e9a09c8b1..c10ef5f6c1b 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -2186,7 +2186,16 @@ impl DeletionStore { Ok(()) } - /// Terminally abort an approved or fenced request before object deletion begins. + /// Terminally abort a request at the reversible pre-destruction boundary. + /// + /// `submitted`, `inventoried`, `approved`, and `fenced` are reversible: + /// nothing tenant-visible has been destroyed, so abort releases the durable + /// request fence over owner listing, unarchive, and owner rotation. Owner + /// admission deliberately has no owner-facing cancellation, so this + /// privileged path is the only recovery when preparation cannot continue. + /// Abort reverses deletion intent, not the owner's archive decision: the + /// community stays archived and the owner restores it explicitly. + /// Stages from `drained` onward have destroyed tenant state and stay closed. pub async fn abort( &self, request_id: Uuid, @@ -2225,7 +2234,10 @@ impl DeletionStore { } if !matches!( request.stage, - DeletionStage::Approved | DeletionStage::Fenced + DeletionStage::Submitted + | DeletionStage::Inventoried + | DeletionStage::Approved + | DeletionStage::Fenced ) { return Err(DbError::DeletionSafety(format!( "deletion {request_id} at stage {} cannot be aborted", @@ -4243,6 +4255,130 @@ mod postgres_tests { ); } + /// Privileged recovery abort at the reversible pre-approval boundary. + /// + /// Owner admission has no owner-facing cancellation. When preparation + /// cannot safely continue, an operator aborts the request; that must + /// release the durable request fence over listing, unarchive, and owner + /// rotation without silently unarchiving the community behind the owner. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn privileged_abort_of_submitted_owner_request_releases_the_request_fence() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let new_owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + let OwnerDeletionAdmission::Accepted(request) = store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request") + else { + panic!("expected accepted request") + }; + assert_eq!(request.stage, DeletionStage::Submitted); + assert_eq!( + db.transfer_ownership(community, &new_owner, &owner) + .await + .expect("fenced transfer"), + TransferResult::DeletionPending + ); + + let aborted = store + .abort( + request_id, + "recovery-operator", + "owner preparation cannot continue", + ) + .await + .expect("privileged abort at the reversible submitted boundary"); + assert_eq!(aborted.stage, DeletionStage::Aborted); + assert_eq!(aborted.aborted_by.as_deref(), Some("recovery-operator")); + + // Abort reverses deletion intent, not the owner's archive decision. + let (deletion_state, archived_at): (String, Option>) = + sqlx::query_as("SELECT deletion_state, archived_at FROM communities WHERE id = $1") + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("community lifecycle after abort"); + assert_eq!(deletion_state, "active"); + assert!( + archived_at.is_some(), + "abort must not unarchive the community on the owner's behalf" + ); + + assert!( + db.list_communities_owned_by(&owner) + .await + .expect("owner list after abort") + .iter() + .any(|row| row.id == community), + "aborting the request must restore the owner's actionable archived row" + ); + let UnarchiveCommunityResult::Unarchived(restored) = db + .unarchive_community_owned_by(&host, &owner) + .await + .expect("unarchive after abort") + else { + panic!("aborting the request must restore owner-authorized unarchive") + }; + assert_eq!(restored.id, community); + assert_eq!( + db.transfer_ownership(community, &new_owner, &owner) + .await + .expect("transfer after abort"), + TransferResult::Transferred { + previous_owner: Some(owner.clone()), + } + ); + } + + /// The reversible boundary stops at `inventoried`. Once execution has + /// destroyed anything, abort must stay closed. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn privileged_abort_spans_only_the_reversible_pre_destruction_boundary() { + let (db, store) = store().await; + let (request, _) = inventoried_request(&db, &store).await; + assert_eq!(request.stage, DeletionStage::Inventoried); + let aborted = store + .abort(request.id, "recovery-operator", "inventory needs recovery") + .await + .expect("privileged abort at the inventoried boundary"); + assert_eq!(aborted.stage, DeletionStage::Aborted); + + for irreversible in [ + DeletionStage::Drained, + DeletionStage::BindingsRemoved, + DeletionStage::PostgresPurged, + DeletionStage::CachePurged, + DeletionStage::LogicallyVerified, + DeletionStage::RetentionPending, + ] { + let (later, _) = inventoried_request(&db, &store).await; + sqlx::query("UPDATE community_deletion_requests SET stage = $2 WHERE id = $1") + .bind(later.id) + .bind(irreversible.to_string()) + .execute(&db.pool) + .await + .expect("advance stage"); + let error = store + .abort(later.id, "recovery-operator", "too late") + .await + .expect_err("abort must stay closed after destruction begins"); + assert!( + error.to_string().contains("cannot be aborted"), + "unexpected error at {irreversible}: {error}" + ); + assert_eq!( + store.get(later.id).await.expect("unchanged request").stage, + irreversible, + "a refused abort must not move the request" + ); + } + } + #[tokio::test] #[ignore = "requires Postgres"] async fn approval_boundary_blocks_claim_until_exact_inventory_is_approved() { From 3edd5948df816051cc32344f96841280c0aa3d9e Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 18:03:39 +0000 Subject: [PATCH 04/10] Reject owner deletion rows that omit provenance The owner branch of community_deletion_owner_provenance tested only the shape of each provenance column. A bare `col ~ '...'` on a NULL column evaluates to NULL, and a CHECK constraint is satisfied by NULL, so `FALSE OR NULL` admitted owner-origin rows with no owner key, no mediating operator, or no acknowledgement version. Reject a NULL in each column before testing its shape. Spell it `NOT (col IS NULL)`: pgschema drops a named CHECK whose body contains `IS NOT NULL` and still exits 0, which would leave the desired-state bootstrap silently unguarded while the migration path stayed correct. Assert one shared case table from both schema sources -- the migration upgrade path and the pgschema desired-state bootstrap -- so the two cannot drift into different owner-provenance guarantees. Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- crates/buzz-db/src/runtime/migration.rs | 16 ++ crates/buzz-db/src/store/deletion.rs | 212 ++++++++++++++++++ ...050_owner_community_deletion_admission.sql | 15 ++ schema/schema.sql | 3 + 4 files changed, 246 insertions(+) diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index a6727c4ac72..7b64547c2f5 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -2498,6 +2498,22 @@ mod postgres_tests { assert_eq!(after, vec![(1, Some(true)), (30_179, None), (30_350, None)]); } + /// Migration-upgrade half of the owner-provenance contract. + /// + /// The desired-state bootstrap half lives in + /// `store::deletion::postgres_tests` and asserts the same shared case + /// table, so `schema/schema.sql` cannot admit owner rows the migration + /// path refuses (or the reverse). + #[tokio::test] + #[ignore = "requires Postgres"] + async fn migrated_schema_enforces_owner_provenance_contract() { + let pool = connect_test_pool().await; + reset_public_schema(&pool).await; + run_migrations(&pool).await.expect("run migrations"); + + crate::store::deletion::owner_provenance_contract::assert_contract(&pool).await; + } + #[tokio::test] #[ignore = "requires Postgres"] async fn run_migrations_applies_consolidated_initial_schema_on_fresh_database() { diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index c10ef5f6c1b..073e9cafd08 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -3605,6 +3605,207 @@ mod tests { } } +/// Shared `community_deletion_owner_provenance` contract cases. +/// +/// The same table is asserted against the migration-upgrade schema +/// (`runtime::migration::postgres_tests`) and the desired-state bootstrap +/// schema (`postgres_tests` below), so the two schema sources cannot drift +/// into different owner-provenance guarantees. +#[cfg(test)] +pub(crate) mod owner_provenance_contract { + use sqlx::{PgPool, Row}; + use uuid::Uuid; + + /// Constraint that must reject every malformed owner-provenance row. + pub(crate) const CONSTRAINT: &str = "community_deletion_owner_provenance"; + + const VALID_OWNER: &str = "1111111111111111111111111111111111111111111111111111111111111111"; + const VALID_OPERATOR: &str = "2222222222222222222222222222222222222222222222222222222222222222"; + + /// One rejected owner-origin row: `(case name, owner, mediator, ack, requested_by)`. + type Case = ( + &'static str, + Option<&'static str>, + Option<&'static str>, + Option, + &'static str, + ); + + /// Every owner-origin row the constraint must refuse. + pub(crate) fn rejected_cases() -> Vec { + vec![ + ( + "missing owner_pubkey", + None, + Some(VALID_OPERATOR), + Some(1), + VALID_OWNER, + ), + ( + "missing mediating_operator_pubkey", + Some(VALID_OWNER), + None, + Some(1), + VALID_OWNER, + ), + ( + "missing acknowledgement_version", + Some(VALID_OWNER), + Some(VALID_OPERATOR), + None, + VALID_OWNER, + ), + ( + "owner_pubkey too short", + Some("abc"), + Some(VALID_OPERATOR), + Some(1), + "abc", + ), + ( + "owner_pubkey uppercase hex", + Some("AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"), + Some(VALID_OPERATOR), + Some(1), + "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA", + ), + ( + "owner_pubkey non-hex", + Some("zzzz111111111111111111111111111111111111111111111111111111111111"), + Some(VALID_OPERATOR), + Some(1), + "zzzz111111111111111111111111111111111111111111111111111111111111", + ), + ( + "mediating_operator_pubkey too short", + Some(VALID_OWNER), + Some("abc"), + Some(1), + VALID_OWNER, + ), + ( + "mediating_operator_pubkey uppercase hex", + Some(VALID_OWNER), + Some("BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB"), + Some(1), + VALID_OWNER, + ), + ( + "acknowledgement_version zero", + Some(VALID_OWNER), + Some(VALID_OPERATOR), + Some(0), + VALID_OWNER, + ), + ( + "acknowledgement_version negative", + Some(VALID_OWNER), + Some(VALID_OPERATOR), + Some(-1), + VALID_OWNER, + ), + ( + "requested_by is not the owner", + Some(VALID_OWNER), + Some(VALID_OPERATOR), + Some(1), + VALID_OPERATOR, + ), + ] + } + + async fn seed_community(pool: &PgPool) -> Uuid { + let host = format!("owner-provenance-{}.example", Uuid::new_v4().simple()); + sqlx::query_scalar::<_, Uuid>("INSERT INTO communities (host) VALUES ($1) RETURNING id") + .bind(&host) + .fetch_one(pool) + .await + .expect("seed owner-provenance community") + } + + async fn insert_owner_row( + pool: &PgPool, + owner: Option<&str>, + mediator: Option<&str>, + acknowledgement: Option, + requested_by: &str, + ) -> Result<(), sqlx::Error> { + let community_id = seed_community(pool).await; + let host: String = sqlx::query("SELECT host FROM communities WHERE id = $1") + .bind(community_id) + .fetch_one(pool) + .await + .expect("seeded host") + .try_get("host") + .expect("host column"); + sqlx::query( + "INSERT INTO community_deletion_requests \ + (id, community_id, community_host, requested_by, request_origin, \ + owner_pubkey, mediating_operator_pubkey, acknowledgement_version) \ + VALUES ($1, $2, $3, $4, 'owner', $5, $6, $7)", + ) + .bind(Uuid::new_v4()) + .bind(community_id) + .bind(host) + .bind(requested_by) + .bind(owner) + .bind(mediator) + .bind(acknowledgement) + .execute(pool) + .await + .map(|_| ()) + } + + /// Assert the live schema refuses every malformed owner row and still + /// admits a well-formed one. + pub(crate) async fn assert_contract(pool: &PgPool) { + for (name, owner, mediator, acknowledgement, requested_by) in rejected_cases() { + let error = insert_owner_row(pool, owner, mediator, acknowledgement, requested_by) + .await + .expect_err(&format!("owner-provenance case must be rejected: {name}")); + let constraint = error + .as_database_error() + .and_then(sqlx::error::DatabaseError::constraint); + assert_eq!( + constraint, + Some(CONSTRAINT), + "case {name:?} must fail {CONSTRAINT}, got: {error}" + ); + } + + // Falsifiability: the constraint must still admit well-formed intent. + insert_owner_row( + pool, + Some(VALID_OWNER), + Some(VALID_OPERATOR), + Some(1), + VALID_OWNER, + ) + .await + .expect("well-formed owner provenance must be accepted"); + + // The operator branch stays exclusive of owner columns. + let community_id = seed_community(pool).await; + let operator_with_owner_columns = sqlx::query( + "INSERT INTO community_deletion_requests \ + (id, community_id, community_host, requested_by, request_origin, owner_pubkey) \ + VALUES ($1, $2, 'operator.example', 'operator', 'operator', $3)", + ) + .bind(Uuid::new_v4()) + .bind(community_id) + .bind(VALID_OWNER) + .execute(pool) + .await; + assert_eq!( + operator_with_owner_columns + .expect_err("operator-origin rows must not carry owner provenance") + .as_database_error() + .and_then(sqlx::error::DatabaseError::constraint), + Some(CONSTRAINT) + ); + } +} + #[cfg(test)] mod postgres_tests { use super::*; @@ -4255,6 +4456,17 @@ mod postgres_tests { ); } + /// Desired-state bootstrap half of the owner-provenance contract. + /// + /// The migration-upgrade half lives in `runtime::migration::postgres_tests` + /// and asserts the same shared case table. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn desired_state_schema_enforces_owner_provenance_contract() { + let (db, _) = store().await; + owner_provenance_contract::assert_contract(&db.pool).await; + } + /// Privileged recovery abort at the reversible pre-approval boundary. /// /// Owner admission has no owner-facing cancellation. When preparation diff --git a/migrations/0050_owner_community_deletion_admission.sql b/migrations/0050_owner_community_deletion_admission.sql index 3fc78a9d92c..25002b32da0 100644 --- a/migrations/0050_owner_community_deletion_admission.sql +++ b/migrations/0050_owner_community_deletion_admission.sql @@ -4,6 +4,18 @@ -- Owner admission supplies that UUID instead of creating a second identity -- column, while these bounded columns distinguish owner intent from the -- deployment operator that mediated it. +-- +-- The owner branch rejects a NULL in every provenance column before testing +-- its shape. A bare `col ~ '...'` on a NULL column yields NULL, and a CHECK is +-- satisfied by NULL, so `FALSE OR NULL` would silently admit an owner-origin +-- row with no owner key, no mediating operator, or no acknowledgement version. +-- +-- The null rejection is spelled `NOT (col IS NULL)` rather than +-- `col IS NOT NULL`: pgschema drops a named CHECK whose body contains +-- `IS NOT NULL` and still exits 0, which would leave the desired-state +-- bootstrap in `schema/schema.sql` silently unguarded while the migration path +-- stayed correct. `store::deletion::owner_provenance_contract` asserts both +-- schema sources against one case table so that divergence fails loudly. SET LOCAL lock_timeout = '5s'; ALTER TABLE community_deletion_requests @@ -19,6 +31,9 @@ ALTER TABLE community_deletion_requests AND acknowledgement_version IS NULL) OR (request_origin = 'owner' + AND NOT (owner_pubkey IS NULL) + AND NOT (mediating_operator_pubkey IS NULL) + AND NOT (acknowledgement_version IS NULL) AND owner_pubkey ~ '^[0-9a-f]{64}$' AND mediating_operator_pubkey ~ '^[0-9a-f]{64}$' AND acknowledgement_version BETWEEN 1 AND 32767 diff --git a/schema/schema.sql b/schema/schema.sql index d18e2eec513..c91cd7375d6 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1280,6 +1280,9 @@ CREATE TABLE community_deletion_requests ( AND acknowledgement_version IS NULL) OR (request_origin = 'owner' + AND NOT (owner_pubkey IS NULL) + AND NOT (mediating_operator_pubkey IS NULL) + AND NOT (acknowledgement_version IS NULL) AND owner_pubkey ~ '^[0-9a-f]{64}$' AND mediating_operator_pubkey ~ '^[0-9a-f]{64}$' AND acknowledgement_version BETWEEN 1 AND 32767 From 2abeed3f65e0d9401e272c2b750c8e763ef99386 Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 18:03:39 +0000 Subject: [PATCH 05/10] Reach each owner-delete rejection guard separately MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The malformed-body coverage sent a single request carrying several invalid fields at once, including an unparseable `request_id`. Request deserialization runs before the host and pubkey guards, so that request was refused by serde and the guards it claimed to cover were never executed — dropping the `owner_pubkey` validation entirely left the test passing. Give each guard its own request whose remaining fields are valid, so the case reaches the guard under test, and assert no deletion intent persists for any refusal. Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- crates/buzz-relay/src/api/operator.rs | 133 +++++++++++++++++++++++++- 1 file changed, 131 insertions(+), 2 deletions(-) diff --git a/crates/buzz-relay/src/api/operator.rs b/crates/buzz-relay/src/api/operator.rs index 20c80efd3d9..ce792c2c2c8 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -790,6 +790,40 @@ mod postgres_tests { .expect("response") } + /// Send a delete-community request with a caller-controlled `Authorization` + /// header so signature-binding failures can be exercised directly. + async fn raw_owner_delete( + state: Arc, + auth: Option, + extra_header: Option<(&str, String)>, + body: String, + ) -> axum::response::Response { + let mut request = Request::builder() + .method("POST") + .uri("/operator/communities/delete") + .header(header::HOST, INGRESS_HOST) + .header(header::CONTENT_TYPE, "application/json"); + if let Some(auth) = auth { + request = request.header(header::AUTHORIZATION, auth); + } + if let Some((name, value)) = extra_header { + request = request.header(name, value); + } + build_router(state) + .oneshot(request.body(Body::from(body)).expect("request")) + .await + .expect("response") + } + + /// The endpoint must not have persisted intent under `request_id`. + async fn assert_no_persisted_request(state: &AppState, request_id: Uuid, case: &str) { + let found = state.db.deletion_store().get(request_id).await; + assert!( + matches!(found, Err(buzz_db::DbError::NotFound(_))), + "{case}: rejected request must not persist deletion intent" + ); + } + async fn provision_community( state: Arc, operator: &Keys, @@ -961,7 +995,14 @@ mod postgres_tests { .await; assert_eq!(protected.status(), StatusCode::CONFLICT); - let malformed = signed_operator_request( + // Each malformed case keeps every unrelated field valid so it reaches + // the guard under test instead of tripping an earlier one, and none of + // them may leave durable intent behind. + let valid_owner = owner.public_key().to_hex(); + let reachable_host = format!("community-{}.example", Uuid::new_v4().simple()); + + let bad_host_id = Uuid::new_v4(); + let bad_host = signed_operator_request( Arc::clone(&state), &operator, "POST", @@ -969,7 +1010,95 @@ mod postgres_tests { Some( serde_json::json!({ "host": "https://not-an-authority.example/path", + "owner_pubkey": valid_owner, + "request_id": bad_host_id, + "acknowledgement_version": 1, + }) + .to_string(), + ), + ) + .await; + assert_eq!(bad_host.status(), StatusCode::BAD_REQUEST); + assert_no_persisted_request(&state, bad_host_id, "unnormalizable host").await; + + let bad_pubkey_id = Uuid::new_v4(); + let bad_pubkey = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some( + serde_json::json!({ + "host": reachable_host, "owner_pubkey": "not-a-pubkey", + "request_id": bad_pubkey_id, + "acknowledgement_version": 1, + }) + .to_string(), + ), + ) + .await; + assert_eq!(bad_pubkey.status(), StatusCode::BAD_REQUEST); + let bad_pubkey_error = read_json(bad_pubkey).await; + assert!( + bad_pubkey_error + .get("error") + .and_then(Value::as_str) + .unwrap_or_default() + .contains("owner_pubkey"), + "invalid owner_pubkey must reach the pubkey guard: {bad_pubkey_error:?}" + ); + assert_no_persisted_request(&state, bad_pubkey_id, "invalid owner_pubkey").await; + + let bad_version_id = Uuid::new_v4(); + let bad_version = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some( + serde_json::json!({ + "host": reachable_host, + "owner_pubkey": valid_owner, + "request_id": bad_version_id, + "acknowledgement_version": 2, + }) + .to_string(), + ), + ) + .await; + assert_eq!(bad_version.status(), StatusCode::BAD_REQUEST); + assert_no_persisted_request(&state, bad_version_id, "unsupported acknowledgement").await; + + let unknown_host_id = Uuid::new_v4(); + let unknown_host = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some( + serde_json::json!({ + "host": reachable_host, + "owner_pubkey": valid_owner, + "request_id": unknown_host_id, + "acknowledgement_version": 1, + }) + .to_string(), + ), + ) + .await; + assert_eq!(unknown_host.status(), StatusCode::NOT_FOUND); + assert_no_persisted_request(&state, unknown_host_id, "unprovisioned host").await; + + let bad_uuid = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some( + serde_json::json!({ + "host": reachable_host, + "owner_pubkey": valid_owner, "request_id": "not-a-uuid", "acknowledgement_version": 1, }) @@ -977,7 +1106,7 @@ mod postgres_tests { ), ) .await; - assert_eq!(malformed.status(), StatusCode::BAD_REQUEST); + assert_eq!(bad_uuid.status(), StatusCode::BAD_REQUEST); let host = format!("community-{}.example", Uuid::new_v4().simple()); assert_eq!( From 1cef1ed116ade2f9812b66ed30fee1b7b6048c99 Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 18:03:39 +0000 Subject: [PATCH 06/10] Cover operator-bound signing and replay on owner delete `/operator/communities/delete` admits an irreversible request, but the endpoint's own tests only exercised the happy path and relied on generic operator-auth coverage elsewhere. Nothing pinned that this route refuses an unsigned caller, the X-Pubkey dev fallback, a non-operator signer, or a signature bound to a different method, URL, or body. Add direct coverage for each binding failure, and for idempotent replay: an identical resubmission converges on the stored request, while a resubmission that changes the owner, the host, or the acknowledgement version is refused and leaves the stored intent untouched. Every rejection asserts no deletion request was persisted. Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- crates/buzz-relay/src/api/operator.rs | 280 ++++++++++++++++++++++++++ 1 file changed, 280 insertions(+) diff --git a/crates/buzz-relay/src/api/operator.rs b/crates/buzz-relay/src/api/operator.rs index ce792c2c2c8..eea2dd7c1f1 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -1145,6 +1145,286 @@ mod postgres_tests { assert_eq!(mismatched.status(), StatusCode::CONFLICT); } + /// Every signature-binding failure mode on the owner-delete endpoint. + /// + /// The endpoint mediates an irreversible request, so a caller that cannot + /// prove operator authority over this exact method, URL, and body must be + /// refused without leaving durable intent behind. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_delete_endpoint_requires_operator_bound_nip98_signature() { + let operator = Keys::generate(); + let outsider = Keys::generate(); + let owner = Keys::generate(); + let Some(state) = operator_test_state(std::slice::from_ref(&operator)).await else { + return; + }; + let host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &host, &owner).await; + + let delete_url = format!("http://{INGRESS_HOST}/operator/communities/delete"); + + // 1. No Authorization header at all. + let unsigned_id = Uuid::new_v4(); + let unsigned = raw_owner_delete( + Arc::clone(&state), + None, + None, + owner_delete_body(&host, &owner, unsigned_id), + ) + .await; + assert_eq!(unsigned.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, unsigned_id, "unsigned").await; + + // 2. X-Pubkey only: the dev fallback is disabled for operator endpoints. + let x_pubkey_id = Uuid::new_v4(); + let x_pubkey_only = raw_owner_delete( + Arc::clone(&state), + None, + Some(("x-pubkey", operator.public_key().to_hex())), + owner_delete_body(&host, &owner, x_pubkey_id), + ) + .await; + assert_eq!(x_pubkey_only.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, x_pubkey_id, "X-Pubkey only").await; + + // 3. Valid NIP-98 from a key that is not an allowlisted operator. + let outsider_id = Uuid::new_v4(); + let outsider_body = owner_delete_body(&host, &owner, outsider_id); + let outsider_response = raw_owner_delete( + Arc::clone(&state), + Some(nip98_auth_header( + &outsider, + &delete_url, + "POST", + Some(outsider_body.as_bytes()), + )), + None, + outsider_body, + ) + .await; + assert_eq!(outsider_response.status(), StatusCode::FORBIDDEN); + assert_no_persisted_request(&state, outsider_id, "non-operator signer").await; + + // 4. Operator signature that omits the payload tag. + let no_payload_id = Uuid::new_v4(); + let no_payload = raw_owner_delete( + Arc::clone(&state), + Some(nip98_auth_header_without_payload( + &operator, + &delete_url, + "POST", + )), + None, + owner_delete_body(&host, &owner, no_payload_id), + ) + .await; + assert_eq!(no_payload.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, no_payload_id, "missing payload tag").await; + + // 5. Payload tag bound to a different body than the one sent. + let signed_id = Uuid::new_v4(); + let tampered_id = Uuid::new_v4(); + let signed_body = owner_delete_body(&host, &owner, signed_id); + let tampered = raw_owner_delete( + Arc::clone(&state), + Some(nip98_auth_header( + &operator, + &delete_url, + "POST", + Some(signed_body.as_bytes()), + )), + None, + owner_delete_body(&host, &owner, tampered_id), + ) + .await; + assert_eq!(tampered.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, tampered_id, "tampered payload").await; + assert_no_persisted_request(&state, signed_id, "tampered payload (signed id)").await; + + // 6. Signature bound to a different operator URL. + let wrong_url_id = Uuid::new_v4(); + let wrong_url_body = owner_delete_body(&host, &owner, wrong_url_id); + let wrong_url = raw_owner_delete( + Arc::clone(&state), + Some(nip98_auth_header( + &operator, + &format!("http://{INGRESS_HOST}/operator/communities"), + "POST", + Some(wrong_url_body.as_bytes()), + )), + None, + wrong_url_body, + ) + .await; + assert_eq!(wrong_url.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, wrong_url_id, "wrong signed URL").await; + + // 7. Signature bound to a different method. + let wrong_method_id = Uuid::new_v4(); + let wrong_method_body = owner_delete_body(&host, &owner, wrong_method_id); + let wrong_method = raw_owner_delete( + Arc::clone(&state), + Some(nip98_auth_header( + &operator, + &delete_url, + "GET", + Some(wrong_method_body.as_bytes()), + )), + None, + wrong_method_body, + ) + .await; + assert_eq!(wrong_method.status(), StatusCode::UNAUTHORIZED); + assert_no_persisted_request(&state, wrong_method_id, "wrong signed method").await; + + // Falsifiability: the same shape, correctly bound, is accepted. + let accepted_id = Uuid::new_v4(); + let accepted = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&host, &owner, accepted_id)), + ) + .await; + assert_eq!(accepted.status(), StatusCode::ACCEPTED); + } + + /// The request UUID is the correlation identity, so a replay that changes + /// any bound field is a conflict rather than a second interpretation. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_delete_replay_rejects_any_changed_field_without_mutating_intent() { + let operator = Keys::generate(); + let owner = Keys::generate(); + let other_owner = Keys::generate(); + let Some(state) = operator_test_state(std::slice::from_ref(&operator)).await else { + return; + }; + let host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &host, &owner).await; + + let other_host = format!("community-{}.example", Uuid::new_v4().simple()); + assert_eq!( + provision_community(Arc::clone(&state), &operator, &other_host, &owner) + .await + .status(), + StatusCode::OK + ); + archive_for_owner_deletion(&state, &other_host, &owner).await; + + let request_id = Uuid::new_v4(); + assert_eq!( + signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&host, &owner, request_id)), + ) + .await + .status(), + StatusCode::ACCEPTED + ); + let admitted = state + .db + .deletion_store() + .get(request_id) + .await + .expect("admitted request"); + + // Exactly one field differs per case; the rest replay verbatim. + let owner_hex = owner.public_key().to_hex(); + let cases: Vec<(&str, StatusCode, String)> = vec![ + ( + "changed owner_pubkey", + StatusCode::CONFLICT, + serde_json::json!({ + "host": host, + "owner_pubkey": other_owner.public_key().to_hex(), + "request_id": request_id, + "acknowledgement_version": 1, + }) + .to_string(), + ), + ( + "changed host", + StatusCode::CONFLICT, + serde_json::json!({ + "host": other_host, + "owner_pubkey": owner_hex, + "request_id": request_id, + "acknowledgement_version": 1, + }) + .to_string(), + ), + ( + "changed acknowledgement_version", + StatusCode::BAD_REQUEST, + serde_json::json!({ + "host": host, + "owner_pubkey": owner_hex, + "request_id": request_id, + "acknowledgement_version": 2, + }) + .to_string(), + ), + ]; + for (case, expected, body) in cases { + let response = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(body), + ) + .await; + assert_eq!(response.status(), expected, "{case}"); + let current = state + .db + .deletion_store() + .get(request_id) + .await + .expect("request still readable"); + assert_eq!(current.community_id, admitted.community_id, "{case}"); + assert_eq!(current.community_host, admitted.community_host, "{case}"); + assert_eq!(current.owner_pubkey, admitted.owner_pubkey, "{case}"); + assert_eq!( + current.acknowledgement_version, admitted.acknowledgement_version, + "{case}" + ); + assert_eq!(current.stage, admitted.stage, "{case}"); + } + + // An identical replay still converges on the same request. + let converged = signed_operator_request( + Arc::clone(&state), + &operator, + "POST", + "/operator/communities/delete", + Some(owner_delete_body(&host, &owner, request_id)), + ) + .await; + assert_eq!(converged.status(), StatusCode::ACCEPTED); + assert_eq!( + read_json(converged).await["request_id"], + request_id.to_string() + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn post_operator_body_requires_payload_tag() { From 8ebe3b881c98f19cfc93da7c1d00d4b0618a356e Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 18:03:39 +0000 Subject: [PATCH 07/10] Record that owner deletion consent is an operator assertion MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The admission path reads as if the relay had verified the owner's consent. It has not: the mediating operator authenticates the owner and collects the acknowledgement out of band, and the request reaching the relay carries only that operator's NIP-98 signature. The relay checks operator authority and current ownership, and stores the owner pubkey, operator pubkey, and acknowledgement version as provenance for the upstream ceremony — never an owner-signed attestation. Also record how manual handoff converges. Owner provenance pins `requested_by` to `owner_pubkey`, so `buzz-admin deletions submit` must repeat the owner pubkey to adopt an admitted request; the operator's own pubkey conflicts with the one-active-request invariant instead. Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- ARCHITECTURE.md | 18 ++++++++++++++++++ crates/buzz-db/src/store/deletion.rs | 12 ++++++++++++ crates/buzz-deletion/src/lib.rs | 6 +++++- crates/buzz-relay/src/api/operator.rs | 7 +++++++ 4 files changed, 42 insertions(+), 1 deletion(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index ede53a4b838..a3fa2af013a 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -26,6 +26,24 @@ management lists suppress the archived row. Replaying the same UUID converges to its current stage; a different UUID conflicts with the existing one-active- request invariant until that request is aborted. +Owner consent on this path is asserted, not proven. The mediating operator +authenticates the owner and collects the deletion acknowledgement out of band, +upstream of the relay; the request itself carries only the operator's NIP-98 +signature. The relay verifies operator authority and that the asserted pubkey +is still the community's owner, then records the owner pubkey, mediating +operator pubkey, and acknowledgement version as durable provenance for that +upstream ceremony. No owner-signed attestation is required or checked, and +owners have no self-service cancellation. Recovery is a privileged abort, +which stays open across the reversible `submitted`, `inventoried`, `approved`, +and `fenced` stages — releasing the request fence while leaving the community +archived — and closes from `drained` onward, once tenant state is destroyed. + +Manual operator handoff converges on an admitted owner request only when +`buzz-admin deletions submit --requested-by` repeats the owner pubkey recorded +on the row. Owner provenance pins `requested_by` to `owner_pubkey`, so passing +the operator's own pubkey does not converge — it conflicts with the existing +one-active-request invariant instead. + Ownership is mutable only while a community is active. Archiving freezes the current owner. Normal transfer and deployment-root legacy convergence take the same community-row lock as owner-deletion admission, then reject archived, diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 073e9cafd08..333b4da5728 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -748,6 +748,12 @@ impl DeletionStore { } /// Persist a request. Only active non-tombstone communities may be submitted. + /// + /// `requested_by` is recorded on a new row and is also the convergence key + /// when a `submitted` request already exists. Owner provenance pins an + /// owner-origin request's `requested_by` to its `owner_pubkey`, so taking + /// one over manually means passing that owner pubkey; the operator's own + /// pubkey conflicts with the existing request instead of converging. pub async fn submit( &self, community_host: &str, @@ -799,6 +805,12 @@ impl DeletionStore { /// correlation/idempotency identity. Replays return the existing request at /// its current stage. This operation only persists intent; it never /// inventories, approves, quiesces, or executes deletion. + /// + /// The owner's consent arrives as the calling operator's assertion: the + /// operator authenticated the owner and collected the acknowledgement + /// upstream. This layer records that provenance and checks that + /// `owner_pubkey` is still the community's owner; it never verifies an + /// owner-signed attestation. pub async fn admit_owner_request( &self, normalized_community_host: &str, diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index 714cf7eaff2..2814c486917 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -227,7 +227,11 @@ pub enum Command { /// Canonical community host. Defaults to RELAY_URL's authority. #[arg(long)] host: Option, - /// Operator identity recorded on the request. + /// Identity recorded on the request. + /// + /// Also the convergence key for an existing `submitted` request: to + /// take over an owner-origin request, pass its owner pubkey. The + /// operator's own pubkey conflicts with it instead of converging. #[arg(long)] requested_by: String, /// Optional reason for the request. diff --git a/crates/buzz-relay/src/api/operator.rs b/crates/buzz-relay/src/api/operator.rs index eea2dd7c1f1..f02fd3e4ea7 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -343,6 +343,13 @@ pub async fn unarchive_community( /// The request UUID is the correlation/idempotency key. Acceptance is a fast /// PostgreSQL-only transaction and returns `202`; inventory, approval, /// quiescing, object-store access, and executor work remain asynchronous. +/// +/// Owner consent is asserted by the operator, not proven to the relay. The +/// operator authenticates the owner and collects the acknowledgement upstream; +/// this request carries only the operator's NIP-98 signature. Authorization is +/// therefore operator authority plus "the asserted pubkey is still the owner". +/// `owner_pubkey` and `acknowledgement_version` are recorded as provenance for +/// that upstream ceremony, not verified as cryptographic owner consent. pub async fn delete_community( State(state): State>, headers: HeaderMap, From e5b3d465748a5eb2d5552ca844c4cd2decdebf8e Mon Sep 17 00:00:00 2001 From: tornquist Date: Fri, 25 Sep 2026 20:27:15 +0000 Subject: [PATCH 08/10] Serialize owner convergence before deletion abort Co-Authored-By: Claude Opus 5 Signed-off-by: tornquist --- crates/buzz-db/src/store/deletion.rs | 202 +++++++++++++++++++++- crates/buzz-db/src/store/relay_members.rs | 1 + 2 files changed, 202 insertions(+), 1 deletion(-) diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 333b4da5728..09752c10093 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -2874,7 +2874,7 @@ async fn lock_community_deletion( Ok(()) } -async fn lock_community_deletion_shared( +pub(crate) async fn lock_community_deletion_shared( tx: &mut Transaction<'_, Postgres>, community: CommunityId, ) -> Result<()> { @@ -4018,6 +4018,111 @@ mod postgres_tests { .expect("release admission fixture serialization"); } + struct OwnerConvergenceGate { + connection: sqlx::pool::PoolConnection, + trigger_name: String, + function_name: String, + first_key: i32, + second_key: i32, + } + + async fn install_owner_convergence_gate( + db: &Db, + community: CommunityId, + owner: &str, + ) -> OwnerConvergenceGate { + let suffix = Uuid::new_v4().simple().to_string(); + let trigger_name = format!("a_owner_convergence_gate_{suffix}"); + let function_name = format!("owner_convergence_gate_fn_{suffix}"); + let gate_id = Uuid::new_v4(); + let first_key = (gate_id.as_u128() as u32 & 0x7fff_ffff) as i32; + let second_key = ((gate_id.as_u128() >> 32) as u32 & 0x7fff_ffff) as i32; + let mut connection = db.pool.acquire().await.expect("acquire gate connection"); + sqlx::query("SELECT pg_advisory_lock($1, $2)") + .bind(first_key) + .bind(second_key) + .execute(&mut *connection) + .await + .expect("hold owner convergence gate"); + sqlx::query(AssertSqlSafe(format!( + "CREATE FUNCTION {function_name}() RETURNS trigger LANGUAGE plpgsql AS $$ \ + BEGIN \ + IF NEW.community_id = '{}'::uuid AND NEW.pubkey = '{}' THEN \ + PERFORM pg_advisory_xact_lock({first_key}, {second_key}); \ + END IF; \ + RETURN NEW; \ + END $$", + community.as_uuid(), + owner + ))) + .execute(&db.pool) + .await + .expect("install owner convergence gate function"); + sqlx::query(AssertSqlSafe(format!( + "CREATE TRIGGER {trigger_name} BEFORE INSERT ON relay_members \ + FOR EACH ROW EXECUTE FUNCTION {function_name}()" + ))) + .execute(&db.pool) + .await + .expect("install owner convergence gate trigger"); + OwnerConvergenceGate { + connection, + trigger_name, + function_name, + first_key, + second_key, + } + } + + async fn wait_for_owner_convergence_gate(db: &Db, gate: &OwnerConvergenceGate) { + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let waiting: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM pg_locks \ + WHERE locktype = 'advisory' AND classid = $1::oid AND objid = $2::oid \ + AND objsubid = 2 AND NOT granted)", + ) + .bind(gate.first_key) + .bind(gate.second_key) + .fetch_one(&db.pool) + .await + .expect("inspect owner convergence gate"); + if waiting { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("owner convergence reached gated member update"); + } + + async fn release_owner_convergence_gate(gate: &mut OwnerConvergenceGate) { + sqlx::query("SELECT pg_advisory_unlock($1, $2)") + .bind(gate.first_key) + .bind(gate.second_key) + .execute(&mut *gate.connection) + .await + .expect("release owner convergence gate"); + } + + async fn remove_owner_convergence_gate(db: &Db, gate: OwnerConvergenceGate) { + sqlx::query(AssertSqlSafe(format!( + "DROP TRIGGER {} ON relay_members", + gate.trigger_name + ))) + .execute(&db.pool) + .await + .expect("remove owner convergence gate trigger"); + sqlx::query(AssertSqlSafe(format!( + "DROP FUNCTION {}()", + gate.function_name + ))) + .execute(&db.pool) + .await + .expect("remove owner convergence gate function"); + } + async fn contender_db(application_name: &str) -> Db { let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") .or_else(|_| std::env::var("DATABASE_URL")) @@ -4203,6 +4308,101 @@ mod postgres_tests { assert_eq!(membership_roles(&db, community).await, before); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn same_owner_convergence_serializes_with_abort_without_deadlock() { + let mut deadlocks = Vec::new(); + for stage in [DeletionStage::Submitted, DeletionStage::Inventoried] { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let OwnerDeletionAdmission::Accepted(request) = store + .admit_owner_request( + &host, + &owner, + "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + 1, + Uuid::new_v4(), + ) + .await + .expect("admit owner deletion") + else { + panic!("expected accepted owner deletion") + }; + if stage == DeletionStage::Inventoried { + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("inventory schema"), + storage: empty_storage_manifest(community), + }; + let inventoried = store + .freeze_inventory(request.id, &inventory) + .await + .expect("freeze inventory"); + assert_eq!(inventoried.stage, DeletionStage::Inventoried); + } + + let mut gate = install_owner_convergence_gate(&db, community, &owner).await; + let convergence_db = + contender_db(&format!("owner-convergence-{}", Uuid::new_v4().simple())).await; + let converging = tokio::spawn({ + let owner = owner.clone(); + async move { convergence_db.provision_owner(community, &owner).await } + }); + wait_for_owner_convergence_gate(&db, &gate).await; + + let abort_application = format!("owner-abort-{}", Uuid::new_v4().simple()); + let abort_store = contender_db(&abort_application).await.deletion_store(); + let aborting = tokio::spawn(async move { + abort_store + .abort(request.id, "operator", "same-owner convergence race") + .await + }); + wait_for_contender_lock(&db, &abort_application).await; + release_owner_convergence_gate(&mut gate).await; + + let (convergence_join, abort_join) = tokio::join!( + tokio::time::timeout(Duration::from_secs(5), converging), + tokio::time::timeout(Duration::from_secs(5), aborting), + ); + remove_owner_convergence_gate(&db, gate).await; + let convergence_result = convergence_join + .expect("same-owner convergence must not hang") + .expect("join same-owner convergence"); + let abort_result = abort_join + .expect("abort must not hang") + .expect("join abort"); + if matches!( + &convergence_result, + Err(DbError::Sqlx(sqlx::Error::Database(error))) + if error.code().as_deref() == Some("40P01") + ) || matches!( + &abort_result, + Err(DbError::Sqlx(sqlx::Error::Database(error))) + if error.code().as_deref() == Some("40P01") + ) { + deadlocks.push(format!( + "{stage}: convergence={convergence_result:?}, abort={abort_result:?}" + )); + continue; + } + assert_eq!( + convergence_result.expect("same-owner convergence"), + ProvisionOwnerResult::Applied + ); + assert_eq!( + abort_result.expect("abort request").stage, + DeletionStage::Aborted + ); + } + assert!( + deadlocks.is_empty(), + "same-owner convergence and abort must serialize without 40P01:\n{}", + deadlocks.join("\n") + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn concurrent_owner_admission_serializes_before_normal_transfer() { diff --git a/crates/buzz-db/src/store/relay_members.rs b/crates/buzz-db/src/store/relay_members.rs index 3bd09c625fe..20309022cbb 100644 --- a/crates/buzz-db/src/store/relay_members.rs +++ b/crates/buzz-db/src/store/relay_members.rs @@ -422,6 +422,7 @@ async fn lock_owner_mutation_admission( proposed_owner: &str, mode: OwnerMutationMode, ) -> Result { + crate::deletion::lock_community_deletion_shared(tx, community).await?; let target = sqlx::query( "SELECT archived_at, deletion_state, deleted_at FROM communities \ WHERE id = $1 FOR UPDATE", From fcd831ea241985eab065233869642d1f891ba183 Mon Sep 17 00:00:00 2001 From: Codex Date: Mon, 28 Sep 2026 15:21:08 +0000 Subject: [PATCH 09/10] Renumber owner deletion admission migration Signed-off-by: Codex Co-authored-by: Codex --- crates/buzz-db/src/runtime/migration.rs | 12 ++++++------ .../runtime/tests/thread_window_postgres_tests.rs | 2 +- ...l => 0051_owner_community_deletion_admission.sql} | 0 3 files changed, 7 insertions(+), 7 deletions(-) rename migrations/{0050_owner_community_deletion_admission.sql => 0051_owner_community_deletion_admission.sql} (100%) diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 7b64547c2f5..9bcc5654eeb 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -705,7 +705,7 @@ mod postgres_tests { assert_eq!(migrations.len(), 50); assert_eq!(migrations[48].version, 49); - assert_eq!(migrations[49].version, 50); + assert_eq!(migrations[49].version, 51); assert!(migrations[48] .sql .as_str() @@ -1770,10 +1770,10 @@ mod postgres_tests { .expect("embedded migration 0029") .sql .as_ref(); - let migration_0050: &str = MIGRATOR + let migration_0051: &str = MIGRATOR .iter() - .find(|migration| migration.version == 50) - .expect("embedded migration 0050") + .find(|migration| migration.version == 51) + .expect("embedded migration 0051") .sql .as_ref(); let workspace_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) @@ -1784,7 +1784,7 @@ mod postgres_tests { .expect("read schema/schema.sql"); let migration = surface(migration_0029); - let owner_admission_migration = surface(migration_0050); + let owner_admission_migration = surface(migration_0051); let schema = surface(&schema_sql); assert_eq!( @@ -1830,7 +1830,7 @@ mod postgres_tests { owner_admission_migration .functions .get("prevent_community_deletion_request_retargeting") - .expect("0050 deletion retargeting guard"), + .expect("0051 deletion retargeting guard"), "schema.sql must carry the latest immutable owner-provenance guard" ); let request_table = schema diff --git a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs index 5ca0185af56..5d51f4348d5 100644 --- a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs +++ b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs @@ -311,7 +311,7 @@ async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer() production_result.is_ok(), "production migrator must preserve ingestion progress: {production_result:?}" ); - assert_eq!(version, 50); + assert_eq!(version, 51); assert_eq!(final_oid, oid, "prebuild must not be replaced"); assert_eq!(count, 4, "all writer witnesses must persist"); } diff --git a/migrations/0050_owner_community_deletion_admission.sql b/migrations/0051_owner_community_deletion_admission.sql similarity index 100% rename from migrations/0050_owner_community_deletion_admission.sql rename to migrations/0051_owner_community_deletion_admission.sql From 1f1b5f4d39f2639053ba58bdd66213609df8dea3 Mon Sep 17 00:00:00 2001 From: Elrond <28d6302a099e5225b02c4155ac4236e4912603df2ab08dbfc2f4fef08ce598c8@buzz.block.builderlab.xyz> Date: Mon, 28 Sep 2026 13:13:09 -0400 Subject: [PATCH 10/10] docs: clarify community deletion abort boundary Signed-off-by: Elrond <28d6302a099e5225b02c4155ac4236e4912603df2ab08dbfc2f4fef08ce598c8@buzz.block.builderlab.xyz> --- ARCHITECTURE.md | 3 ++- crates/buzz-db/src/store/deletion.rs | 7 ++++--- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index a3fa2af013a..2de723c5e3e 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -36,7 +36,8 @@ upstream ceremony. No owner-signed attestation is required or checked, and owners have no self-service cancellation. Recovery is a privileged abort, which stays open across the reversible `submitted`, `inventoried`, `approved`, and `fenced` stages — releasing the request fence while leaving the community -archived — and closes from `drained` onward, once tenant state is destroyed. +archived — and closes from `drained` onward, when tenant-state destruction may +have begun. Manual operator handoff converges on an admitted owner request only when `buzz-admin deletions submit --requested-by` repeats the owner pubkey recorded diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 09752c10093..cdd13d94c81 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -2207,7 +2207,8 @@ impl DeletionStore { /// privileged path is the only recovery when preparation cannot continue. /// Abort reverses deletion intent, not the owner's archive decision: the /// community stays archived and the owner restores it explicitly. - /// Stages from `drained` onward have destroyed tenant state and stay closed. + /// Stages from `drained` onward stay closed because tenant-state destruction + /// may have begun. pub async fn abort( &self, request_id: Uuid, @@ -4758,8 +4759,8 @@ mod postgres_tests { ); } - /// The reversible boundary stops at `inventoried`. Once execution has - /// destroyed anything, abort must stay closed. + /// The reversible boundary extends through `fenced`. From `drained` + /// onward, destruction may have begun, so abort must stay closed. #[tokio::test] #[ignore = "requires Postgres"] async fn privileged_abort_spans_only_the_reversible_pre_destruction_boundary() {