am
2026-08-13 11:04:18 -07:00
parent 30d9609691
commit 520f0feea1
+212 -76
View File
@@ -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<String>,
/// 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<chrono::Utc>,
observation_completed_at: chrono::DateTime<chrono::Utc>,
@@ -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<i32> {
async fn run_with_services(command: Command, services: Services) -> Result<i32> {
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::<u64>()
);
match run_preview(&services, host.as_deref(), include_keys).await {
Err(error) if error.downcast_ref::<OutputClosed>().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<i32>
}
}
async fn run_preview(services: &Services, host: Option<&str>, include_keys: bool) -> Result<i32> {
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<i32> {
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::<u64>()
);
}
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<String> {
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<FrozenInventory> {
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<StorageManifest> {
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<StorageManifest> {
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<PreviewRedis> {
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<usize> {
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::<OutputClosed>().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::<serde_json::Value>(line).expect("valid record"))
.collect::<Vec<_>>();
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);
}