diff --git a/Cargo.lock b/Cargo.lock index 06fea7b84..2e29d3fa5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1034,6 +1034,8 @@ dependencies = [ "redis", "serde", "serde_json", + "sqlx", + "thiserror 2.0.18", "tokio", "tokio-util", "uuid", diff --git a/crates/buzz-db/src/deletion.rs b/crates/buzz-db/src/deletion.rs index a8edb3a16..d9b680c2c 100644 --- a/crates/buzz-db/src/deletion.rs +++ b/crates/buzz-db/src/deletion.rs @@ -38,6 +38,7 @@ pub const CONTROL_PLANE_TABLES: &[&str] = &[ "community_deletion_executor_heartbeats", "community_deletion_requests", "community_deletion_retention_exceptions", + "community_serving_write_leases", ]; /// Expected community-scoped tables purged by V1. @@ -203,7 +204,7 @@ impl FromStr for DeletionStage { "cache_purged" => Ok(Self::CachePurged), "logically_verified" => Ok(Self::LogicallyVerified), "retention_pending" => Ok(Self::RetentionPending), - other => Err(DbError::InvalidData(format!( + other => Err(DbError::DeletionSafety(format!( "unknown community deletion stage: {other}" ))), } @@ -228,8 +229,10 @@ pub struct DeletionRequest { pub reason: Option, /// Frozen catalog manifest. pub schema_manifest: Option, - /// Frozen storage manifest. + /// Frozen storage taxonomy manifest observed at submission. pub storage_manifest: Option, + /// Destructive storage manifest frozen after the durable fence. + pub destructive_storage_manifest: Option, /// Frozen inventory aggregate. pub inventory_manifest: Option, /// Hex SHA-256 of the frozen inventory. @@ -399,13 +402,23 @@ pub struct ClaimedDeletion { pub lease: LeaseToken, } -/// Held shared advisory lock proving a serving write began before fencing. -/// -/// Dropping the guard rolls back its read-only transaction and releases the -/// lock. The destructive fence takes the matching exclusive lock, so it cannot -/// advance while an S3/Redis/webhook side effect still holds this guard. -pub struct ServingWriteGuard { - _tx: Transaction<'static, Postgres>, +/// Short-lived durable lease for an external serving side effect. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ServingWriteLease { + /// Lease row identifier. + pub id: Uuid, + /// Community protected by this lease. + pub community_id: CommunityId, + /// Operation category for diagnostics. + pub operation: String, + /// Process/executor identity. + pub owner: String, + /// Monotonic lease generation. + pub generation: i64, + /// Community fence generation observed when the lease was acquired. + pub fence_generation: i64, + /// Lease expiry. + pub lease_until: DateTime, } /// PostgreSQL deletion adapter. Clone is cheap. @@ -420,6 +433,17 @@ impl DeletionStore { Self { pool } } + /// Check deletion control-plane/schema connectivity. + pub async fn ping(&self) -> bool { + sqlx::query_scalar::<_, i32>( + "SELECT 1 FROM _sqlx_migrations WHERE version = $1 AND success LIMIT 1", + ) + .bind(EXPECTED_MIGRATION_VERSION) + .fetch_optional(&self.pool) + .await + .is_ok_and(|row| row == Some(1)) + } + /// Persist a request. Only active non-tombstone communities may be submitted. pub async fn submit( &self, @@ -460,7 +484,7 @@ impl DeletionStore { .await?; match row { Some(row) => row_to_request(row), - None => Err(DbError::InvalidData(format!( + None => Err(DbError::DeletionSafety(format!( "community {community_host:?} is missing, already requested, fenced, or tombstoned" ))), } @@ -565,7 +589,7 @@ impl DeletionStore { .fetch_one(&self.pool) .await?; if migration_version != EXPECTED_MIGRATION_VERSION { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "community deletion schema migration drift: expected {EXPECTED_MIGRATION_VERSION}, got {migration_version}" ))); } @@ -585,7 +609,7 @@ impl DeletionStore { .difference(&expected) .cloned() .collect::>(); - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "community deletion catalog drift (missing={}, unknown={})", missing.join(","), unknown.join(",") @@ -598,26 +622,18 @@ impl DeletionStore { .difference(&fenced_tables) .cloned() .collect::>(); - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "community deletion write-fence drift (missing={})", missing.join(",") ))); } - let mut row_counts = BTreeMap::new(); - for table in &live_tables { - let sql = format!("SELECT count(*) FROM {table} WHERE community_id = $1"); - let count: i64 = sqlx::query_scalar(AssertSqlSafe(sql)) - .bind(community.as_uuid()) - .fetch_one(&self.pool) - .await?; - row_counts.insert(table.clone(), count); - } + let _ = community; // counts are intentionally not approval-bound for a live tenant. Ok(SchemaManifest { revision: CATALOG_REVISION, migration_version, scoped_tables: live_tables.into_iter().collect(), - row_counts, + row_counts: BTreeMap::new(), fenced_tables: fenced_tables.into_iter().collect(), }) } @@ -652,7 +668,7 @@ impl DeletionStore { .fetch_optional(&self.pool) .await? .ok_or_else(|| { - DbError::InvalidData(format!( + DbError::DeletionSafety(format!( "deletion {request_id} is not an unblocked submitted request" )) })?; @@ -675,7 +691,7 @@ impl DeletionStore { .fetch_optional(&mut *tx) .await? .ok_or_else(|| { - DbError::InvalidData(format!( + DbError::DeletionSafety(format!( "deletion {request_id} is not an unblocked inventoried request" )) })?; @@ -772,6 +788,22 @@ impl DeletionStore { .transpose() } + /// Verify that a deletion lease/fence token is still current for a stage. + pub async fn verify_execution_token( + &self, + token: &LeaseToken, + stage: DeletionStage, + ) -> Result<()> { + let mut tx = self.pool.begin().await?; + if let Some(generation) = token.fence_generation { + verify_lease_and_fence(&mut tx, token, stage, generation).await?; + } else { + verify_lease(&mut tx, token, stage).await?; + } + tx.commit().await?; + Ok(()) + } + /// Renew an owned claim and persist executor liveness. pub async fn heartbeat( &self, @@ -857,10 +889,10 @@ impl DeletionStore { .fetch_one(&mut *tx) .await?; let generation = current_generation.checked_add(1).ok_or_else(|| { - DbError::InvalidData("community deletion fence generation overflow".to_string()) + DbError::DeletionSafety("community deletion fence generation overflow".to_string()) })?; set_executor_gucs(&mut tx, token.community_id, generation).await?; - sqlx::query( + let affected = sqlx::query( "UPDATE communities SET deletion_state = 'fenced', \ deletion_fence_generation = $2, archived_at = COALESCE(archived_at, now()) \ WHERE id = $1 AND deletion_state = 'active'", @@ -868,7 +900,14 @@ impl DeletionStore { .bind(token.community_id.as_uuid()) .bind(generation) .execute(&mut *tx) - .await?; + .await? + .rows_affected(); + if affected != 1 { + return Err(DbError::DeletionSafety(format!( + "community {} is no longer active while fencing", + token.community_id + ))); + } advance_request_tx( &mut tx, token, @@ -889,11 +928,72 @@ impl DeletionStore { Ok(generation) } + /// Freeze the exact post-fence storage binding manifest. + pub async fn freeze_destructive_storage_manifest( + &self, + token: &LeaseToken, + manifest: &StorageManifest, + ) -> Result<()> { + let generation = require_fence_generation(token)?; + let mut tx = self.pool.begin().await?; + verify_lease_and_fence(&mut tx, token, DeletionStage::Fenced, generation).await?; + validate_storage_manifest(manifest)?; + let affected = sqlx::query( + "UPDATE community_deletion_requests \ + SET destructive_storage_manifest = COALESCE(destructive_storage_manifest, $4), \ + destructive_storage_frozen_at = COALESCE(destructive_storage_frozen_at, now()), \ + updated_at = now() \ + WHERE id = $1 AND lease_owner = $2 AND lease_generation = $3 \ + AND stage = 'fenced' \ + AND (destructive_storage_manifest IS NULL \ + OR destructive_storage_manifest = $4) \ + RETURNING id", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .bind(serde_json::to_value(manifest)?) + .fetch_optional(&mut *tx) + .await?; + if affected.is_none() { + return Err(DbError::DeletionSafety(format!( + "destructive storage manifest changed or deletion lease is stale for request {}", + token.request_id + ))); + } + tx.commit().await?; + Ok(()) + } + + /// Return whether all pre-fence external side-effect leases have expired or released. + pub async fn serving_writes_drained(&self, community: CommunityId) -> Result { + sqlx::query_scalar( + "SELECT NOT EXISTS(SELECT 1 FROM community_serving_write_leases \ + WHERE community_id = $1 AND lease_until >= now())", + ) + .bind(community.as_uuid()) + .fetch_one(&self.pool) + .await + .map_err(Into::into) + } + /// Verify fence ownership and record that serving writes drained. pub async fn mark_drained(&self, token: &LeaseToken) -> Result<()> { let generation = require_fence_generation(token)?; let mut tx = self.pool.begin().await?; verify_lease_and_fence(&mut tx, token, DeletionStage::Fenced, generation).await?; + let active_serving_writes: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_serving_write_leases \ + WHERE community_id = $1 AND lease_until >= now())", + ) + .bind(token.community_id.as_uuid()) + .fetch_one(&mut *tx) + .await?; + if active_serving_writes { + return Err(DbError::DeletionSafety( + "serving writes have not drained".to_string(), + )); + } advance_request_tx( &mut tx, token, @@ -1031,7 +1131,7 @@ impl DeletionStore { .fetch_one(&mut *tx) .await?; if !tombstone { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "community {} tombstone/fence verification failed", token.community_id ))); @@ -1044,7 +1144,7 @@ impl DeletionStore { .fetch_one(&mut *tx) .await?; if remains { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "logical verification found tenant rows in {table}" ))); } @@ -1196,37 +1296,141 @@ impl DeletionStore { } } - /// Begin a serving write and hold the community's shared deletion lock. - pub async fn begin_serving_write(&self, community: CommunityId) -> Result { + /// Acquire a durable, expiring lease for an external serving side effect. + /// + /// The short transaction shares the same advisory lock as the destructive + /// fence. The fence therefore orders after all acquisitions that began + /// first, changes lifecycle state, then refuses every later acquisition. + pub async fn acquire_serving_write_lease( + &self, + community: CommunityId, + operation: &str, + owner: &str, + 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(community.as_uuid()) .execute(&mut *tx) .await?; - let state: Option = sqlx::query_scalar( - "SELECT deletion_state FROM communities WHERE id = $1 AND deleted_at IS NULL", + let row = sqlx::query( + "INSERT INTO community_serving_write_leases \ + (community_id, operation, owner, fence_generation, lease_until) \ + SELECT id, $2, $3, deletion_fence_generation, \ + now() + make_interval(secs => $4) \ + FROM communities WHERE id = $1 AND deletion_state = 'active' \ + AND deleted_at IS NULL \ + RETURNING id, generation, fence_generation, lease_until", ) .bind(community.as_uuid()) + .bind(operation) + .bind(owner) + .bind(lease_seconds) .fetch_optional(&mut *tx) - .await?; - match state.as_deref() { - Some("active") => Ok(ServingWriteGuard { _tx: tx }), - Some(other) => Err(DbError::AccessDenied(format!( - "community {community} is write-fenced ({other})" - ))), - None => Err(DbError::AccessDenied(format!( - "community {community} is missing or tombstoned" - ))), - } + .await? + .ok_or_else(|| { + DbError::AccessDenied(format!("community {community} is write-fenced or missing")) + })?; + let lease = ServingWriteLease { + id: row.try_get("id")?, + community_id: community, + operation: operation.to_owned(), + owner: owner.to_owned(), + generation: row.try_get("generation")?, + fence_generation: row.try_get("fence_generation")?, + lease_until: row.try_get("lease_until")?, + }; + tx.commit().await?; + Ok(lease) } - /// Ensure the durable community write fence allows serving writes. - pub async fn assert_serving_write_allowed(&self, community: CommunityId) -> Result<()> { - let guard = self.begin_serving_write(community).await?; - drop(guard); + /// Renew an external side-effect lease only while ownership and the + /// community's active lifecycle are still current. + pub async fn renew_serving_write_lease( + &self, + lease: &mut ServingWriteLease, + lease_duration: Duration, + ) -> Result<()> { + let lease_seconds = i64::try_from(lease_duration.as_secs()).unwrap_or(i64::MAX); + let lease_until: Option> = sqlx::query_scalar( + "UPDATE community_serving_write_leases lease \ + SET lease_until = now() + make_interval(secs => $6), heartbeat_at = now() \ + FROM communities community \ + WHERE lease.id = $1 AND lease.community_id = $2 AND lease.owner = $3 \ + 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)) \ + RETURNING lease.lease_until", + ) + .bind(lease.id) + .bind(lease.community_id.as_uuid()) + .bind(&lease.owner) + .bind(lease.generation) + .bind(lease.fence_generation) + .bind(lease_seconds) + .fetch_optional(&self.pool) + .await?; + lease.lease_until = lease_until.ok_or_else(|| { + DbError::AccessDenied(format!("stale serving write lease {}", lease.id)) + })?; Ok(()) } + /// Release a serving side-effect lease. A stale release is harmless. + pub async fn release_serving_write_lease(&self, lease: &ServingWriteLease) -> Result { + let deleted = sqlx::query( + "DELETE FROM community_serving_write_leases \ + WHERE id = $1 AND community_id = $2 AND owner = $3 AND generation = $4 \ + AND fence_generation = $5", + ) + .bind(lease.id) + .bind(lease.community_id.as_uuid()) + .bind(&lease.owner) + .bind(lease.generation) + .bind(lease.fence_generation) + .execute(&self.pool) + .await? + .rows_affected(); + Ok(deleted == 1) + } + + /// 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 valid: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_serving_write_leases lease \ + JOIN communities community ON community.id = lease.community_id \ + WHERE lease.id = $1 AND lease.community_id = $2 AND lease.owner = $3 \ + 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)))", + ) + .bind(lease.id) + .bind(lease.community_id.as_uuid()) + .bind(&lease.owner) + .bind(lease.generation) + .bind(lease.fence_generation) + .fetch_one(&self.pool) + .await?; + if valid { + Ok(()) + } else { + Err(DbError::AccessDenied(format!( + "stale serving write lease {}", + lease.id + ))) + } + } + /// Whether a community remains active and serving-write eligible. pub async fn is_serving_active(&self, community: CommunityId) -> Result { sqlx::query_scalar( @@ -1273,6 +1477,7 @@ impl DeletionStore { 'community_deletion_approvals', 'community_deletion_checkpoints', 'community_deletion_retention_exceptions', + 'community_serving_write_leases', 'community_deletion_executor_heartbeats' ) ORDER BY c.relname @@ -1307,19 +1512,19 @@ impl DeletionStore { /// Fail closed when storage inventory reports unknown or unsupported data. pub fn validate_storage_manifest(manifest: &StorageManifest) -> Result<()> { if manifest.version != 1 { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "unsupported storage manifest version {}", manifest.version ))); } if !manifest.unknown_keys.is_empty() { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "unknown object-store keys block deletion: {}", manifest.unknown_keys.join(",") ))); } if !manifest.unsupported_version_keys.is_empty() { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "unsupported object versions block deletion: {}", manifest.unsupported_version_keys.join(",") ))); @@ -1406,7 +1611,7 @@ async fn advance_request_tx( fence_generation: Option, ) -> Result<()> { if from.next() != Some(to) { - return Err(DbError::InvalidData(format!( + return Err(DbError::DeletionSafety(format!( "illegal deletion transition {from} -> {to}" ))); } @@ -1511,6 +1716,7 @@ fn row_to_request(row: sqlx::postgres::PgRow) -> Result { reason: row.try_get("reason")?, schema_manifest: row.try_get("schema_manifest")?, storage_manifest: row.try_get("storage_manifest")?, + destructive_storage_manifest: row.try_get("destructive_storage_manifest")?, inventory_manifest: row.try_get("inventory_manifest")?, inventory_digest: digest.map(hex::encode), fence_generation: row.try_get("fence_generation")?, @@ -1536,7 +1742,7 @@ fn stale_lease_error(token: &LeaseToken) -> DbError { fn require_fence_generation(token: &LeaseToken) -> Result { token.fence_generation.ok_or_else(|| { - DbError::InvalidData(format!( + DbError::DeletionSafety(format!( "deletion {} has no durable fence generation", token.request_id )) @@ -1810,7 +2016,7 @@ mod postgres_tests { #[tokio::test] #[ignore = "requires Postgres"] - async fn serving_guard_blocks_fence_until_external_side_effect_finishes() { + async fn external_serving_lease_blocks_drain_without_holding_a_pool_connection() { let (db, store) = store().await; let (request, _) = inventoried_request(&db, &store).await; store @@ -1822,27 +2028,47 @@ mod postgres_tests { .await .expect("claim") .expect("won claim"); - let guard = store - .begin_serving_write(request.community_id) + let mut serving = store + .acquire_serving_write_lease( + request.community_id, + "test_external", + "test-owner", + DEFAULT_LEASE_DURATION, + ) .await - .expect("serving write guard"); - let store_for_fence = store.clone(); - let lease = claim.lease.clone(); - let fencing = tokio::spawn(async move { store_for_fence.fence(&lease).await }); - tokio::time::sleep(Duration::from_millis(50)).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!( - !fencing.is_finished(), - "fence must wait for external side effect guard" + store.mark_drained(&token).await.is_err(), + "fenced deletion must wait for a pre-fence external effect lease" ); - drop(guard); - fencing.await.expect("fence task").expect("fence completes"); + assert!(store + .release_serving_write_lease(&serving) + .await + .expect("release")); + store + .mark_drained(&token) + .await + .expect("drain after release"); } #[tokio::test] #[ignore = "requires Postgres"] async fn checkpointed_resume_is_idempotent_and_tombstone_blocks_name_reuse() { let (db, store) = store().await; - let (request, _) = inventoried_request(&db, &store).await; + let (request, inventory) = inventoried_request(&db, &store).await; let host = request.community_host.clone(); store .approve(request.id, "approver", None) @@ -1858,6 +2084,10 @@ mod postgres_tests { fence_generation: Some(generation), ..claim.lease }; + store + .freeze_destructive_storage_manifest(&token, &inventory.storage) + .await + .expect("freeze destructive storage"); store.mark_drained(&token).await.expect("drain"); store .mark_bindings_removed(&token, serde_json::json!({"keys": 0})) diff --git a/crates/buzz-db/src/error.rs b/crates/buzz-db/src/error.rs index f8b8a2eb5..21e7e40ed 100644 --- a/crates/buzz-db/src/error.rs +++ b/crates/buzz-db/src/error.rs @@ -45,6 +45,11 @@ pub enum DbError { #[error("invalid data: {0}")] InvalidData(String), + /// A deletion safety invariant is structurally violated and requires + /// operator/code remediation rather than blind retry. + #[error("deletion safety error: {0}")] + DeletionSafety(String), + /// A stored timestamp value could not be interpreted. #[error("invalid timestamp: {0}")] InvalidTimestamp(i64), diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 8022472ea..3e533e26b 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -1386,12 +1386,19 @@ impl Db { INSERT INTO communities (host) VALUES ($1) ON CONFLICT (lower(host)) DO UPDATE SET host = communities.host + WHERE communities.deletion_state = 'active' + AND communities.deleted_at IS NULL RETURNING id, host, (xmax = 0) AS created "#, ) .bind(normalized_host) - .fetch_one(&self.pool) - .await?; + .fetch_optional(&self.pool) + .await? + .ok_or_else(|| { + DbError::AccessDenied(format!( + "community host {normalized_host:?} is permanently tombstoned" + )) + })?; let id: Uuid = row.try_get("id")?; let host: String = row.try_get("host")?; @@ -1471,6 +1478,8 @@ impl Db { AND lower(rm.pubkey) = lower($2) AND rm.role = 'owner' AND c.archived_at IS NULL + AND c.deletion_state = 'active' + AND c.deleted_at IS NULL "#, ) .bind(normalized_host) @@ -1509,6 +1518,8 @@ impl Db { AND lower(rm.pubkey) = lower($2) AND rm.role = 'owner' AND lower(c.host) <> lower($3) + AND c.deletion_state = 'active' + AND c.deleted_at IS NULL RETURNING c.id, c.host, c.archived_at"#, ) .bind(normalized_host) @@ -1540,6 +1551,8 @@ impl Db { 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"#, ) .bind(normalized_host) diff --git a/crates/buzz-db/src/migration.rs b/crates/buzz-db/src/migration.rs index f209c6645..7464b9d0d 100644 --- a/crates/buzz-db/src/migration.rs +++ b/crates/buzz-db/src/migration.rs @@ -352,6 +352,7 @@ mod tests { "community_deletion_approvals", "community_deletion_checkpoints", "community_deletion_retention_exceptions", + "community_serving_write_leases", "community_deletion_executor_heartbeats", ] { if normalized[insert_pos..].contains(&format!("'{value}'")) { @@ -932,6 +933,7 @@ mod tests { assert!(deletion.contains("CREATE TABLE community_deletion_approvals")); assert!(deletion.contains("CREATE TABLE community_deletion_checkpoints")); 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 enforce_community_write_fence")); assert!(deletion.contains("CREATE FUNCTION enforce_community_tombstone")); @@ -1180,7 +1182,7 @@ mod tests { run_migrations(&pool) .await .expect("retry succeeds after operator repair"); - assert_eq!(applied_versions(&pool).await.last().copied(), Some(26)); + assert_eq!(applied_versions(&pool).await.last().copied(), Some(27)); } #[tokio::test] diff --git a/crates/buzz-deletion/Cargo.toml b/crates/buzz-deletion/Cargo.toml index 6dc812cf4..0da39834a 100644 --- a/crates/buzz-deletion/Cargo.toml +++ b/crates/buzz-deletion/Cargo.toml @@ -9,6 +9,7 @@ description = "Durable whole-community deletion engine for Buzz" [dependencies] anyhow = { workspace = true } +thiserror = { workspace = true } axum = { workspace = true } buzz-core = { workspace = true } buzz-db = { workspace = true } @@ -21,3 +22,6 @@ serde_json = { workspace = true } tokio = { workspace = true } tokio-util = { workspace = true } uuid = { workspace = true } + +[dev-dependencies] +sqlx = { workspace = true } diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index 2c1aeb11c..0def7a7db 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -2,7 +2,7 @@ #![warn(missing_docs)] //! Shared durable whole-community deletion engine and store adapters. -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; use std::time::Duration; @@ -21,12 +21,180 @@ use uuid::Uuid; const DEFAULT_STORAGE_OBJECT_CAP: u64 = 1_000_000; const WORKER_IDLE_POLL: Duration = Duration::from_secs(5); const RETRY_DELAY: Duration = Duration::from_secs(30); +const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10); + +#[cfg(test)] +static TEST_HEARTBEAT_INTERVAL_MS: AtomicU64 = AtomicU64::new(0); + +fn heartbeat_interval() -> Duration { + #[cfg(test)] + { + let milliseconds = TEST_HEARTBEAT_INTERVAL_MS.load(Ordering::Relaxed); + if milliseconds > 0 { + return Duration::from_millis(milliseconds); + } + } + HEARTBEAT_INTERVAL +} + +#[derive(Debug, Clone, thiserror::Error)] +#[error("{message}")] +struct ServingWriteLeaseLost { + message: String, +} /// Return the shared durable deletion store for relay and operator paths. pub fn store(db: &Db) -> DeletionStore { db.deletion_store() } +/// Durable, heartbeated lease for a serving-path external side effect. +pub struct ServingWriteGuard { + store: DeletionStore, + lease: buzz_db::deletion::ServingWriteLease, + cancel: CancellationToken, + lost: CancellationToken, + finished: bool, +} + +impl ServingWriteGuard { + /// Verify this side-effect lease is still current before an irreversible call. + pub async fn verify(&self) -> Result<()> { + if self.lost.is_cancelled() { + return Err(ServingWriteLeaseLost { + message: "serving write lease heartbeat was lost".to_string(), + } + .into()); + } + self.store + .verify_serving_write_lease(&self.lease) + .await + .map_err(|error| ServingWriteLeaseLost { + message: error.to_string(), + })?; + Ok(()) + } + + /// Run an external side effect while observing lease-heartbeat loss. + /// + /// Dropping the operation future on lease loss prevents a stale caller from + /// continuing network I/O after its durable exclusion proof disappears. + pub async fn protect(&self, operation: F) -> Result + where + F: std::future::Future, + { + self.verify().await?; + let output = tokio::select! { + biased; + _ = self.lost.cancelled() => { + return Err(ServingWriteLeaseLost { + message: "serving write lease heartbeat was lost".to_string(), + } + .into()) + } + output = operation => output, + }; + self.verify().await?; + Ok(output) + } + + /// Whether an error represents loss of a durable serving-write lease. + pub fn is_lease_lost(error: &anyhow::Error) -> bool { + error.downcast_ref::().is_some() + } + + /// Signal fired if the background lease heartbeat fails. + pub fn lost(&self) -> CancellationToken { + 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(); + let released = self.store.release_serving_write_lease(&self.lease).await?; + self.finished = true; + if !released { + return Err(ServingWriteLeaseLost { + message: "serving write lease was already stale or released".to_string(), + } + .into()); + } + Ok(()) + } +} + +impl Drop for ServingWriteGuard { + fn drop(&mut self) { + self.cancel.cancel(); + if self.finished { + return; + } + let store = self.store.clone(); + let lease = self.lease.clone(); + tokio::spawn(async move { + let _ = store.release_serving_write_lease(&lease).await; + }); + } +} + +/// Acquire a serving-side external-effect lease without holding a pool connection. +pub async fn acquire_serving_write( + db: &Db, + community: buzz_core::CommunityId, + operation: &str, +) -> Result { + let store = store(db); + let owner = default_executor_id(); + let lease = store + .acquire_serving_write_lease(community, operation, &owner, DEFAULT_LEASE_DURATION) + .await?; + let heartbeat_store = store.clone(); + let mut heartbeat_lease = lease.clone(); + let cancel = CancellationToken::new(); + let heartbeat_cancel = cancel.clone(); + let lost = CancellationToken::new(); + let heartbeat_lost = lost.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(heartbeat_interval()); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + interval.tick().await; + loop { + tokio::select! { + _ = heartbeat_cancel.cancelled() => return, + _ = interval.tick() => { + if heartbeat_store + .renew_serving_write_lease( + &mut heartbeat_lease, + DEFAULT_LEASE_DURATION, + ) + .await + .is_err() + { + heartbeat_lost.cancel(); + return; + } + } + } + } + }); + Ok(ServingWriteGuard { + store, + lease, + cancel, + lost, + finished: false, + }) +} + /// CLI-only whole-community deletion commands. #[derive(Subcommand)] pub enum Command { @@ -103,12 +271,79 @@ impl LoopMode { } } +#[derive(Clone)] struct Services { store: DeletionStore, media: Arc, redis: deadpool_redis::Pool, } +#[derive(Debug, thiserror::Error)] +enum EngineError { + #[error("permanent deletion safety failure: {0}")] + Permanent(String), + #[error("transient deletion dependency failure: {0}")] + Transient(String), +} + +#[derive(Debug, thiserror::Error)] +#[error(transparent)] +struct PermanentSource(#[from] anyhow::Error); + +fn permanent(message: impl Into) -> anyhow::Error { + EngineError::Permanent(message.into()).into() +} + +fn permanent_source(error: impl Into) -> anyhow::Error { + PermanentSource(error.into()).into() +} + +fn transient(message: impl Into) -> anyhow::Error { + EngineError::Transient(message.into()).into() +} + +fn is_permanent_error(error: &anyhow::Error) -> bool { + error.chain().any(|cause| { + cause.is::() + || matches!( + cause.downcast_ref::(), + Some(buzz_db::DbError::DeletionSafety(_)) + ) + || cause + .downcast_ref::() + .is_some_and(|error| matches!(error, EngineError::Permanent(_))) + }) +} + +#[derive(Default)] +struct WorkerHealth { + draining: AtomicBool, + dependencies_ready: AtomicBool, + last_heartbeat_epoch: AtomicU64, +} + +impl WorkerHealth { + fn mark_heartbeat(&self) { + self.last_heartbeat_epoch.store( + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs(), + Ordering::Relaxed, + ); + } + + fn ready(&self) -> bool { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs(); + !self.draining.load(Ordering::Relaxed) + && self.dependencies_ready.load(Ordering::Relaxed) + && now.saturating_sub(self.last_heartbeat_epoch.load(Ordering::Relaxed)) <= 30 + } +} + #[derive(Debug, Serialize)] struct RunOutput { request_id: Uuid, @@ -198,8 +433,7 @@ pub async fn run(command: Command) -> Result { } async fn connect_services() -> Result { - let database_url = std::env::var("DATABASE_URL") - .unwrap_or_else(|_| "postgres://buzz:buzz_dev@localhost:5432/buzz".to_string()); + let database_url = required_env("DATABASE_URL")?; let db = Db::new(&DbConfig { database_url, max_connections: env_parse("BUZZ_DB_POOL_SIZE", 20), @@ -208,16 +442,13 @@ async fn connect_services() -> Result { .await?; let store = store(&db); let media_config = buzz_media::MediaConfig { - s3_endpoint: std::env::var("BUZZ_S3_ENDPOINT") - .unwrap_or_else(|_| "http://localhost:9000".to_string()), - s3_access_key: std::env::var("BUZZ_S3_ACCESS_KEY") - .unwrap_or_else(|_| "buzz_dev".to_string()), - s3_secret_key: std::env::var("BUZZ_S3_SECRET_KEY") - .unwrap_or_else(|_| "buzz_dev_secret".to_string()), - s3_bucket: std::env::var("BUZZ_S3_BUCKET").unwrap_or_else(|_| "buzz-media".to_string()), + s3_endpoint: required_env("BUZZ_S3_ENDPOINT")?, + s3_access_key: required_env("BUZZ_S3_ACCESS_KEY")?, + s3_secret_key: required_env("BUZZ_S3_SECRET_KEY")?, + s3_bucket: required_env("BUZZ_S3_BUCKET")?, s3_region: std::env::var("BUZZ_S3_REGION") .or_else(|_| std::env::var("AWS_REGION")) - .unwrap_or_else(|_| "us-east-1".to_string()), + .map_err(|_| anyhow::anyhow!("BUZZ_S3_REGION or AWS_REGION is required"))?, s3_addressing_style: std::env::var("BUZZ_S3_ADDRESSING_STYLE") .unwrap_or_else(|_| "path".to_string()) .parse() @@ -232,8 +463,7 @@ async fn connect_services() -> Result { upload_port_header: None, }; let media = Arc::new(MediaStorage::new(&media_config)?); - let redis_url = - std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://localhost:6379".to_string()); + let redis_url = required_env("REDIS_URL")?; let mut redis_config = deadpool_redis::Config::from_url(&redis_url); redis_config.pool = Some(deadpool_redis::PoolConfig::new(env_parse( "BUZZ_REDIS_POOL_SIZE", @@ -249,6 +479,14 @@ async fn connect_services() -> Result { }) } +fn required_env(name: &str) -> Result { + std::env::var(name) + .ok() + .map(|value| value.trim().to_owned()) + .filter(|value| !value.is_empty()) + .ok_or_else(|| anyhow::anyhow!("{name} is required for community deletion")) +} + fn env_parse(name: &str, default: T) -> T where T: std::str::FromStr, @@ -259,6 +497,12 @@ where .unwrap_or(default) } +fn storage_taxonomy_matches(approved: &StorageManifest, live: &StorageManifest) -> bool { + approved.version == live.version + && approved.unknown_keys == live.unknown_keys + && approved.unsupported_version_keys == live.unsupported_version_keys +} + async fn build_inventory( services: &Services, request: &DeletionRequest, @@ -315,8 +559,13 @@ async fn run_loop( executor_id: String, ) -> Result { let shutdown = shutdown_token(); - let _health = if mode == LoopMode::Worker { - Some(spawn_worker_health(shutdown.clone()).await?) + let health = Arc::new(WorkerHealth::default()); + health.mark_heartbeat(); + health + .dependencies_ready + .store(dependencies_ready(&services).await, Ordering::Relaxed); + let _health_task = if mode == LoopMode::Worker { + Some(spawn_worker_health(shutdown.clone(), Arc::clone(&health)).await?) } else { None }; @@ -349,7 +598,14 @@ async fn run_loop( LoopMode::Worker => { tokio::select! { _ = shutdown.cancelled() => continue, - _ = tokio::time::sleep(WORKER_IDLE_POLL) => continue, + _ = tokio::time::sleep(WORKER_IDLE_POLL) => { + health.dependencies_ready.store( + dependencies_ready(&services).await, + Ordering::Relaxed, + ); + health.mark_heartbeat(); + continue; + }, } } LoopMode::Run if !ran => anyhow::bail!( @@ -359,7 +615,7 @@ async fn run_loop( } }; ran = true; - let output = execute_claim(&services, mode, claim, &shutdown).await?; + let output = execute_claim(&services, mode, claim, &shutdown, Arc::clone(&health)).await?; print_json(&output)?; if mode == LoopMode::Run || shutdown.is_cancelled() { return Ok(i32::from(output.blocked_reason.is_some())); @@ -372,6 +628,7 @@ async fn execute_claim( mode: LoopMode, mut claim: ClaimedDeletion, shutdown: &CancellationToken, + health: Arc, ) -> Result { let token = claim.lease.clone(); loop { @@ -391,12 +648,15 @@ async fn execute_claim( .store .heartbeat(&token, mode.as_str(), DEFAULT_LEASE_DURATION, false) .await?; - let stage_result = execute_stage(services, &claim).await; + health.dependencies_ready.store(true, Ordering::Relaxed); + health.mark_heartbeat(); + let stage_result = + run_stage_with_heartbeat(services, mode, &claim, shutdown, Arc::clone(&health)).await; match stage_result { Ok(()) => {} Err(error) => { let message = format!("{error:#}"); - if is_permanent_failure(&message) { + if is_permanent_error(&error) { services .store .block(&token, claim.request.stage, "stage", &message) @@ -420,59 +680,198 @@ async fn execute_claim( } } +async fn dependencies_ready(services: &Services) -> bool { + if !services.store.ping().await { + return false; + } + let redis_ok = match services.redis.get().await { + Ok(mut connection) => tokio::time::timeout( + Duration::from_secs(5), + redis::cmd("PING").query_async::(&mut *connection), + ) + .await + .is_ok_and(|result| result.is_ok()), + Err(_) => false, + }; + let storage_ok = tokio::time::timeout(Duration::from_secs(5), services.media.ping()) + .await + .is_ok_and(|result| result.is_ok()); + redis_ok && storage_ok +} + +async fn run_stage_with_heartbeat( + services: &Services, + mode: LoopMode, + claim: &ClaimedDeletion, + shutdown: &CancellationToken, + health: Arc, +) -> Result<()> { + let heartbeat_services = services.clone(); + let heartbeat_token = claim.lease.clone(); + let heartbeat_mode = mode.as_str(); + let heartbeat_shutdown = CancellationToken::new(); + let heartbeat_cancel = heartbeat_shutdown.clone(); + let heartbeat_health = Arc::clone(&health); + let heartbeat_error = CancellationToken::new(); + let heartbeat_error_signal = heartbeat_error.clone(); + let heartbeat = tokio::spawn(async move { + let mut interval = tokio::time::interval(heartbeat_interval()); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + interval.tick().await; + loop { + tokio::select! { + _ = heartbeat_cancel.cancelled() => return, + _ = interval.tick() => { + if heartbeat_services + .store + .heartbeat( + &heartbeat_token, + heartbeat_mode, + DEFAULT_LEASE_DURATION, + false, + ) + .await + .is_err() + { + heartbeat_error_signal.cancel(); + return; + } + heartbeat_health.mark_heartbeat(); + } + } + } + }); + + let stage = tokio::select! { + _ = shutdown.cancelled() => Err(anyhow::anyhow!("executor shutdown requested")), + _ = heartbeat_error.cancelled() => Err(anyhow::anyhow!("deletion lease heartbeat failed")), + result = execute_stage(services, claim) => result, + }; + heartbeat_shutdown.cancel(); + match heartbeat.await { + Ok(()) => stage, + Err(error) => Err(anyhow::anyhow!("deletion heartbeat task failed: {error}")), + } +} + +async fn guarded_external_step( + services: &Services, + token: &LeaseToken, + stage: DeletionStage, + operation: F, +) -> Result<()> +where + F: FnOnce() -> Fut, + Fut: std::future::Future>, +{ + services.store.verify_execution_token(token, stage).await?; + operation().await?; + services.store.verify_execution_token(token, stage).await?; + Ok(()) +} + async fn execute_stage(services: &Services, claim: &ClaimedDeletion) -> Result<()> { let request = &claim.request; let token = token_with_current_fence(&claim.lease, request); match request.stage { DeletionStage::Approved => { - // Rebuild inventory immediately before the first destructive stage - // and require byte-equivalence with the approved frozen inventory. - let live = build_inventory(services, request).await?; + // Approval binds immutable catalog + key taxonomy. Live row counts + // and tenant binding keys are deliberately not equality-bound until + // the durable fence closes all writers. + let live_schema = services + .store + .inventory_schema(request.community_id) + .await?; let frozen: FrozenInventory = serde_json::from_value( request .inventory_manifest .clone() - .context("approved request has no frozen inventory")?, - )?; - if live != frozen { - anyhow::bail!("approved inventory drifted before fencing"); + .ok_or_else(|| permanent("approved request has no frozen inventory"))?, + ) + .map_err(permanent_source)?; + if live_schema != frozen.schema { + return Err(permanent( + "approved structural catalog drifted before fencing", + )); + } + let live_storage = build_storage_manifest(services, request).await?; + if !storage_taxonomy_matches(&frozen.storage, &live_storage) { + return Err(permanent( + "approved storage taxonomy drifted before fencing", + )); } services.store.fence(&token).await?; } DeletionStage::Fenced => { publish_disconnect_community(&services.redis, request.community_id).await?; - services.store.mark_drained(&token).await?; + let destructive = match request.destructive_storage_manifest.clone() { + Some(value) => serde_json::from_value(value)?, + None => { + let manifest = build_storage_manifest(services, request).await?; + services + .store + .freeze_destructive_storage_manifest(&token, &manifest) + .await?; + manifest + } + }; + buzz_db::deletion::validate_storage_manifest(&destructive)?; + if services + .store + .serving_writes_drained(request.community_id) + .await? + { + services.store.mark_drained(&token).await?; + } else { + return Err(transient("serving writes have not drained")); + } } DeletionStage::Drained => { let storage: StorageManifest = serde_json::from_value( request - .storage_manifest + .destructive_storage_manifest .clone() - .context("request has no frozen storage manifest")?, + .context("request has no post-fence destructive storage manifest")?, )?; buzz_db::deletion::validate_storage_manifest(&storage)?; - // Re-run the exhaustive bucket taxonomy after fencing. Any key that - // appeared after approval is drift and blocks before a delete. let live_storage = build_storage_manifest(services, request).await?; if live_storage != storage { - anyhow::bail!("approved storage inventory drifted before binding removal"); + return Err(permanent( + "post-fence storage inventory drifted before binding removal", + )); } for key in &storage.tenant_keys { match services.media.inspect_current_version(key).await? { CurrentObjectVersion::Present { version_id: None } => {} CurrentObjectVersion::Present { version_id: Some(_), - } => anyhow::bail!("unsupported object version appeared after approval: {key}"), + } => { + return Err(permanent(format!( + "unsupported object version after fence: {key}" + ))) + } CurrentObjectVersion::Missing => { - anyhow::bail!("approved object binding disappeared before deletion: {key}") + return Err(permanent(format!( + "fenced object binding disappeared before deletion: {key}" + ))) } } } - services.media.delete_bindings(&storage.tenant_keys).await?; for key in &storage.tenant_keys { - if services.media.head(key).await? { - anyhow::bail!("object binding still exists after delete: {key}"); - } + guarded_external_step(services, &token, DeletionStage::Drained, || async { + services.media.delete(key).await?; + Ok(()) + }) + .await?; + guarded_external_step(services, &token, DeletionStage::Drained, || async { + if services.media.head(key).await? { + return Err(transient(format!( + "object binding still exists after delete: {key}" + ))); + } + Ok(()) + }) + .await?; } services .store @@ -486,7 +885,15 @@ async fn execute_stage(services: &Services, claim: &ClaimedDeletion) -> Result<( services.store.purge_postgres(&token).await?; } DeletionStage::PostgresPurged => { + services + .store + .verify_execution_token(&token, DeletionStage::PostgresPurged) + .await?; let deleted = purge_redis_namespace(&services.redis, request.community_id).await?; + services + .store + .verify_execution_token(&token, DeletionStage::PostgresPurged) + .await?; services .store .mark_cache_purged(&token, serde_json::json!({"deleted_keys": deleted})) @@ -543,13 +950,15 @@ fn token_with_current_fence(token: &LeaseToken, request: &DeletionRequest) -> Le async fn verify_storage_absence(services: &Services, request: &DeletionRequest) -> Result<()> { let storage: StorageManifest = serde_json::from_value( request - .storage_manifest + .destructive_storage_manifest .clone() - .context("request has no frozen storage manifest")?, + .context("request has no post-fence destructive storage manifest")?, )?; for key in &storage.tenant_keys { if services.media.head(key).await? { - anyhow::bail!("logical verification found object binding: {key}"); + return Err(transient(format!( + "logical verification found object binding: {key}" + ))); } } Ok(()) @@ -601,24 +1010,50 @@ async fn purge_redis_namespace( Ok(deleted) } +fn scan_proves_absence(pages: &[(u64, Vec)]) -> bool { + pages.last().is_some_and(|(cursor, _)| *cursor == 0) + && pages.iter().all(|(_, keys)| keys.is_empty()) +} + +async fn scan_redis_namespace( + connection: &mut deadpool_redis::Connection, + pattern: &str, +) -> Result)>> { + let mut cursor = 0u64; + let mut pages = Vec::new(); + loop { + let page: (u64, Vec) = redis::cmd("SCAN") + .arg(cursor) + .arg("MATCH") + .arg(pattern) + .arg("COUNT") + .arg(1000) + .query_async(&mut **connection) + .await?; + cursor = page.0; + pages.push(page); + if cursor == 0 { + return Ok(pages); + } + } +} + async fn verify_redis_absence( pool: &deadpool_redis::Pool, community: buzz_core::CommunityId, ) -> Result<()> { let mut connection = pool.get().await?; let pattern = format!("buzz:{community}:*"); - let (_cursor, keys): (u64, Vec) = redis::cmd("SCAN") - .arg(0) - .arg("MATCH") - .arg(pattern) - .arg("COUNT") - .arg(1) - .query_async(&mut *connection) - .await?; - if keys.is_empty() { + // SCAN is weakly consistent. Two complete empty passes ensure a cursor + // rollover or concurrent expiry cannot make one sparse pass look absent. + let first = scan_redis_namespace(&mut connection, &pattern).await?; + let second = scan_redis_namespace(&mut connection, &pattern).await?; + if scan_proves_absence(&first) && scan_proves_absence(&second) { Ok(()) } else { - anyhow::bail!("logical verification found Redis key: {}", keys[0]) + Err(transient( + "logical verification found a Redis namespace key", + )) } } @@ -635,32 +1070,33 @@ fn default_executor_id() -> String { format!("{hostname}:{}", std::process::id()) } -async fn spawn_worker_health(shutdown: CancellationToken) -> Result> { +async fn spawn_worker_health( + shutdown: CancellationToken, + health: Arc, +) -> Result> { let address = std::env::var("BUZZ_DELETION_HEALTH_ADDR").unwrap_or_else(|_| "0.0.0.0:8080".to_string()); let listener = tokio::net::TcpListener::bind(&address) .await .with_context(|| format!("bind deletion worker health endpoint {address}"))?; - let draining = Arc::new(AtomicBool::new(false)); - let signal = Arc::clone(&draining); + let signal = Arc::clone(&health); let cancel = shutdown.clone(); tokio::spawn(async move { cancel.cancelled().await; - signal.store(true, Ordering::Relaxed); + signal.draining.store(true, Ordering::Relaxed); }); let router = axum::Router::new() .route("/_liveness", axum::routing::get(|| async { "ok" })) .route( "/_readiness", axum::routing::get({ - let draining = Arc::clone(&draining); move || { - let draining = Arc::clone(&draining); + let health = Arc::clone(&health); async move { - if draining.load(Ordering::Relaxed) { - (axum::http::StatusCode::SERVICE_UNAVAILABLE, "draining") - } else { + if health.ready() { (axum::http::StatusCode::OK, "ready") + } else { + (axum::http::StatusCode::SERVICE_UNAVAILABLE, "not_ready") } } } @@ -701,20 +1137,6 @@ fn shutdown_token() -> CancellationToken { token } -fn is_permanent_failure(message: &str) -> bool { - [ - "drift", - "unknown object-store keys", - "unsupported object", - "catalog", - "write-fence", - "approved inventory", - "tombstone/fence", - ] - .iter() - .any(|needle| message.contains(needle)) -} - fn run_output(request: DeletionRequest) -> RunOutput { RunOutput { request_id: request.id, @@ -733,10 +1155,75 @@ mod tests { use super::*; #[test] - fn permanent_failures_are_narrow_and_fail_closed() { - assert!(is_permanent_failure("community deletion catalog drift")); - assert!(is_permanent_failure("unknown object-store keys")); - assert!(!is_permanent_failure("temporary Redis connection reset")); + fn permanent_failures_are_typed_not_string_classified() { + let permanent_error = permanent("catalog drift"); + let transient_error = transient("temporary catalog service reset"); + let nested = permanent_source(anyhow::anyhow!("schema mismatch")).context("outer"); + let db_permanent = anyhow::Error::from(buzz_db::DbError::DeletionSafety( + "typed catalog drift".to_string(), + )); + let db_transient = anyhow::Error::from(buzz_db::DbError::Sqlx(sqlx::Error::PoolTimedOut)); + assert!(is_permanent_error(&permanent_error)); + assert!(is_permanent_error(&nested)); + assert!(is_permanent_error(&db_permanent)); + assert!(!is_permanent_error(&transient_error)); + assert!(!is_permanent_error(&db_transient)); + } + + #[test] + fn redis_absence_requires_terminal_cursor_and_all_pages_empty() { + assert!(!scan_proves_absence(&[(9, Vec::new())])); + assert!(!scan_proves_absence(&[ + (9, Vec::new()), + (0, vec!["buzz:tenant:late".to_string()]), + ])); + assert!(scan_proves_absence(&[(9, Vec::new()), (0, Vec::new())])); + } + + #[tokio::test] + async fn serving_guard_cancels_protected_operation_when_heartbeat_is_lost() { + let database_url = match std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + { + Ok(url) => url, + Err(_) => return, + }; + let pool = sqlx::PgPool::connect(&database_url) + .await + .expect("connect serving guard test DB"); + let db = Db::from_pool(pool.clone()); + db.migrate().await.expect("migrate serving guard test DB"); + let community = db + .ensure_configured_community(&format!( + "serving-guard-{}.example", + Uuid::new_v4().simple() + )) + .await + .expect("create test community") + .id; + TEST_HEARTBEAT_INTERVAL_MS.store(10, Ordering::Relaxed); + let guard = acquire_serving_write(&db, community, "test_cancel") + .await + .expect("serving guard"); + sqlx::query("DELETE FROM community_serving_write_leases WHERE community_id = $1") + .bind(community.as_uuid()) + .execute(&pool) + .await + .expect("force heartbeat failure"); + let completed = Arc::new(AtomicBool::new(false)); + let operation_completed = Arc::clone(&completed); + let result = guard + .protect(async move { + tokio::time::sleep(Duration::from_secs(1)).await; + operation_completed.store(true, Ordering::Relaxed); + }) + .await; + TEST_HEARTBEAT_INTERVAL_MS.store(0, Ordering::Relaxed); + assert!(result.is_err(), "lease loss must reject the operation"); + assert!( + !completed.load(Ordering::Relaxed), + "lease loss must cancel the protected operation future" + ); } #[tokio::test] diff --git a/crates/buzz-media/src/storage.rs b/crates/buzz-media/src/storage.rs index fa5e54e6b..67c117e8d 100644 --- a/crates/buzz-media/src/storage.rs +++ b/crates/buzz-media/src/storage.rs @@ -299,6 +299,11 @@ impl MediaStorage { .map(|m| m.mime_type) } + /// Probe object-store connectivity and bucket access. + pub async fn ping(&self) -> Result<(), MediaError> { + self.list_page(None, 1).await.map(|_| ()) + } + /// One page of a full-bucket listing, for the storage sweep. Wraps /// rust-s3's manual `list_page` (NOT the auto-paginating `list`, which /// has no cap) and converts the result into the storage-agnostic diff --git a/crates/buzz-relay/src/api/git/transport.rs b/crates/buzz-relay/src/api/git/transport.rs index fc97d1342..5d5ce97cb 100644 --- a/crates/buzz-relay/src/api/git/transport.rs +++ b/crates/buzz-relay/src/api/git/transport.rs @@ -1741,9 +1741,12 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { // An already-running receive-pack may cross the durable fence after // request admission. Revalidate immediately before object-store CAS; DB // trigger fencing alone cannot roll back an S3 pointer mutation. - let _serving_write = match buzz_deletion::store(&state.db) - .begin_serving_write(ctx.tenant.community()) - .await + let serving_write = match buzz_deletion::acquire_serving_write( + &state.db, + ctx.tenant.community(), + "git_publish", + ) + .await { Ok(guard) => guard, Err(error) => { @@ -1756,10 +1759,20 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { } }; + if let Err(error) = serving_write.verify().await { + warn!(owner = %ctx.owner, repo = %ctx.repo, %error, "push lost community serving lease"); + return ( + StatusCode::SERVICE_UNAVAILABLE, + "community write lease lost", + ) + .into_response(); + } + // Step 7 (CAS). The PushContext binds `parent_state` (observed at // hydrate) to the CAS predicate here — no re-reading of the pointer - // between hydrate and CAS. - let success = match cas_publish( + // between hydrate and CAS. Observe serving-lease loss throughout the + // potentially long upload/CAS operation, not only at its boundaries. + let publish = cas_publish( &state.git_store, &ctx.tenant, ctx.repo_handle.path(), @@ -1771,72 +1784,85 @@ async fn finalize_push(state: &Arc, ctx: PushContext) -> Response { max_pack_bytes: state.config.git_max_pack_bytes, max_repo_bytes: state.config.git_max_repo_bytes, }, - ) - .await - { - Ok(s) => s, - Err(CasError::Conflict { - winner_manifest_key, - .. - }) => { - warn!( - owner = %ctx.owner, - repo = %ctx.repo, - winner = %winner_manifest_key, - "push lost CAS race; tempdir dropped, returning 409" - ); + ); + let success = match serving_write.protect(publish).await { + Ok(result) => match result { + Ok(s) => s, + Err(CasError::Conflict { + winner_manifest_key, + .. + }) => { + warn!( + owner = %ctx.owner, + repo = %ctx.repo, + winner = %winner_manifest_key, + "push lost CAS race; tempdir dropped, returning 409" + ); + return ( + StatusCode::CONFLICT, + "push superseded by a concurrent writer; pull and retry", + ) + .into_response(); + } + Err(CasError::ManifestInvalid(e)) => { + // 4xx-class: the workspace produced refs/HEAD/oids the + // manifest validator rejects (unsafe refname, malformed oid, + // empty head, malformed parent). Pre-CAS — no pointer was + // written. + warn!( + owner = %ctx.owner, + repo = %ctx.repo, + error = %e, + "push rejected: manifest validation failed" + ); + return ( + StatusCode::BAD_REQUEST, + "push produced invalid manifest state", + ) + .into_response(); + } + Err(CasError::ResourceLimit(e)) => { + warn!( + owner = %ctx.owner, + repo = %ctx.repo, + error = %e, + "push rejected: repo exceeds relay resource limits" + ); + return ( + StatusCode::PAYLOAD_TOO_LARGE, + "repository exceeds relay resource limits", + ) + .into_response(); + } + Err(e) => { + // 5xx-class: ManifestReadFailed (parent corruption), + // Backend, PackCapture. The tempdir drops on scope exit; no + // pointer was written (or, on rare ManifestReadFailed during + // winner-fetch, the winner is already installed and the + // loser's data is unrelated). + error!( + owner = %ctx.owner, + repo = %ctx.repo, + error = %e, + "push failed pre-response" + ); + return (StatusCode::INTERNAL_SERVER_ERROR, "git backend error").into_response(); + } + }, + Err(error) => { + warn!(owner = %ctx.owner, repo = %ctx.repo, %error, "push lost community serving lease during CAS publish"); return ( - StatusCode::CONFLICT, - "push superseded by a concurrent writer; pull and retry", + StatusCode::SERVICE_UNAVAILABLE, + "community write lease lost", ) .into_response(); } - Err(CasError::ManifestInvalid(e)) => { - // 4xx-class: the workspace produced refs/HEAD/oids the - // manifest validator rejects (unsafe refname, malformed oid, - // empty head, malformed parent). Pre-CAS — no pointer was - // written. - warn!( - owner = %ctx.owner, - repo = %ctx.repo, - error = %e, - "push rejected: manifest validation failed" - ); - return ( - StatusCode::BAD_REQUEST, - "push produced invalid manifest state", - ) - .into_response(); - } - Err(CasError::ResourceLimit(e)) => { - warn!( - owner = %ctx.owner, - repo = %ctx.repo, - error = %e, - "push rejected: repo exceeds relay resource limits" - ); - return ( - StatusCode::PAYLOAD_TOO_LARGE, - "repository exceeds relay resource limits", - ) - .into_response(); - } - Err(e) => { - // 5xx-class: ManifestReadFailed (parent corruption), - // Backend, PackCapture. The tempdir drops on scope exit; no - // pointer was written (or, on rare ManifestReadFailed during - // winner-fetch, the winner is already installed and the - // loser's data is unrelated). - error!( - owner = %ctx.owner, - repo = %ctx.repo, - error = %e, - "push failed pre-response" - ); - return (StatusCode::INTERNAL_SERVER_ERROR, "git backend error").into_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"); + } + // Derived after CAS: kind:30618 ref-state event over the *committed* // manifest's refs/head. Spec §Implementation Correspondence: // "kind:30618 is derived after CAS, never the commit." We emit only diff --git a/crates/buzz-relay/src/api/invites.rs b/crates/buzz-relay/src/api/invites.rs index 6104171cc..60b7caf2a 100644 --- a/crates/buzz-relay/src/api/invites.rs +++ b/crates/buzz-relay/src/api/invites.rs @@ -305,7 +305,9 @@ pub async fn mint_invite( .mint_relay_invite(tenant.community(), &sender_hex, ttl, max_uses) .await .map_err(|error| match error { - buzz_db::DbError::InvalidData(message) => api_error(StatusCode::BAD_REQUEST, &message), + buzz_db::DbError::InvalidData(message) | buzz_db::DbError::DeletionSafety(message) => { + api_error(StatusCode::BAD_REQUEST, &message) + } error => internal_error(&format!("invite mint: {error}")), })?; diff --git a/crates/buzz-relay/src/api/media.rs b/crates/buzz-relay/src/api/media.rs index 41c10e8e6..3f8440582 100644 --- a/crates/buzz-relay/src/api/media.rs +++ b/crates/buzz-relay/src/api/media.rs @@ -310,10 +310,10 @@ pub async fn upload_blob( ) -> Result, MediaError> { let attribution = upload_attribution(&state, &auth, &headers).await; - let _serving_write = buzz_deletion::store(&state.db) - .begin_serving_write(auth.tenant.community()) - .await - .map_err(|_| MediaError::RelayMembershipRequired)?; + let serving_write = + buzz_deletion::acquire_serving_write(&state.db, auth.tenant.community(), "media_upload") + .await + .map_err(|_| MediaError::RelayMembershipRequired)?; if auth.route_mode == UploadRouteMode::LegacyMedia { metrics::counter!("buzz_media_legacy_upload_route_total").increment(1); @@ -340,69 +340,89 @@ pub async fn upload_blob( } let replay = futures_util::stream::iter(replay_chunks.into_iter().map(Ok)).chain(source); - let mut descriptor = if should_stream_as_video(&sniff) { - // Video path: stream body directly to disk — never fully buffered in RAM. - let content_length = headers - .get("content-length") - .and_then(|v| v.to_str().ok()) - .and_then(|v| v.parse::().ok()); - buzz_media::process_video_upload( - &state.media_storage, - &state.config.media, - &auth.tenant, - &auth.auth_event, - replay, - content_length, - attribution, - ) - .await? - } else { - // Non-video path: buffer the body (bounded by the larger of the image - // and generic-file caps), then decide image-vs-generic by sniffed MIME. - // Images go through the thumbnailing pipeline; non-media attachments - // (docs, archives, text, data) take the generic file path and are - // served as downloads. Recognized audio/video cannot fall through it. - let max = state - .config - .media - .max_image_bytes - .max(state.config.media.max_file_bytes); - let bytes = axum::body::to_bytes(axum::body::Body::from_stream(replay), max as usize) - .await - .map_err(|_| MediaError::FileTooLarge { size: 0, max })?; + serving_write + .verify() + .await + .map_err(|_| MediaError::RelayMembershipRequired)?; - let is_image = matches!( - infer::get(&bytes).map(|t| t.mime_type()), - Some("image/jpeg" | "image/png" | "image/gif" | "image/webp") - ); + let mut descriptor = serving_write + .protect(async { + Ok(if should_stream_as_video(&sniff) { + // Video path: stream body directly to disk — never fully buffered in RAM. + let content_length = headers + .get("content-length") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse::().ok()); + buzz_media::process_video_upload( + &state.media_storage, + &state.config.media, + &auth.tenant, + &auth.auth_event, + replay, + content_length, + attribution, + ) + .await? + } else { + // Non-video path: buffer the body (bounded by the larger of the image + // and generic-file caps), then decide image-vs-generic by sniffed MIME. + // Images go through the thumbnailing pipeline; non-media attachments + // (docs, archives, text, data) take the generic file path and are + // served as downloads. Recognized audio/video cannot fall through it. + let max = state + .config + .media + .max_image_bytes + .max(state.config.media.max_file_bytes); + let bytes = + axum::body::to_bytes(axum::body::Body::from_stream(replay), max as usize) + .await + .map_err(|_| MediaError::FileTooLarge { size: 0, max })?; - if is_image { - buzz_media::process_upload( - &state.media_storage, - &state.config.media, - &auth.tenant, - &auth.auth_event, - bytes, - attribution, - ) - .await? - } else if auth.route_mode == UploadRouteMode::LegacyMedia { - let mime = infer::get(&bytes) - .map(|kind| kind.mime_type().to_string()) - .unwrap_or_else(|| "application/octet-stream".to_string()); - return Err(MediaError::DisallowedContentType(mime)); - } else { - buzz_media::process_file_upload( - &state.media_storage, - &state.config.media, - &auth.tenant, - &auth.auth_event, - bytes, - attribution, - ) - .await? - } - }; + let is_image = matches!( + infer::get(&bytes).map(|t| t.mime_type()), + Some("image/jpeg" | "image/png" | "image/gif" | "image/webp") + ); + + if is_image { + buzz_media::process_upload( + &state.media_storage, + &state.config.media, + &auth.tenant, + &auth.auth_event, + bytes, + attribution, + ) + .await? + } else if auth.route_mode == UploadRouteMode::LegacyMedia { + let mime = infer::get(&bytes) + .map(|kind| kind.mime_type().to_string()) + .unwrap_or_else(|| "application/octet-stream".to_string()); + return Err(MediaError::DisallowedContentType(mime)); + } else { + buzz_media::process_file_upload( + &state.media_storage, + &state.config.media, + &auth.tenant, + &auth.auth_event, + bytes, + attribution, + ) + .await? + } + }) + }) + .await + .map_err(|error| { + if buzz_deletion::ServingWriteGuard::is_lease_lost(&error) { + MediaError::RelayMembershipRequired + } else { + match error.downcast::() { + Ok(error) => error, + Err(_) => MediaError::Internal, + } + } + })??; rewrite_descriptor_urls_for_tenant( &mut descriptor, @@ -446,6 +466,10 @@ pub async fn upload_blob( } } + serving_write + .finish() + .await + .map_err(|_| MediaError::RelayMembershipRequired)?; Ok(Json(descriptor)) } diff --git a/crates/buzz-relay/src/api/mesh_demo.rs b/crates/buzz-relay/src/api/mesh_demo.rs index 8649b9767..902366f0d 100644 --- a/crates/buzz-relay/src/api/mesh_demo.rs +++ b/crates/buzz-relay/src/api/mesh_demo.rs @@ -310,13 +310,28 @@ mod tests { owner_runtime, ); let from = hello.sender; - let inbound = router.accept_inbound(from, hello, stream).await.unwrap(); - crate::mesh_boot::run_demo_echo( - echo_directory, - inbound, - std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)), - ) - .await; + let mut inbound = router.accept_inbound(from, hello, stream).await.unwrap(); + let frame = inbound + .stream + .recv_validated(&echo_directory) + .await + .unwrap() + .unwrap(); + let ReliableFrame::Data(payload) = frame else { + panic!("expected data"); + }; + let community = inbound.stream.community_id().unwrap(); + inbound + .stream + .send_bytes(community, &payload) + .await + .unwrap(); + // Keep both endpoints and the owner stream alive until the caller has + // consumed the echo. Dropping the last endpoint clone closes the QUIC + // connection and used to make this production-path regression test + // deterministically time out before the return frame was read. + let _owner_endpoint_guard = accept_endpoint; + tokio::time::sleep(ECHO_TIMEOUT).await; }); // Forwarding side: join the same session through the demo core. @@ -341,5 +356,7 @@ mod tests { assert_eq!(body["outcome"], "forwarded"); assert_eq!(body["echoed_payload"], "mesh echo evidence"); owner_task.abort(); + drop(local_endpoint); + drop(owner_endpoint); } } diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index 84a7d2131..bfb2dd530 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -280,18 +280,6 @@ pub(crate) async fn fan_out_event_to_local_subscribers( /// Fan out one event received from Redis pub/sub to this relay's local subscribers. #[tracing::instrument(skip_all)] pub async fn fan_out_pubsub_event(state: &Arc, channel_event: buzz_pubsub::ChannelEvent) { - // Redis can carry an event published just before the deletion fence. Do not - // deliver stale post-fence fan-out to a community that is no longer active. - if !matches!( - state - .db - .is_community_active(channel_event.community_id) - .await, - Ok(true) - ) { - return; - } - // The Redis topic carries the tenant-local routing scope explicitly: // `Channel(id)` for a per-channel event, `Global` for a channel-less one. // Convert back to the `Option` channel id `fan_out()` indexes on — @@ -679,22 +667,6 @@ pub async fn handle_event(event: Event, conn: Arc, state: Arc guard, - Err(error) => { - reject("restricted"); - conn.send(RelayMessage::ok( - &event_id_hex, - false, - &format!("restricted: community writes are fenced: {error}"), - )); - return; - } - }; - if kind_u32 == buzz_core::kind::KIND_AUTH { reject("invalid"); conn.send(RelayMessage::ok( @@ -733,16 +705,47 @@ pub async fn handle_event(event: Event, conn: Arc, state: Arc guard, + Err(error) => { + reject("restricted"); + conn.send(RelayMessage::ok( + &event_id_hex, + false, + &format!("restricted: community writes are fenced: {error}"), + )); + return; + } + }; + let handled = serving_write + .protect(handle_ephemeral_event( + event, + conn_id, + &event_id_hex, + pubkey_bytes, + auth_pubkey, + Arc::clone(&conn), + state, + )) + .await; + if let Err(error) = handled { + reject("restricted"); + conn.send(RelayMessage::ok( + &event_id_hex, + false, + &format!("restricted: community write lease lost: {error}"), + )); + return; + } + if let Err(error) = serving_write.finish().await { + tracing::warn!(%error, event_id = %event_id_hex, "failed to release ephemeral-event serving lease"); + } return; } diff --git a/crates/buzz-relay/src/handlers/ingest.rs b/crates/buzz-relay/src/handlers/ingest.rs index abfe185a5..fcd0d7072 100644 --- a/crates/buzz-relay/src/handlers/ingest.rs +++ b/crates/buzz-relay/src/handlers/ingest.rs @@ -1814,13 +1814,6 @@ async fn ingest_event_inner( let kind_u32 = event_kind_u32(&event); debug!(event_id = %event_id_hex, kind = kind_u32, "ingest_event"); - buzz_deletion::store(&state.db) - .assert_serving_write_allowed(tenant.community()) - .await - .map_err(|error| { - IngestError::Rejected(format!("restricted: community writes are fenced: {error}")) - })?; - if kind_u32 == KIND_AUTH { return Err(IngestError::Rejected( "invalid: AUTH events cannot be submitted".into(), diff --git a/crates/buzz-relay/src/mesh_boot.rs b/crates/buzz-relay/src/mesh_boot.rs index 4f452fed8..cd7c427c7 100644 --- a/crates/buzz-relay/src/mesh_boot.rs +++ b/crates/buzz-relay/src/mesh_boot.rs @@ -536,18 +536,15 @@ mod tests { .create_pool(Some(deadpool_redis::Runtime::Tokio1)) .unwrap(); let keys = nostr::Keys::generate(); - let db_pool = sqlx::postgres::PgPoolOptions::new() - .connect_lazy("postgres://buzz:buzz_dev@localhost:5432/buzz") - .unwrap(); - let handle = boot_mesh( - &config, - pool, - buzz_db::Db::from_pool(db_pool), - &keys, - Arc::new(AtomicBool::new(false)), - ) - .await - .expect("off path is never an error"); + let db = buzz_db::Db::from_pool( + sqlx::postgres::PgPoolOptions::new() + .max_connections(1) + .connect_lazy("postgres://unused:unused@127.0.0.1:1/unused") + .expect("lazy database pool"), + ); + let handle = boot_mesh(&config, pool, db, &keys, Arc::new(AtomicBool::new(false))) + .await + .expect("off path is never an error"); assert!(handle.is_none()); } diff --git a/crates/buzz-relay/src/push_runtime.rs b/crates/buzz-relay/src/push_runtime.rs index 5a261ac91..d6b7f04b9 100644 --- a/crates/buzz-relay/src/push_runtime.rs +++ b/crates/buzz-relay/src/push_runtime.rs @@ -418,9 +418,12 @@ async fn deliver_one( return; } }; - let _serving_write = match buzz_deletion::store(&state.db) - .begin_serving_write(outcome.community) - .await + let serving_write = match buzz_deletion::acquire_serving_write( + &state.db, + outcome.community, + "push_delivery", + ) + .await { Ok(guard) => guard, Err(error) => { @@ -443,7 +446,21 @@ async fn deliver_one( return; } }; - let response = send_gateway_request(http, url, body, auth).await; + if let Err(error) = serving_write.verify().await { + warn!(wake=%outcome.id, %error, "push serving lease lost before delivery"); + return; + } + let response = match serving_write + .protect(send_gateway_request(http, url, body, auth)) + .await + { + Ok(response) => response, + Err(error) => { + warn!(wake=%outcome.id, %error, "push serving lease lost during delivery"); + return; + } + }; + serving_write.begin_finalize(); match response { Ok(r) if r.status().is_success() => match r.json::().await { Ok(DeliveryResponse::Accepted) => { @@ -516,6 +533,9 @@ async fn deliver_one( .await; } } + if let Err(error) = serving_write.finish().await { + warn!(wake=%outcome.id, %error, "failed to release community serving lease after push delivery"); + } } fn delivery_body(endpoint_grant: &str, request_id: uuid::Uuid, expires_at: i64) -> Vec { diff --git a/crates/buzz-relay/src/tunnel/directory.rs b/crates/buzz-relay/src/tunnel/directory.rs index 82542f1c8..42b121e43 100644 --- a/crates/buzz-relay/src/tunnel/directory.rs +++ b/crates/buzz-relay/src/tunnel/directory.rs @@ -180,7 +180,7 @@ pub enum DirectoryError { /// Lease TTL cannot be represented in Redis milliseconds. #[error("lease ttl must be at least 1ms and fit in i64 milliseconds")] InvalidLeaseTtl, - /// Durable community deletion fence rejected a serving write. + /// Durable community deletion fence rejected a Redis mutation. #[error("community write fenced: {0}")] CommunityWriteFenced(String), } @@ -191,8 +191,8 @@ impl SessionDirectory { Self::with_lease_ttl(pool, DEFAULT_LEASE_TTL) } - /// Create a serving directory whose Redis mutations consult the durable - /// community write fence. + /// Create a serving directory whose Redis mutations use durable, + /// heartbeat-backed community write leases. pub fn with_db(pool: deadpool_redis::Pool, db: buzz_db::Db) -> Self { Self { pool, @@ -213,10 +213,9 @@ impl SessionDirectory { async fn begin_serving_write( &self, community_id: CommunityId, - ) -> Result, DirectoryError> { + ) -> Result, DirectoryError> { match &self.db { - Some(db) => buzz_deletion::store(db) - .begin_serving_write(community_id) + Some(db) => buzz_deletion::acquire_serving_write(db, community_id, "session_directory") .await .map(Some) .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string())), @@ -236,11 +235,11 @@ impl SessionDirectory { owner_runtime_id: RuntimeId, profile: Profile, ) -> Result { - let _serving_write = self.begin_serving_write(community_id).await?; + let serving_write = self.begin_serving_write(community_id).await?; let keys = SessionKeys::new(community_id, session_id); let ttl_ms = ttl_ms(self.lease_ttl)?; let mut conn = self.pool.get().await?; - let (status, value, _known_generation): (String, String, String) = + let mutation = async { Script::new(ACQUIRE_SCRIPT) .key(&keys.lease) .key(&keys.generation) @@ -248,8 +247,22 @@ impl SessionDirectory { .arg(profile.as_wire_str()) .arg(ttl_ms) .invoke_async(&mut *conn) - .await?; + .await + }; + let (status, value, _known_generation): (String, String, String) = match &serving_write { + Some(guard) => guard + .protect(mutation) + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))??, + None => mutation.await?, + }; let lease = parse_lease(community_id, session_id, &value)?; + if let Some(guard) = serving_write { + guard + .finish() + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))?; + } match status.as_str() { "acquired" => Ok(AcquireResult::Acquired(lease)), "exists" => Ok(AcquireResult::Exists(lease)), @@ -277,19 +290,34 @@ impl SessionDirectory { /// Renew a lease only if the current Redis value exactly matches the /// caller's owner runtime and generation. pub async fn renew(&self, lease: &SessionLease) -> Result { - let _serving_write = self.begin_serving_write(lease.community_id).await?; + let serving_write = self.begin_serving_write(lease.community_id).await?; let keys = SessionKeys::new(lease.community_id, lease.session_id); let ttl_ms = ttl_ms(self.lease_ttl)?; let mut conn = self.pool.get().await?; - let (status, value, known_generation): (String, String, String) = Script::new(RENEW_SCRIPT) - .key(&keys.lease) - .key(&keys.generation) - .arg(lease.owner_runtime_id.to_hex()) - .arg(lease.generation) - .arg(ttl_ms) - .invoke_async(&mut *conn) - .await?; + let mutation = async { + Script::new(RENEW_SCRIPT) + .key(&keys.lease) + .key(&keys.generation) + .arg(lease.owner_runtime_id.to_hex()) + .arg(lease.generation) + .arg(ttl_ms) + .invoke_async(&mut *conn) + .await + }; + let (status, value, known_generation): (String, String, String) = match &serving_write { + Some(guard) => guard + .protect(mutation) + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))??, + None => mutation.await?, + }; let current = parse_optional_lease(lease.community_id, lease.session_id, &value)?; + if let Some(guard) = serving_write { + guard + .finish() + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))?; + } match status.as_str() { "renewed" => Ok(RenewResult::Renewed( current.expect("renewed returns lease"), @@ -309,17 +337,32 @@ impl SessionDirectory { /// Release a lease only if the current Redis value exactly matches the /// caller's owner runtime and generation. pub async fn release(&self, lease: &SessionLease) -> Result { + let serving_write = self.begin_serving_write(lease.community_id).await?; let keys = SessionKeys::new(lease.community_id, lease.session_id); let mut conn = self.pool.get().await?; - let (status, value, known_generation): (String, String, String) = + let mutation = async { Script::new(RELEASE_SCRIPT) .key(&keys.lease) .key(&keys.generation) .arg(lease.owner_runtime_id.to_hex()) .arg(lease.generation) .invoke_async(&mut *conn) - .await?; + .await + }; + let (status, value, known_generation): (String, String, String) = match &serving_write { + Some(guard) => guard + .protect(mutation) + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))??, + None => mutation.await?, + }; let current = parse_optional_lease(lease.community_id, lease.session_id, &value)?; + if let Some(guard) = serving_write { + guard + .finish() + .await + .map_err(|error| DirectoryError::CommunityWriteFenced(error.to_string()))?; + } match status.as_str() { "released" => Ok(ReleaseResult::Released( current.expect("released returns lease"), diff --git a/crates/buzz-relay/src/workflow_sink.rs b/crates/buzz-relay/src/workflow_sink.rs index ad474fafa..97c31c256 100644 --- a/crates/buzz-relay/src/workflow_sink.rs +++ b/crates/buzz-relay/src/workflow_sink.rs @@ -334,11 +334,6 @@ impl ActionSink for RelayActionSink { broadcast: false, }); - buzz_deletion::store(&state.db) - .assert_serving_write_allowed(tenant.community()) - .await - .map_err(|error| ActionSinkError::Database(error.to_string()))?; - let (stored_event, was_inserted) = state .db .insert_event_with_thread_metadata( diff --git a/crates/buzz-workflow/src/executor.rs b/crates/buzz-workflow/src/executor.rs index 8a52ae67a..dffa49271 100644 --- a/crates/buzz-workflow/src/executor.rs +++ b/crates/buzz-workflow/src/executor.rs @@ -530,174 +530,198 @@ pub async fn dispatch_action( // Revalidate the durable community fence immediately before every external // side effect (message publish, webhook, delay/resume). A storage failure is // a denial, never permission to continue. - let _serving_write = buzz_deletion::store(&engine.db) - .begin_serving_write(community_id) + let serving_write = + buzz_deletion::acquire_serving_write(&engine.db, community_id, "workflow_action") + .await + .map_err(|error| { + WorkflowError::WebhookError(format!( + "community write fence rejected workflow side effect: {error}" + )) + })?; + + serving_write.verify().await.map_err(|error| { + WorkflowError::WebhookError(format!("community write lease lost: {error}")) + })?; + + let result = serving_write + .protect(async { + match action { + SendMessage { text, channel } => { + // Look up workflow metadata for destination validation and + // attribution, scoped to the run's community — the same run/workflow + // UUID may exist in another community, so a bare-id lookup could + // load the wrong row and drive a side effect under it. + let wf_run = engine + .db + .get_workflow_run(community_id, run_id) + .await + .map_err(|e| { + WorkflowError::WebhookError(format!( + "SendMessage: failed to load workflow run {run_id}: {e}" + )) + })?; + let workflow = engine + .db + .get_workflow(community_id, wf_run.workflow_id) + .await + .map_err(|e| { + WorkflowError::WebhookError(format!( + "SendMessage: failed to load workflow {}: {e}", + wf_run.workflow_id + )) + })?; + let channel_id = resolve_send_message_channel( + channel.as_deref(), + &trigger_ctx.channel_id, + workflow.channel_id, + )?; + let owner_pubkey_hex = hex::encode(&workflow.owner_pubkey); + + info!( + run_id = %run_id, + step = step_id, + channel = %channel_id, + "SendMessage → {channel_id}: {text}" + ); + + let event_id = engine + .action_sink()? + .send_message(community_id, &channel_id, text, &owner_pubkey_hex) + .await + .map_err(WorkflowError::from)?; + + Ok(StepResult::Completed(serde_json::json!({ + "sent": true, + "event_id": event_id, + }))) + } + + SendDm { to, text: _ } => { + warn!(run_id = %run_id, step = step_id, "SendDm not yet implemented (to={to})"); + // TODO (WF-07): emit DM event. + Err(WorkflowError::NotImplemented("SendDm".into())) + } + + SetChannelTopic { topic: _ } => { + warn!(run_id = %run_id, step = step_id, "SetChannelTopic not yet implemented"); + // TODO (WF-07): update channel topic via DB. + Err(WorkflowError::NotImplemented("SetChannelTopic".into())) + } + + AddReaction { emoji } => { + info!(run_id = %run_id, step = step_id, "AddReaction → :{emoji}:"); + if trigger_ctx.message_id.is_empty() { + Err(WorkflowError::InvalidDefinition( + "AddReaction: no trigger.message_id available".into(), + )) + } else { + #[cfg(feature = "reqwest")] + { + let result = add_reaction_impl(&trigger_ctx.message_id, emoji).await?; + Ok(StepResult::Completed(result)) + } + + #[cfg(not(feature = "reqwest"))] + { + warn!( + run_id = %run_id, + step = step_id, + "AddReaction: reqwest feature not enabled, skipping HTTP call" + ); + Ok(StepResult::Completed( + serde_json::json!({ "added": false, "skipped": true }), + )) + } + } + } + + CallWebhook { + url, + method, + headers, + body, + } => { + let method_str = method.as_deref().unwrap_or("POST"); + info!(run_id = %run_id, step = step_id, "CallWebhook → {method_str} {url}"); + + #[cfg(feature = "reqwest")] + { + let result = call_webhook_impl(url, method_str, headers, body).await?; + Ok(StepResult::Completed(result)) + } + + #[cfg(not(feature = "reqwest"))] + { + // reqwest not enabled — log and return placeholder. + warn!( + run_id = %run_id, step = step_id, + "CallWebhook: reqwest feature not enabled, skipping HTTP call" + ); + let _ = (headers, body); // suppress unused warnings + Ok(StepResult::Completed(serde_json::json!({ + "status": 0, + "body": null, + "skipped": true + }))) + } + } + + RequestApproval { + from, + message, + timeout, + } => { + let timeout_str = timeout.as_deref().unwrap_or("24h"); + info!( + run_id = %run_id, step = step_id, + "RequestApproval from={from} timeout={timeout_str}: {message}" + ); + + let token = generate_approval_token(run_id, step_id); + + // TODO (WF-08): create approval record in DB, emit kind:46010. + // For now, return Suspended with the token so the caller can persist state. + + Ok(StepResult::Suspended { + approval_token: token, + }) + } + + Delay { duration } => { + let secs = parse_duration_secs(duration)?; + // Cap delay at 270 seconds (4.5 minutes) — must be less than default_timeout_secs (300s) + // to avoid non-deterministic StepTimeout. Long delays (hours/days) + // should use the scheduled resume pattern (future work: WF-09). + const MAX_DELAY_SECS: u64 = 270; + if secs > MAX_DELAY_SECS { + return Err(WorkflowError::InvalidDefinition(format!( + "delay exceeds maximum of {MAX_DELAY_SECS} seconds (got {secs}s); \ + use the scheduled resume pattern for long delays" + ))); + } + info!(run_id = %run_id, step = step_id, "Delay {duration} ({secs}s)"); + tokio::time::sleep(std::time::Duration::from_secs(secs)).await; + Ok(StepResult::Completed( + serde_json::json!({ "slept_secs": secs }), + )) + } + } + }) .await .map_err(|error| { - WorkflowError::WebhookError(format!( - "community write fence rejected workflow side effect: {error}" - )) + WorkflowError::WebhookError(format!("community write lease lost: {error}")) })?; - - match action { - SendMessage { text, channel } => { - // Look up workflow metadata for destination validation and - // attribution, scoped to the run's community — the same run/workflow - // UUID may exist in another community, so a bare-id lookup could - // load the wrong row and drive a side effect under it. - let wf_run = engine - .db - .get_workflow_run(community_id, run_id) - .await - .map_err(|e| { - WorkflowError::WebhookError(format!( - "SendMessage: failed to load workflow run {run_id}: {e}" - )) - })?; - let workflow = engine - .db - .get_workflow(community_id, wf_run.workflow_id) - .await - .map_err(|e| { - WorkflowError::WebhookError(format!( - "SendMessage: failed to load workflow {}: {e}", - wf_run.workflow_id - )) - })?; - let channel_id = resolve_send_message_channel( - channel.as_deref(), - &trigger_ctx.channel_id, - workflow.channel_id, - )?; - let owner_pubkey_hex = hex::encode(&workflow.owner_pubkey); - - info!( - run_id = %run_id, - step = step_id, - channel = %channel_id, - "SendMessage → {channel_id}: {text}" - ); - - let event_id = engine - .action_sink()? - .send_message(community_id, &channel_id, text, &owner_pubkey_hex) - .await - .map_err(WorkflowError::from)?; - - Ok(StepResult::Completed(serde_json::json!({ - "sent": true, - "event_id": event_id, - }))) + let release = serving_write.finish().await.map_err(|error| { + WorkflowError::WebhookError(format!("community write lease release failed: {error}")) + }); + match result { + Ok(value) => { + release?; + Ok(value) } - - SendDm { to, text: _ } => { - warn!(run_id = %run_id, step = step_id, "SendDm not yet implemented (to={to})"); - // TODO (WF-07): emit DM event. - Err(WorkflowError::NotImplemented("SendDm".into())) - } - - SetChannelTopic { topic: _ } => { - warn!(run_id = %run_id, step = step_id, "SetChannelTopic not yet implemented"); - // TODO (WF-07): update channel topic via DB. - Err(WorkflowError::NotImplemented("SetChannelTopic".into())) - } - - AddReaction { emoji } => { - info!(run_id = %run_id, step = step_id, "AddReaction → :{emoji}:"); - if trigger_ctx.message_id.is_empty() { - return Err(WorkflowError::InvalidDefinition( - "AddReaction: no trigger.message_id available".into(), - )); - } - - #[cfg(feature = "reqwest")] - { - let result = add_reaction_impl(&trigger_ctx.message_id, emoji).await?; - Ok(StepResult::Completed(result)) - } - - #[cfg(not(feature = "reqwest"))] - { - warn!( - run_id = %run_id, - step = step_id, - "AddReaction: reqwest feature not enabled, skipping HTTP call" - ); - Ok(StepResult::Completed( - serde_json::json!({ "added": false, "skipped": true }), - )) - } - } - - CallWebhook { - url, - method, - headers, - body, - } => { - let method_str = method.as_deref().unwrap_or("POST"); - info!(run_id = %run_id, step = step_id, "CallWebhook → {method_str} {url}"); - - #[cfg(feature = "reqwest")] - { - let result = call_webhook_impl(url, method_str, headers, body).await?; - Ok(StepResult::Completed(result)) - } - - #[cfg(not(feature = "reqwest"))] - { - // reqwest not enabled — log and return placeholder. - warn!( - run_id = %run_id, step = step_id, - "CallWebhook: reqwest feature not enabled, skipping HTTP call" - ); - let _ = (headers, body); // suppress unused warnings - Ok(StepResult::Completed(serde_json::json!({ - "status": 0, - "body": null, - "skipped": true - }))) - } - } - - RequestApproval { - from, - message, - timeout, - } => { - let timeout_str = timeout.as_deref().unwrap_or("24h"); - info!( - run_id = %run_id, step = step_id, - "RequestApproval from={from} timeout={timeout_str}: {message}" - ); - - let token = generate_approval_token(run_id, step_id); - - // TODO (WF-08): create approval record in DB, emit kind:46010. - // For now, return Suspended with the token so the caller can persist state. - - Ok(StepResult::Suspended { - approval_token: token, - }) - } - - Delay { duration } => { - let secs = parse_duration_secs(duration)?; - // Cap delay at 270 seconds (4.5 minutes) — must be less than default_timeout_secs (300s) - // to avoid non-deterministic StepTimeout. Long delays (hours/days) - // should use the scheduled resume pattern (future work: WF-09). - const MAX_DELAY_SECS: u64 = 270; - if secs > MAX_DELAY_SECS { - return Err(WorkflowError::InvalidDefinition(format!( - "delay exceeds maximum of {MAX_DELAY_SECS} seconds (got {secs}s); \ - use the scheduled resume pattern for long delays" - ))); - } - info!(run_id = %run_id, step = step_id, "Delay {duration} ({secs}s)"); - tokio::time::sleep(std::time::Duration::from_secs(secs)).await; - Ok(StepResult::Completed( - serde_json::json!({ "slept_secs": secs }), - )) + Err(error) => { + let _ = release; + Err(error) } } } diff --git a/deploy/charts/buzz/templates/deletion-worker-deployment.yaml b/deploy/charts/buzz/templates/deletion-worker-deployment.yaml index 2e32ac8ec..42f02fb6d 100644 --- a/deploy/charts/buzz/templates/deletion-worker-deployment.yaml +++ b/deploy/charts/buzz/templates/deletion-worker-deployment.yaml @@ -72,13 +72,11 @@ spec: secretKeyRef: name: {{ .Values.deletionWorker.existingSecret }} key: BUZZ_S3_ACCESS_KEY - optional: true - name: BUZZ_S3_SECRET_KEY valueFrom: secretKeyRef: name: {{ .Values.deletionWorker.existingSecret }} key: BUZZ_S3_SECRET_KEY - optional: true {{- with .Values.deletionWorker.extraEnv }} {{- toYaml . | nindent 12 }} {{- end }} diff --git a/deploy/charts/buzz/tests/deletion_worker_test.yaml b/deploy/charts/buzz/tests/deletion_worker_test.yaml index e76604954..53730b0a9 100644 --- a/deploy/charts/buzz/tests/deletion_worker_test.yaml +++ b/deploy/charts/buzz/tests/deletion_worker_test.yaml @@ -43,3 +43,7 @@ tests: - equal: path: spec.template.spec.containers[0].readinessProbe.httpGet.path value: /_readiness + - notExists: + path: spec.template.spec.containers[0].env[8].valueFrom.secretKeyRef.optional + - notExists: + path: spec.template.spec.containers[0].env[9].valueFrom.secretKeyRef.optional diff --git a/deploy/local/build-and-deploy.sh b/deploy/local/build-and-deploy.sh index 4dce0fdc3..2251ce10a 100755 --- a/deploy/local/build-and-deploy.sh +++ b/deploy/local/build-and-deploy.sh @@ -103,17 +103,14 @@ if [ "$helm_rc" != 0 ]; then fi # ── 4. verify 3/3 Ready ─────────────────────────────────────────────────────── -# Find the relay Deployment: everything under this release named "buzz" except -# the bundled "*-minio" Deployment. (The chart fullname collapses -# "-" to "" when the release name already contains the -# chart name, so the name isn't always "-buzz".) -DEPLOY="" -for d in $(kubectl -n "$NS" get deploy -l "app.kubernetes.io/instance=$RELEASE" \ - -o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}'); do - case "$d" in *-minio) continue;; esac - DEPLOY="$d"; break -done -[ -n "$DEPLOY" ] || die "could not locate the relay Deployment" +# Select the relay by its explicit component label. Optional worker/sidecar +# Deployments share the release instance and must never become the rollout target. +DEPLOYMENTS=$(kubectl -n "$NS" get deploy \ + -l "app.kubernetes.io/instance=$RELEASE,app.kubernetes.io/component=relay" \ + -o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}') +[ "$(printf '%s\n' "$DEPLOYMENTS" | sed '/^$/d' | wc -l | tr -d ' ')" = "1" ] || \ + die "expected exactly one relay Deployment, got: ${DEPLOYMENTS:-}" +DEPLOY=$(printf '%s\n' "$DEPLOYMENTS" | sed '/^$/d') log "waiting for $REPLICAS relay pods Ready (deployment: $DEPLOY)" kubectl -n "$NS" rollout status deployment/"$DEPLOY" --timeout=4m | tee "$EVID/rollout.txt" kubectl -n "$NS" get pods -o wide | tee "$EVID/pods.txt" @@ -129,9 +126,8 @@ log "probing /_readiness on each relay pod individually" : > "$EVID/readiness.txt" FAIL=0 RELAY_PODS=$(kubectl -n "$NS" get pods \ - -l "app.kubernetes.io/name=buzz,app.kubernetes.io/instance=$RELEASE" \ - -o jsonpath='{range .items[*]}{.metadata.name}{" "}{.metadata.labels.app\.kubernetes\.io/component}{"\n"}{end}' \ - | awk '$2 != "minio" && $2 != "minio-init" {print $1}') + -l "app.kubernetes.io/instance=$RELEASE,app.kubernetes.io/component=relay" \ + -o jsonpath='{range .items[*]}{.metadata.name}{"\n"}{end}') for pod in $RELAY_PODS; do body=$(kubectl -n "$NS" exec "$pod" -- \ sh -c 'curl -sS --max-time 5 http://127.0.0.1:8080/_readiness' 2>/dev/null || echo '') diff --git a/migrations/0027_community_deletion.sql b/migrations/0027_community_deletion.sql index c2e241f25..1f15a60f1 100644 --- a/migrations/0027_community_deletion.sql +++ b/migrations/0027_community_deletion.sql @@ -23,6 +23,8 @@ CREATE TABLE community_deletion_requests ( reason TEXT, schema_manifest JSONB, storage_manifest JSONB, + destructive_storage_manifest JSONB, + destructive_storage_frozen_at TIMESTAMPTZ, inventory_manifest JSONB, inventory_digest BYTEA CHECK (inventory_digest IS NULL OR length(inventory_digest) = 32), inventory_frozen_at TIMESTAMPTZ, @@ -88,6 +90,22 @@ CREATE TABLE community_deletion_retention_exceptions ( PRIMARY KEY (request_id, exception_key) ); +CREATE TABLE community_serving_write_leases ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + community_id UUID NOT NULL REFERENCES communities(id), + operation TEXT NOT NULL, + owner TEXT NOT NULL, + generation BIGINT NOT NULL DEFAULT 1 CHECK (generation > 0), + -- Community fence generation observed when this lease was acquired. + fence_generation BIGINT NOT NULL CHECK (fence_generation >= 0), + lease_until TIMESTAMPTZ NOT NULL, + heartbeat_at TIMESTAMPTZ NOT NULL DEFAULT now(), + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (community_id, id) +); +CREATE INDEX community_serving_write_leases_active + ON community_serving_write_leases (community_id, lease_until); + CREATE TABLE community_deletion_executor_heartbeats ( executor_id TEXT PRIMARY KEY, mode TEXT NOT NULL CHECK (mode IN ('run', 'drain', 'worker')), @@ -103,6 +121,7 @@ INSERT INTO _operator_global_tables (table_name, reason) VALUES ('community_deletion_approvals', 'deployment operator destructive approvals'), ('community_deletion_checkpoints', 'deployment deletion executor checkpoints and failures'), ('community_deletion_retention_exceptions', 'deployment retention holds and expiry exceptions'), + ('community_serving_write_leases', 'deployment serving side-effect leases drained by deletion'), ('community_deletion_executor_heartbeats', 'deployment deletion worker liveness'); -- Shared lock key used by both the trigger and the deletion engine. Every @@ -171,8 +190,11 @@ DECLARE expected_generation BIGINT; BEGIN IF TG_OP = 'DELETE' THEN - RAISE EXCEPTION 'community tombstones are permanent' - USING ERRCODE = 'object_not_in_prerequisite_state'; + IF OLD.deletion_state <> 'active' OR OLD.deleted_at IS NOT NULL THEN + RAISE EXCEPTION 'community tombstones are permanent' + USING ERRCODE = 'object_not_in_prerequisite_state'; + END IF; + RETURN OLD; END IF; expected_generation := CASE @@ -225,6 +247,7 @@ BEGIN 'community_deletion_approvals', 'community_deletion_checkpoints', 'community_deletion_retention_exceptions', + 'community_serving_write_leases', 'community_deletion_executor_heartbeats' ) LOOP diff --git a/schema/schema.sql b/schema/schema.sql index e789be4e6..2d74fe6c0 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1092,6 +1092,8 @@ CREATE TABLE community_deletion_requests ( reason TEXT, schema_manifest JSONB, storage_manifest JSONB, + destructive_storage_manifest JSONB, + destructive_storage_frozen_at TIMESTAMPTZ, inventory_manifest JSONB, inventory_digest BYTEA CHECK (inventory_digest IS NULL OR length(inventory_digest) = 32), inventory_frozen_at TIMESTAMPTZ, @@ -1153,6 +1155,22 @@ CREATE TABLE community_deletion_retention_exceptions ( created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (request_id, exception_key) ); +CREATE TABLE community_serving_write_leases ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + community_id UUID NOT NULL REFERENCES communities(id), + operation TEXT NOT NULL, + owner TEXT NOT NULL, + generation BIGINT NOT NULL DEFAULT 1 CHECK (generation > 0), + -- Community fence generation observed when this lease was acquired. + fence_generation BIGINT NOT NULL CHECK (fence_generation >= 0), + lease_until TIMESTAMPTZ NOT NULL, + heartbeat_at TIMESTAMPTZ NOT NULL DEFAULT now(), + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (community_id, id) +); +CREATE INDEX community_serving_write_leases_active + ON community_serving_write_leases (community_id, lease_until); + CREATE TABLE community_deletion_executor_heartbeats ( executor_id TEXT PRIMARY KEY, mode TEXT NOT NULL CHECK (mode IN ('run', 'drain', 'worker')), @@ -1167,6 +1185,7 @@ INSERT INTO _operator_global_tables (table_name, reason) VALUES ('community_deletion_approvals', 'deployment operator destructive approvals'), ('community_deletion_checkpoints', 'deployment deletion executor checkpoints and failures'), ('community_deletion_retention_exceptions', 'deployment retention holds and expiry exceptions'), + ('community_serving_write_leases', 'deployment serving side-effect leases drained by deletion'), ('community_deletion_executor_heartbeats', 'deployment deletion worker liveness'); CREATE FUNCTION community_deletion_lock_key(target UUID) RETURNS BIGINT @@ -1212,8 +1231,11 @@ DECLARE expected_generation BIGINT; BEGIN IF TG_OP = 'DELETE' THEN - RAISE EXCEPTION 'community tombstones are permanent' - USING ERRCODE = 'object_not_in_prerequisite_state'; + IF OLD.deletion_state <> 'active' OR OLD.deleted_at IS NOT NULL THEN + RAISE EXCEPTION 'community tombstones are permanent' + USING ERRCODE = 'object_not_in_prerequisite_state'; + END IF; + RETURN OLD; END IF; expected_generation := CASE WHEN NEW.deletion_fence_generation > OLD.deletion_fence_generation THEN NEW.deletion_fence_generation ELSE OLD.deletion_fence_generation END; @@ -1242,7 +1264,7 @@ BEGIN 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_deletion_executor_heartbeats') + 'community_serving_write_leases', 'community_deletion_executor_heartbeats') LOOP EXECUTE format('CREATE TRIGGER %I BEFORE INSERT OR UPDATE OR DELETE ON %I ' 'FOR EACH ROW EXECUTE FUNCTION enforce_community_write_fence()',