From e00dee77e23c4713a8e46d72b82daa5a94b985fd Mon Sep 17 00:00:00 2001 From: npub1dccv64krpcpse5cmkzfeh998cftungyatw3djt8jwdw6g43f7fyqzzmrf7 <6e30cd56c30e030cd31bb0939b94a7c257c9a09d5ba2d92cf2735da45629f248@buzz.block.builderlab.xyz> Date: Sun, 2 Aug 2026 13:55:16 -0700 Subject: [PATCH] fix: preserve admitted git writes through quiescing Co-authored-by: npub1dccv64krpcpse5cmkzfeh998cftungyatw3djt8jwdw6g43f7fyqzzmrf7 <6e30cd56c30e030cd31bb0939b94a7c257c9a09d5ba2d92cf2735da45629f248@buzz.block.builderlab.xyz> Signed-off-by: npub1dccv64krpcpse5cmkzfeh998cftungyatw3djt8jwdw6g43f7fyqzzmrf7 <6e30cd56c30e030cd31bb0939b94a7c257c9a09d5ba2d92cf2735da45629f248@buzz.block.builderlab.xyz> --- .github/workflows/ci.yml | 15 +++++ crates/buzz-db/src/deletion.rs | 74 +++++++++++++++++++++- crates/buzz-db/src/lib.rs | 14 ++-- crates/buzz-deletion/src/lib.rs | 63 +++++++++--------- crates/buzz-relay/src/api/git/transport.rs | 4 +- deploy/charts/buzz/README.md | 7 +- migrations/0027_community_deletion.sql | 33 ++++++++++ schema/schema.sql | 33 ++++++++++ 8 files changed, 195 insertions(+), 48 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 60507182d..0a8b252ee 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -353,6 +353,7 @@ jobs: cargo nextest archive \ --cargo-profile ci \ -p buzz-db \ + -p buzz-deletion \ -p buzz-relay \ -p buzz-test-client \ --lib \ @@ -684,6 +685,20 @@ jobs: done cat /tmp/buzz-relay.log exit 1 + - name: Community deletion production-path regressions + run: | + cargo nextest run \ + --archive-file target/ci/backend-integration-tests.tar.zst \ + -E '(package(buzz-deletion) and test(/tests::/)) or (package(buzz-relay) and test(/finalize_push_(holds_serving_lease_through_post_cas_publication|db_failure_after_cas_is_not_success_and_releases_lease)/))' \ + --run-ignored ignored-only + env: + DATABASE_URL: postgres://buzz:${{ env.BUZZ_TEST_POSTGRES_PASSWORD }}@localhost:5432/buzz + BUZZ_TEST_S3_ENDPOINT: http://localhost:9000 + BUZZ_TEST_S3_ACCESS_KEY: buzz_dev + BUZZ_TEST_S3_SECRET_KEY: buzz_dev_secret + BUZZ_TEST_S3_BUCKET: buzz-media + BUZZ_TEST_S3_REGION: us-east-1 + - name: Invite security tests run: | cargo nextest run \ diff --git a/crates/buzz-db/src/deletion.rs b/crates/buzz-db/src/deletion.rs index c9df6d147..72a52ac2b 100644 --- a/crates/buzz-db/src/deletion.rs +++ b/crates/buzz-db/src/deletion.rs @@ -1577,6 +1577,60 @@ impl DeletionStore { } } + /// Take the shared community deletion lock inside an existing transaction + /// and authorize a final mutation under an already-admitted serving lease. + /// + /// The lease is checked in the same transaction as the mutation. During + /// quiescing, only this exact unexpired lease and fence generation may + /// finish; active communities continue to accept the admitted write too. + pub async fn guard_transaction_with_serving_lease( + &self, + tx: &mut Transaction<'_, Postgres>, + lease: &ServingWriteLease, + ) -> Result<()> { + sqlx::query("SELECT pg_advisory_xact_lock_shared(community_deletion_lock_key($1))") + .bind(lease.community_id.as_uuid()) + .execute(&mut **tx) + .await?; + let valid: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_serving_write_leases lease \ + JOIN communities community ON community.id = lease.community_id \ + 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 IN ('active', 'quiescing') \ + AND community.deletion_fence_generation = lease.fence_generation)", + ) + .bind(lease.id) + .bind(lease.community_id.as_uuid()) + .bind(&lease.owner) + .bind(lease.generation) + .bind(lease.fence_generation) + .fetch_one(&mut **tx) + .await?; + if !valid { + return Err(DbError::AccessDenied(format!( + "stale serving write lease {}", + lease.id + ))); + } + sqlx::query( + "SELECT set_config('buzz.serving_write_community', $1, true), \ + set_config('buzz.serving_write_lease_id', $2, true), \ + set_config('buzz.serving_write_owner', $3, true), \ + set_config('buzz.serving_write_generation', $4, true), \ + set_config('buzz.serving_write_fence_generation', $5, true)", + ) + .bind(lease.community_id.to_string()) + .bind(lease.id.to_string()) + .bind(&lease.owner) + .bind(lease.generation.to_string()) + .bind(lease.fence_generation.to_string()) + .execute(&mut **tx) + .await?; + Ok(()) + } + /// Acquire a durable, expiring lease for an external serving side effect. /// /// The short transaction shares the same advisory lock as the destructive @@ -1929,13 +1983,14 @@ async fn verify_lease( ) -> Result<()> { let valid: bool = sqlx::query_scalar( "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ - WHERE id = $1 AND stage = $2 AND lease_owner = $3 \ + WHERE id = $1 AND community_id = $5 AND stage = $2 AND lease_owner = $3 \ AND lease_generation = $4 AND lease_until >= now() AND blocked_at IS NULL)", ) .bind(token.request_id) .bind(stage.to_string()) .bind(&token.owner) .bind(token.generation) + .bind(token.community_id.as_uuid()) .fetch_one(&mut **tx) .await?; if valid { @@ -1954,7 +2009,8 @@ async fn verify_lease_and_fence( let valid: bool = sqlx::query_scalar( "SELECT EXISTS(SELECT 1 FROM community_deletion_requests request \ JOIN communities community ON community.id = request.community_id \ - WHERE request.id = $1 AND request.stage = $2 AND request.lease_owner = $3 \ + WHERE request.id = $1 AND request.community_id = $6 \ + AND request.stage = $2 AND request.lease_owner = $3 \ AND request.lease_generation = $4 AND request.lease_until >= now() \ AND request.blocked_at IS NULL AND request.fence_generation = $5 \ AND community.deletion_state IN ('fenced', 'tombstone') \ @@ -1965,6 +2021,7 @@ async fn verify_lease_and_fence( .bind(&token.owner) .bind(token.generation) .bind(fence_generation) + .bind(token.community_id.as_uuid()) .fetch_one(&mut **tx) .await?; if valid { @@ -2419,6 +2476,19 @@ mod postgres_tests { store.fence(&stale).await.is_err(), "stale lease must reject" ); + let mut wrong_community = claim.lease.clone(); + wrong_community.community_id = db + .ensure_configured_community(&format!( + "wrong-lease-community-{}.example", + Uuid::new_v4().simple() + )) + .await + .expect("create unrelated community") + .id; + assert!( + store.begin_quiescing(&wrong_community).await.is_err(), + "a lease token must remain bound to its durable request community" + ); store.begin_quiescing(&claim.lease).await.expect("quiesce"); let generation = store.fence(&claim.lease).await.expect("fence"); diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 6b9a5afd9..92a2ac942 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -1664,19 +1664,19 @@ impl Db { Ok(result) } - /// Insert an event while holding the admitted serving-write lease's - /// community ordering lock through commit. + /// Insert an event while holding and validating an admitted serving-write + /// lease under the community ordering lock through commit. /// /// External side effects use a durable lease rather than one long-lived DB - /// transaction. Their final database mutation must still share the - /// acquire/fence advisory-lock order so a fence cannot commit between the - /// caller's lease verification and this insert. + /// transaction. Their final database mutation presents that exact lease so + /// it may finish during quiescing without admitting any new serving work. pub async fn insert_event_with_serving_write_guard( &self, - community_id: CommunityId, + lease: &deletion::ServingWriteLease, event: &nostr::Event, channel_id: Option, ) -> Result<(StoredEvent, bool)> { + let community_id = lease.community_id; let kind_u16 = event.kind.as_u16(); let kind_u32 = u32::from(kind_u16); if kind_u32 == buzz_core::kind::KIND_AUTH { @@ -1688,7 +1688,7 @@ impl Db { let mut tx = self.pool.begin().await?; self.deletion_store() - .guard_transaction(&mut tx, community_id) + .guard_transaction_with_serving_lease(&mut tx, lease) .await?; let result = event::insert_event_with_thread_metadata_tx( &mut tx, diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index 1aab912d0..3281deabf 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -125,6 +125,11 @@ impl ServingWriteGuard { self.lost.clone() } + /// The durable lease token presented to a final database mutation. + pub fn lease(&self) -> &buzz_db::deletion::ServingWriteLease { + &self.lease + } + /// Release the lease after the side effect completes. pub async fn finish(mut self) -> Result<()> { self.cancel.cancel(); @@ -1379,10 +1384,10 @@ fn print_json(value: &impl Serialize) -> Result<()> { mod tests { use super::*; - async fn claimed_test_deletion(prefix: &str) -> Option<(Services, ClaimedDeletion)> { + async fn claimed_test_deletion(prefix: &str) -> (Services, ClaimedDeletion) { let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") .or_else(|_| std::env::var("DATABASE_URL")) - .ok()?; + .expect("BUZZ_TEST_DATABASE_URL or DATABASE_URL is required"); let pool = sqlx::PgPool::connect(&database_url) .await .expect("connect deletion engine test DB"); @@ -1453,23 +1458,23 @@ mod tests { .create_pool(Some(deadpool_redis::Runtime::Tokio1)) .expect("construct unused Redis pool"), }; - Some((services, claim)) + (services, claim) } - fn deletion_test_media_storage() -> Option> { + fn deletion_test_media_storage() -> Arc { let endpoint = std::env::var("BUZZ_TEST_S3_ENDPOINT") .or_else(|_| std::env::var("BUZZ_S3_ENDPOINT")) - .ok()?; + .expect("BUZZ_TEST_S3_ENDPOINT or BUZZ_S3_ENDPOINT is required"); let access_key = std::env::var("BUZZ_TEST_S3_ACCESS_KEY") .or_else(|_| std::env::var("BUZZ_S3_ACCESS_KEY")) - .ok()?; + .expect("BUZZ_TEST_S3_ACCESS_KEY or BUZZ_S3_ACCESS_KEY is required"); let secret_key = std::env::var("BUZZ_TEST_S3_SECRET_KEY") .or_else(|_| std::env::var("BUZZ_S3_SECRET_KEY")) - .ok()?; + .expect("BUZZ_TEST_S3_SECRET_KEY or BUZZ_S3_SECRET_KEY is required"); let bucket = std::env::var("BUZZ_TEST_S3_BUCKET") .or_else(|_| std::env::var("BUZZ_S3_BUCKET")) - .ok()?; - Some(Arc::new( + .expect("BUZZ_TEST_S3_BUCKET or BUZZ_S3_BUCKET is required"); + Arc::new( MediaStorage::new(&buzz_media::MediaConfig { s3_endpoint: endpoint, s3_access_key: access_key, @@ -1489,7 +1494,7 @@ mod tests { upload_port_header: None, }) .expect("construct deletion test media service"), - )) + ) } #[test] @@ -1500,14 +1505,10 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres and S3-compatible storage"] async fn drained_stage_reconciles_first_object_deleted_before_checkpoint() { - let Some((mut services, claim)) = claimed_test_deletion("deletion-first-missing").await - else { - return; - }; - let Some(media) = deletion_test_media_storage() else { - return; - }; + let (mut services, claim) = claimed_test_deletion("deletion-first-missing").await; + let media = deletion_test_media_storage(); services.media = media; let object_key = format!( @@ -1667,10 +1668,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres"] async fn final_storage_verification_rejects_late_target_binding() { - let Some((services, claim)) = claimed_test_deletion("deletion-late-binding").await else { - return; - }; + let (services, claim) = claimed_test_deletion("deletion-late-binding").await; let community = claim.request.community_id; let late_key = format!("_meta/{community}/{}.json", "a".repeat(64)); let error = verify_storage_absence_from_objects( @@ -1689,10 +1689,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres"] async fn stale_lease_during_failure_recording_is_lost_ownership() { - let Some((services, claim)) = claimed_test_deletion("deletion-stale-record").await else { - return; - }; + let (services, claim) = claimed_test_deletion("deletion-stale-record").await; let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") .or_else(|_| std::env::var("DATABASE_URL")) .expect("test database URL"); @@ -1718,13 +1717,11 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres"] async fn serving_guard_cancels_protected_operation_when_heartbeat_is_lost() { - let database_url = match std::env::var("BUZZ_TEST_DATABASE_URL") + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") .or_else(|_| std::env::var("DATABASE_URL")) - { - Ok(url) => url, - Err(_) => return, - }; + .expect("BUZZ_TEST_DATABASE_URL or DATABASE_URL is required"); let pool = sqlx::PgPool::connect(&database_url) .await .expect("connect serving guard test DB"); @@ -1764,10 +1761,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres"] async fn guarded_external_step_rejects_preexisting_heartbeat_loss_without_polling_operation() { - let Some((services, claim)) = claimed_test_deletion("deletion-heartbeat").await else { - return; - }; + let (services, claim) = claimed_test_deletion("deletion-heartbeat").await; let heartbeat_lost = CancellationToken::new(); heartbeat_lost.cancel(); let polled = Arc::new(AtomicBool::new(false)); @@ -1798,10 +1794,9 @@ mod tests { } #[tokio::test] + #[ignore = "requires Postgres"] async fn shutdown_during_stage_releases_claim_without_recording_retry() { - let Some((services, claim)) = claimed_test_deletion("deletion-shutdown").await else { - return; - }; + let (services, claim) = claimed_test_deletion("deletion-shutdown").await; let request_id = claim.request.id; let retry_count = claim.request.retry_count; let shutdown = CancellationToken::new(); diff --git a/crates/buzz-relay/src/api/git/transport.rs b/crates/buzz-relay/src/api/git/transport.rs index fa7791d1b..a9cc5bb43 100644 --- a/crates/buzz-relay/src/api/git/transport.rs +++ b/crates/buzz-relay/src/api/git/transport.rs @@ -1933,13 +1933,13 @@ async fn finalize_push_inner( } else { state .db - .insert_event_with_serving_write_guard(ctx.tenant.community(), &event, None) + .insert_event_with_serving_write_guard(serving_write.lease(), &event, None) .await }; #[cfg(not(test))] let insert_result = state .db - .insert_event_with_serving_write_guard(ctx.tenant.community(), &event, None) + .insert_event_with_serving_write_guard(serving_write.lease(), &event, None) .await; match insert_result { Ok((stored, true)) => { diff --git a/deploy/charts/buzz/README.md b/deploy/charts/buzz/README.md index 2ecc5d444..9ec31282f 100644 --- a/deploy/charts/buzz/README.md +++ b/deploy/charts/buzz/README.md @@ -201,9 +201,10 @@ Track application p99 for events/media/Git/workflows alongside the lease gauges; investigate rising dead tuples or expired rows before raising worker scale. Object removal is checkpointed per key and each attempt deletes at most -`deletionWorker.storageDeleteBatchSize` (default 100). After recorded deletion -progress, missing keys resume as already-complete; a missing key before any -checkpoint still blocks as unexplained drift. Changed/versioned keys fail closed. +`deletionWorker.storageDeleteBatchSize` (default 100). Deletion is resumable +from durable checkpoints. Missing keys are recorded as `already_missing` and +processing continues, including the crash window where a DELETE committed +before its first PostgreSQL checkpoint. Changed/versioned keys fail closed. Final verification performs a fresh complete target inventory before advancing. ## Device pairing relay diff --git a/migrations/0027_community_deletion.sql b/migrations/0027_community_deletion.sql index 72e99d1a2..1b074ecf4 100644 --- a/migrations/0027_community_deletion.sql +++ b/migrations/0027_community_deletion.sql @@ -161,6 +161,12 @@ DECLARE generation BIGINT; executor_community TEXT; executor_generation TEXT; + serving_community TEXT; + serving_lease_id TEXT; + serving_owner TEXT; + serving_generation TEXT; + serving_fence_generation TEXT; + serving_lease_valid BOOLEAN := false; BEGIN -- Nullable operator-attribution rows without a tenant are unrelated. IF target IS NULL THEN @@ -186,6 +192,33 @@ BEGIN RETURN; END IF; + -- A serving mutation admitted before quiescing may finish only while its + -- exact durable lease remains current and bound to this fence generation. + serving_community := current_setting('buzz.serving_write_community', true); + serving_lease_id := current_setting('buzz.serving_write_lease_id', true); + serving_owner := current_setting('buzz.serving_write_owner', true); + serving_generation := current_setting('buzz.serving_write_generation', true); + serving_fence_generation := current_setting('buzz.serving_write_fence_generation', true); + IF lifecycle IN ('active', 'quiescing') + AND serving_community = target::TEXT + AND serving_lease_id ~ '^[0-9a-fA-F-]{36}$' + AND serving_generation ~ '^[0-9]+$' + AND serving_fence_generation ~ '^[0-9]+$' + AND serving_fence_generation::BIGINT = generation THEN + SELECT EXISTS( + SELECT 1 FROM community_serving_write_leases lease + WHERE lease.id = serving_lease_id::UUID + AND lease.community_id = target + AND lease.owner = serving_owner + AND lease.generation = serving_generation::BIGINT + AND lease.fence_generation = serving_fence_generation::BIGINT + AND lease.lease_until >= now() + ) INTO serving_lease_valid; + IF serving_lease_valid THEN + RETURN; + END IF; + END IF; + IF lifecycle <> 'active' THEN RAISE EXCEPTION 'community write fenced: community % generation %', target, generation USING ERRCODE = 'object_not_in_prerequisite_state'; diff --git a/schema/schema.sql b/schema/schema.sql index 184e297e1..ee7c9ee95 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1214,6 +1214,12 @@ DECLARE generation BIGINT; executor_community TEXT; executor_generation TEXT; + serving_community TEXT; + serving_lease_id TEXT; + serving_owner TEXT; + serving_generation TEXT; + serving_fence_generation TEXT; + serving_lease_valid BOOLEAN := false; BEGIN -- Nullable operator-attribution rows without a tenant are unrelated. IF target IS NULL THEN @@ -1239,6 +1245,33 @@ BEGIN RETURN; END IF; + -- A serving mutation admitted before quiescing may finish only while its + -- exact durable lease remains current and bound to this fence generation. + serving_community := current_setting('buzz.serving_write_community', true); + serving_lease_id := current_setting('buzz.serving_write_lease_id', true); + serving_owner := current_setting('buzz.serving_write_owner', true); + serving_generation := current_setting('buzz.serving_write_generation', true); + serving_fence_generation := current_setting('buzz.serving_write_fence_generation', true); + IF lifecycle IN ('active', 'quiescing') + AND serving_community = target::TEXT + AND serving_lease_id ~ '^[0-9a-fA-F-]{36}$' + AND serving_generation ~ '^[0-9]+$' + AND serving_fence_generation ~ '^[0-9]+$' + AND serving_fence_generation::BIGINT = generation THEN + SELECT EXISTS( + SELECT 1 FROM community_serving_write_leases lease + WHERE lease.id = serving_lease_id::UUID + AND lease.community_id = target + AND lease.owner = serving_owner + AND lease.generation = serving_generation::BIGINT + AND lease.fence_generation = serving_fence_generation::BIGINT + AND lease.lease_until >= now() + ) INTO serving_lease_valid; + IF serving_lease_valid THEN + RETURN; + END IF; + END IF; + IF lifecycle <> 'active' THEN RAISE EXCEPTION 'community write fenced: community % generation %', target, generation USING ERRCODE = 'object_not_in_prerequisite_state';