fix: close community deletion serving-write races

Order fencing behind admitted serving-write leases, keep external write
heartbeats live through durable finalization, and make Git publication hold
the community lock through the relay event commit. Harden tenant trigger
coverage for cross-community updates and future tables, then fail startup
and readiness closed on live catalog drift.

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