mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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>
This commit is contained in:
parent
1cfa319d04
commit
e00dee77e2
@@ -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 \
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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<Uuid>,
|
||||
) -> 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,
|
||||
|
||||
@@ -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<Arc<MediaStorage>> {
|
||||
fn deletion_test_media_storage() -> Arc<MediaStorage> {
|
||||
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();
|
||||
|
||||
@@ -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)) => {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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';
|
||||
|
||||
@@ -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';
|
||||
|
||||
Reference in New Issue
Block a user