From 19c5c6d9df92856778580acadf152e9e6296b6fd Mon Sep 17 00:00:00 2001 From: Atish Patel Date: Thu, 16 Jul 2026 17:35:02 -0500 Subject: [PATCH] fix(media): correct upload policy accounting Co-authored-by: Codex --- crates/buzz-media/src/upload.rs | 116 ++++++++++++++++++++++------- crates/buzz-relay/src/api/media.rs | 33 +++++++- crates/buzz-relay/src/router.rs | 4 +- 3 files changed, 121 insertions(+), 32 deletions(-) diff --git a/crates/buzz-media/src/upload.rs b/crates/buzz-media/src/upload.rs index b30828a5c..647eae6ab 100644 --- a/crates/buzz-media/src/upload.rs +++ b/crates/buzz-media/src/upload.rs @@ -37,6 +37,8 @@ pub struct StreamingIngestInput<'a> { pub ctx: &'a TenantContext, /// Verified Blossom authorization event for the source bytes. pub auth_event: &'a nostr::Event, + /// Source hash claimed by the `X-SHA-256` request header. + pub claimed_source_hash: &'a str, /// Optional request content length for an early size rejection. pub content_length: Option, /// Optional moderation attribution recorded after durable publication. @@ -58,6 +60,7 @@ pub async fn process_streaming_ingest( config, ctx, auth_event, + claimed_source_hash, content_length, attribution, mode, @@ -81,24 +84,22 @@ pub async fn process_streaming_ingest( let source_path = source.path().to_path_buf(); let (source_hash, source_size, sniff) = stream_source(body_stream, &source_path, max_bytes).await?; + if source_hash != claimed_source_hash { + return Err(MediaError::HashMismatch); + } - let auth = auth_event.clone(); - let auth_hash = source_hash.clone(); - let bound_host = ctx.host().to_string(); - tokio::task::spawn_blocking(move || match mode { - UploadRouteMode::Media => { - verify_blossom_media_auth(&auth, &auth_hash, Some(&bound_host), 3600) - } - UploadRouteMode::Upload | UploadRouteMode::Legacy => { - verify_blossom_upload_auth(&auth, &auth_hash, Some(&bound_host), 3600) - } - }) - .await - .map_err(|_| MediaError::Internal)??; + // Authenticate the computed source hash before invoking any media parser. + // The permissive bound is narrowed below once content probing identifies + // whether this is a long-running video upload or a short-lived token class. + verify_source_auth(mode, auth_event, &source_hash, ctx.host(), 3_600).await?; let source_probe = crate::sanitize::probe_media(&source_path, &sniff, config).await?; let is_media = source_probe.is_some(); enforce_route_policy(mode, source_probe.as_ref())?; + let max_auth_age = auth_max_age_secs(source_probe.as_ref()); + if max_auth_age < 3_600 { + verify_source_auth(mode, auth_event, &source_hash, ctx.host(), max_auth_age).await?; + } enforce_source_size(source_probe.as_ref(), source_size, config)?; @@ -170,6 +171,8 @@ pub async fn process_streaming_ingest( if sidecar_exists && blob_exists { let meta = storage.get_sidecar(ctx, &output_hash).await?; if let Some(attribution) = &attribution { + let (record_source_hash, record_source_size, record_source_mime) = + source_record_fields(source_probe.as_ref(), &source_hash, source_size); record_upload_event( storage, ctx, @@ -180,13 +183,9 @@ pub async fn process_streaming_ingest( ext: &ext, mime: &mime, size: output_size, - source_sha256: Some(source_hash.as_str()), - source_size: Some(source_size), - source_mime: Some( - source_probe - .as_ref() - .map_or(mime.as_str(), |probe| probe.mime.as_str()), - ), + source_sha256: record_source_hash, + source_size: record_source_size, + source_mime: record_source_mime, sanitization_policy: is_media.then_some(1), tool_versions: is_media.then(crate::sanitize::tool_versions).flatten(), uploaded_at: chrono::Utc::now().timestamp(), @@ -244,6 +243,8 @@ pub async fn process_streaming_ingest( }; if let Some(attribution) = &attribution { + let (record_source_hash, record_source_size, record_source_mime) = + source_record_fields(source_probe.as_ref(), &source_hash, source_size); record_upload_event( storage, ctx, @@ -254,13 +255,9 @@ pub async fn process_streaming_ingest( ext: &ext, mime: &mime, size: output_size, - source_sha256: Some(source_hash.as_str()), - source_size: Some(source_size), - source_mime: Some( - source_probe - .as_ref() - .map_or(mime.as_str(), |probe| probe.mime.as_str()), - ), + source_sha256: record_source_hash, + source_size: record_source_size, + source_mime: record_source_mime, sanitization_policy: is_media.then_some(1), tool_versions: is_media.then(crate::sanitize::tool_versions).flatten(), uploaded_at, @@ -293,6 +290,46 @@ pub async fn process_streaming_ingest( )) } +async fn verify_source_auth( + mode: UploadRouteMode, + auth_event: &nostr::Event, + source_hash: &str, + bound_host: &str, + max_age_secs: u64, +) -> Result<(), MediaError> { + let auth = auth_event.clone(); + let auth_hash = source_hash.to_string(); + let bound_host = bound_host.to_string(); + tokio::task::spawn_blocking(move || match mode { + UploadRouteMode::Media => { + verify_blossom_media_auth(&auth, &auth_hash, Some(&bound_host), max_age_secs) + } + UploadRouteMode::Upload | UploadRouteMode::Legacy => { + verify_blossom_upload_auth(&auth, &auth_hash, Some(&bound_host), max_age_secs) + } + }) + .await + .map_err(|_| MediaError::Internal)? +} + +fn auth_max_age_secs(source_probe: Option<&crate::sanitize::MediaProbe>) -> u64 { + match source_probe.map(|probe| probe.class) { + Some(crate::sanitize::MediaClass::Video) => 3_600, + Some(crate::sanitize::MediaClass::Image | crate::sanitize::MediaClass::Audio) | None => 600, + } +} + +fn source_record_fields<'a>( + source_probe: Option<&'a crate::sanitize::MediaProbe>, + source_hash: &'a str, + source_size: u64, +) -> (Option<&'a str>, Option, Option<&'a str>) { + match source_probe { + Some(probe) => (Some(source_hash), Some(source_size), Some(&probe.mime)), + None => (None, None, None), + } +} + fn enforce_route_policy( mode: UploadRouteMode, source_probe: Option<&crate::sanitize::MediaProbe>, @@ -1160,6 +1197,31 @@ mod tests { )); } + #[test] + fn authorization_freshness_is_class_specific() { + let image = image_probe(); + let mut video = image.clone(); + video.class = crate::sanitize::MediaClass::Video; + let mut audio = image.clone(); + audio.class = crate::sanitize::MediaClass::Audio; + + assert_eq!(auth_max_age_secs(Some(&image)), 600); + assert_eq!(auth_max_age_secs(Some(&audio)), 600); + assert_eq!(auth_max_age_secs(None), 600); + assert_eq!(auth_max_age_secs(Some(&video)), 3_600); + } + + #[test] + fn exact_upload_records_omit_transformation_fields() { + assert_eq!(source_record_fields(None, "source", 42), (None, None, None)); + + let image = image_probe(); + assert_eq!( + source_record_fields(Some(&image), "source", 42), + (Some("source"), Some(42), Some("image/jpeg")) + ); + } + #[test] fn test_build_descriptor_no_meta() { // When meta is None, all optional fields should be None. diff --git a/crates/buzz-relay/src/api/media.rs b/crates/buzz-relay/src/api/media.rs index 60a9223e2..36656efaf 100644 --- a/crates/buzz-relay/src/api/media.rs +++ b/crates/buzz-relay/src/api/media.rs @@ -320,6 +320,7 @@ pub async fn upload_blob( config: &state.config.media, ctx: &auth.tenant, auth_event: &auth.auth_event, + claimed_source_hash: &auth.claimed_hash, content_length, attribution, mode: auth.mode, @@ -350,7 +351,8 @@ pub async fn upload_blob( // Audit via bounded channel — same pattern as event audit. let desc = descriptor.clone(); - let sanitized = desc.sha256 != auth.claimed_hash; + let sanitization_applied = route_applies_sanitization(auth.mode); + let transformed = desc.sha256 != auth.claimed_hash; if let Err(e) = state .audit_tx .send(NewAuditEntry { @@ -361,12 +363,15 @@ pub async fn upload_blob( detail: serde_json::json!({ "sha256": desc.sha256, "source_sha256": auth.claimed_hash, - "sanitized": sanitized, - "sanitization_policy": sanitized.then_some(1), + "sanitized": sanitization_applied, + "transformed": transformed, + "sanitization_policy": sanitization_applied.then_some(1), "outcome": "accepted", "output_size": desc.size, "output_mime": desc.mime_type, - "tool_versions": buzz_media::sanitize::tool_versions(), + "tool_versions": sanitization_applied + .then(buzz_media::sanitize::tool_versions) + .flatten(), }), }) .await @@ -378,6 +383,13 @@ pub async fn upload_blob( Ok(Json(descriptor)) } +fn route_applies_sanitization(mode: buzz_media::UploadRouteMode) -> bool { + matches!( + mode, + buzz_media::UploadRouteMode::Media | buzz_media::UploadRouteMode::Legacy + ) +} + pub(crate) fn media_base_url_for_tenant(config_relay_url: &str, tenant_host: &str) -> String { let scheme = if config_relay_url.trim_start().starts_with("wss://") || config_relay_url.trim_start().starts_with("https://") @@ -1183,6 +1195,19 @@ mod tests { ); } + #[test] + fn audit_records_sanitizer_application_from_route_policy() { + assert!(route_applies_sanitization( + buzz_media::UploadRouteMode::Media + )); + assert!(route_applies_sanitization( + buzz_media::UploadRouteMode::Legacy + )); + assert!(!route_applies_sanitization( + buzz_media::UploadRouteMode::Upload + )); + } + #[test] fn rewrite_descriptor_urls_for_tenant_replaces_global_media_host() { let hash = "a".repeat(64); diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 70bd3ec5b..d9f390dd7 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -34,7 +34,9 @@ pub fn build_router(state: Arc) -> Router { .config .media .max_image_bytes - .max(state.config.media.max_video_bytes) as usize; + .max(state.config.media.max_video_bytes) + .max(state.config.media.max_audio_bytes) + .max(state.config.media.max_file_bytes) as usize; let media_router = Router::new() .route("/media", put(api::media::upload_blob)) .route("/upload", put(api::media::upload_blob))