diff --git a/crates/buzz-deletion/src/lib.rs b/crates/buzz-deletion/src/lib.rs index 31c1f5613..49d983559 100644 --- a/crates/buzz-deletion/src/lib.rs +++ b/crates/buzz-deletion/src/lib.rs @@ -2,6 +2,7 @@ #![warn(missing_docs)] //! Shared durable whole-community deletion engine and store adapters. +use std::io::{self, Write}; #[cfg(test)] use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; @@ -223,7 +224,7 @@ pub enum Command { /// Canonical community host. Defaults to RELAY_URL's authority. #[arg(long)] host: Option, - /// Stream concrete object-store and Redis keys as NDJSON before the summary. + /// Stream concrete object-store and Redis keys as NDJSON before a required completion record. #[arg(long)] include_keys: bool, }, @@ -366,6 +367,8 @@ fn is_permanent_error(error: &anyhow::Error) -> bool { #[derive(Debug, Serialize)] struct PreviewOutput { record: &'static str, + observation_id: Uuid, + complete: bool, advisory: bool, observation_started_at: chrono::DateTime, observation_completed_at: chrono::DateTime, @@ -385,13 +388,15 @@ struct PreviewStorage { #[derive(Debug, Serialize)] struct PreviewRedis { - observed_key_count: u64, + observed_scan_entry_count: u64, + count_semantics: &'static str, keys_included: bool, } #[derive(Debug, Serialize)] struct PreviewKey<'a> { record: &'static str, + observation_id: Uuid, store: &'static str, key: &'a str, } @@ -474,53 +479,10 @@ pub async fn run(command: Command) -> Result { async fn run_with_services(command: Command, services: Services) -> Result { match command { Command::Preview { host, include_keys } => { - let observation_started_at = chrono::Utc::now(); - let relay_url = std::env::var("RELAY_URL").ok(); - let host = resolve_submit_host(host.as_deref(), relay_url.as_deref())?; - let (community_id, community_host) = services.store.preview_target(&host).await?; - let postgres = services.store.inventory_schema(community_id).await?; - let mut object_key_count = 0u64; - let object_store = enumerate_tenant_prefixes( - &services, - community_id, - None, - None, - include_keys.then_some(&mut object_key_count), - ) - .await?; - if include_keys { - debug_assert_eq!( - object_key_count, - object_store - .prefixes - .iter() - .map(|prefix| prefix.object_count) - .sum::() - ); + match run_preview(&services, host.as_deref(), include_keys).await { + Err(error) if error.downcast_ref::().is_some() => Ok(0), + result => result, } - let redis = - preview_redis_namespace(&services.redis, community_id, include_keys).await?; - let output = PreviewOutput { - record: "deletion_preview_summary", - advisory: true, - observation_started_at, - observation_completed_at: chrono::Utc::now(), - warning: "live state may change after this read-only observation; execution freezes its authoritative manifest only after fence and drain", - community_id: community_id.to_string(), - community_host, - postgres, - object_store: PreviewStorage { - manifest: object_store, - keys_included: include_keys, - }, - redis, - }; - if include_keys { - print_ndjson(&output)?; - } else { - print_json(&output)?; - } - Ok(0) } Command::Submit { host, @@ -592,6 +554,74 @@ async fn run_with_services(command: Command, services: Services) -> Result } } +async fn run_preview(services: &Services, host: Option<&str>, include_keys: bool) -> Result { + let mut output = io::stdout(); + run_preview_to(services, host, include_keys, &mut output).await +} + +async fn run_preview_to( + services: &Services, + host: Option<&str>, + include_keys: bool, + output: &mut (dyn Write + Send), +) -> Result { + let observation_id = Uuid::new_v4(); + let observation_started_at = chrono::Utc::now(); + let relay_url = std::env::var("RELAY_URL").ok(); + let host = resolve_submit_host(host, relay_url.as_deref())?; + let (community_id, community_host) = services.store.preview_target(&host).await?; + let postgres = services.store.inventory_schema(community_id).await?; + let mut object_key_count = 0u64; + let object_store = enumerate_tenant_prefixes( + services, + community_id, + None, + None, + include_keys.then_some((observation_id, &mut *output, &mut object_key_count)), + ) + .await?; + if include_keys { + debug_assert_eq!( + object_key_count, + object_store + .prefixes + .iter() + .map(|prefix| prefix.object_count) + .sum::() + ); + } + let redis = preview_redis_namespace( + &services.redis, + community_id, + observation_id, + include_keys.then_some(&mut *output), + ) + .await?; + let result = PreviewOutput { + record: "deletion_preview_complete", + observation_id, + complete: true, + advisory: true, + observation_started_at, + observation_completed_at: chrono::Utc::now(), + warning: "live state may change after this read-only observation; Redis SCAN entries may be duplicated or omitted during concurrent writes; output without this completion record is invalid", + community_id: community_id.to_string(), + community_host, + postgres, + object_store: PreviewStorage { + manifest: object_store, + keys_included: include_keys, + }, + redis, + }; + if include_keys { + write_ndjson(output, &result)?; + } else { + write_json(output, &result)?; + } + Ok(0) +} + fn resolve_submit_host(host: Option<&str>, relay_url: Option<&str>) -> Result { if let Some(host) = host { let host = host.trim(); @@ -742,10 +772,20 @@ fn validate_storage_ownership(request: &DeletionRequest, manifest: &StorageManif Ok(()) } +async fn ensure_bucket_is_unversioned(services: &Services) -> Result<()> { + if services.media.bucket_versioning_detected().await? { + return Err(permanent( + "bucket versioning detected; deletion cannot prove logical absence with delete markers", + )); + } + Ok(()) +} + async fn build_inventory( services: &Services, request: &DeletionRequest, ) -> Result { + ensure_bucket_is_unversioned(services).await?; let schema = services .store .inventory_schema(request.community_id) @@ -789,13 +829,8 @@ async fn enumerate_tenant_prefixes( community: buzz_core::CommunityId, heartbeat_lost: Option<&CancellationToken>, mut sink: Option<&mut ChunkSink<'_>>, - mut preview_key_count: Option<&mut u64>, + mut preview_output: Option<(Uuid, &mut (dyn Write + Send), &mut u64)>, ) -> Result { - if services.media.bucket_versioning_detected().await? { - return Err(permanent( - "bucket versioning detected; deletion cannot prove logical absence with delete markers", - )); - } let community_uuid = *community.as_uuid(); let chunk_keys = manifest_chunk_keys(); let mut prefixes = Vec::new(); @@ -819,13 +854,17 @@ async fn enumerate_tenant_prefixes( } digest.fold(&key)?; total_bytes = total_bytes.saturating_add(size); - if let Some(key_count) = preview_key_count.as_deref_mut() { - print_ndjson(&PreviewKey { - record: "deletion_preview_key", - store: "object_store", - key: &key, - })?; - *key_count = key_count.saturating_add(1); + if let Some((observation_id, output, key_count)) = preview_output.as_mut() { + write_ndjson( + &mut **output, + &PreviewKey { + record: "deletion_preview_key", + observation_id: *observation_id, + store: "object_store", + key: &key, + }, + )?; + **key_count = key_count.saturating_add(1); } if let Some(sink) = sink.as_deref_mut() { sink.buffered.push(key); @@ -873,6 +912,7 @@ async fn freeze_destructive_manifest( heartbeat_lost: &CancellationToken, ) -> Result { validate_frozen_inventory(request)?; + ensure_bucket_is_unversioned(services).await?; // A prior interrupted freeze may have left partial chunks; they were // never bound to a committed manifest, so rewrite them from scratch. services.store.clear_manifest_key_chunks(token).await?; @@ -1391,15 +1431,24 @@ async fn publish_disconnect_community( Ok(()) } -fn preview_redis_page(include_keys: bool, keys: &[String], key_count: &mut u64) -> Result<()> { +fn preview_redis_page( + output: &mut Option<&mut (dyn Write + Send)>, + observation_id: Uuid, + keys: &[String], + key_count: &mut u64, +) -> Result<()> { *key_count = key_count.saturating_add(keys.len() as u64); - if include_keys { + if let Some(output) = output.as_deref_mut() { for key in keys { - print_ndjson(&PreviewKey { - record: "deletion_preview_key", - store: "redis", - key, - })?; + write_ndjson( + &mut *output, + &PreviewKey { + record: "deletion_preview_key", + observation_id, + store: "redis", + key, + }, + )?; } } Ok(()) @@ -1408,7 +1457,8 @@ fn preview_redis_page(include_keys: bool, keys: &[String], key_count: &mut u64) async fn preview_redis_namespace( pool: &deadpool_redis::Pool, community: buzz_core::CommunityId, - include_keys: bool, + observation_id: Uuid, + mut output: Option<&mut (dyn Write + Send)>, ) -> Result { let mut connection = pool.get().await?; let pattern = format!("buzz:{community}:*"); @@ -1423,15 +1473,16 @@ async fn preview_redis_namespace( .arg(1000) .query_async(&mut *connection) .await?; - preview_redis_page(include_keys, &keys, &mut key_count)?; + preview_redis_page(&mut output, observation_id, &keys, &mut key_count)?; if next == 0 { break; } cursor = next; } Ok(PreviewRedis { - observed_key_count: key_count, - keys_included: include_keys, + observed_scan_entry_count: key_count, + count_semantics: "Redis SCAN observations; concurrent keyspace changes may produce duplicate or omitted entries", + keys_included: output.is_some(), }) } @@ -1572,9 +1623,30 @@ fn run_output(request: DeletionRequest) -> RunOutput { } } -fn print_ndjson(value: &impl Serialize) -> Result<()> { - println!("{}", serde_json::to_string(value)?); - Ok(()) +#[derive(Debug, thiserror::Error)] +#[error("stdout consumer closed the output stream")] +struct OutputClosed; + +fn write_bytes(output: &mut (dyn Write + Send), bytes: &[u8]) -> Result<()> { + output.write_all(bytes).map_err(|error| { + if error.kind() == io::ErrorKind::BrokenPipe { + OutputClosed.into() + } else { + error.into() + } + }) +} + +fn write_ndjson(output: &mut (dyn Write + Send), value: &impl Serialize) -> Result<()> { + let mut bytes = serde_json::to_vec(value)?; + bytes.push(b'\n'); + write_bytes(output, &bytes) +} + +fn write_json(output: &mut (dyn Write + Send), value: &impl Serialize) -> Result<()> { + let mut bytes = serde_json::to_vec_pretty(value)?; + bytes.push(b'\n'); + write_bytes(output, &bytes) } fn print_json(value: &impl Serialize) -> Result<()> { @@ -1586,11 +1658,75 @@ fn print_json(value: &impl Serialize) -> Result<()> { mod tests { use super::*; + struct BrokenPipeWriter; + + impl Write for BrokenPipeWriter { + fn write(&mut self, _buffer: &[u8]) -> io::Result { + Err(io::Error::from(io::ErrorKind::BrokenPipe)) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + #[test] + fn preview_ndjson_treats_a_closed_consumer_as_output_closed() { + let error = write_ndjson( + &mut BrokenPipeWriter, + &PreviewKey { + record: "deletion_preview_key", + observation_id: Uuid::nil(), + store: "redis", + key: "buzz:test:key", + }, + ) + .expect_err("broken pipe must stop streaming"); + assert!(error.downcast_ref::().is_some()); + } + + #[test] + fn preview_ndjson_binds_keys_to_a_terminal_completion_record() { + let observation_id = Uuid::new_v4(); + let mut output = Vec::new(); + write_ndjson( + &mut output, + &PreviewKey { + record: "deletion_preview_key", + observation_id, + store: "object_store", + key: "tenant/key", + }, + ) + .expect("write key record"); + write_ndjson( + &mut output, + &serde_json::json!({ + "record": "deletion_preview_complete", + "observation_id": observation_id, + "complete": true, + }), + ) + .expect("write completion record"); + + let records = String::from_utf8(output) + .expect("UTF-8 NDJSON") + .lines() + .map(|line| serde_json::from_str::(line).expect("valid record")) + .collect::>(); + assert_eq!(records.len(), 2); + assert_eq!(records[0]["observation_id"], records[1]["observation_id"]); + assert_eq!(records[1]["record"], "deletion_preview_complete"); + assert_eq!(records[1]["complete"], true); + } + #[test] fn preview_redis_count_saturates_without_retaining_keys() { let keys = vec!["buzz:a:one".to_string(), "buzz:a:two".to_string()]; let mut count = u64::MAX - 1; - preview_redis_page(false, &keys, &mut count).expect("count preview page"); + let mut output = None; + preview_redis_page(&mut output, Uuid::nil(), &keys, &mut count) + .expect("count preview page"); assert_eq!(count, u64::MAX); }