diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index ed34ee4c8d3..8aecbd61cc4 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -14,6 +14,43 @@ 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. + +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, 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 +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, +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/lib.rs b/crates/buzz-db/src/lib.rs index 27beed51736..3dac0b1ba2f 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -73,7 +73,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 1b210ed0f7e..66e3645915d 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -705,12 +705,17 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 50); + assert_eq!(migrations.len(), 51); assert_eq!(migrations[48].version, 49); + assert_eq!(migrations[50].version, 51); assert!(migrations[48] .sql .as_str() .contains("idx_thread_metadata_window")); + assert!(migrations[50] + .sql + .as_str() + .contains("community_deletion_owner_provenance")); assert_eq!(migrations[0].version, 1); assert_eq!(&*migrations[0].description, "initial schema"); assert!(migrations[0] @@ -1820,6 +1825,12 @@ mod postgres_tests { .expect("embedded migration 0029") .sql .as_ref(); + let migration_0051: &str = MIGRATOR + .iter() + .find(|migration| migration.version == 51) + .expect("embedded migration 0051") + .sql + .as_ref(); let workspace_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) .parent() .and_then(std::path::Path::parent) @@ -1828,6 +1839,7 @@ mod postgres_tests { .expect("read schema/schema.sql"); let migration = surface(migration_0029); + let owner_admission_migration = surface(migration_0051); let schema = surface(&schema_sql); assert_eq!( @@ -1856,13 +1868,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("0051 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 @@ -2512,6 +2553,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/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/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..cdd13d94c81 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 { @@ -697,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, @@ -742,6 +799,137 @@ 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. + /// + /// 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, + 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( @@ -2010,7 +2198,17 @@ 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 stay closed because tenant-state destruction + /// may have begun. pub async fn abort( &self, request_id: Uuid, @@ -2049,7 +2247,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", @@ -2674,7 +2875,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<()> { @@ -3112,6 +3313,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")?, @@ -3413,10 +3618,216 @@ 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::*; - use crate::{CreateCommunityWithOwnerResult, Db, DbConfig}; + use crate::{ + 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") @@ -3485,6 +3896,914 @@ 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) + } + + 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"); + } + + 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")) + .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() { + 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_eq!( + db.transfer_ownership(created.id, &replacement, &owner) + .await + .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(_) + )); + } + + #[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("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 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() { + 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] + #[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" + ); + } + + /// 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 + /// 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 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() { + 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() { diff --git a/crates/buzz-db/src/store/relay_members.rs b/crates/buzz-db/src/store/relay_members.rs index ecde3924fef..20309022cbb 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,93 @@ 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 { + 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", + ) + .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 +479,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 +520,7 @@ async fn bootstrap_owner_with_operation( .await?; tx.commit().await?; - Ok(()) + Ok(ProvisionOwnerResult::Applied) } /// The result of a transfer-ownership attempt. @@ -444,12 +540,27 @@ 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. /// Enforced atomically inside the transfer transaction so concurrent /// transfers to the same recipient cannot both pass the limit. 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. @@ -500,13 +611,17 @@ 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 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 /// `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,18 +648,28 @@ pub async fn transfer_ownership( ) .await?; - // 2. 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?; @@ -570,7 +695,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 +710,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 +721,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", @@ -798,7 +923,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, @@ -1341,6 +1470,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-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 1c912e599e0..64095decc1a 100644 --- a/crates/buzz-relay/src/api/operator.rs +++ b/crates/buzz-relay/src/api/operator.rs @@ -293,7 +293,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) @@ -314,6 +318,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>, @@ -398,12 +411,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(), @@ -413,6 +437,113 @@ 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. +/// +/// 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, + 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>, @@ -535,6 +666,18 @@ 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, + "community deletion is pending", + )); + } buzz_db::relay_members::TransferResult::LimitReached => { return Err(api_error( StatusCode::CONFLICT, @@ -798,6 +941,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, @@ -813,6 +990,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") @@ -882,6 +1082,500 @@ 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); + + // 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", + "/operator/communities/delete", + 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, + }) + .to_string(), + ), + ) + .await; + assert_eq!(bad_uuid.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); + } + + /// 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() { @@ -1378,6 +2072,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; } diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 3e96c02c13d..9993ea4c8ad 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -330,6 +330,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), diff --git a/migrations/0051_owner_community_deletion_admission.sql b/migrations/0051_owner_community_deletion_admission.sql new file mode 100644 index 00000000000..25002b32da0 --- /dev/null +++ b/migrations/0051_owner_community_deletion_admission.sql @@ -0,0 +1,81 @@ +-- 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. +-- +-- 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 + 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 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 + 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 1ba4bb0399b..fdd382f4635 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,21 @@ 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 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 + AND requested_by = owner_pubkey) + ), UNIQUE (id, community_id, inventory_digest) ); CREATE UNIQUE INDEX community_deletion_requests_active_community @@ -1293,7 +1313,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 +1324,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