diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index d30b9c9a2f5..c11db7c7d32 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -18,8 +18,9 @@ Deployment-root community management uses operator-signed NIP-98 HTTP requests. `POST /operator/communities/delete` accepts only an exact normalized, archived community whose asserted pubkey is still its owner. The caller supplies the request UUID as the stable correlation/idempotency identity; the durable row -records owner intent, mediating operator, and acknowledgement version. Admission -returns `202` at the `submitted` stage and performs no inventory, approval, +records operator-attested owner intent, mediating operator, and acknowledgement +version. Admission returns `202` at the `submitted` stage and performs no +inventory, approval, quiescing, object-store access, or deletion execution synchronously. While that non-aborted request exists, unarchive and ownership transfer conflict and owner management lists suppress the archived row. Replaying the same UUID converges @@ -45,6 +46,19 @@ on the row. Owner provenance pins `requested_by` to `owner_pubkey`, so passing the operator's own pubkey does not converge — it conflicts with the existing one-active-request invariant instead. +The privileged one-shot `buzz-admin deletions drain` process gives already- +approved work priority. When none is ready, it may claim only an +operator-attested owner-origin `submitted` request under the same durable +generation lease used for execution, inventory it with lease-loss cancellation, +and atomically freeze the inventory plus a digest-bound `owner_automatic` +approval. The mediating operator remains the approval actor; the owner +acknowledgement is pre-inventory +intent, not a claim that the owner reviewed the digest. The retained lease then +enters the unchanged approved-request executor. Operator-origin requests never +auto-progress and still require explicit inventory and approval. Owner +admission still has no owner-facing cancellation or grace period; the +privileged abort described above remains the recovery path. + Ownership is mutable only while a community is active. Archiving freezes the current owner. Normal transfer and deployment-root legacy convergence take the same community-row lock as owner-deletion admission, then reject archived, @@ -762,6 +776,10 @@ Postgres/Redis clients and S3 client; it does not call relay HTTP. Durable requests, leases, retry timing, and checkpoints in Postgres are the handoff and execution authority, so Kubernetes uses `Forbid` concurrency and zero Job retries rather than introducing a second retry system. +The same drain first claims runnable approved work and, only when none exists, +may prepare one owner-origin submission. Inventory, automatic approval, and +execution share one generation lease and the existing retry/block/checkpoint +records; there is no preparation worker, command, queue, or retry authority. --- diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 9a003339da0..8d71a48daa0 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -705,9 +705,12 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 52); + assert_eq!(migrations.len(), 53); assert_eq!(migrations[48].version, 49); + assert_eq!(migrations[49].version, 50); assert_eq!(migrations[50].version, 51); + assert_eq!(migrations[51].version, 52); + assert_eq!(migrations[52].version, 53); assert!(migrations[48] .sql .as_str() @@ -716,6 +719,11 @@ mod postgres_tests { .sql .as_str() .contains("community_deletion_owner_provenance")); + assert!(migrations[52].sql.as_str().contains("approval_origin")); + assert!(migrations[52] + .sql + .as_str() + .contains("community_deletion_requests_owner_preparable")); assert_eq!(migrations[0].version, 1); assert_eq!(&*migrations[0].description, "initial schema"); assert!(migrations[0] @@ -1722,6 +1730,34 @@ mod postgres_tests { ); } + #[test] + fn owner_deletion_auto_approval_migration_matches_desired_schema() { + let migration = MIGRATOR + .iter() + .find(|migration| migration.version == 53) + .expect("embedded migration 0053") + .sql + .as_ref() + .to_ascii_lowercase(); + let workspace_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) + .parent() + .and_then(std::path::Path::parent) + .expect("workspace root"); + let schema = std::fs::read_to_string(workspace_root.join("schema/schema.sql")) + .expect("read schema/schema.sql") + .to_ascii_lowercase(); + + for sql in [&migration, &schema] { + assert!(sql.contains("approval_origin text not null default 'operator'")); + assert!(sql.contains("approval_origin in ('operator', 'owner_automatic')")); + assert!(sql.contains("'submitted', 'approved', 'fenced'")); + assert!(sql.contains("community_deletion_requests_owner_preparable")); + assert!(sql.contains("request_origin = 'owner'")); + assert!(sql.contains("stage = 'submitted'")); + } + assert!(migration.contains("set local lock_timeout = '5s'")); + } + /// Structural parity between migration 0029's deletion surface and the /// desired-state bootstrap schema (`schema/schema.sql`). /// @@ -1866,13 +1902,50 @@ mod postgres_tests { .tables .get(table) .unwrap_or_else(|| panic!("schema.sql is missing deletion table {table}")); - if table != "community_deletion_requests" { + if table != "community_deletion_requests" && table != "community_deletion_approvals" { assert_eq!( in_schema, definition, "schema.sql definition of {table} drifted from migration 0029" ); } } + let migration_approval_table = migration + .tables + .get("community_deletion_approvals") + .expect("0029 approval table"); + let schema_approval_table = schema + .tables + .get("community_deletion_approvals") + .expect("schema.sql approval table"); + for invariant in [ + "inventory_digest bytea not null check (length(inventory_digest) = 32)", + "foreign key (request_id, community_id, inventory_digest) references community_deletion_requests(id, community_id, inventory_digest) on delete restrict", + ] { + assert!( + migration_approval_table.contains(invariant), + "0029 deletion approvals are missing {invariant}" + ); + assert!( + schema_approval_table.contains(invariant), + "schema.sql deletion approvals are missing {invariant}" + ); + } + let migration_request_table = migration + .tables + .get("community_deletion_requests") + .expect("0029 deletion request table"); + for request_table in [ + migration_request_table, + schema + .tables + .get("community_deletion_requests") + .expect("schema.sql deletion request table"), + ] { + assert!( + request_table.contains("unique (id, community_id, inventory_digest)"), + "deletion requests must expose the exact composite approval target" + ); + } for (function, definition) in &migration.functions { let in_schema = schema .functions diff --git a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs index 82b7d2fb0a6..c6d38474e1d 100644 --- a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs +++ b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs @@ -311,7 +311,7 @@ async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer() production_result.is_ok(), "production migrator must preserve ingestion progress: {production_result:?}" ); - assert_eq!(version, 52); + assert_eq!(version, 53); assert_eq!(final_oid, oid, "prebuild must not be replaced"); assert_eq!(count, 4, "all writer witnesses must persist"); } diff --git a/crates/buzz-db/src/store/deletion.rs b/crates/buzz-db/src/store/deletion.rs index 23b6d43049c..09c17ee1656 100644 --- a/crates/buzz-db/src/store/deletion.rs +++ b/crates/buzz-db/src/store/deletion.rs @@ -241,11 +241,11 @@ pub struct DeletionRequest { pub retry_stage: Option, /// Legacy display identity that submitted the request. pub requested_by: String, - /// Whether the request originated from an operator or authenticated owner intent. + /// Whether the request originated from an operator or operator-attested owner intent. pub request_origin: DeletionRequestOrigin, - /// Current owner identity authenticated at owner-request admission. + /// Current owner identity asserted by the mediating operator at admission. pub owner_pubkey: Option, - /// Deployment operator that mediated the authenticated owner intent. + /// Deployment operator that attested to the owner intent. pub mediating_operator_pubkey: Option, /// Owner-facing destructive-action acknowledgement contract version. pub acknowledgement_version: Option, @@ -303,7 +303,7 @@ pub struct DeletionRequest { pub enum DeletionRequestOrigin { /// Request was submitted directly by a deployment operator. Operator, - /// Request records authenticated owner intent mediated by an operator. + /// Request records operator-attested owner intent. Owner, } @@ -321,7 +321,7 @@ impl FromStr for DeletionRequestOrigin { } } -/// Result of atomically admitting authenticated owner deletion intent. +/// Result of atomically admitting operator-attested owner deletion intent. #[derive(Debug, Clone, PartialEq)] pub enum OwnerDeletionAdmission { /// A new request was created or the stable request UUID converged to its row. @@ -610,12 +610,38 @@ pub struct DeletionApproval { pub inventory_digest: String, /// Approving operator identity. pub approved_by: String, + /// Bounded provenance for the approval decision. + pub approval_origin: DeletionApprovalOrigin, /// Optional approval note. pub note: Option, /// Approval timestamp. pub approved_at: DateTime, } +/// Durable provenance class for an inventory approval. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum DeletionApprovalOrigin { + /// A deployment operator explicitly approved the inventory. + Operator, + /// Privileged policy approved operator-attested owner intent automatically. + OwnerAutomatic, +} + +impl FromStr for DeletionApprovalOrigin { + type Err = DbError; + + fn from_str(value: &str) -> std::result::Result { + match value { + "operator" => Ok(Self::Operator), + "owner_automatic" => Ok(Self::OwnerAutomatic), + other => Err(DbError::InvalidData(format!( + "unknown community deletion approval origin: {other}" + ))), + } + } +} + /// Monotonic lease token required by every execution mutation. #[derive(Debug, Clone, PartialEq, Eq)] pub struct LeaseToken { @@ -959,7 +985,7 @@ impl DeletionStore { pub async fn inspect(&self, request_id: Uuid) -> Result { let request = self.get(request_id).await?; let approval_row = sqlx::query( - "SELECT inventory_digest, approved_by, note, approved_at \ + "SELECT inventory_digest, approved_by, approval_origin, note, approved_at \ FROM community_deletion_approvals WHERE request_id = $1", ) .bind(request_id) @@ -970,6 +996,7 @@ impl DeletionStore { Ok::(DeletionApproval { inventory_digest: hex::encode(row.try_get::, _>("inventory_digest")?), approved_by: row.try_get("approved_by")?, + approval_origin: row.try_get::("approval_origin")?.parse()?, note: row.try_get("note")?, approved_at: row.try_get("approved_at")?, }) @@ -1187,7 +1214,7 @@ impl DeletionStore { .bind(schema) .bind(storage) .bind(frozen) - .bind(digest) + .bind(&digest) .fetch_optional(&self.pool) .await? .ok_or_else(|| { @@ -1229,8 +1256,8 @@ impl DeletionStore { } sqlx::query( "INSERT INTO community_deletion_approvals \ - (request_id, community_id, inventory_digest, approved_by, note) \ - VALUES ($1, $2, $3, $4, $5)", + (request_id, community_id, inventory_digest, approved_by, approval_origin, note) \ + VALUES ($1, $2, $3, $4, 'operator', $5)", ) .bind(request_id) .bind(community_id) @@ -1276,6 +1303,255 @@ impl DeletionStore { self.claim(None, owner, lease_duration).await } + /// Claim the oldest due operator-attested owner submission for preparation. + pub async fn claim_next_owner_submission( + &self, + owner: &str, + lease_duration: Duration, + ) -> Result> { + self.claim_owner_submission(None, owner, lease_duration) + .await + } + + /// Claim one due operator-attested owner submission by request id. + pub async fn claim_specific_owner_submission( + &self, + request_id: Uuid, + owner: &str, + lease_duration: Duration, + ) -> Result> { + self.claim_owner_submission(Some(request_id), owner, lease_duration) + .await + } + + async fn claim_owner_submission( + &self, + request_id: Option, + owner: &str, + lease_duration: Duration, + ) -> Result> { + let lease_seconds = i64::try_from(lease_duration.as_secs()).unwrap_or(i64::MAX); + let mut tx = self.pool.begin().await?; + let row = sqlx::query( + r#"WITH candidate AS ( + SELECT id FROM community_deletion_requests + WHERE ($1::uuid IS NULL OR id = $1) + AND request_origin = 'owner' AND acknowledgement_version = $4 + AND stage = 'submitted' + AND blocked_at IS NULL AND next_attempt_at <= now() + AND (lease_until IS NULL OR lease_until < now()) + ORDER BY created_at, id + FOR UPDATE SKIP LOCKED LIMIT 1 + ) + UPDATE community_deletion_requests request + SET lease_owner = $2, lease_generation = lease_generation + 1, + lease_until = now() + make_interval(secs => $3), + attempts = attempts + 1, updated_at = now() + FROM candidate WHERE request.id = candidate.id + RETURNING request.*"#, + ) + .bind(request_id) + .bind(owner) + .bind(lease_seconds) + .bind(OWNER_DELETION_ACKNOWLEDGEMENT_VERSION) + .fetch_optional(&mut *tx) + .await?; + tx.commit().await?; + let Some(row) = row else { + return Ok(None); + }; + let request = row_to_request(row)?; + let lease = LeaseToken { + request_id: request.id, + owner: owner.to_owned(), + generation: request.lease_generation, + community_id: request.community_id, + fence_generation: request.fence_generation, + }; + Ok(Some(ClaimedDeletion { request, lease })) + } + + /// Renew a live owner-submission preparation lease. + pub async fn heartbeat_owner_submission( + &self, + token: &LeaseToken, + executor_mode: &str, + lease_duration: Duration, + draining: bool, + ) -> Result<()> { + let lease_seconds = i64::try_from(lease_duration.as_secs()).unwrap_or(i64::MAX); + let mut tx = self.pool.begin().await?; + verify_owner_submission_lease(&mut tx, token).await?; + let affected = sqlx::query( + "UPDATE community_deletion_requests SET lease_until = \ + now() + make_interval(secs => $4), updated_at = now() \ + WHERE id = $1 AND lease_owner = $2 AND lease_generation = $3", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .bind(lease_seconds) + .execute(&mut *tx) + .await? + .rows_affected(); + if affected != 1 { + return Err(stale_lease_error(token)); + } + upsert_executor_heartbeat(&mut tx, token, executor_mode, draining).await?; + tx.commit().await?; + Ok(()) + } + + /// Atomically freeze inventory and approve operator-attested owner intent. + pub async fn complete_owner_preparation( + &self, + token: &LeaseToken, + inventory: &FrozenInventory, + ) -> Result { + validate_storage_manifest(&inventory.storage)?; + let digest = inventory.digest()?; + let schema = serde_json::to_value(&inventory.schema)?; + let storage = serde_json::to_value(&inventory.storage)?; + let frozen = serde_json::to_value(inventory)?; + let mut tx = self.pool.begin().await?; + + lock_community_deletion_shared(&mut tx, token.community_id).await?; + let community = sqlx::query( + "SELECT archived_at, deletion_state, deleted_at FROM communities \ + WHERE id = $1 FOR UPDATE", + ) + .bind(token.community_id.as_uuid()) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| { + DbError::DeletionSafety(format!( + "owner deletion {} community is missing before automatic approval", + token.request_id + )) + })?; + let current_owners: Vec = sqlx::query_scalar( + "SELECT pubkey FROM relay_members \ + WHERE community_id = $1 AND role = 'owner' ORDER BY pubkey FOR UPDATE", + ) + .bind(token.community_id.as_uuid()) + .fetch_all(&mut *tx) + .await?; + let request_row = + sqlx::query("SELECT * FROM community_deletion_requests WHERE id = $1 FOR UPDATE") + .bind(token.request_id) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| stale_lease_error(token))?; + let request = row_to_request(request_row)?; + let lease_live: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests \ + WHERE id = $1 AND lease_owner = $2 AND lease_generation = $3 \ + AND lease_until >= now())", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .fetch_one(&mut *tx) + .await?; + let lease_matches = lease_live + && request.community_id == token.community_id + && request.request_origin == DeletionRequestOrigin::Owner + && request.acknowledgement_version == Some(OWNER_DELETION_ACKNOWLEDGEMENT_VERSION) + && request.lease_owner.as_deref() == Some(token.owner.as_str()) + && request.lease_generation == token.generation + && request.blocked_reason.is_none(); + if !lease_matches { + return Err(stale_lease_error(token)); + } + let archived_at: Option> = community.try_get("archived_at")?; + let deletion_state: String = community.try_get("deletion_state")?; + let deleted_at: Option> = community.try_get("deleted_at")?; + let owner_authority_matches = request.owner_pubkey.as_ref().is_some_and(|owner| { + current_owners.len() == 1 && current_owners.first() == Some(owner) + }); + if archived_at.is_none() + || deletion_state != "active" + || deleted_at.is_some() + || !owner_authority_matches + { + return Err(DbError::DeletionSafety(format!( + "owner deletion {} community is no longer archived under the admitted owner", + token.request_id + ))); + } + if request.stage == DeletionStage::Approved { + let approval: Option<(Vec, String, String)> = sqlx::query_as( + "SELECT inventory_digest, approved_by, approval_origin \ + FROM community_deletion_approvals WHERE request_id = $1", + ) + .bind(token.request_id) + .fetch_optional(&mut *tx) + .await?; + let expected_operator = request.mediating_operator_pubkey.as_deref(); + let digest_hex = hex::encode(&digest); + let converged = request.inventory_digest.as_deref() == Some(digest_hex.as_str()) + && approval + .as_ref() + .is_some_and(|(approved_digest, approved_by, origin)| { + approved_digest.as_slice() == digest + && Some(approved_by.as_str()) == expected_operator + && origin == "owner_automatic" + }); + if converged { + tx.commit().await?; + return Ok(request); + } + return Err(DbError::DeletionSafety(format!( + "deletion {} automatic preparation does not match frozen approval evidence", + token.request_id + ))); + } + if request.stage != DeletionStage::Submitted { + return Err(stale_lease_error(token)); + } + let mediating_operator = request.mediating_operator_pubkey.clone().ok_or_else(|| { + DbError::DeletionSafety(format!( + "owner deletion {} is missing mediating operator provenance", + token.request_id + )) + })?; + sqlx::query( + "UPDATE community_deletion_requests SET stage = 'inventoried', \ + schema_manifest = $2, storage_manifest = $3, inventory_manifest = $4, \ + inventory_digest = $5, inventory_frozen_at = now(), updated_at = now() \ + WHERE id = $1", + ) + .bind(token.request_id) + .bind(schema) + .bind(storage) + .bind(frozen) + .bind(&digest) + .execute(&mut *tx) + .await?; + sqlx::query( + "INSERT INTO community_deletion_approvals \ + (request_id, community_id, inventory_digest, approved_by, approval_origin) \ + VALUES ($1, $2, $3, $4, 'owner_automatic')", + ) + .bind(token.request_id) + .bind(token.community_id.as_uuid()) + .bind(digest) + .bind(mediating_operator) + .execute(&mut *tx) + .await?; + let row = sqlx::query( + "UPDATE community_deletion_requests SET stage = 'approved', \ + retry_count = 0, retry_stage = NULL, next_attempt_at = now(), \ + last_error = NULL, last_error_at = NULL, updated_at = now() \ + WHERE id = $1 AND stage = 'inventoried' RETURNING *", + ) + .bind(token.request_id) + .fetch_one(&mut *tx) + .await?; + tx.commit().await?; + row_to_request(row) + } + async fn claim( &self, request_id: Option, @@ -1406,19 +1682,7 @@ impl DeletionStore { if affected != 1 { return Err(stale_lease_error(token)); } - sqlx::query( - "INSERT INTO community_deletion_executor_heartbeats \ - (executor_id, mode, request_id, draining) VALUES ($1, $2, $3, $4) \ - ON CONFLICT (executor_id) DO UPDATE SET mode = EXCLUDED.mode, \ - request_id = EXCLUDED.request_id, heartbeat_at = now(), \ - draining = EXCLUDED.draining, stopped_at = NULL", - ) - .bind(&token.owner) - .bind(executor_mode) - .bind(token.request_id) - .bind(draining) - .execute(&mut *tx) - .await?; + upsert_executor_heartbeat(&mut tx, token, executor_mode, draining).await?; tx.commit().await?; Ok(()) } @@ -2202,6 +2466,81 @@ impl DeletionStore { Ok(()) } + /// Persist a retryable owner-inventory failure and release its preparation lease. + pub async fn record_owner_preparation_retry( + &self, + token: &LeaseToken, + unit_key: &str, + error: &str, + retry_after: Duration, + ) -> Result<()> { + let bounded = bound_text(error, 4096); + let retry_seconds = i64::try_from(retry_after.as_secs()).unwrap_or(i64::MAX); + let mut tx = self.pool.begin().await?; + verify_owner_submission_lease(&mut tx, token).await?; + let (retry_count, retry_stage): (i32, Option) = sqlx::query_as( + "SELECT retry_count, retry_stage FROM community_deletion_requests \ + WHERE id = $1 FOR UPDATE", + ) + .bind(token.request_id) + .fetch_one(&mut *tx) + .await?; + let consecutive_retries = if retry_stage.as_deref() == Some("submitted") { + retry_count.saturating_add(1) + } else { + 1 + }; + let exhausted = consecutive_retries >= 8; + checkpoint_failed_tx(&mut tx, token, DeletionStage::Submitted, unit_key, &bounded).await?; + sqlx::query( + "UPDATE community_deletion_requests SET retry_count = $5, \ + retry_stage = 'submitted', last_error = $4, last_error_at = now(), \ + next_attempt_at = CASE WHEN $6 THEN next_attempt_at \ + ELSE now() + make_interval(secs => $7) END, \ + blocked_at = CASE WHEN $6 THEN now() ELSE blocked_at END, \ + blocked_reason = CASE WHEN $6 THEN $4 ELSE blocked_reason END, \ + lease_owner = NULL, lease_until = NULL, updated_at = now() \ + WHERE id = $1 AND lease_owner = $2 AND lease_generation = $3", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .bind(&bounded) + .bind(consecutive_retries) + .bind(exhausted) + .bind(retry_seconds) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(()) + } + + /// Permanently block unsafe owner inventory preparation and release its lease. + pub async fn block_owner_preparation( + &self, + token: &LeaseToken, + unit_key: &str, + error: &str, + ) -> Result<()> { + let bounded = bound_text(error, 4096); + let mut tx = self.pool.begin().await?; + verify_owner_submission_lease(&mut tx, token).await?; + checkpoint_failed_tx(&mut tx, token, DeletionStage::Submitted, unit_key, &bounded).await?; + sqlx::query( + "UPDATE community_deletion_requests SET blocked_at = now(), blocked_reason = $4, \ + last_error = $4, last_error_at = now(), lease_owner = NULL, lease_until = NULL, \ + updated_at = now() WHERE id = $1 AND lease_owner = $2 AND lease_generation = $3", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .bind(&bounded) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(()) + } + /// Terminally abort a request at the reversible pre-destruction boundary. /// /// `submitted`, `inventoried`, `approved`, and `fenced` are reversible: @@ -3121,31 +3460,80 @@ fn validate_manifest_key_chunks( Ok(()) } -async fn verify_lease( +async fn upsert_executor_heartbeat( tx: &mut Transaction<'_, Postgres>, token: &LeaseToken, - stage: DeletionStage, + executor_mode: &str, + draining: bool, ) -> Result<()> { - let valid: bool = sqlx::query_scalar( - "SELECT EXISTS(SELECT 1 FROM community_deletion_requests request \ - JOIN community_deletion_approvals approval ON approval.request_id = request.id \ - AND approval.community_id = request.community_id \ - AND approval.inventory_digest = request.inventory_digest \ - WHERE request.id = $1 AND request.community_id = $5 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)", + sqlx::query( + "INSERT INTO community_deletion_executor_heartbeats \ + (executor_id, mode, request_id, draining) VALUES ($1, $2, $3, $4) \ + ON CONFLICT (executor_id) DO UPDATE SET mode = EXCLUDED.mode, \ + request_id = EXCLUDED.request_id, heartbeat_at = now(), \ + draining = EXCLUDED.draining, stopped_at = 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) + .bind(executor_mode) + .bind(token.request_id) + .bind(draining) + .execute(&mut **tx) .await?; - if valid { - Ok(()) - } else { - Err(stale_lease_error(token)) + Ok(()) +} + +async fn verify_owner_submission_lease( + tx: &mut Transaction<'_, Postgres>, + token: &LeaseToken, +) -> Result<()> { + let valid = sqlx::query_scalar::<_, Uuid>( + "SELECT request.id FROM community_deletion_requests request \ + WHERE request.id = $1 AND request.community_id = $4 \ + AND request.request_origin = 'owner' AND request.stage = 'submitted' \ + AND request.acknowledgement_version = $5 \ + AND request.lease_owner = $2 AND request.lease_generation = $3 \ + AND request.lease_until >= now() AND request.blocked_at IS NULL \ + FOR UPDATE", + ) + .bind(token.request_id) + .bind(&token.owner) + .bind(token.generation) + .bind(token.community_id.as_uuid()) + .bind(OWNER_DELETION_ACKNOWLEDGEMENT_VERSION) + .fetch_optional(&mut **tx) + .await?; + if valid.is_some() { + Ok(()) + } else { + Err(stale_lease_error(token)) + } +} + +async fn verify_lease( + tx: &mut Transaction<'_, Postgres>, + token: &LeaseToken, + stage: DeletionStage, +) -> Result<()> { + let valid: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM community_deletion_requests request \ + JOIN community_deletion_approvals approval ON approval.request_id = request.id \ + AND approval.community_id = request.community_id \ + AND approval.inventory_digest = request.inventory_digest \ + WHERE request.id = $1 AND request.community_id = $5 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)", + ) + .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 { + Ok(()) + } else { + Err(stale_lease_error(token)) } } @@ -4763,6 +5151,77 @@ mod postgres_tests { ); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn privileged_abort_fences_a_live_owner_preparation_lease() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let claim = store + .claim_specific_owner_submission(request_id, "preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(community), + }; + let (lease_owner, lease_generation, lease_is_live): (Option, i64, bool) = + sqlx::query_as( + "SELECT lease_owner, lease_generation, lease_until >= now() \ + FROM community_deletion_requests WHERE id = $1", + ) + .bind(request_id) + .fetch_one(&db.pool) + .await + .expect("read live preparation lease"); + assert_eq!(lease_owner.as_deref(), Some("preparer")); + assert_eq!(lease_generation, claim.lease.generation); + assert!(lease_is_live, "preparation lease must be live before abort"); + + let aborted = store + .abort(request_id, "recovery-operator", "cancel live preparation") + .await + .expect("abort live preparation"); + assert_eq!(aborted.stage, DeletionStage::Aborted); + assert_eq!(aborted.lease_generation, claim.lease.generation + 1); + assert!(aborted.lease_owner.is_none()); + assert!(aborted.lease_until.is_none()); + + let heartbeat_error = store + .heartbeat_owner_submission(&claim.lease, "drain", DEFAULT_LEASE_DURATION, false) + .await + .expect_err("aborted preparation lease cannot heartbeat"); + assert!(is_stale_deletion_lease(&heartbeat_error)); + let completion_error = store + .complete_owner_preparation(&claim.lease, &inventory) + .await + .expect_err("aborted preparation lease cannot approve"); + assert!(is_stale_deletion_lease(&completion_error)); + + let request = store.get(request_id).await.expect("load aborted request"); + assert_eq!(request.stage, DeletionStage::Aborted); + assert!(request.inventory_digest.is_none()); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM community_deletion_approvals WHERE request_id = $1", + ) + .bind(request_id) + .fetch_one(&db.pool) + .await + .expect("count automatic approvals"), + 0 + ); + } + /// The reversible boundary extends through `fenced`. From `drained` /// onward, destruction may have begun, so abort must stay closed. #[tokio::test] @@ -4808,6 +5267,474 @@ mod postgres_tests { } } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_preparation_claim_is_owner_only_concurrent_and_reclaimable() { + let (db, store) = store().await; + let (owner_host, owner, _) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + let OwnerDeletionAdmission::Accepted(owner_request) = store + .admit_owner_request(&owner_host, &owner, operator, 1, request_id) + .await + .expect("admit owner request") + else { + panic!("expected accepted owner request") + }; + let operator_host = format!("operator-delete-{}.example", Uuid::new_v4().simple()); + db.ensure_configured_community(&operator_host) + .await + .expect("create operator community"); + let operator_request = store + .submit(&operator_host, "manual-operator", None) + .await + .expect("submit manual request"); + + let (first, second) = tokio::join!( + store.claim_specific_owner_submission( + owner_request.id, + "preparer-a", + DEFAULT_LEASE_DURATION, + ), + store.claim_specific_owner_submission( + owner_request.id, + "preparer-b", + DEFAULT_LEASE_DURATION, + ), + ); + let claims = [first.expect("first claim"), second.expect("second claim")]; + assert_eq!(claims.iter().filter(|claim| claim.is_some()).count(), 1); + let claim = claims.into_iter().flatten().next().expect("one winner"); + assert_eq!(claim.request.id, owner_request.id); + assert_eq!(claim.request.request_origin, DeletionRequestOrigin::Owner); + assert_ne!(claim.request.id, operator_request.id); + assert!(store + .claim_specific_owner_submission( + operator_request.id, + "operator-preparer", + DEFAULT_LEASE_DURATION, + ) + .await + .expect("operator request selection") + .is_none()); + + store + .heartbeat_owner_submission(&claim.lease, "drain", DEFAULT_LEASE_DURATION, false) + .await + .expect("heartbeat owner preparation"); + sqlx::query( + "UPDATE community_deletion_requests SET lease_until = now() - interval '1 second' WHERE id = $1", + ) + .bind(owner_request.id) + .execute(&db.pool) + .await + .expect("expire preparation lease"); + let successor = store + .claim_specific_owner_submission(owner_request.id, "preparer-c", DEFAULT_LEASE_DURATION) + .await + .expect("reclaim expired preparation") + .expect("expired owner preparation is reclaimable"); + assert_eq!(successor.request.id, owner_request.id); + assert!(successor.lease.generation > claim.lease.generation); + assert!(store + .heartbeat_owner_submission(&claim.lease, "drain", DEFAULT_LEASE_DURATION, false) + .await + .is_err()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn blocked_owner_submission_cannot_be_claimed() { + let (db, store) = store().await; + let (host, owner, _) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + sqlx::query( + "UPDATE community_deletion_requests SET blocked_at = now(), \ + blocked_reason = 'operator hold' WHERE id = $1", + ) + .bind(request_id) + .execute(&db.pool) + .await + .expect("block owner submission before claim"); + + assert!(store + .claim_specific_owner_submission( + request_id, + "specific-preparer", + DEFAULT_LEASE_DURATION + ) + .await + .expect("specific claim selection") + .is_none()); + assert!(store + .claim_next_owner_submission("next-preparer", DEFAULT_LEASE_DURATION) + .await + .expect("next claim selection") + .is_none()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn blocked_claimed_owner_submission_cannot_heartbeat_or_approve() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let claim = store + .claim_specific_owner_submission(request_id, "preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + sqlx::query( + "UPDATE community_deletion_requests SET blocked_at = now(), \ + blocked_reason = 'operator hold' WHERE id = $1", + ) + .bind(request_id) + .execute(&db.pool) + .await + .expect("block owner submission while retaining lease"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(community), + }; + + assert!(store + .heartbeat_owner_submission(&claim.lease, "drain", DEFAULT_LEASE_DURATION, false) + .await + .is_err()); + assert!(store + .complete_owner_preparation(&claim.lease, &inventory) + .await + .is_err()); + assert_eq!( + store + .get(request_id) + .await + .expect("load blocked request") + .stage, + DeletionStage::Submitted + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM community_deletion_approvals WHERE request_id = $1", + ) + .bind(request_id) + .fetch_one(&db.pool) + .await + .expect("count automatic approvals"), + 0 + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_preparation_atomically_approves_exact_inventory_and_converges() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let claim = store + .claim_specific_owner_submission(request_id, "preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(community), + }; + + let approved = store + .complete_owner_preparation(&claim.lease, &inventory) + .await + .expect("complete owner preparation"); + assert_eq!(approved.stage, DeletionStage::Approved); + assert_eq!( + approved.inventory_digest, + Some(hex::encode(inventory.digest().unwrap())) + ); + let replay = store + .complete_owner_preparation(&claim.lease, &inventory) + .await + .expect("ambiguous commit replay converges"); + assert_eq!(replay, approved); + let inspection = store.inspect(request_id).await.expect("inspect approval"); + let approval = inspection.approval.expect("automatic approval evidence"); + assert_eq!( + approval.approval_origin, + DeletionApprovalOrigin::OwnerAutomatic + ); + assert_eq!(approval.approved_by, operator); + assert!(store + .claim_specific(request_id, "other-executor", DEFAULT_LEASE_DURATION) + .await + .expect("existing execution claim remains approval-bound") + .is_none()); + store + .verify_execution_token(&claim.lease, DeletionStage::Approved) + .await + .expect("retained lease is immediately execution eligible"); + + let changed = FrozenInventory { + schema: inventory.schema.clone(), + storage: StorageManifest { + version: inventory.storage.version, + prefixes: inventory + .storage + .prefixes + .iter() + .cloned() + .map(|mut prefix| { + prefix.total_bytes += 1; + prefix + }) + .collect(), + }, + }; + assert!(store + .complete_owner_preparation(&claim.lease, &changed) + .await + .is_err()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_preparation_rechecks_archived_current_owner_before_automatic_approval() { + enum AuthorityDrift { + Unarchived, + OwnerChanged, + } + + let mut failures = Vec::new(); + for drift in [AuthorityDrift::Unarchived, AuthorityDrift::OwnerChanged] { + let label = match drift { + AuthorityDrift::Unarchived => "unarchived", + AuthorityDrift::OwnerChanged => "owner-changed", + }; + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let claim = store + .claim_specific_owner_submission(request_id, "preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(community), + }; + + match drift { + AuthorityDrift::Unarchived => { + sqlx::query("UPDATE communities SET archived_at = NULL WHERE id = $1") + .bind(community.as_uuid()) + .execute(&db.pool) + .await + .expect("simulate stale archive authority"); + } + AuthorityDrift::OwnerChanged => { + sqlx::query( + "UPDATE relay_members SET role = 'member' \ + WHERE community_id = $1 AND pubkey = $2 AND role = 'owner'", + ) + .bind(community.as_uuid()) + .bind(&owner) + .execute(&db.pool) + .await + .expect("simulate stale owner authority"); + } + } + + let completion = store + .complete_owner_preparation(&claim.lease, &inventory) + .await; + let request = store.get(request_id).await.expect("load request"); + let approval_count = sqlx::query_scalar::<_, i64>( + "SELECT count(*) FROM community_deletion_approvals WHERE request_id = $1", + ) + .bind(request_id) + .fetch_one(&db.pool) + .await + .expect("count automatic approvals"); + if completion.is_ok() + || request.stage != DeletionStage::Submitted + || request.inventory_digest.is_some() + || approval_count != 0 + { + failures.push(format!( + "{label}: completion={completion:?}, stage={}, inventory_frozen={}, approvals={approval_count}", + request.stage, + request.inventory_digest.is_some(), + )); + } + } + + assert!( + failures.is_empty(), + "stale owner authority reached automatic approval: {failures:#?}" + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn stale_owner_preparation_generation_cannot_freeze_or_approve() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let stale = store + .claim_specific_owner_submission(request_id, "stale-preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + sqlx::query( + "UPDATE community_deletion_requests SET lease_until = now() - interval '1 second' WHERE id = $1", + ) + .bind(stale.request.id) + .execute(&db.pool) + .await + .expect("expire stale lease"); + let successor = store + .claim_specific_owner_submission( + request_id, + "successor-preparer", + DEFAULT_LEASE_DURATION, + ) + .await + .expect("reclaim owner request") + .expect("expired request is reclaimable"); + let inventory = FrozenInventory { + schema: store + .inventory_schema(community) + .await + .expect("schema inventory"), + storage: empty_storage_manifest(community), + }; + assert!(is_stale_deletion_lease( + &store + .complete_owner_preparation(&stale.lease, &inventory) + .await + .expect_err("stale generation cannot commit") + )); + let approved = store + .complete_owner_preparation(&successor.lease, &inventory) + .await + .expect("successor commits preparation"); + assert_eq!(approved.stage, DeletionStage::Approved); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn owner_preparation_retry_block_and_privileged_abort_are_recoverable() { + let (db, store) = store().await; + let (host, owner, community) = archived_owned_community(&db).await; + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let request_id = Uuid::new_v4(); + store + .admit_owner_request(&host, &owner, operator, 1, request_id) + .await + .expect("admit owner request"); + let claim = store + .claim_specific_owner_submission(request_id, "preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + store + .record_owner_preparation_retry( + &claim.lease, + "inventory", + "temporary object-store failure", + Duration::ZERO, + ) + .await + .expect("record preparation retry"); + let retried = store.get(request_id).await.expect("load retried request"); + assert_eq!(retried.retry_stage, Some(DeletionStage::Submitted)); + assert_eq!(retried.retry_count, 1); + assert!(retried.lease_owner.is_none()); + + let claim = store + .claim_specific_owner_submission(request_id, "preparer-2", DEFAULT_LEASE_DURATION) + .await + .expect("reclaim owner request") + .expect("retried request is due"); + store + .block_owner_preparation(&claim.lease, "inventory", "unsafe storage taxonomy") + .await + .expect("block preparation"); + assert_eq!( + store + .get(request_id) + .await + .expect("load blocked request") + .blocked_reason + .as_deref(), + Some("unsafe storage taxonomy") + ); + let aborted = store + .abort( + request_id, + "recovery-operator", + "cannot safely enumerate storage", + ) + .await + .expect("abort submitted request"); + assert_eq!(aborted.stage, DeletionStage::Aborted); + assert_eq!( + sqlx::query_scalar::<_, String>("SELECT deletion_state FROM communities WHERE id = $1") + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("community lifecycle"), + "active" + ); + assert!(sqlx::query_scalar::<_, Option>>( + "SELECT archived_at FROM communities WHERE id = $1" + ) + .bind(community.as_uuid()) + .fetch_one(&db.pool) + .await + .expect("community archive state") + .is_some()); + assert!(matches!( + store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("fresh owner request after recovery abort"), + OwnerDeletionAdmission::Accepted(_) + )); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn approval_boundary_blocks_claim_until_exact_inventory_is_approved() { @@ -4856,6 +5783,16 @@ mod postgres_tests { .await .expect("approve"); assert_eq!(approved.stage, DeletionStage::Approved); + assert_eq!( + store + .inspect(request.id) + .await + .expect("inspect manual approval") + .approval + .expect("manual approval evidence") + .approval_origin, + DeletionApprovalOrigin::Operator + ); assert_eq!( approved.inventory_digest, Some(hex::encode(inventory.digest().unwrap())) diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index 2814c486917..46cdad8265b 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -290,7 +290,7 @@ pub enum Command { #[arg(long)] executor_id: Option, }, - /// Drain the currently runnable deletion queue, then exit. + /// Drain runnable work, preparing operator-attested owner submissions when idle. Drain { /// Executor identity (defaults to hostname/pid). #[arg(long)] @@ -450,7 +450,7 @@ async fn run_with_services(command: Command, services: Services) -> Result .store .submit(&host, &requested_by, reason.as_deref()) .await?; - let inventory = build_inventory(&services, &request).await?; + let inventory = build_inventory(&services, &request, None).await?; let request = services .store .freeze_inventory(request.id, &inventory) @@ -683,12 +683,13 @@ fn validate_storage_ownership(request: &DeletionRequest, manifest: &StorageManif async fn build_inventory( services: &Services, request: &DeletionRequest, + heartbeat_lost: Option<&CancellationToken>, ) -> Result { let schema = services .store .inventory_schema(request.community_id) .await?; - let storage = enumerate_tenant_prefixes(services, request, None, None).await?; + let storage = enumerate_tenant_prefixes(services, request, heartbeat_lost, None).await?; Ok(FrozenInventory { schema, storage }) } @@ -977,7 +978,15 @@ async fn run_loop( .await? } }; - let Some(claim) = claim else { + let claim = if claim.is_none() && mode == LoopMode::Drain && request_id.is_none() { + services + .store + .claim_next_owner_submission(&executor_id, DEFAULT_LEASE_DURATION) + .await? + } else { + claim + }; + let Some(mut claim) = claim else { if mode == LoopMode::Run && !ran { anyhow::bail!( "deletion request is not runnable, is blocked, or is leased by another executor" @@ -986,6 +995,34 @@ async fn run_loop( return Ok(0); }; ran = true; + if claim.request.stage == DeletionStage::Submitted { + let preparation_request = claim.request.clone(); + let preparation_services = &services; + match prepare_owner_claim_with( + &services, + mode, + claim, + &shutdown, + HEARTBEAT_INTERVAL, + |heartbeat_lost| async move { + build_inventory( + preparation_services, + &preparation_request, + Some(&heartbeat_lost), + ) + .await + }, + ) + .await? + { + OwnerPreparationOutcome::Prepared(prepared) => claim = *prepared, + OwnerPreparationOutcome::Stopped => return Ok(0), + OwnerPreparationOutcome::Failed(output) => { + print_json(&output)?; + return Ok(1); + } + } + } let output = execute_claim(&services, mode, claim, &shutdown).await?; print_json(&output)?; let failed = output.last_error.is_some() || output.blocked_reason.is_some(); @@ -995,6 +1032,173 @@ async fn run_loop( } } +enum OwnerPreparationOutcome { + Prepared(Box), + Stopped, + Failed(RunOutput), +} + +async fn record_owner_preparation_failure( + services: &Services, + token: &LeaseToken, + error: &anyhow::Error, +) -> Result<()> { + let message = format!("{error:#}"); + let result = if is_permanent_error(error) { + services + .store + .block_owner_preparation(token, "inventory", &message) + .await + } else { + services + .store + .record_owner_preparation_retry(token, "inventory", &message, RETRY_DELAY) + .await + }; + match result { + Ok(()) => Ok(()), + Err(error) if buzz_db::deletion::is_stale_deletion_lease(&error) => Ok(()), + Err(error) => Err(error.into()), + } +} + +async fn prepare_owner_claim_with( + services: &Services, + mode: LoopMode, + claim: ClaimedDeletion, + shutdown: &CancellationToken, + heartbeat_period: Duration, + build: F, +) -> Result +where + F: FnOnce(CancellationToken) -> Fut, + Fut: Future>, +{ + let token = claim.lease.clone(); + if let Err(error) = services + .store + .heartbeat_owner_submission(&token, mode.as_str(), DEFAULT_LEASE_DURATION, false) + .await + { + if buzz_db::deletion::is_stale_deletion_lease(&error) { + let request = services.store.get(token.request_id).await?; + return Ok(OwnerPreparationOutcome::Failed(run_output(request))); + } + return Err(error.into()); + } + + let heartbeat_services = services.clone(); + let heartbeat_token = token.clone(); + let heartbeat_mode = mode.as_str(); + let heartbeat_shutdown = CancellationToken::new(); + let heartbeat_cancel = heartbeat_shutdown.clone(); + let heartbeat_error = CancellationToken::new(); + let heartbeat_error_signal = heartbeat_error.clone(); + let heartbeat = tokio::spawn(async move { + let mut interval = tokio::time::interval(heartbeat_period); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + interval.tick().await; + loop { + tokio::select! { + _ = heartbeat_cancel.cancelled() => return, + _ = interval.tick() => { + if heartbeat_services + .store + .heartbeat_owner_submission( + &heartbeat_token, + heartbeat_mode, + DEFAULT_LEASE_DURATION, + false, + ) + .await + .is_err() + { + heartbeat_error_signal.cancel(); + return; + } + } + } + } + }); + + enum PreparationStage { + Prepared(Box), + Stopped, + Failed(anyhow::Error), + } + let preparation = async { + let inventory = build(heartbeat_error.clone()).await?; + services + .store + .complete_owner_preparation(&token, &inventory) + .await + .map_err(Into::into) + }; + let stage = tokio::select! { + biased; + _ = shutdown.cancelled() => PreparationStage::Stopped, + _ = heartbeat_error.cancelled() => PreparationStage::Failed(DeletionLeaseLost.into()), + result = preparation => match result { + Ok(request) => PreparationStage::Prepared(Box::new(request)), + Err(error) => PreparationStage::Failed(error), + }, + }; + heartbeat_shutdown.cancel(); + let stage = match heartbeat.await { + Ok(()) => stage, + Err(error) => PreparationStage::Failed(anyhow::anyhow!( + "owner deletion preparation heartbeat task failed: {error}" + )), + }; + + match stage { + PreparationStage::Prepared(request) => Ok(OwnerPreparationOutcome::Prepared(Box::new( + ClaimedDeletion { + request: *request, + lease: token, + }, + ))), + PreparationStage::Stopped => { + let _ = services + .store + .heartbeat_owner_submission(&token, mode.as_str(), DEFAULT_LEASE_DURATION, true) + .await; + services + .store + .stop_executor(Some(&token), &token.owner) + .await?; + Ok(OwnerPreparationOutcome::Stopped) + } + PreparationStage::Failed(error) => { + let request = services.store.get(token.request_id).await?; + if request.stage == DeletionStage::Approved + && request.lease_owner.as_deref() == Some(token.owner.as_str()) + && request.lease_generation == token.generation + && services + .store + .verify_execution_token(&token, DeletionStage::Approved) + .await + .is_ok() + { + return Ok(OwnerPreparationOutcome::Prepared(Box::new( + ClaimedDeletion { + request, + lease: token, + }, + ))); + } + if request.lease_owner.as_deref() != Some(token.owner.as_str()) + || request.lease_generation != token.generation + { + return Ok(OwnerPreparationOutcome::Failed(run_output(request))); + } + record_owner_preparation_failure(services, &token, &error).await?; + let request = services.store.get(token.request_id).await?; + Ok(OwnerPreparationOutcome::Failed(run_output(request))) + } + } +} + async fn stop_claim_executor( services: &Services, mode: LoopMode, @@ -1725,6 +1929,344 @@ mod postgres_tests { (db, services, claim) } + async fn owner_preparation_fixture( + prefix: &str, + ) -> (Db, Services, DeletionRequest, FrozenInventory) { + let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") + .or_else(|_| std::env::var("DATABASE_URL")) + .expect("BUZZ_TEST_DATABASE_URL or DATABASE_URL is required"); + let pool = sqlx::PgPool::connect(&database_url) + .await + .expect("connect owner preparation test DB"); + let db = Db::from_pool(pool); + if std::env::var("BUZZ_TEST_SCHEMA_MODE").as_deref() != Ok("desired") { + db.migrate().await.expect("migrate deletion engine test DB"); + } + let store = db.deletion_store(); + let host = format!("{prefix}-{}.example", Uuid::new_v4().simple()); + let owner = format!("{}{}", Uuid::new_v4().simple(), Uuid::new_v4().simple()); + let buzz_db::CreateCommunityWithOwnerResult::Created(community) = db + .create_community_with_owner(&host, &owner) + .await + .expect("create owner preparation community") + else { + panic!("expected fresh owner preparation community") + }; + db.archive_community_owned_by(&host, &owner, "protected.example") + .await + .expect("archive owner preparation community") + .expect("owned community"); + let operator = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + let buzz_db::deletion::OwnerDeletionAdmission::Accepted(request) = store + .admit_owner_request(&host, &owner, operator, 1, Uuid::new_v4()) + .await + .expect("admit owner request") + else { + panic!("owner request must be accepted") + }; + let inventory = FrozenInventory { + schema: store + .inventory_schema(community.id) + .await + .expect("inventory owner schema"), + storage: empty_storage_manifest(community.id), + }; + let services = Services { + store, + media: Arc::new( + MediaStorage::new(&buzz_media::MediaConfig { + s3_endpoint: "http://127.0.0.1:1".to_string(), + s3_access_key: "unused".to_string(), + s3_secret_key: "unused".to_string(), + s3_bucket: "unused".to_string(), + s3_region: "us-east-1".to_string(), + s3_addressing_style: buzz_media::S3AddressingStyle::Path, + max_image_bytes: 1, + max_gif_bytes: 1, + max_video_bytes: 1, + max_file_bytes: 1, + public_base_url: "http://localhost/media".to_string(), + upload_records_enabled: false, + upload_ip_header: None, + upload_port_header: None, + }) + .expect("construct unused media service"), + ), + redis: deadpool_redis::Config::from_url("redis://127.0.0.1:1") + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("construct unused Redis pool"), + }; + (db, services, *request, inventory) + } + + async fn claimed_owner_preparation( + prefix: &str, + ) -> (Db, Services, ClaimedDeletion, FrozenInventory) { + let (db, services, request, inventory) = owner_preparation_fixture(prefix).await; + let claim = services + .store + .claim_specific_owner_submission(request.id, "test-preparer", DEFAULT_LEASE_DURATION) + .await + .expect("claim owner request") + .expect("owner request is preparable"); + (db, services, claim, inventory) + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn drain_loop_prioritizes_approved_work_then_dispatches_owner_submission() { + let (db, services, owner_request, _) = + owner_preparation_fixture("owner-preparation-dispatch").await; + let approved_host = format!("approved-{}.example", Uuid::new_v4().simple()); + let approved_community = db + .ensure_configured_community(&approved_host) + .await + .expect("create approved community"); + let approved = services + .store + .submit(&approved_host, "manual-operator", None) + .await + .expect("submit approved request"); + let inventory = FrozenInventory { + schema: services + .store + .inventory_schema(approved_community.id) + .await + .expect("inventory approved community"), + storage: empty_storage_manifest(approved_community.id), + }; + services + .store + .freeze_inventory(approved.id, &inventory) + .await + .expect("freeze approved request"); + services + .store + .approve(approved.id, "manual-approver", None) + .await + .expect("approve request"); + + assert_eq!( + run_loop( + services.clone(), + LoopMode::Drain, + None, + "priority-executor".to_string(), + ) + .await + .expect("approved work failure remains typed"), + 1 + ); + let approved_after = services + .store + .get(approved.id) + .await + .expect("approved request after drain"); + assert!( + approved_after.attempts > 0, + "approved work must be claimed first" + ); + let owner_after_priority = services + .store + .get(owner_request.id) + .await + .expect("owner request after approved work"); + assert_eq!(owner_after_priority.stage, DeletionStage::Submitted); + assert_eq!(owner_after_priority.attempts, 0); + assert!(owner_after_priority.lease_owner.is_none()); + + assert_eq!( + run_loop( + services.clone(), + LoopMode::Drain, + None, + "owner-dispatch-executor".to_string(), + ) + .await + .expect("owner inventory failure remains typed"), + 1 + ); + let owner_after_dispatch = services + .store + .get(owner_request.id) + .await + .expect("owner request after dispatch"); + assert_eq!(owner_after_dispatch.attempts, 1); + assert_eq!( + owner_after_dispatch.retry_stage, + Some(DeletionStage::Submitted) + ); + assert!(owner_after_dispatch.lease_owner.is_none()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn drain_owner_preparation_approves_for_existing_execution_only() { + let (db, services, claim, inventory) = + claimed_owner_preparation("owner-preparation-success").await; + let operator_host = format!("manual-{}.example", Uuid::new_v4().simple()); + db.ensure_configured_community(&operator_host) + .await + .expect("create manual community"); + let manual = services + .store + .submit(&operator_host, "manual-operator", None) + .await + .expect("submit manual request"); + + let outcome = prepare_owner_claim_with( + &services, + LoopMode::Drain, + claim, + &CancellationToken::new(), + HEARTBEAT_INTERVAL, + |_| async move { Ok(inventory) }, + ) + .await + .expect("prepare owner claim"); + let OwnerPreparationOutcome::Prepared(prepared) = outcome else { + panic!("owner preparation should produce an approved execution claim") + }; + assert_eq!(prepared.request.stage, DeletionStage::Approved); + services + .store + .verify_execution_token(&prepared.lease, DeletionStage::Approved) + .await + .expect("prepared claim enters unchanged execution boundary"); + assert_eq!( + services + .store + .get(manual.id) + .await + .expect("manual request") + .stage, + DeletionStage::Submitted + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn drain_owner_preparation_persists_transient_and_permanent_failures() { + let (_, services, claim, _) = claimed_owner_preparation("owner-preparation-failure").await; + let request_id = claim.request.id; + let outcome = prepare_owner_claim_with( + &services, + LoopMode::Drain, + claim, + &CancellationToken::new(), + HEARTBEAT_INTERVAL, + |_| async { Err(transient("temporary inventory failure")) }, + ) + .await + .expect("record transient preparation failure"); + assert!(matches!(outcome, OwnerPreparationOutcome::Failed(_))); + let retried = services + .store + .get(request_id) + .await + .expect("retried request"); + assert_eq!(retried.retry_stage, Some(DeletionStage::Submitted)); + assert_eq!(retried.retry_count, 1); + + let (_, blocked_services, claim, _) = + claimed_owner_preparation("owner-preparation-permanent").await; + let blocked_request_id = claim.request.id; + let outcome = prepare_owner_claim_with( + &blocked_services, + LoopMode::Drain, + claim, + &CancellationToken::new(), + HEARTBEAT_INTERVAL, + |_| async { Err(permanent("unsafe inventory taxonomy")) }, + ) + .await + .expect("record permanent preparation failure"); + assert!(matches!(outcome, OwnerPreparationOutcome::Failed(_))); + assert!(blocked_services + .store + .get(blocked_request_id) + .await + .expect("blocked request") + .blocked_reason + .is_some()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn drain_owner_preparation_rejects_stale_lease_before_inventory() { + let (_, services, claim, inventory) = + claimed_owner_preparation("owner-preparation-stale").await; + services + .store + .stop_executor(Some(&claim.lease), &claim.lease.owner) + .await + .expect("release preparation lease"); + let built = Arc::new(AtomicBool::new(false)); + let observed = Arc::clone(&built); + let outcome = prepare_owner_claim_with( + &services, + LoopMode::Drain, + claim, + &CancellationToken::new(), + HEARTBEAT_INTERVAL, + move |_| async move { + observed.store(true, Ordering::SeqCst); + Ok(inventory) + }, + ) + .await + .expect("stale preparation converges without mutation"); + assert!(matches!(outcome, OwnerPreparationOutcome::Failed(_))); + assert!(!built.load(Ordering::SeqCst)); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn drain_owner_preparation_cancels_inventory_after_lease_loss() { + struct DropSignal(Arc); + + impl Drop for DropSignal { + fn drop(&mut self) { + self.0.store(true, Ordering::SeqCst); + } + } + + let (_, services, claim, _) = claimed_owner_preparation("owner-preparation-cancel").await; + let token = claim.lease.clone(); + let dropped = Arc::new(AtomicBool::new(false)); + let observed = Arc::clone(&dropped); + let shutdown = CancellationToken::new(); + let (inventory_started_tx, inventory_started_rx) = tokio::sync::oneshot::channel(); + let preparation = prepare_owner_claim_with( + &services, + LoopMode::Drain, + claim, + &shutdown, + Duration::from_millis(10), + move |_| async move { + let _drop_signal = DropSignal(observed); + inventory_started_tx + .send(()) + .expect("signal inventory started"); + std::future::pending::>().await + }, + ); + let revoke = async { + inventory_started_rx.await.expect("inventory started"); + services + .store + .stop_executor(Some(&token), &token.owner) + .await + .expect("revoke preparation lease"); + }; + let (outcome, ()) = tokio::join!(preparation, revoke); + assert!(matches!( + outcome.expect("lease loss is typed control flow"), + OwnerPreparationOutcome::Failed(_) + )); + assert!(dropped.load(Ordering::SeqCst)); + } + fn env_of<'a>(set: &'a [(&'a str, &'a str)]) -> impl Fn(&str) -> Option + use<'a> { move |name| { set.iter() diff --git a/docs/operator-community-deletion.md b/docs/operator-community-deletion.md index 61d8fd89456..0a3b1de5148 100644 --- a/docs/operator-community-deletion.md +++ b/docs/operator-community-deletion.md @@ -5,9 +5,13 @@ Buzz executes whole-community deletion through the typed, one-shot that command as a Kubernetes CronJob; it does not call relay HTTP and it does not add another queue or retry service. -Postgres remains the handoff and source of truth. A run claims only requests -that the deletion store considers runnable, heartbeats the existing lease, and -resumes from durable checkpoints. `concurrencyPolicy: Forbid` prevents scheduled +Postgres remains the handoff and source of truth. A run gives already-approved +work priority. When none is ready, it may claim an operator-attested +owner-origin request at `submitted`, build the existing bounded inventory, and +atomically freeze that inventory with a digest-bound `owner_automatic` +approval. The same lease then enters the unchanged executor and resumes from +durable checkpoints. +Operator-origin requests never auto-progress. `concurrencyPolicy: Forbid` prevents scheduled pod overlap, `backoffLimit: 0` prevents Kubernetes Job retries, and the deletion store remains authoritative when a pod exits, reaches its deadline, or is replaced. @@ -70,9 +74,11 @@ reports their objects with the `null` version id. ## Runbook -1. Confirm database migrations are current and the deletion request has crossed - the explicit inventory and approval boundary with - `buzz-admin deletions inspect `. +1. Confirm database migrations are current. For operator-origin requests, + confirm explicit inventory and operator approval with + `buzz-admin deletions inspect `. For owner-origin requests, + expect the drain to record `approval_origin: owner_automatic`; `approved_by` + is the immutable mediating operator, not the owner. 2. Confirm the selected Secret contains the required keys and the S3 principal has version-list and exact-version delete permissions. 3. Enable the CronJob and inspect its rendered command and environment before @@ -97,6 +103,10 @@ reports their objects with the `null` version id. 7. If a run fails or times out, fix the recorded dependency or permission failure. Do not add Kubernetes retries: the next scheduled drain consults the durable retry/checkpoint state and resumes only when the store allows it. +7. Use `buzz-admin deletions abort` as privileged recovery while a request is + still at `submitted` or `inventoried` when safe preparation cannot continue. + Abort preserves the archived community and immutable request evidence. An + operator may also `unblock` a remediated preparation failure. ## Deadlines, termination, and the retry budget @@ -133,12 +143,13 @@ killed process. Expect up to roughly a lease duration of delay before the request is runnable again; do not raise the grace period expecting a clean handoff. -The current owner self-serve relay admission records an owner-origin request at -`submitted` and intentionally performs no inventory or approval synchronously. -The deletion engine rejects `submitted` and `inventoried` requests at its -explicit approval boundary. Automating the privileged inventory/approval step -is therefore a separate control-plane slice; enabling this CronJob alone does -not make a newly accepted owner request destructive. +Owner self-serve relay admission still records only a `submitted` row and does +no inventory, approval, S3 work, or execution synchronously. A successful drain +has no human approval step or cooling-off period: operator-attested owner intent +is prepared automatically under privileged policy and becomes immediately +eligible for execution. Transient preparation failures use the existing retry +schedule; permanent or exhausted failures block durably. Owner-facing +admission has no cancellation endpoint. The chart has no existing PrometheusRule or provider-neutral CronJob alert integration. Operators must alert on failed/missed Jobs and long-running active diff --git a/migrations/0053_owner_deletion_auto_approval.sql b/migrations/0053_owner_deletion_auto_approval.sql new file mode 100644 index 00000000000..139ffb0bcf2 --- /dev/null +++ b/migrations/0053_owner_deletion_auto_approval.sql @@ -0,0 +1,22 @@ +-- Owner-origin deletion requests may be prepared automatically by the +-- privileged one-shot deletion drain. Manual approvals retain operator +-- provenance and the existing exact inventory-digest foreign key. +SET LOCAL lock_timeout = '5s'; + +ALTER TABLE community_deletion_approvals + ADD COLUMN approval_origin TEXT NOT NULL DEFAULT 'operator' + CHECK (approval_origin IN ('operator', 'owner_automatic')); + +ALTER TABLE community_deletion_requests + DROP CONSTRAINT community_deletion_requests_retry_stage_check, + ADD CONSTRAINT community_deletion_requests_retry_stage_check + CHECK (retry_stage IS NULL OR retry_stage IN ( + 'submitted', 'approved', 'fenced', 'drained', 'bindings_removed', + 'postgres_purged', 'cache_purged', 'logically_verified' + )); + +CREATE INDEX community_deletion_requests_owner_preparable + ON community_deletion_requests (next_attempt_at, created_at) + WHERE request_origin = 'owner' + AND stage = 'submitted' + AND blocked_at IS NULL; diff --git a/schema/schema.sql b/schema/schema.sql index c8f294c40ef..317044b0b56 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1256,7 +1256,7 @@ CREATE TABLE community_deletion_requests ( attempts INTEGER NOT NULL DEFAULT 0 CHECK (attempts >= 0), retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), retry_stage TEXT CHECK (retry_stage IS NULL OR retry_stage IN ( - 'approved', 'fenced', 'drained', 'bindings_removed', + 'submitted', 'approved', 'fenced', 'drained', 'bindings_removed', 'postgres_purged', 'cache_purged', 'logically_verified' )), next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), @@ -1304,12 +1304,19 @@ CREATE INDEX community_deletion_requests_runnable 'postgres_purged', 'cache_purged', 'logically_verified'); CREATE INDEX community_deletion_requests_lease ON community_deletion_requests (lease_until) WHERE lease_owner IS NOT NULL; +CREATE INDEX community_deletion_requests_owner_preparable + ON community_deletion_requests (next_attempt_at, created_at) + WHERE request_origin = 'owner' + AND stage = 'submitted' + AND blocked_at IS NULL; CREATE TABLE community_deletion_approvals ( request_id UUID PRIMARY KEY, community_id UUID NOT NULL, inventory_digest BYTEA NOT NULL CHECK (length(inventory_digest) = 32), approved_by TEXT NOT NULL, + approval_origin TEXT NOT NULL DEFAULT 'operator' + CHECK (approval_origin IN ('operator', 'owner_automatic')), note TEXT, approved_at TIMESTAMPTZ NOT NULL DEFAULT now(), FOREIGN KEY (request_id, community_id, inventory_digest)