From e120231d705fe3c2fe1d8e386e7da600c02b330b Mon Sep 17 00:00:00 2001 From: npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc <5288f082d91aff2de5bc8be23fbfa01ca2e3859cac57a91b7a3fa9f12349505c@buzz.block.builderlab.xyz> Date: Fri, 31 Jul 2026 19:12:16 -0700 Subject: [PATCH] fix: close community deletion serving-write races Order fencing behind admitted serving-write leases, keep external write heartbeats live through durable finalization, and make Git publication hold the community lock through the relay event commit. Harden tenant trigger coverage for cross-community updates and future tables, then fail startup and readiness closed on live catalog drift. Co-authored-by: npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc <5288f082d91aff2de5bc8be23fbfa01ca2e3859cac57a91b7a3fa9f12349505c@buzz.block.builderlab.xyz> Signed-off-by: npub122y0pqkertljmedu303rl0aqrj3w8pvu43t6jxm6875lzg6f2pwqegc3xc <5288f082d91aff2de5bc8be23fbfa01ca2e3859cac57a91b7a3fa9f12349505c@buzz.block.builderlab.xyz> --- crates/buzz-db/src/deletion.rs | 246 ++++++++++--- crates/buzz-db/src/error.rs | 14 + crates/buzz-db/src/event.rs | 2 +- crates/buzz-db/src/lib.rs | 48 +++ crates/buzz-db/src/migration.rs | 299 ++++++++++++++++ crates/buzz-deletion/src/lib.rs | 23 +- crates/buzz-relay/src/api/git/transport.rs | 390 +++++++++++++++++++-- crates/buzz-relay/src/main.rs | 6 + crates/buzz-relay/src/push_runtime.rs | 1 - crates/buzz-relay/src/router.rs | 26 +- migrations/0027_community_deletion.sql | 126 +++++-- schema/schema.sql | 142 ++++++-- 12 files changed, 1183 insertions(+), 140 deletions(-) diff --git a/crates/buzz-db/src/deletion.rs b/crates/buzz-db/src/deletion.rs index f69d46d96..71070874e 100644 --- a/crates/buzz-db/src/deletion.rs +++ b/crates/buzz-db/src/deletion.rs @@ -581,8 +581,12 @@ impl DeletionStore { }) } - /// Build and validate a live PostgreSQL schema inventory. - pub async fn inventory_schema(&self, community: CommunityId) -> Result { + /// Validate the live scoped-table and write-fence catalog before serving. + /// + /// This uses the same exact-set checks as deletion inventory but does not + /// require a community row. Schema drift therefore fails startup/readiness + /// even before an operator submits a deletion request. + pub async fn validate_catalog(&self) -> Result<()> { let migration_version: i64 = sqlx::query_scalar( "SELECT COALESCE(max(version), 0) FROM _sqlx_migrations WHERE success", ) @@ -594,12 +598,12 @@ impl DeletionStore { ))); } - let live_tables = self.live_scoped_tables().await?; let expected = EXPECTED_SCOPED_TABLES .iter() .copied() .map(str::to_owned) .collect::>(); + let live_tables = self.live_scoped_tables().await?; if live_tables != expected { let missing = expected .difference(&live_tables) @@ -622,16 +626,28 @@ impl DeletionStore { .difference(&fenced_tables) .cloned() .collect::>(); + let unknown = fenced_tables + .difference(&expected) + .cloned() + .collect::>(); return Err(DbError::DeletionSafety(format!( - "community deletion write-fence drift (missing={})", - missing.join(",") + "community deletion write-fence drift (missing={}, unknown={})", + missing.join(","), + unknown.join(",") ))); } + Ok(()) + } + /// Build and validate a live PostgreSQL schema inventory. + pub async fn inventory_schema(&self, community: CommunityId) -> Result { + self.validate_catalog().await?; + let live_tables = self.live_scoped_tables().await?; + let fenced_tables = self.live_fenced_tables().await?; let _ = community; // counts are intentionally not approval-bound for a live tenant. Ok(SchemaManifest { revision: CATALOG_REVISION, - migration_version, + migration_version: EXPECTED_MIGRATION_VERSION, scoped_tables: live_tables.into_iter().collect(), row_counts: BTreeMap::new(), fenced_tables: fenced_tables.into_iter().collect(), @@ -882,6 +898,23 @@ impl DeletionStore { .bind(token.community_id.as_uuid()) .execute(&mut *tx) .await?; + let active_serving_writes = sqlx::query( + "SELECT count(*)::BIGINT AS active_count, \ + COALESCE(array_agg(DISTINCT operation ORDER BY operation), ARRAY[]::TEXT[]) AS operations \ + FROM community_serving_write_leases \ + WHERE community_id = $1 AND lease_until >= now()", + ) + .bind(token.community_id.as_uuid()) + .fetch_one(&mut *tx) + .await?; + let active_count: i64 = active_serving_writes.try_get("active_count")?; + if active_count > 0 { + return Err(DbError::ServingWritesNotDrained { + community_id: *token.community_id.as_uuid(), + active_count, + operations: active_serving_writes.try_get("operations")?, + }); + } let current_generation: i64 = sqlx::query_scalar( "SELECT deletion_fence_generation FROM communities WHERE id = $1 FOR UPDATE", ) @@ -1353,6 +1386,11 @@ impl DeletionStore { lease_duration: Duration, ) -> Result<()> { let lease_seconds = i64::try_from(lease_duration.as_secs()).unwrap_or(i64::MAX); + let mut tx = self.pool.begin().await?; + sqlx::query("SELECT pg_advisory_xact_lock_shared(community_deletion_lock_key($1))") + .bind(lease.community_id.as_uuid()) + .execute(&mut *tx) + .await?; let lease_until: Option> = sqlx::query_scalar( "UPDATE community_serving_write_leases lease \ SET lease_until = now() + make_interval(secs => $6), heartbeat_at = now() \ @@ -1361,10 +1399,9 @@ impl DeletionStore { AND lease.generation = $4 AND lease.fence_generation = $5 \ AND lease.lease_until >= now() \ AND community.id = lease.community_id \ - AND ((community.deletion_state = 'active' \ - AND community.deletion_fence_generation = lease.fence_generation) \ - OR (community.deletion_state = 'fenced' \ - AND community.deletion_fence_generation = lease.fence_generation + 1)) \ + AND community.deletion_state = 'active' \ + AND community.deleted_at IS NULL \ + AND community.deletion_fence_generation = lease.fence_generation \ RETURNING lease.lease_until", ) .bind(lease.id) @@ -1373,11 +1410,13 @@ impl DeletionStore { .bind(lease.generation) .bind(lease.fence_generation) .bind(lease_seconds) - .fetch_optional(&self.pool) + .fetch_optional(&mut *tx) .await?; - lease.lease_until = lease_until.ok_or_else(|| { + let lease_until = lease_until.ok_or_else(|| { DbError::AccessDenied(format!("stale serving write lease {}", lease.id)) })?; + tx.commit().await?; + lease.lease_until = lease_until; Ok(()) } @@ -1402,6 +1441,11 @@ impl DeletionStore { /// Check that an external side-effect lease and the community's active /// lifecycle are still current. pub async fn verify_serving_write_lease(&self, lease: &ServingWriteLease) -> Result<()> { + let mut tx = self.pool.begin().await?; + sqlx::query("SELECT pg_advisory_xact_lock_shared(community_deletion_lock_key($1))") + .bind(lease.community_id.as_uuid()) + .execute(&mut *tx) + .await?; let valid: bool = sqlx::query_scalar( "SELECT EXISTS(SELECT 1 FROM community_serving_write_leases lease \ JOIN communities community ON community.id = lease.community_id \ @@ -1409,19 +1453,18 @@ impl DeletionStore { AND lease.generation = $4 AND lease.fence_generation = $5 \ AND lease.lease_until >= now() \ AND community.deleted_at IS NULL \ - AND ((community.deletion_state = 'active' \ - AND community.deletion_fence_generation = lease.fence_generation) \ - OR (community.deletion_state = 'fenced' \ - AND community.deletion_fence_generation = lease.fence_generation + 1)))", + AND community.deletion_state = 'active' \ + AND community.deletion_fence_generation = lease.fence_generation)", ) .bind(lease.id) .bind(lease.community_id.as_uuid()) .bind(&lease.owner) .bind(lease.generation) .bind(lease.fence_generation) - .fetch_one(&self.pool) + .fetch_one(&mut *tx) .await?; if valid { + tx.commit().await?; Ok(()) } else { Err(DbError::AccessDenied(format!( @@ -1472,14 +1515,7 @@ impl DeletionStore { AND NOT c.relispartition AND a.attname = 'community_id' AND NOT a.attisdropped - AND c.relname NOT IN ( - 'community_deletion_requests', - 'community_deletion_approvals', - 'community_deletion_checkpoints', - 'community_deletion_retention_exceptions', - 'community_serving_write_leases', - 'community_deletion_executor_heartbeats' - ) + AND NOT community_write_fence_excluded_table(c.relname) ORDER BY c.relname "#, ) @@ -1500,6 +1536,12 @@ impl DeletionStore { AND NOT trigger.tgisinternal AND NOT c.relispartition AND procedure.proname = 'enforce_community_write_fence' + AND trigger.tgenabled = 'O' + AND (trigger.tgtype & 1) = 1 + AND (trigger.tgtype & 2) = 2 + AND (trigger.tgtype & 4) = 4 + AND (trigger.tgtype & 8) = 8 + AND (trigger.tgtype & 16) = 16 ORDER BY c.relname "#, ) @@ -2026,6 +2068,127 @@ mod postgres_tests { ); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn active_serving_lease_retries_fence_without_lifecycle_change() { + let (db, store) = store().await; + let (request, _) = inventoried_request(&db, &store).await; + store + .approve(request.id, "approver", None) + .await + .expect("approve"); + let claim = store + .claim_specific(request.id, "executor", DEFAULT_LEASE_DURATION) + .await + .expect("claim") + .expect("won claim"); + let serving = store + .acquire_serving_write_lease( + request.community_id, + "test_external", + "test-owner", + DEFAULT_LEASE_DURATION, + ) + .await + .expect("serving lease"); + + let error = store + .fence(&claim.lease) + .await + .expect_err("active serving lease must make fence retry"); + assert!(matches!( + error, + DbError::ServingWritesNotDrained { + active_count: 1, + ref operations, + .. + } if operations == &["test_external"] + )); + let unchanged = store.get(request.id).await.expect("request after retry"); + assert_eq!(unchanged.stage, DeletionStage::Approved); + assert_eq!(unchanged.fence_generation, None); + assert!(store + .is_serving_active(request.community_id) + .await + .expect("community remains active")); + + assert!(store + .release_serving_write_lease(&serving) + .await + .expect("release")); + let generation = store.fence(&claim.lease).await.expect("retry fence"); + assert_eq!(generation, 1); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn serving_lease_acquisition_and_fence_have_exactly_one_winner() { + let (db, store) = store().await; + let (request, _) = inventoried_request(&db, &store).await; + store + .approve(request.id, "approver", None) + .await + .expect("approve"); + let claim = store + .claim_specific(request.id, "executor", DEFAULT_LEASE_DURATION) + .await + .expect("claim") + .expect("won claim"); + + let mut lock_holder = db.pool.acquire().await.expect("lock connection"); + sqlx::query("BEGIN") + .execute(&mut *lock_holder) + .await + .expect("begin lock transaction"); + sqlx::query("SELECT pg_advisory_xact_lock(community_deletion_lock_key($1))") + .bind(request.community_id.as_uuid()) + .execute(&mut *lock_holder) + .await + .expect("hold ordering lock"); + + let acquire_store = store.clone(); + let community = request.community_id; + let acquire = tokio::spawn(async move { + acquire_store + .acquire_serving_write_lease( + community, + "race_acquire", + "race-owner", + DEFAULT_LEASE_DURATION, + ) + .await + }); + let fence_store = store.clone(); + let fence_token = claim.lease.clone(); + let fence = tokio::spawn(async move { fence_store.fence(&fence_token).await }); + sqlx::query("COMMIT") + .execute(&mut *lock_holder) + .await + .expect("release ordering lock"); + + let lease_result = acquire.await.expect("acquire task"); + let fence_result = fence.await.expect("fence task"); + match (lease_result, fence_result) { + (Ok(lease), Err(DbError::ServingWritesNotDrained { .. })) => { + assert!(store + .is_serving_active(request.community_id) + .await + .expect("active after lease wins")); + assert!(store + .release_serving_write_lease(&lease) + .await + .expect("release winner")); + } + (Err(DbError::AccessDenied(_)), Ok(_)) => { + assert!(!store + .is_serving_active(request.community_id) + .await + .expect("fenced after deletion wins")); + } + other => panic!("invalid acquisition/fence race outcome: {other:?}"), + } + } + #[tokio::test] #[ignore = "requires Postgres"] async fn external_serving_lease_blocks_drain_without_holding_a_pool_connection() { @@ -2040,7 +2203,7 @@ mod postgres_tests { .await .expect("claim") .expect("won claim"); - let mut serving = store + let serving = store .acquire_serving_write_lease( request.community_id, "test_external", @@ -2049,31 +2212,26 @@ mod postgres_tests { ) .await .expect("serving lease"); - let generation = store.fence(&claim.lease).await.expect("fence"); - store - .renew_serving_write_lease(&mut serving, DEFAULT_LEASE_DURATION) - .await - .expect("pre-fence side effect continues heartbeating while fenced"); - store - .verify_serving_write_lease(&serving) - .await - .expect("pre-fence side effect remains valid while draining"); - let token = LeaseToken { - fence_generation: Some(generation), - ..claim.lease - }; - assert!( - store.mark_drained(&token).await.is_err(), - "fenced deletion must wait for a pre-fence external effect lease" - ); + assert!(matches!( + store.fence(&claim.lease).await, + Err(DbError::ServingWritesNotDrained { .. }) + )); assert!(store .release_serving_write_lease(&serving) .await .expect("release")); + let generation = store + .fence(&claim.lease) + .await + .expect("fence after release"); + let token = LeaseToken { + fence_generation: Some(generation), + ..claim.lease + }; store .mark_drained(&token) .await - .expect("drain after release"); + .expect("drain immediately after successful fence"); } #[tokio::test] diff --git a/crates/buzz-db/src/error.rs b/crates/buzz-db/src/error.rs index 21e7e40ed..593eea1cc 100644 --- a/crates/buzz-db/src/error.rs +++ b/crates/buzz-db/src/error.rs @@ -45,6 +45,20 @@ pub enum DbError { #[error("invalid data: {0}")] InvalidData(String), + /// A serving write admitted before the lifecycle transition is still live. + /// This is an ordinary retryable drain condition, not a safety violation. + #[error( + "community {community_id} still has {active_count} active serving write lease(s): {operations:?}" + )] + ServingWritesNotDrained { + /// Community whose lifecycle transition must retry. + community_id: uuid::Uuid, + /// Number of currently unexpired serving-write leases. + active_count: i64, + /// Distinct operation categories holding those leases. + operations: Vec, + }, + /// A deletion safety invariant is structurally violated and requires /// operator/code remediation rather than blind retry. #[error("deletion safety error: {0}")] diff --git a/crates/buzz-db/src/event.rs b/crates/buzz-db/src/event.rs index a670a1340..007beb945 100644 --- a/crates/buzz-db/src/event.rs +++ b/crates/buzz-db/src/event.rs @@ -1111,7 +1111,7 @@ pub struct ThreadMetadataParams<'a> { pub broadcast: bool, } -async fn insert_event_with_thread_metadata_tx( +pub(crate) async fn insert_event_with_thread_metadata_tx( tx: &mut Transaction<'_, Postgres>, community_id: CommunityId, event: &Event, diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 3e533e26b..cdb01b6c8 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -1020,6 +1020,11 @@ impl Db { sqlx::query("SELECT 1").execute(&self.pool).await.is_ok() } + /// Validate the live community-deletion tenant catalog and write fences. + pub async fn validate_deletion_catalog(&self) -> Result<()> { + self.deletion_store().validate_catalog().await + } + /// Returns pool utilisation stats for metrics emission. /// /// `size` — total connections (idle + active) @@ -1654,6 +1659,49 @@ impl Db { Ok(result) } + /// Insert an event while holding the admitted serving-write lease's + /// community ordering lock through commit. + /// + /// External side effects use a durable lease rather than one long-lived DB + /// transaction. Their final database mutation must still share the + /// acquire/fence advisory-lock order so a fence cannot commit between the + /// caller's lease verification and this insert. + pub async fn insert_event_with_serving_write_guard( + &self, + community_id: CommunityId, + event: &nostr::Event, + channel_id: Option, + ) -> Result<(StoredEvent, bool)> { + let kind_u16 = event.kind.as_u16(); + let kind_u32 = u32::from(kind_u16); + if kind_u32 == buzz_core::kind::KIND_AUTH { + return Err(DbError::AuthEventRejected); + } + if buzz_core::kind::is_ephemeral(kind_u32) { + return Err(DbError::EphemeralEventRejected(kind_u16)); + } + + let mut tx = self.pool.begin().await?; + self.deletion_store() + .guard_transaction(&mut tx, community_id) + .await?; + let result = event::insert_event_with_thread_metadata_tx( + &mut tx, + community_id, + event, + channel_id, + None, + ) + .await?; + tx.commit().await?; + if result.1 { + if let Err(e) = insert_mentions(&self.pool, community_id, event, channel_id).await { + tracing::warn!(event_id = %event.id, "Failed to insert mentions: {e}"); + } + } + Ok(result) + } + /// Queries events matching the given filter parameters. /// /// Always reads from the WRITER pool. If the result influences a write diff --git a/crates/buzz-db/src/migration.rs b/crates/buzz-db/src/migration.rs index 7464b9d0d..8dac21f94 100644 --- a/crates/buzz-db/src/migration.rs +++ b/crates/buzz-db/src/migration.rs @@ -495,6 +495,106 @@ mod tests { .collect() } + fn migrations_missing_community_write_fence_attachment() -> Vec { + let mut migrations: Vec<_> = MIGRATOR.iter().collect(); + migrations.sort_by_key(|migration| migration.version); + migrations + .into_iter() + .filter(|migration| migration.version > 27) + .flat_map(|migration| { + let statements = split_sql_statements(migration.sql.as_str()); + let attachments = statements + .iter() + .filter_map(|statement| { + let normalized = normalize_sql(statement); + if !normalized.contains("attach_community_write_fence(") { + return None; + } + let call = normalized.split_once("attach_community_write_fence(")?.1; + let table = call + .split(')') + .next()? + .trim() + .split("::") + .next()? + .trim_matches(|ch| ch == '\'' || ch == '"') + .rsplit('.') + .next()? + .to_owned(); + (!table.is_empty()).then_some(table) + }) + .collect::>(); + statements.into_iter().filter_map(move |statement| { + let normalized = normalize_sql(&statement); + let table = if normalized.starts_with("create table") + && normalized.contains("community_id") + && !normalized.contains(" partition of ") + { + identifier_after_keyword(&statement, "create table") + } else if normalized.starts_with("alter table") + && normalized.contains("add") + && normalized.contains("community_id") + { + identifier_after_keyword(&statement, "alter table") + } else { + None + }?; + (!attachments.contains(&table)).then(|| { + format!( + "migration {} introduces {table}.community_id without attach_community_write_fence('{table}'::regclass)", + migration.version + ) + }) + }) + }) + .collect() + } + + fn community_write_fence_attachment_violations(sql: &str) -> Vec { + let statements = split_sql_statements(sql); + let attachments = statements + .iter() + .filter_map(|statement| { + let normalized = normalize_sql(statement); + if !normalized.contains("attach_community_write_fence(") { + return None; + } + let call = normalized.split_once("attach_community_write_fence(")?.1; + let table = call + .split(')') + .next()? + .trim() + .split("::") + .next()? + .trim_matches(|ch| ch == '\'' || ch == '"') + .rsplit('.') + .next()? + .to_owned(); + (!table.is_empty()).then_some(table) + }) + .collect::>(); + statements + .into_iter() + .filter_map(|statement| { + let normalized = normalize_sql(&statement); + let table = if normalized.starts_with("create table") + && normalized.contains("community_id") + && !normalized.contains(" partition of ") + { + identifier_after_keyword(&statement, "create table") + } else if normalized.starts_with("alter table") + && normalized.contains("add") + && normalized.contains("community_id") + { + identifier_after_keyword(&statement, "alter table") + } else { + None + }?; + (!attachments.contains(&table)).then_some(table) + }) + .collect() + } + fn scoped_constraint_lints(sql: &str, scoped_tables: &BTreeSet) -> Vec { let mut constraints = table_constraints(sql, scoped_tables); constraints.extend(alter_table_constraints(sql, scoped_tables)); @@ -935,13 +1035,54 @@ mod tests { assert!(deletion.contains("CREATE TABLE community_deletion_retention_exceptions")); assert!(deletion.contains("CREATE TABLE community_serving_write_leases")); assert!(deletion.contains("CREATE TABLE community_deletion_executor_heartbeats")); + assert!(deletion.contains("CREATE FUNCTION assert_community_write_allowed")); assert!(deletion.contains("CREATE FUNCTION enforce_community_write_fence")); + assert!(deletion.contains("CREATE FUNCTION attach_community_write_fence")); + assert!(deletion.contains("community_write_fence_excluded_table")); assert!(deletion.contains("CREATE FUNCTION enforce_community_tombstone")); assert!(deletion.contains("community tombstones are permanent")); assert!(deletion.contains("_operator_global_tables")); assert!(deletion.contains("'submitted', 'inventoried', 'approved', 'fenced', 'drained'")); } + #[test] + fn post_deletion_migrations_attach_every_new_community_write_fence() { + let violations = migrations_missing_community_write_fence_attachment(); + assert!( + violations.is_empty(), + "migrations after 0027 must explicitly attach every new community write fence:\n{}", + violations.join("\n") + ); + } + + #[test] + fn community_write_fence_attachment_lint_covers_create_and_alter() { + let missing = r#" + CREATE TABLE created_late ( + community_id UUID NOT NULL, + id UUID NOT NULL + ); + CREATE TABLE altered_late (id UUID NOT NULL); + ALTER TABLE altered_late ADD COLUMN community_id UUID NOT NULL; + "#; + assert_eq!( + community_write_fence_attachment_violations(missing), + vec!["created_late", "altered_late"] + ); + + let attached = r#" + CREATE TABLE created_late ( + community_id UUID NOT NULL, + id UUID NOT NULL + ); + SELECT attach_community_write_fence('created_late'::regclass); + CREATE TABLE altered_late (id UUID NOT NULL); + ALTER TABLE altered_late ADD COLUMN community_id UUID NOT NULL; + SELECT attach_community_write_fence('altered_late'::regclass); + "#; + assert!(community_write_fence_attachment_violations(attached).is_empty()); + } + #[test] fn migration_lint_detects_tables_missing_community_id_by_default() { let sql = r#" @@ -1304,5 +1445,163 @@ mod tests { search_expression.contains("ELSE NULL::tsvector"), "fresh installs must default non-allowlisted kinds to NULL: {search_expression}" ); + + let active_a = uuid::Uuid::new_v4(); + let active_b = uuid::Uuid::new_v4(); + let to_fence = uuid::Uuid::new_v4(); + for (community, label) in [ + (active_a, "active-a"), + (active_b, "active-b"), + (to_fence, "to-fence"), + ] { + sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") + .bind(community) + .bind(format!("late-fence-{label}-{}.example", community.simple())) + .execute(&pool) + .await + .expect("insert late-table test community"); + } + sqlx::query( + "CREATE TABLE late_created_scoped (\ + community_id UUID NOT NULL, id BIGINT PRIMARY KEY, value TEXT NOT NULL\ + )", + ) + .execute(&pool) + .await + .expect("create late scoped table"); + sqlx::query("SELECT attach_community_write_fence('late_created_scoped'::regclass)") + .execute(&pool) + .await + .expect("attach late create fence"); + sqlx::query("CREATE TABLE late_altered_scoped (id BIGINT PRIMARY KEY)") + .execute(&pool) + .await + .expect("create table before late alter"); + sqlx::query("ALTER TABLE late_altered_scoped ADD COLUMN community_id UUID NOT NULL") + .execute(&pool) + .await + .expect("add late community id"); + sqlx::query("SELECT attach_community_write_fence('late_altered_scoped'::regclass)") + .execute(&pool) + .await + .expect("attach late alter fence"); + let attached: Vec = sqlx::query_scalar( + "SELECT c.relname FROM pg_trigger trigger \ + JOIN pg_class c ON c.oid = trigger.tgrelid \ + JOIN pg_proc procedure ON procedure.oid = trigger.tgfoid \ + WHERE c.relname IN ('late_created_scoped', 'late_altered_scoped') \ + AND procedure.proname = 'enforce_community_write_fence' \ + AND NOT trigger.tgisinternal ORDER BY c.relname", + ) + .fetch_all(&pool) + .await + .expect("read late trigger catalog"); + assert_eq!(attached, vec!["late_altered_scoped", "late_created_scoped"]); + let malformed_fence_triggers: i64 = sqlx::query_scalar( + "SELECT count(*)::BIGINT FROM pg_trigger trigger \ + JOIN pg_class c ON c.oid = trigger.tgrelid \ + JOIN pg_proc procedure ON procedure.oid = trigger.tgfoid \ + WHERE c.relname IN ('late_created_scoped', 'late_altered_scoped') \ + AND procedure.proname = 'enforce_community_write_fence' \ + AND NOT trigger.tgisinternal \ + AND (trigger.tgenabled <> 'O' OR (trigger.tgtype & 31) <> 31)", + ) + .fetch_one(&pool) + .await + .expect("validate late trigger mode and operations"); + assert_eq!(malformed_fence_triggers, 0); + + sqlx::query( + "INSERT INTO late_created_scoped (community_id, id, value) \ + VALUES ($1, 1, 'same'), ($2, 2, 'source-fenced'), \ + ($1, 3, 'destination-fenced'), ($1, 4, 'opposite-a'), \ + ($3, 5, 'opposite-b')", + ) + .bind(active_a) + .bind(to_fence) + .bind(active_b) + .execute(&pool) + .await + .expect("seed late table while communities active"); + sqlx::query("UPDATE late_created_scoped SET value = 'same-ok' WHERE id = 1") + .execute(&pool) + .await + .expect("same-tenant active update"); + sqlx::query("UPDATE late_created_scoped SET community_id = $1 WHERE id = 1") + .bind(active_b) + .execute(&pool) + .await + .expect("active-to-active update"); + + let mut fence_connection = pool.acquire().await.expect("fence connection"); + sqlx::query("BEGIN") + .execute(&mut *fence_connection) + .await + .expect("begin direct fence"); + sqlx::query( + "SELECT set_config('buzz.deletion_executor_community', $1, true), \ + set_config('buzz.deletion_fence_generation', '1', true)", + ) + .bind(to_fence.to_string()) + .execute(&mut *fence_connection) + .await + .expect("authorize direct fence"); + sqlx::query( + "UPDATE communities SET deletion_state = 'fenced', \ + deletion_fence_generation = 1, archived_at = now() WHERE id = $1", + ) + .bind(to_fence) + .execute(&mut *fence_connection) + .await + .expect("fence test destination"); + sqlx::query("COMMIT") + .execute(&mut *fence_connection) + .await + .expect("commit direct fence"); + + let active_to_fenced = + sqlx::query("UPDATE late_created_scoped SET community_id = $1 WHERE id = 3") + .bind(to_fence) + .execute(&pool) + .await + .expect_err("active to fenced destination must fail"); + assert!(active_to_fenced + .to_string() + .contains("community write fenced")); + let fenced_to_active = + sqlx::query("UPDATE late_created_scoped SET community_id = $1 WHERE id = 2") + .bind(active_a) + .execute(&pool) + .await + .expect_err("fenced source to active destination must fail"); + assert!(fenced_to_active + .to_string() + .contains("community write fenced")); + let row_locations: Vec<(i64, uuid::Uuid)> = sqlx::query_as( + "SELECT id, community_id FROM late_created_scoped WHERE id IN (2, 3) ORDER BY id", + ) + .fetch_all(&pool) + .await + .expect("failed moves preserve row location"); + assert_eq!(row_locations, vec![(2, to_fence), (3, active_a)]); + + let move_a = sqlx::query("UPDATE late_created_scoped SET community_id = $1 WHERE id = 4") + .bind(active_b) + .execute(&pool); + let move_b = sqlx::query("UPDATE late_created_scoped SET community_id = $1 WHERE id = 5") + .bind(active_a) + .execute(&pool); + tokio::time::timeout(std::time::Duration::from_secs(2), async { + let (a, b) = tokio::join!(move_a, move_b); + a.expect("opposite active move A"); + b.expect("opposite active move B"); + }) + .await + .expect("opposite cross-tenant updates must not deadlock"); + + sqlx::query("DROP TABLE late_created_scoped, late_altered_scoped") + .execute(&pool) + .await + .expect("drop late-table fixtures"); } } diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index a77e00d7d..372eb1368 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -116,15 +116,6 @@ impl ServingWriteGuard { self.lost.clone() } - /// Begin finalization after the side effect itself has succeeded. - /// - /// This stops lease renewal without releasing the durable row. Call - /// [`Self::finish`] immediately after the result's database finalization so - /// deletion cannot overtake the side effect before its outcome is recorded. - pub fn begin_finalize(&self) { - self.cancel.cancel(); - } - /// Release the lease after the side effect completes. pub async fn finish(mut self) -> Result<()> { self.cancel.cancel(); @@ -871,7 +862,19 @@ async fn execute_stage( "approved storage taxonomy drifted before fencing", )); } - services.store.fence(&token).await?; + match services.store.fence(&token).await { + Ok(_) => {} + Err(buzz_db::DbError::ServingWritesNotDrained { + active_count, + operations, + .. + }) => { + return Err(transient(format!( + "serving writes not drained before fence: count={active_count}, operations={operations:?}" + ))); + } + Err(error) => return Err(error.into()), + } } DeletionStage::Fenced => { services diff --git a/crates/buzz-relay/src/api/git/transport.rs b/crates/buzz-relay/src/api/git/transport.rs index 5d5ce97cb..b73b9c9dd 100644 --- a/crates/buzz-relay/src/api/git/transport.rs +++ b/crates/buzz-relay/src/api/git/transport.rs @@ -1698,6 +1698,21 @@ pub(crate) struct PushContext { pub repo_handle: HydratedRepo, } +#[derive(Default)] +struct FinalizePushHooks { + #[cfg(test)] + post_cas_gate: Option>, + #[cfg(test)] + fail_ref_state_insert: bool, +} + +#[cfg(test)] +#[derive(Default)] +struct PostCasGate { + reached: tokio::sync::Notify, + resume: tokio::sync::Notify, +} + /// Finalize a push request: CAS-commit the new state into the object /// store, derive kind:30618 from the committed manifest, and only then /// build the success response. @@ -1708,6 +1723,17 @@ pub(crate) struct PushContext { /// constructor of a push 2xx, so the seam is structural (not by /// convention). async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { + finalize_push_inner(state, ctx, &FinalizePushHooks::default()).await +} + +async fn finalize_push_inner( + state: &Arc, + ctx: PushContext, + hooks: &FinalizePushHooks, +) -> Response { + #[cfg(not(test))] + let _ = hooks; + // The push fence, part 0 — **a rejected push publishes nothing.** // // `ctx.pack.ok` is false when git aborted the ref updates: either the @@ -1859,8 +1885,10 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { } }; - if let Err(error) = serving_write.finish().await { - warn!(owner = %ctx.owner, repo = %ctx.repo, %error, "failed to release community serving lease after push"); + #[cfg(test)] + if let Some(gate) = &hooks.post_cas_gate { + gate.reached.notify_one(); + gate.resume.notified().await; } // Derived after CAS: kind:30618 ref-state event over the *committed* @@ -1886,7 +1914,7 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { (Some(before), Some(after)) => before != after, _ => true, // first push (parent None) or impossible-shape after key → publish }; - if manifest_changed { + let publication_result: Result<(), String> = if manifest_changed { let inputs = RefStateInputs { repo_id: &ctx.repo_id, head: &success.manifest.head, @@ -1897,11 +1925,23 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { Ok(event) => { // Relay-signed kind:30618 belongs to the same server-resolved // tenant as the git request that committed the pointer. - match state + #[cfg(test)] + let insert_result = if hooks.fail_ref_state_insert { + Err(buzz_db::DbError::InvalidData( + "injected kind:30618 insert failure".to_string(), + )) + } else { + state + .db + .insert_event_with_serving_write_guard(ctx.tenant.community(), &event, None) + .await + }; + #[cfg(not(test))] + let insert_result = state .db - .insert_event(ctx.tenant.community(), &event, None) - .await - { + .insert_event_with_serving_write_guard(ctx.tenant.community(), &event, None) + .await; + match insert_result { Ok((stored, true)) => { // Routed through the guarded send path for uniformity; // the access gate no-ops for this globally-scoped @@ -1918,6 +1958,7 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { manifest = %success.manifest_key, "kind:30618 published (derived after CAS)" ); + Ok(()) } Ok((_, false)) => { info!( @@ -1925,26 +1966,41 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { repo = %ctx.repo_id, "kind:30618 deduplicated by relay db" ); + Ok(()) } - Err(e) => { - warn!( - owner = %ctx.owner, - repo = %ctx.repo_id, - error = %e, - "kind:30618 insert failed; push remains durable in object store" - ); - } + Err(error) => Err(format!("kind:30618 insert failed: {error}")), } } - Err(e) => { - warn!( - owner = %ctx.owner, - repo = %ctx.repo_id, - error = %e, - "kind:30618 build failed; push remains durable in object store" - ); - } + Err(error) => Err(format!("kind:30618 build failed: {error}")), } + } else { + Ok(()) + }; + + // The admitted serving write spans the complete publication attempt. Fence + // acquisition cannot overtake the pointer CAS, durable 30618 insert, or + // local fan-out attempt; only now may the lease be released. + if let Err(error) = serving_write.finish().await { + warn!(owner = %ctx.owner, repo = %ctx.repo, %error, "failed to release community serving lease after push publication"); + return ( + StatusCode::SERVICE_UNAVAILABLE, + "community write lease lost during publication", + ) + .into_response(); + } + if let Err(error) = publication_result { + error!( + owner = %ctx.owner, + repo = %ctx.repo_id, + manifest = %success.manifest_key, + %error, + "push pointer committed but kind:30618 publication failed" + ); + return ( + StatusCode::INTERNAL_SERVER_ERROR, + "push committed but ref-state publication failed; retry", + ) + .into_response(); } // Only now — after CAS commit and (optional) 30618 emission — build @@ -1973,12 +2029,14 @@ pub fn git_router(state: Arc) -> Router { #[cfg(test)] mod track_c_tests { use super::*; + use crate::api::git::hydrate::{hydrate_for_write, HydrationOptions}; use crate::api::git::manifest::Manifest; use buzz_core::CommunityId; use nostr::{EventBuilder, Keys, Kind, Tag}; use std::collections::BTreeMap; use std::io::Write; use std::process::Output; + use tempfile::TempDir; fn oid_sha1() -> String { "cb09a769da1c01f458fa6959d4e8eded38fac8d3".to_string() @@ -2119,6 +2177,292 @@ mod track_c_tests { assert!(remote.join("refs/heads/master").exists()); } + async fn run_finalize_git(repo: &Path, args: &[&str]) -> std::process::Output { + let mut command = Command::new("git"); + command.current_dir(repo).args(args); + harden_git_env(&mut command); + let output = command.output().await.expect("spawn git"); + assert!( + output.status.success(), + "git {args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + output + } + + async fn finalize_test_state() -> (Arc, sqlx::PgPool) { + const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; // sadscan:disable np.postgres.1 + let mut config = crate::config::Config::from_env().expect("default config loads"); + config.require_relay_membership = false; + config.redis_url = "redis://127.0.0.1:1".to_string(); + config.database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .unwrap_or_else(|_| TEST_DB_URL.to_string()); + let pool = sqlx::PgPool::connect(&config.database_url) + .await + .expect("connect test DB"); + let db = buzz_db::Db::from_pool(pool.clone()); + db.migrate().await.expect("migrate test DB"); + let redis_pool = deadpool_redis::Config::from_url(&config.redis_url) + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("redis pool"); + let pubsub = Arc::new( + buzz_pubsub::PubSubManager::new(&config.redis_url, redis_pool.clone()) + .await + .expect("pubsub manager"), + ); + let audit = buzz_audit::AuditService::new(pool.clone()); + let auth = buzz_auth::AuthService::new(config.auth.clone()); + let search = buzz_search::SearchService::new(pool.clone()); + let workflow_engine = Arc::new(buzz_workflow::WorkflowEngine::new( + db.clone(), + buzz_workflow::WorkflowConfig::default(), + )); + let media_storage = buzz_media::MediaStorage::new(&config.media).expect("media storage"); + let (state, _audit_shutdown) = AppState::new( + config, + db, + redis_pool, + audit, + pubsub, + auth, + search, + workflow_engine, + Keys::generate(), + media_storage, + ); + (Arc::new(state), pool) + } + + async fn approved_deletion( + state: &AppState, + host: &str, + ) -> ( + buzz_db::deletion::DeletionRequest, + buzz_db::deletion::ClaimedDeletion, + ) { + use buzz_db::deletion::{FrozenInventory, StorageManifest, DEFAULT_LEASE_DURATION}; + + let store = state.db.deletion_store(); + let request = store + .submit(host, "git-finalize-test", Some("post-CAS lease regression")) + .await + .expect("submit deletion"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(request.community_id) + .await + .expect("schema inventory"), + storage: StorageManifest { + version: 1, + tenant_keys: Vec::new(), + git_pointer_keys: Vec::new(), + media_sidecar_keys: Vec::new(), + media_upload_keys: Vec::new(), + retained_shared_cas_keys: Vec::new(), + unknown_keys: Vec::new(), + unsupported_version_keys: Vec::new(), + }, + }; + store + .freeze_inventory(request.id, &inventory) + .await + .expect("freeze inventory"); + store + .approve(request.id, "git-finalize-test", None) + .await + .expect("approve deletion"); + let claim = store + .claim_specific(request.id, "git-finalize-test", DEFAULT_LEASE_DURATION) + .await + .expect("claim deletion") + .expect("won deletion claim"); + (request, claim) + } + + async fn pushed_context( + state: &AppState, + community: CommunityId, + host: &str, + owner: String, + repo: String, + pusher: nostr::PublicKey, + scratch: &Path, + ) -> PushContext { + let tenant = TenantContext::resolved(community, host); + let (hydrated, parent_state) = hydrate_for_write( + &state.git_store, + &tenant, + &owner, + &repo, + HydrationOptions { + pack_cache: &state.git_pack_cache, + scratch_dir: scratch, + max_pack_bytes: 1024 * 1024, + max_repo_bytes: 2 * 1024 * 1024, + }, + ) + .await + .expect("hydrate empty test repo"); + let source = scratch.join("source"); + tokio::fs::create_dir(&source) + .await + .expect("source directory"); + run_finalize_git(&source, &["init", "--quiet", "--initial-branch=main"]).await; + run_finalize_git(&source, &["config", "user.email", "finalize@test"]).await; + run_finalize_git(&source, &["config", "user.name", "finalize"]).await; + tokio::fs::write(source.join("file.txt"), b"committed\n") + .await + .expect("write source file"); + run_finalize_git(&source, &["add", "file.txt"]).await; + run_finalize_git(&source, &["commit", "--quiet", "-m", "committed"]).await; + let remote = hydrated.path().to_str().expect("hydrated path utf8"); + run_finalize_git(&source, &["push", "--quiet", remote, "main"]).await; + + PushContext { + pack: PackOutput { + stdout: b"push-ok".to_vec(), + ok: true, + }, + parent_state, + owner, + repo: repo.clone(), + repo_id: repo, + pusher, + tenant, + repo_handle: hydrated, + } + } + + #[tokio::test] + #[ignore = "requires Postgres and MinIO"] + async fn finalize_push_holds_serving_lease_through_post_cas_publication() { + let (state, pool) = finalize_test_state().await; + let host = format!("git-finalize-{}.example", uuid::Uuid::new_v4().simple()); + let community = state + .db + .ensure_configured_community(&host) + .await + .expect("create test community") + .id; + let (request, claim) = approved_deletion(&state, &host).await; + let scratch = TempDir::new().expect("scratch"); + let owner = format!("owner-{}", uuid::Uuid::new_v4().simple()); + let repo = format!("repo-{}", uuid::Uuid::new_v4().simple()); + let ctx = pushed_context( + &state, + community, + &host, + owner, + repo.clone(), + Keys::generate().public_key(), + scratch.path(), + ) + .await; + let gate = Arc::new(PostCasGate::default()); + let hooks = FinalizePushHooks { + post_cas_gate: Some(Arc::clone(&gate)), + fail_ref_state_insert: false, + }; + let finalize_state = Arc::clone(&state); + let finalize = + tokio::spawn(async move { finalize_push_inner(&finalize_state, ctx, &hooks).await }); + + gate.reached.notified().await; + let error = state + .db + .deletion_store() + .fence(&claim.lease) + .await + .expect_err("post-CAS serving lease must block fence"); + assert!(matches!( + error, + buzz_db::DbError::ServingWritesNotDrained { .. } + )); + assert!(state + .db + .deletion_store() + .is_serving_active(community) + .await + .expect("community remains active")); + + gate.resume.notify_one(); + let response = finalize.await.expect("finalize task"); + assert_eq!(response.status(), StatusCode::OK); + let mut query = buzz_db::event::EventQuery::for_community(community); + query.kinds = Some(vec![30_618]); + query.d_tag = Some(repo); + let events = state.db.query_events(&query).await.expect("query 30618"); + assert_eq!(events.len(), 1, "kind:30618 must be durable before release"); + assert!(state + .db + .deletion_store() + .serving_writes_drained(community) + .await + .expect("serving lease released")); + let generation = state + .db + .deletion_store() + .fence(&claim.lease) + .await + .expect("fence after publication"); + assert_eq!(generation, 1); + assert_eq!( + state + .db + .deletion_store() + .get(request.id) + .await + .expect("fenced request") + .stage, + buzz_db::deletion::DeletionStage::Fenced + ); + drop(state); + pool.close().await; + } + + #[tokio::test] + #[ignore = "requires Postgres and MinIO"] + async fn finalize_push_db_failure_after_cas_is_not_success_and_releases_lease() { + let (state, pool) = finalize_test_state().await; + let host = format!( + "git-finalize-fail-{}.example", + uuid::Uuid::new_v4().simple() + ); + let community = state + .db + .ensure_configured_community(&host) + .await + .expect("create test community") + .id; + let scratch = TempDir::new().expect("scratch"); + let ctx = pushed_context( + &state, + community, + &host, + format!("owner-{}", uuid::Uuid::new_v4().simple()), + format!("repo-{}", uuid::Uuid::new_v4().simple()), + Keys::generate().public_key(), + scratch.path(), + ) + .await; + let hooks = FinalizePushHooks { + post_cas_gate: None, + fail_ref_state_insert: true, + }; + + let response = finalize_push_inner(&state, ctx, &hooks).await; + assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR); + assert!(state + .db + .deletion_store() + .serving_writes_drained(community) + .await + .expect("serving lease released on failure")); + drop(state); + pool.close().await; + } + /// A gzip-encoded request body is transparently inflated before it /// reaches the git subprocess. Git's smart-HTTP client gzips the /// upload-pack/receive-pack request body past a size threshold (fires diff --git a/crates/buzz-relay/src/main.rs b/crates/buzz-relay/src/main.rs index 122f470b4..a866ba31d 100644 --- a/crates/buzz-relay/src/main.rs +++ b/crates/buzz-relay/src/main.rs @@ -201,6 +201,12 @@ async fn main() -> anyhow::Result<()> { error!("Failed to ensure partitions: {e}"); } + db.validate_deletion_catalog().await.map_err(|e| { + error!("Community deletion catalog validation failed: {e}"); + anyhow::anyhow!("Community deletion catalog is unsafe: {e}") + })?; + info!("Community deletion catalog and write fences verified"); + // Freshness fence probe: cursor pages route to the replica only for // history the probe has verified as fully replayed. Deliberately AFTER // the migration decision: spawn_fence_probe first verifies the diff --git a/crates/buzz-relay/src/push_runtime.rs b/crates/buzz-relay/src/push_runtime.rs index d6b7f04b9..4946b248c 100644 --- a/crates/buzz-relay/src/push_runtime.rs +++ b/crates/buzz-relay/src/push_runtime.rs @@ -460,7 +460,6 @@ async fn deliver_one( return; } }; - serving_write.begin_finalize(); match response { Ok(r) if r.status().is_success() => match r.json::().await { Ok(DeliveryResponse::Accepted) => { diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 400ed1dfe..f0daa5a73 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -376,22 +376,30 @@ async fn readiness_handler(State(state): State>) -> impl IntoRespo } let check = async { - let (pg_ok, redis_ok) = tokio::join!(state.db.ping(), async { - state.redis_pool.get().await.is_ok() - },); - (pg_ok, redis_ok) + let (pg_ok, redis_ok, deletion_catalog_ok) = tokio::join!( + state.db.ping(), + async { state.redis_pool.get().await.is_ok() }, + async { state.db.validate_deletion_catalog().await.is_ok() }, + ); + (pg_ok, redis_ok, deletion_catalog_ok) }; - let (pg_ok, redis_ok) = tokio::time::timeout(Duration::from_secs(2), check) - .await - .unwrap_or((false, false)); + let (pg_ok, redis_ok, deletion_catalog_ok) = + tokio::time::timeout(Duration::from_secs(2), check) + .await + .unwrap_or((false, false, false)); - if pg_ok && redis_ok { + if pg_ok && redis_ok && deletion_catalog_ok { (StatusCode::OK, Json(json!({"status": "ready"}))).into_response() } else { ( StatusCode::SERVICE_UNAVAILABLE, - Json(json!({"status": "not_ready", "postgres": pg_ok, "redis": redis_ok})), + Json(json!({ + "status": "not_ready", + "postgres": pg_ok, + "redis": redis_ok, + "deletion_catalog": deletion_catalog_ok + })), ) .into_response() } diff --git a/migrations/0027_community_deletion.sql b/migrations/0027_community_deletion.sql index 1f15a60f1..b31690f53 100644 --- a/migrations/0027_community_deletion.sql +++ b/migrations/0027_community_deletion.sql @@ -133,24 +133,32 @@ LANGUAGE SQL IMMUTABLE STRICT PARALLEL SAFE AS $$ SELECT hashtextextended('buzz-community-deletion:' || target::text, 0) $$; -CREATE FUNCTION enforce_community_write_fence() RETURNS TRIGGER +-- Keep the deletion control plane writable while its target tenant is fenced. +-- This predicate is the single SQL source of truth used by attachment and live +-- catalog validation. +CREATE FUNCTION community_write_fence_excluded_table(target NAME) RETURNS BOOLEAN +LANGUAGE SQL IMMUTABLE STRICT PARALLEL SAFE AS $$ + SELECT target::TEXT = ANY (ARRAY[ + 'community_deletion_requests', + 'community_deletion_approvals', + 'community_deletion_checkpoints', + 'community_deletion_retention_exceptions', + 'community_serving_write_leases', + 'community_deletion_executor_heartbeats' + ]::TEXT[]) +$$; + +CREATE FUNCTION assert_community_write_allowed(target UUID) RETURNS VOID LANGUAGE plpgsql AS $$ DECLARE - target UUID; lifecycle TEXT; generation BIGINT; executor_community TEXT; executor_generation TEXT; BEGIN - IF TG_OP = 'INSERT' THEN - target := NEW.community_id; - ELSE - target := OLD.community_id; - END IF; - -- Nullable operator-attribution rows without a tenant are unrelated. IF target IS NULL THEN - RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; + RETURN; END IF; PERFORM pg_advisory_xact_lock_shared(community_deletion_lock_key(target)); @@ -163,18 +171,43 @@ BEGIN USING ERRCODE = 'object_not_in_prerequisite_state'; END IF; + -- Authorization is evaluated independently for every community checked. executor_community := current_setting('buzz.deletion_executor_community', true); executor_generation := current_setting('buzz.deletion_fence_generation', true); - IF executor_community = target::text + IF executor_community = target::TEXT AND executor_generation ~ '^[0-9]+$' AND executor_generation::BIGINT = generation THEN - RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; + RETURN; END IF; IF lifecycle <> 'active' THEN RAISE EXCEPTION 'community write fenced: community % generation %', target, generation USING ERRCODE = 'object_not_in_prerequisite_state'; END IF; +END +$$; + +CREATE FUNCTION enforce_community_write_fence() RETURNS TRIGGER +LANGUAGE plpgsql AS $$ +BEGIN + IF TG_OP = 'INSERT' THEN + PERFORM assert_community_write_allowed(NEW.community_id); + ELSIF TG_OP = 'DELETE' THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id IS NOT DISTINCT FROM NEW.community_id THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id IS NULL THEN + PERFORM assert_community_write_allowed(NEW.community_id); + ELSIF NEW.community_id IS NULL THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id < NEW.community_id THEN + PERFORM assert_community_write_allowed(OLD.community_id); + PERFORM assert_community_write_allowed(NEW.community_id); + ELSE + PERFORM assert_community_write_allowed(NEW.community_id); + PERFORM assert_community_write_allowed(OLD.community_id); + END IF; + RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; END $$; @@ -224,16 +257,60 @@ CREATE TRIGGER communities_deletion_tombstone BEFORE UPDATE OR DELETE ON communities FOR EACH ROW EXECUTE FUNCTION enforce_community_tombstone(); +-- Attach the universal fence to one community-scoped relation. Future +-- migrations must invoke this helper explicitly after CREATE/ALTER introduces +-- community_id; the migration lint enforces that contract. +CREATE FUNCTION attach_community_write_fence(target REGCLASS) RETURNS VOID +LANGUAGE plpgsql AS $$ +DECLARE + relation_name NAME; +BEGIN + SELECT c.relname + INTO relation_name + FROM pg_class c + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE c.oid = target + AND n.nspname = 'public' + AND c.relkind IN ('r', 'p') + AND NOT c.relispartition; + IF NOT FOUND THEN + RAISE EXCEPTION 'community write fence target % is not a public table', target + USING ERRCODE = 'wrong_object_type'; + END IF; + IF community_write_fence_excluded_table(relation_name) THEN + RETURN; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_attribute + WHERE attrelid = target AND attname = 'community_id' AND NOT attisdropped + ) THEN + RAISE EXCEPTION 'community write fence target % has no community_id', target + USING ERRCODE = 'undefined_column'; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_trigger + WHERE tgrelid = target + AND tgname = 'community_write_fence_' || relation_name + AND NOT tgisinternal + ) THEN + EXECUTE format( + 'CREATE TRIGGER %I BEFORE INSERT OR UPDATE OR DELETE ON %s ' + 'FOR EACH ROW EXECUTE FUNCTION enforce_community_write_fence()', + 'community_write_fence_' || relation_name, + target + ); + END IF; +END +$$; + -- Attach the universal fence to every existing table carrying community_id, -- including deployment-private sidecars whose community_id is provenance. --- Deletion control-plane tables are excluded: they must remain writable while --- the target tenant is fenced. DO $$ DECLARE - table_name TEXT; + target REGCLASS; BEGIN - FOR table_name IN - SELECT c.relname + FOR target IN + SELECT c.oid::REGCLASS FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace JOIN pg_attribute a ON a.attrelid = c.oid @@ -242,21 +319,10 @@ BEGIN AND NOT c.relispartition AND a.attname = 'community_id' AND NOT a.attisdropped - AND c.relname NOT IN ( - 'community_deletion_requests', - 'community_deletion_approvals', - 'community_deletion_checkpoints', - 'community_deletion_retention_exceptions', - 'community_serving_write_leases', - 'community_deletion_executor_heartbeats' - ) + AND NOT community_write_fence_excluded_table(c.relname) + ORDER BY c.oid::REGCLASS::TEXT LOOP - EXECUTE format( - 'CREATE TRIGGER %I BEFORE INSERT OR UPDATE OR DELETE ON %I ' - 'FOR EACH ROW EXECUTE FUNCTION enforce_community_write_fence()', - 'community_write_fence_' || table_name, - table_name - ); + PERFORM attach_community_write_fence(target); END LOOP; END $$; diff --git a/schema/schema.sql b/schema/schema.sql index 2d74fe6c0..651ac932d 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1192,37 +1192,85 @@ CREATE FUNCTION community_deletion_lock_key(target UUID) RETURNS BIGINT LANGUAGE SQL IMMUTABLE STRICT PARALLEL SAFE AS $$ SELECT hashtextextended('buzz-community-deletion:' || target::text, 0) $$; -CREATE FUNCTION enforce_community_write_fence() RETURNS TRIGGER +-- Keep the deletion control plane writable while its target tenant is fenced. +-- This predicate is the single SQL source of truth used by attachment and live +-- catalog validation. +CREATE FUNCTION community_write_fence_excluded_table(target NAME) RETURNS BOOLEAN +LANGUAGE SQL IMMUTABLE STRICT PARALLEL SAFE AS $$ + SELECT target::TEXT = ANY (ARRAY[ + 'community_deletion_requests', + 'community_deletion_approvals', + 'community_deletion_checkpoints', + 'community_deletion_retention_exceptions', + 'community_serving_write_leases', + 'community_deletion_executor_heartbeats' + ]::TEXT[]) +$$; + +CREATE FUNCTION assert_community_write_allowed(target UUID) RETURNS VOID LANGUAGE plpgsql AS $$ DECLARE - target UUID; lifecycle TEXT; generation BIGINT; executor_community TEXT; executor_generation TEXT; BEGIN - IF TG_OP = 'INSERT' THEN target := NEW.community_id; ELSE target := OLD.community_id; END IF; - IF target IS NULL THEN RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; END IF; + -- Nullable operator-attribution rows without a tenant are unrelated. + IF target IS NULL THEN + RETURN; + END IF; + PERFORM pg_advisory_xact_lock_shared(community_deletion_lock_key(target)); - SELECT deletion_state, deletion_fence_generation INTO lifecycle, generation - FROM communities WHERE id = target; + SELECT deletion_state, deletion_fence_generation + INTO lifecycle, generation + FROM communities + WHERE id = target; IF NOT FOUND THEN RAISE EXCEPTION 'community write rejected: community % is missing', target USING ERRCODE = 'object_not_in_prerequisite_state'; END IF; + + -- Authorization is evaluated independently for every community checked. executor_community := current_setting('buzz.deletion_executor_community', true); executor_generation := current_setting('buzz.deletion_fence_generation', true); - IF executor_community = target::text AND executor_generation ~ '^[0-9]+$' + IF executor_community = target::TEXT + AND executor_generation ~ '^[0-9]+$' AND executor_generation::BIGINT = generation THEN - RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; + RETURN; END IF; + IF lifecycle <> 'active' THEN RAISE EXCEPTION 'community write fenced: community % generation %', target, generation USING ERRCODE = 'object_not_in_prerequisite_state'; END IF; +END +$$; + +CREATE FUNCTION enforce_community_write_fence() RETURNS TRIGGER +LANGUAGE plpgsql AS $$ +BEGIN + IF TG_OP = 'INSERT' THEN + PERFORM assert_community_write_allowed(NEW.community_id); + ELSIF TG_OP = 'DELETE' THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id IS NOT DISTINCT FROM NEW.community_id THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id IS NULL THEN + PERFORM assert_community_write_allowed(NEW.community_id); + ELSIF NEW.community_id IS NULL THEN + PERFORM assert_community_write_allowed(OLD.community_id); + ELSIF OLD.community_id < NEW.community_id THEN + PERFORM assert_community_write_allowed(OLD.community_id); + PERFORM assert_community_write_allowed(NEW.community_id); + ELSE + PERFORM assert_community_write_allowed(NEW.community_id); + PERFORM assert_community_write_allowed(OLD.community_id); + END IF; + RETURN CASE WHEN TG_OP = 'DELETE' THEN OLD ELSE NEW END; END $$; + CREATE FUNCTION enforce_community_tombstone() RETURNS TRIGGER LANGUAGE plpgsql AS $$ DECLARE @@ -1253,22 +1301,72 @@ END $$; CREATE TRIGGER communities_deletion_tombstone BEFORE UPDATE OR DELETE ON communities FOR EACH ROW EXECUTE FUNCTION enforce_community_tombstone(); -DO $$ -DECLARE table_name TEXT; +-- Attach the universal fence to one community-scoped relation. Future +-- migrations must invoke this helper explicitly after CREATE/ALTER introduces +-- community_id; the migration lint enforces that contract. +CREATE FUNCTION attach_community_write_fence(target REGCLASS) RETURNS VOID +LANGUAGE plpgsql AS $$ +DECLARE + relation_name NAME; BEGIN - FOR table_name IN - SELECT c.relname FROM pg_class c - JOIN pg_namespace n ON n.oid = c.relnamespace - JOIN pg_attribute a ON a.attrelid = c.oid - WHERE n.nspname = 'public' AND c.relkind IN ('r', 'p') AND NOT c.relispartition - AND a.attname = 'community_id' AND NOT a.attisdropped - AND c.relname NOT IN ('community_deletion_requests', 'community_deletion_approvals', - 'community_deletion_checkpoints', 'community_deletion_retention_exceptions', - 'community_serving_write_leases', 'community_deletion_executor_heartbeats') - LOOP - EXECUTE format('CREATE TRIGGER %I BEFORE INSERT OR UPDATE OR DELETE ON %I ' + SELECT c.relname + INTO relation_name + FROM pg_class c + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE c.oid = target + AND n.nspname = 'public' + AND c.relkind IN ('r', 'p') + AND NOT c.relispartition; + IF NOT FOUND THEN + RAISE EXCEPTION 'community write fence target % is not a public table', target + USING ERRCODE = 'wrong_object_type'; + END IF; + IF community_write_fence_excluded_table(relation_name) THEN + RETURN; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_attribute + WHERE attrelid = target AND attname = 'community_id' AND NOT attisdropped + ) THEN + RAISE EXCEPTION 'community write fence target % has no community_id', target + USING ERRCODE = 'undefined_column'; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_trigger + WHERE tgrelid = target + AND tgname = 'community_write_fence_' || relation_name + AND NOT tgisinternal + ) THEN + EXECUTE format( + 'CREATE TRIGGER %I BEFORE INSERT OR UPDATE OR DELETE ON %s ' 'FOR EACH ROW EXECUTE FUNCTION enforce_community_write_fence()', - 'community_write_fence_' || table_name, table_name); + 'community_write_fence_' || relation_name, + target + ); + END IF; +END +$$; + +-- Attach the universal fence to every existing table carrying community_id, +-- including deployment-private sidecars whose community_id is provenance. +DO $$ +DECLARE + target REGCLASS; +BEGIN + FOR target IN + SELECT c.oid::REGCLASS + FROM pg_class c + JOIN pg_namespace n ON n.oid = c.relnamespace + JOIN pg_attribute a ON a.attrelid = c.oid + WHERE n.nspname = 'public' + AND c.relkind IN ('r', 'p') + AND NOT c.relispartition + AND a.attname = 'community_id' + AND NOT a.attisdropped + AND NOT community_write_fence_excluded_table(c.relname) + ORDER BY c.oid::REGCLASS::TEXT + LOOP + PERFORM attach_community_write_fence(target); END LOOP; END $$;