fix(media): correct upload policy accounting

Co-authored-by: Codex <noreply@openai.com>
This commit is contained in:
Atish Patel
2026-07-16 17:35:02 -05:00
co-authored by Codex
parent fe6477509c
commit 19c5c6d9df
3 changed files with 121 additions and 32 deletions
+89 -27
View File
@@ -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<u64>,
/// 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<u64>, 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.
+29 -4
View File
@@ -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);
+3 -1
View File
@@ -34,7 +34,9 @@ pub fn build_router(state: Arc<AppState>) -> 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))