Compare commits

...
Author SHA1 Message Date
Hampus 7d2c51a57f fix(media-proxy): budget native image work by pod memory (#3385) 2026-10-11 02:07:19 +02:00
20 changed files with 663 additions and 133 deletions

No files matched your search

+4
View File
@@ -14,6 +14,10 @@ pub const MAX_VIDEO_PACKETS_FOR_THUMBNAIL: usize = 512;
pub const DEFAULT_IMAGE_SIZE: u32 = 128;
pub const MAX_ANIMATED_FRAMES_DEFAULT: u32 = 20_000;
pub const MAX_ANIMATED_TOTAL_PIXELS_DEFAULT: usize = 4 * MAX_MEDIA_IMAGE_PIXELS_DEFAULT;
pub const MAX_JXL_DECODE_PIXELS: usize = 4096 * 4096;
pub const MAX_STATIC_OUTPUT_PIXELS: usize = 6144 * 4096;
pub const SVG_RASTER_MAX_DIMENSION: u32 = 4096;
pub const LARGE_SOURCE_LOG_PIXELS: usize = 4096 * 4096;
const _: () = assert!(MAX_ANIMATED_FRAMES_DEFAULT >= 20_000);
pub const IMAGE_SIZES: &[u32] = &[
@@ -1,8 +1,9 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use super::native_runtime::clear_vips_error;
use super::{MediaError, native_status_error};
use super::{ImageDimensions, MediaError, native_status_error};
use crate::{
constants,
image_transform::{ImageOptions, ResizeMode},
media_limits::MediaLimits,
native,
@@ -55,6 +56,51 @@ pub(crate) fn validate_dimensions_u32(
Ok(())
}
pub(crate) fn validate_jxl_decode_pixels(
sniffed_mime: &str,
dims: ImageDimensions,
) -> Result<(), MediaError> {
if sniffed_mime != "image/jxl" {
return Ok(());
}
let pixels = (dims.width as usize)
.saturating_mul(dims.height as usize)
.saturating_mul(dims.pages.max(1) as usize);
if pixels > constants::MAX_JXL_DECODE_PIXELS {
return Err(MediaError::InvalidImageDimensions);
}
Ok(())
}
pub(crate) fn static_output_pixels(options: &ImageOptions, dims: ImageDimensions) -> usize {
let source_width = dims.width as f64;
let source_height = dims.height as f64;
let source_pixels = dims.width as usize * dims.height as usize;
if options.wants_cover_crop()
&& let (Some(width), Some(height)) = (options.width, options.height)
{
return (width as usize * height as usize).min(source_pixels);
}
let width_scale = options
.width
.map_or(1.0, |width| width as f64 / source_width);
let height_scale = options
.height
.map_or(1.0, |height| height as f64 / source_height);
let scale = width_scale.min(height_scale).min(1.0);
((source_width * scale).round() * (source_height * scale).round()) as usize
}
pub(crate) fn validate_static_output_pixels(
options: &ImageOptions,
dims: ImageDimensions,
) -> Result<(), MediaError> {
if static_output_pixels(options, dims) > constants::MAX_STATIC_OUTPUT_PIXELS {
return Err(MediaError::InvalidImageDimensions);
}
Ok(())
}
pub(crate) fn validate_dimensions(
media_limits: &MediaLimits,
width: c_int,
@@ -3,7 +3,7 @@
use super::MediaError;
use super::av_metadata::{AVMetadata, AVProbe, NSFW_PREVIEW_MAX_DIMENSION, probe_av_metadata};
use super::image_probe::probe_image_dims;
use super::loaded_image::validate_dimensions_u32;
use super::loaded_image::{validate_dimensions_u32, validate_jxl_decode_pixels};
use super::nsfw_processing::{
NSFWScanPreparation, NSFWScanSource, classify_nsfw_buffers, nsfw_scan_buffers,
};
@@ -22,6 +22,7 @@ use sha2::{Digest, Sha256};
pub struct MetadataOptions {
pub placeholder: bool,
pub nsfw: NSFWPolicy,
pub deadline_ms: Option<i64>,
}
impl Default for MetadataOptions {
@@ -29,6 +30,7 @@ impl Default for MetadataOptions {
Self {
placeholder: true,
nsfw: NSFWPolicy::Disabled,
deadline_ms: None,
}
}
}
@@ -50,12 +52,12 @@ struct MetadataResponse {
nsfw_probability: f32,
}
struct MetadataBlocking {
pub struct PreparedMetadata {
response: MetadataResponse,
nsfw_scan: Option<NSFWScanPreparation>,
}
impl MetadataBlocking {
impl PreparedMetadata {
fn take_nsfw_scan(&mut self) -> Option<NSFWScanPreparation> {
self.nsfw_scan.take()
}
@@ -128,7 +130,7 @@ fn probe_av_metadata_without_requiring_a_frame(
}
}
fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlocking, MediaError> {
fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<PreparedMetadata, MediaError> {
let MetadataBlockingRequest {
media_limits,
metrics,
@@ -145,7 +147,9 @@ fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlo
let initial_category = mime::category(sniffed.mime).ok_or(MediaError::UnsupportedMediaType)?;
let is_image = initial_category == mime::Category::Image;
let dims = if is_image {
Some(probe_image_dims(media_limits, input)?)
let dims = probe_image_dims(media_limits, input)?;
validate_jxl_decode_pixels(sniffed.mime, dims)?;
Some(dims)
} else {
None
};
@@ -190,7 +194,7 @@ fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlo
let placeholder = if options.placeholder {
let hash = if is_image {
optional_thumbhash(
encode_thumbhash(media_limits, input, None),
encode_thumbhash(media_limits, input, options.deadline_ms),
"image_metadata",
)
} else {
@@ -219,7 +223,7 @@ fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlo
frame_count: frames_count.max(sniffed.frames),
input,
duration_seconds: av_probe.as_ref().and_then(|probe| probe.duration_seconds),
deadline_ms: None,
deadline_ms: options.deadline_ms,
})?
} else {
None
@@ -239,7 +243,7 @@ fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlo
let format = metadata_format(content_type);
let content_hash = hex::encode(Sha256::digest(input));
Ok(MetadataBlocking {
Ok(PreparedMetadata {
response: MetadataResponse {
content_type: content_type.to_owned(),
size: input.len(),
@@ -257,13 +261,37 @@ fn metadata_blocking(request: MetadataBlockingRequest<'_>) -> Result<MetadataBlo
})
}
fn metadata_finalize(prepared: MetadataBlocking, verdict: NSFWClassification) -> MetadataResponse {
fn metadata_finalize(prepared: PreparedMetadata, verdict: NSFWClassification) -> MetadataResponse {
let mut response = prepared.response;
response.nsfw = verdict.is_nsfw;
response.nsfw_probability = verdict.probability;
response
}
pub fn prepare_metadata(
input: &[u8],
options: &MetadataOptions,
media_limits: &MediaLimits,
metrics: &TransformMetrics,
) -> Result<PreparedMetadata, MediaError> {
metadata_blocking(MetadataBlockingRequest {
media_limits,
metrics,
input,
options,
})
}
pub async fn finish_metadata(
mut prepared: PreparedMetadata,
nsfw_client: &NSFWClient,
) -> Result<String, MediaError> {
let scan = prepared.take_nsfw_scan();
let verdict = classify_nsfw_buffers(nsfw_client, scan).await?;
let response = metadata_finalize(prepared, verdict);
serde_json::to_string(&response).map_err(|_| MediaError::MediaEncodeFailed)
}
pub async fn metadata_json_with_options(
input: &[u8],
_filename: &str,
@@ -272,14 +300,6 @@ pub async fn metadata_json_with_options(
nsfw_client: &NSFWClient,
metrics: &TransformMetrics,
) -> Result<String, MediaError> {
let mut prepared = metadata_blocking(MetadataBlockingRequest {
media_limits,
metrics,
input,
options: &options,
})?;
let scan = prepared.take_nsfw_scan();
let verdict = classify_nsfw_buffers(nsfw_client, scan).await?;
let response = metadata_finalize(prepared, verdict);
serde_json::to_string(&response).map_err(|_| MediaError::MediaEncodeFailed)
let prepared = prepare_metadata(input, &options, media_limits, metrics)?;
finish_metadata(prepared, nsfw_client).await
}
+6 -1
View File
@@ -18,6 +18,7 @@ mod encoding;
mod image_probe;
mod loaded_image;
mod metadata;
mod native_cost;
pub(crate) mod native_runtime;
mod nsfw_processing;
mod placeholder;
@@ -32,7 +33,11 @@ mod tests;
pub use av_metadata::{
AVMetadata, AVMetadataFrame, AVProbe, NSFW_PREVIEW_MAX_DIMENSION, probe_av_metadata,
};
pub use metadata::{MetadataOptions, metadata_json_with_options};
pub use metadata::{
MetadataOptions, PreparedMetadata, finish_metadata, metadata_json_with_options,
prepare_metadata,
};
pub use native_cost::{AV_NATIVE_COST_BYTES, image_transform_cost, metadata_cost};
pub use nsfw_processing::encode_static_image_for_nsfw;
pub use transform::transform_image;
pub use video_thumbnail::{
@@ -0,0 +1,72 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use super::MediaError;
use super::image_probe::probe_image_dims;
use super::loaded_image::static_output_pixels;
use super::native_runtime::ensure_vips_init;
use crate::{
image_quality::ImageQuality, image_transform::ImageOptions, media_limits::MediaLimits, mime,
output_format::OutputFormat,
};
pub const AV_NATIVE_COST_BYTES: usize = 256 << 20;
const ANIMATED_CANVAS_BUFFERS: usize = 8;
const RGBA_BYTES_PER_PIXEL: usize = 4;
fn decode_bytes_per_pixel(sniffed_mime: &str) -> usize {
match sniffed_mime {
"image/jxl" => 40,
"image/webp" => 20,
"image/avif" | "image/heic" | "image/heif" => 12,
"image/jpeg" => 2,
"image/svg+xml" => 0,
_ => 8,
}
}
fn encode_bytes_per_pixel(format: OutputFormat, quality: ImageQuality) -> usize {
match (format, quality) {
(OutputFormat::WebP, ImageQuality::Lossless | ImageQuality::Auto) => 48,
(OutputFormat::WebP, _) => 28,
(OutputFormat::JPEG, _) => 4,
_ => 8,
}
}
pub fn image_transform_cost(
input: &[u8],
options: &ImageOptions,
media_limits: &MediaLimits,
) -> Result<usize, MediaError> {
ensure_vips_init()?;
let sniffed_mime = mime::sniff(input).mime;
let dims = probe_image_dims(media_limits, input)?;
let source_pixels = dims.width as usize * dims.height as usize;
let output_pixels = static_output_pixels(options, dims);
let encode = encode_bytes_per_pixel(options.format, options.quality);
let pixel_cost = if options.is_animated() && dims.pages > 1 {
source_pixels.max(output_pixels) * RGBA_BYTES_PER_PIXEL * ANIMATED_CANVAS_BUFFERS
+ output_pixels * encode
} else {
source_pixels * decode_bytes_per_pixel(sniffed_mime) + output_pixels * encode
};
Ok(input.len().saturating_add(pixel_cost))
}
pub fn metadata_cost(input: &[u8], media_limits: &MediaLimits) -> Result<usize, MediaError> {
let sniffed_mime = mime::sniff(input).mime;
match mime::category(sniffed_mime) {
Some(mime::Category::Image) => {
ensure_vips_init()?;
let dims = probe_image_dims(media_limits, input)?;
let source_pixels = dims.width as usize * dims.height as usize;
Ok(input
.len()
.saturating_add(source_pixels * decode_bytes_per_pixel(sniffed_mime)))
}
Some(mime::Category::Video | mime::Category::Audio) => {
Ok(input.len().saturating_add(AV_NATIVE_COST_BYTES))
}
_ => Ok(input.len()),
}
}
@@ -64,6 +64,7 @@ fn metadata_json_returns_unavailable_when_nsfw_service_fails() {
MetadataOptions {
placeholder: false,
nsfw: NSFWPolicy::enabled(0.85).expect("valid threshold"),
deadline_ms: None,
},
&test_media_limits(),
&client,
@@ -237,3 +238,23 @@ fn metadata_succeeds_without_a_placeholder_when_thumbhash_generation_fails() {
assert_eq!(64, value["height"]);
assert_eq!(None, value.get("placeholder"));
}
#[test]
fn metadata_refuses_jxl_above_the_decode_pixel_cap() {
let err = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap()
.block_on(async {
metadata_json_with_options(
crate::test_fixtures::JXL_4100_LOSSY_ALPHA,
"big.jxl",
MetadataOptions::default(),
&test_media_limits(),
&NSFWClient::disabled(),
&TransformMetrics::new(),
)
.await
.unwrap_err()
});
assert_eq!(MediaError::InvalidImageDimensions, err);
}
@@ -1,14 +1,16 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use super::super::ImageOptions;
use super::super::image_probe::load_image;
use super::super::native_runtime::ensure_vips_init;
use super::super::transform::source_supports_pages;
use super::super::{ImageOptions, MediaError, image_transform_cost};
use super::fixtures::{animated_mode, metadata_value, transform_image};
use crate::{
mime, native,
output_format::OutputFormat,
test_fixtures::{synthetic_bmp, synthetic_png},
test_fixtures::{
JXL_64_LOSSY_ALPHA, JXL_4100_LOSSY_ALPHA, png_dimensions, synthetic_bmp, synthetic_png,
},
};
#[test]
@@ -103,3 +105,105 @@ fn source_supports_pages_matches_libvips_loader_list() {
assert!(!source_supports_pages("image/bmp"));
assert!(!source_supports_pages("application/octet-stream"));
}
#[test]
fn refuses_jxl_above_the_decode_pixel_cap_before_decoding() {
for width in [None, Some(128)] {
let err = transform_image(
JXL_4100_LOSSY_ALPHA,
&ImageOptions {
width,
format: OutputFormat::WebP,
..Default::default()
},
)
.unwrap_err();
assert_eq!(MediaError::InvalidImageDimensions, err);
}
}
#[test]
fn transforms_jxl_within_the_decode_pixel_cap() {
let out = transform_image(
JXL_64_LOSSY_ALPHA,
&ImageOptions {
width: Some(32),
format: OutputFormat::WebP,
..Default::default()
},
)
.unwrap();
assert_eq!("image/webp", out.content_type);
}
#[test]
fn refuses_full_size_svg_raster_above_the_static_output_cap() {
let svg = br#"<svg xmlns="http://www.w3.org/2000/svg" width="16383" height="16383"></svg>"#;
let err = transform_image(
svg,
&ImageOptions {
format: OutputFormat::WebP,
..Default::default()
},
)
.unwrap_err();
assert_eq!(MediaError::InvalidImageDimensions, err);
let err = transform_image(
svg,
&ImageOptions {
width: Some(16384),
format: OutputFormat::WebP,
..Default::default()
},
)
.unwrap_err();
assert_eq!(MediaError::InvalidImageDimensions, err);
}
#[test]
fn rasterizes_large_svg_into_a_bounded_box() {
let svg = br#"<svg xmlns="http://www.w3.org/2000/svg" width="16383" height="16383"></svg>"#;
let out = transform_image(
svg,
&ImageOptions {
width: Some(64),
height: Some(64),
format: OutputFormat::PNG,
..Default::default()
},
)
.unwrap();
assert_eq!(Some((64, 64)), png_dimensions(&out.bytes));
}
#[test]
fn transform_cost_charges_the_full_source_decode_for_jxl() {
let cost = image_transform_cost(
JXL_4100_LOSSY_ALPHA,
&ImageOptions {
width: Some(128),
format: OutputFormat::WebP,
..Default::default()
},
&super::fixtures::test_media_limits(),
)
.unwrap();
assert!(cost >= 4100 * 4100 * 40, "cost {cost}");
}
#[test]
fn transform_cost_charges_svg_at_its_output_size() {
let svg = br#"<svg xmlns="http://www.w3.org/2000/svg" width="16383" height="16383"></svg>"#;
let cost = image_transform_cost(
svg,
&ImageOptions {
width: Some(64),
height: Some(64),
format: OutputFormat::WebP,
..Default::default()
},
&super::fixtures::test_media_limits(),
)
.unwrap();
assert!(cost < 1 << 20, "cost {cost}");
}
@@ -8,10 +8,13 @@ use super::encoding::{
VipsEncodeRequest, anim_limits_from_options, encode_vips_image,
try_transform_animated_webp_direct,
};
use super::image_probe::{animated_probe_from_image, load_image, probe_animated, try_decode_bmp};
use super::image_probe::{
animated_probe_from_image, load_image, probe_animated, probe_image_dims, try_decode_bmp,
};
use super::loaded_image::{
normalize_vips_image_to_uchar, page_height, resize_loaded_image, resize_loaded_image_by_scale,
validate_dimensions_u32, validate_vips_image,
validate_dimensions_u32, validate_jxl_decode_pixels, validate_static_output_pixels,
validate_vips_image,
};
use super::native_runtime::{clear_vips_error, ensure_vips_init, last_vips_error};
use super::transform_plan::{
@@ -77,6 +80,54 @@ fn heif_source_may_be_a_sequence(sniffed: mime::SniffInfo) -> bool {
}
}
fn ensure_decode_cost_admissible(
input: &[u8],
sniffed_mime: &str,
options: &ImageOptions,
media_limits: &MediaLimits,
) -> Result<(), MediaError> {
let animated = options.is_animated();
if animated && sniffed_mime != "image/jxl" {
return Ok(());
}
let dims = probe_image_dims(media_limits, input)?;
let source_pixels = dims.width as usize * dims.height as usize;
if source_pixels >= constants::LARGE_SOURCE_LOG_PIXELS {
tracing::info!(
mime = sniffed_mime,
input_bytes = input.len(),
width = dims.width,
height = dims.height,
pages = dims.pages,
out_width = options.width.unwrap_or(0),
out_height = options.height.unwrap_or(0),
animated,
"large native transform starting"
);
}
let admitted = validate_jxl_decode_pixels(sniffed_mime, dims).and_then(|()| {
if animated {
Ok(())
} else {
validate_static_output_pixels(options, dims)
}
});
if admitted.is_err() {
tracing::warn!(
mime = sniffed_mime,
input_bytes = input.len(),
width = dims.width,
height = dims.height,
pages = dims.pages,
out_width = options.width.unwrap_or(0),
out_height = options.height.unwrap_or(0),
animated,
"native transform refused by decode cost limits"
);
}
admitted
}
pub fn transform_image(
input: &[u8],
options: &ImageOptions,
@@ -99,6 +150,7 @@ pub fn transform_image(
ensure_vips_init()?;
let sniffed = mime::sniff(input);
let animated = options.is_animated();
ensure_decode_cost_admissible(input, sniffed.mime, options, media_limits)?;
let format = effective_transform_format(sniffed.mime, options.format, animated);
let animated_avif = sniffed.mime == "image/avif" && sniffed.animated;
let full_canvas_animation = animated
+5 -2
View File
@@ -23,7 +23,7 @@ use fluxer_common::attachment_url_signature::is_signature_parameter_name;
use rand::RngExt;
use stage::{StageTimingSnapshot, StageTimings};
use std::{future::Future, sync::Arc, time::Instant};
use tracing::{Level, event};
use tracing::{Instrument, Level, event};
const ID_ALPHABET: &[u8] = b"0123456789ABCDEFGHJKMNPQRSTVWXYZ";
const ID_LEN: usize = 12;
@@ -179,7 +179,10 @@ where
} = observation;
let stages = Arc::new(StageTimings::default());
let started = Instant::now();
let response = stage::scope(Arc::clone(&stages), future).await;
let span = tracing::info_span!("request", req = %id.as_str(), path = %path);
let response = stage::scope(Arc::clone(&stages), future)
.instrument(span)
.await;
let elapsed_ms = metrics::duration_millis(started.elapsed());
let StageTimingSnapshot {
fetch_ms,
@@ -7,7 +7,7 @@ pub(in crate::server) use failure::MediaFailure;
pub(in crate::server) use input::{MediaInput, MediaInputLimit, load_media_input};
use crate::{
constants::AssetExtension,
constants::{self, AssetExtension},
image_quality::ImageQuality,
image_transform::AnimationMode,
media_process, mime,
@@ -16,7 +16,7 @@ use crate::{
server::{
format_policy::image_extension_from_filename,
state::AppState,
transform::execution::{run_transform, transform_error_is_timeout},
transform::execution::{deadline_instant, run_transform, transform_error_is_timeout},
},
};
use bytes::Bytes;
@@ -43,21 +43,39 @@ pub(in crate::server) async fn resolve_metadata(
app.media.nsfw().record_declined_scan();
NSFWPolicy::Disabled
};
let json = media_process::metadata_json_with_options(
&input.data,
&input.filename,
media_process::MetadataOptions {
placeholder: true,
nsfw,
},
&app.media.limits(),
app.media.nsfw(),
&app.metrics.transform(),
)
.await
.map_err(|err| MediaFailure::MetadataExtractionFailed {
let runtime = app.media.transforms();
let options = media_process::MetadataOptions {
placeholder: true,
nsfw,
deadline_ms: runtime.transform_deadline_ms(),
};
let deadline = deadline_instant(options.deadline_ms);
let media_limits = app.media.limits();
let transform_metrics = app.metrics.transform();
let data = input.data.clone();
let cost_data = input.data.clone();
let extraction_failed = |err: &dyn std::fmt::Debug| MediaFailure::MetadataExtractionFailed {
detail: format!("filename={} err={err:?}", input.filename),
})?;
};
let prepared = runtime
.tasks()
.run_native(
deadline,
move || Ok(media_process::metadata_cost(&cost_data, &media_limits)?),
move || {
Ok(media_process::prepare_metadata(
&data,
&options,
&media_limits,
&transform_metrics,
)?)
},
)
.await
.map_err(|err| extraction_failed(&err))?;
let json = media_process::finish_metadata(prepared, app.media.nsfw())
.await
.map_err(|err| extraction_failed(&err))?;
Ok(MetadataOutput {
metadata: serde_json::from_str(&json).unwrap_or_else(|_| serde_json::json!({})),
data: include_data.then_some(input.data),
@@ -74,6 +92,8 @@ async fn rasterize_metadata_svg(
input: LoadedMediaInput,
) -> Result<LoadedMediaInput, MediaFailure> {
let options = media_process::ImageOptions {
width: Some(constants::SVG_RASTER_MAX_DIMENSION),
height: Some(constants::SVG_RASTER_MAX_DIMENSION),
format: OutputFormat::WebP,
quality: ImageQuality::Lossless,
animation: AnimationMode::Static,
@@ -18,20 +18,20 @@ use std::{
time::Instant,
};
use tokio::sync::{Notify, oneshot};
use tracing::error;
use tracing::{error, warn};
pub(in crate::server) struct NativeTaskExecutorSettings {
pub(in crate::server) max_native_transforms: usize,
pub(in crate::server) worker_queue_capacity: usize,
pub(in crate::server) decoded_bytes_per_transform: usize,
pub(in crate::server) memory_budget_bytes: usize,
pub(in crate::server) native_metrics: Arc<NativeTransformMetrics>,
pub(in crate::server) transform_metrics: Arc<TransformMetrics>,
}
pub(in crate::server) struct NativeTaskExecutor {
native_transforms: TimedSemaphore,
decoded_bytes: ByteBudget,
decoded_bytes_per_transform: usize,
memory: ByteBudget,
memory_budget_bytes: usize,
tasks: NativeTaskTracker,
native_metrics: Arc<NativeTransformMetrics>,
transform_metrics: Arc<TransformMetrics>,
@@ -70,34 +70,33 @@ impl NativeTaskExecutor {
let NativeTaskExecutorSettings {
max_native_transforms,
worker_queue_capacity,
decoded_bytes_per_transform,
memory_budget_bytes,
native_metrics,
transform_metrics,
} = settings;
assert!(decoded_bytes_per_transform > 0);
let decoded_bytes_capacity = decoded_bytes_per_transform
.checked_mul(max_native_transforms)
.expect("native decoded image budget must not overflow");
assert!(memory_budget_bytes > 0);
Self {
native_transforms: TimedSemaphore::with_queue_capacity(
max_native_transforms,
worker_queue_capacity,
),
decoded_bytes: ByteBudget::new(decoded_bytes_capacity),
decoded_bytes_per_transform,
memory: ByteBudget::new(memory_budget_bytes),
memory_budget_bytes,
tasks: NativeTaskTracker::default(),
native_metrics,
transform_metrics,
}
}
pub(in crate::server) async fn run_native<T, F>(
pub(in crate::server) async fn run_native<T, C, F>(
&self,
deadline: Option<Instant>,
cost: C,
work: F,
) -> anyhow::Result<T>
where
T: Send + 'static,
C: FnOnce() -> anyhow::Result<usize> + Send + 'static,
F: FnOnce() -> anyhow::Result<T> + Send + 'static,
{
let admission = match self.native_transforms.try_admit() {
@@ -116,13 +115,24 @@ impl NativeTaskExecutor {
self.native_metrics
.observe_wait(metrics::duration_millis(wait_started.elapsed()));
let permit = permit?;
let decoded_bytes = self
.decoded_bytes
.try_reserve(self.decoded_bytes_per_transform)
.ok_or(MediaError::AllocationFailed)?;
let memory = self.memory.clone();
let capacity = self.memory_budget_bytes;
let span = tracing::Span::current();
self.run_task(deadline, move || {
let _span = span.enter();
let _permit = permit;
let _decoded_bytes = decoded_bytes;
let bytes = cost()?;
if bytes > capacity {
warn!(
bytes,
capacity, "native work refused, larger than the memory budget"
);
return Err(MediaError::InvalidImageDimensions.into());
}
let Some(_reserved) = memory.try_reserve(bytes) else {
warn!(bytes, capacity, "native work refused, memory budget in use");
return Err(MediaError::AllocationFailed.into());
};
work()
})
.await
@@ -266,7 +276,7 @@ mod tests {
NativeTaskExecutor::new(NativeTaskExecutorSettings {
max_native_transforms: permits,
worker_queue_capacity: queue,
decoded_bytes_per_transform: 1024,
memory_budget_bytes: 1024,
native_metrics: metrics.native_transform(),
transform_metrics: metrics.transform(),
})
@@ -288,6 +298,7 @@ mod tests {
let error = executor
.run_native(
Some(Instant::now() + Duration::from_millis(30)),
|| Ok(0),
move || {
released.blocking_recv();
let _ = finished.blocking_send(());
@@ -338,7 +349,11 @@ mod tests {
let metrics = metrics::Metrics::new();
let executor = executor(&metrics, 1, 0);
let value = executor
.run_native(Some(Instant::now() + Duration::from_secs(30)), || Ok(7u32))
.run_native(
Some(Instant::now() + Duration::from_secs(30)),
|| Ok(0),
|| Ok(7u32),
)
.await
.expect("the task completes");
assert_eq!(7, value);
@@ -361,15 +376,19 @@ mod tests {
let holder = Arc::clone(&executor);
let held = tokio::spawn(async move {
holder
.run_native(None, move || {
released.blocking_recv();
Ok(())
})
.run_native(
None,
|| Ok(0),
move || {
released.blocking_recv();
Ok(())
},
)
.await
});
tokio::time::sleep(Duration::from_millis(50)).await;
let error = executor
.run_native(None, || Ok(()))
.run_native(None, || Ok(0), || Ok(()))
.await
.expect_err("the admission queue is full");
assert_eq!(
@@ -393,7 +412,7 @@ mod tests {
let executor = executor(&metrics, 2, 4);
executor.begin_shutdown();
let error = executor
.run_native(None, || Ok(()))
.run_native(None, || Ok(0), || Ok(()))
.await
.expect_err("a closed executor admits nothing");
assert_eq!(
@@ -402,4 +421,59 @@ mod tests {
);
executor.wait_for_shutdown().await;
}
#[tokio::test]
async fn work_larger_than_the_memory_budget_is_refused_before_it_runs() {
let metrics = metrics::Metrics::new();
let executor = executor(&metrics, 1, 0);
let error = executor
.run_native(
None,
|| Ok(1025),
|| -> anyhow::Result<()> { panic!("refused work must not run") },
)
.await
.expect_err("the job exceeds the whole budget");
assert_eq!(
Some(&MediaError::InvalidImageDimensions),
error.downcast_ref::<MediaError>()
);
}
#[tokio::test]
async fn work_that_does_not_fit_beside_a_running_job_fails_fast_and_the_budget_frees() {
let metrics = metrics::Metrics::new();
let executor = Arc::new(executor(&metrics, 2, 0));
let (release, mut released) = mpsc::channel::<()>(1);
let (started, mut has_started) = mpsc::channel::<()>(1);
let holder = Arc::clone(&executor);
let held = tokio::spawn(async move {
holder
.run_native(
None,
|| Ok(1000),
move || {
let _ = started.blocking_send(());
released.blocking_recv();
Ok(())
},
)
.await
});
has_started.recv().await;
let error = executor
.run_native(None, || Ok(100), || Ok(()))
.await
.expect_err("the budget is in use");
assert_eq!(
Some(&MediaError::AllocationFailed),
error.downcast_ref::<MediaError>()
);
drop(release);
held.await.expect("held task").expect("held work");
executor
.run_native(None, || Ok(1024), || Ok(()))
.await
.expect("the released budget admits a job of the full size");
}
}
@@ -3,6 +3,7 @@
use crate::{
byte_budget::BudgetedBytes,
constants,
image_quality::ImageQuality,
image_transform::ResizeMode,
media_process, mime,
output_format::OutputFormat,
@@ -20,7 +21,9 @@ use crate::{
media_response,
},
state::AppState,
transform::execution::{run_transform, transform_error_is_timeout},
transform::execution::{
VideoTransformOptions, run_transform, run_video_transform, transform_error_is_timeout,
},
},
};
use axum::{
@@ -176,12 +179,23 @@ pub(in crate::server) async fn thumbnail_handler(
Err(err) => return storage_error_response(&req.upload_filename, err),
};
let media = if mime::category(&object.content_type) == Some(mime::Category::Video) {
match media_process::extract_video_thumbnail(
&object.data,
OutputFormat::WebP,
&app.media.limits(),
) {
let options = VideoTransformOptions {
format: OutputFormat::WebP,
width: None,
height: None,
quality: ImageQuality::Auto,
deadline_ms: app.media.transforms().transform_deadline_ms(),
};
match run_video_transform(app.media.transforms(), object.data.clone(), options).await {
Ok(media) => media,
Err(err) if transform_error_is_timeout(&err) => {
return text_with_source(
StatusCode::GATEWAY_TIMEOUT,
"Gateway Timeout",
"video_thumbnail_timeout",
err,
);
}
Err(err) => {
return text_with_source(
StatusCode::BAD_REQUEST,
@@ -253,11 +267,14 @@ pub(in crate::server) async fn frames_handler(
Ok(input) => input,
Err(failure) => return failure.into_response(),
};
match media_process::extract_video_thumbnail(
&input.data,
OutputFormat::JPEG,
&app.media.limits(),
) {
let options = VideoTransformOptions {
format: OutputFormat::JPEG,
width: None,
height: None,
quality: ImageQuality::Auto,
deadline_ms: app.media.transforms().transform_deadline_ms(),
};
match run_video_transform(app.media.transforms(), input.data, options).await {
Ok(frame) => {
let encoded = general_purpose::STANDARD.encode(frame.bytes);
json_response(
@@ -683,11 +700,15 @@ mod tests {
.media
.transforms()
.tasks()
.run_native(None, move || {
let _ = started.blocking_send(());
let _ = released.blocking_recv();
Ok(())
})
.run_native(
None,
|| Ok(0),
move || {
let _ = started.blocking_send(());
let _ = released.blocking_recv();
Ok(())
},
)
.await
});
has_started
@@ -272,6 +272,8 @@ async fn serve_stored_svg_rasterized(
let format = OutputFormat::WebP;
let quality = ImageQuality::Lossless;
let options = media_process::ImageOptions {
width: Some(constants::SVG_RASTER_MAX_DIMENSION),
height: Some(constants::SVG_RASTER_MAX_DIMENSION),
format,
quality,
animation: AnimationMode::Static,
@@ -282,8 +284,8 @@ async fn serve_stored_svg_rasterized(
route: TransformRoute::Stored,
asset_kind: None,
cache_identity,
width: None,
height: None,
width: options.width,
height: options.height,
format,
quality: Some(quality),
animated: false,
@@ -35,17 +35,32 @@ pub(in crate::server) async fn run_transform(
let deadline = deadline_instant(options.deadline_ms);
let media_limits = runtime.limits();
let transform_metrics = runtime.metrics();
let cost_data = data.clone();
let timed = runtime
.tasks()
.run_native(deadline, move || {
let started = Instant::now();
let media =
media_process::transform_image(&data, &options, &media_limits, &transform_metrics)?;
Ok(TimedMedia {
media,
elapsed_ms: metrics::duration_millis(started.elapsed()),
})
})
.run_native(
deadline,
move || {
Ok(media_process::image_transform_cost(
&cost_data,
&options,
&media_limits,
)?)
},
move || {
let started = Instant::now();
let media = media_process::transform_image(
&data,
&options,
&media_limits,
&transform_metrics,
)?;
Ok(TimedMedia {
media,
elapsed_ms: metrics::duration_millis(started.elapsed()),
})
},
)
.await?;
runtime.metrics().observe_image_duration(timed.elapsed_ms);
request_log::record_stage(Stage::Transform, timed.elapsed_ms);
@@ -67,34 +82,42 @@ pub(in crate::server) async fn run_video_transform(
let deadline = deadline_instant(deadline_ms);
let media_limits = runtime.limits();
let transform_metrics = runtime.metrics();
let cost_bytes = data
.len()
.saturating_add(media_process::AV_NATIVE_COST_BYTES);
let timed = runtime
.tasks()
.run_native(deadline, move || {
let started = Instant::now();
let thumbnail = media_process::extract_video_thumbnail(&data, format, &media_limits)?;
let media = if width.is_none() && height.is_none() {
thumbnail
} else {
media_process::transform_image(
&thumbnail.bytes,
&ImageOptions {
width,
height,
format,
quality,
animation: AnimationMode::Static,
deadline_ms,
..Default::default()
},
&media_limits,
&transform_metrics,
)?
};
Ok(TimedMedia {
media,
elapsed_ms: metrics::duration_millis(started.elapsed()),
})
})
.run_native(
deadline,
move || Ok(cost_bytes),
move || {
let started = Instant::now();
let thumbnail =
media_process::extract_video_thumbnail(&data, format, &media_limits)?;
let media = if width.is_none() && height.is_none() {
thumbnail
} else {
media_process::transform_image(
&thumbnail.bytes,
&ImageOptions {
width,
height,
format,
quality,
animation: AnimationMode::Static,
deadline_ms,
..Default::default()
},
&media_limits,
&transform_metrics,
)?
};
Ok(TimedMedia {
media,
elapsed_ms: metrics::duration_millis(started.elapsed()),
})
},
)
.await?;
runtime.metrics().observe_video_duration(timed.elapsed_ms);
request_log::record_stage(Stage::Transform, timed.elapsed_ms);
+49 -9
View File
@@ -59,8 +59,12 @@ use axum::{
use bytes::Bytes;
use std::{collections::HashMap, sync::Arc};
const DECODED_BYTES_PER_PIXEL: usize = 4;
const DECODED_PIXEL_BUFFERS_PER_TRANSFORM: usize = 2;
const NATIVE_MEMORY_HEADROOM_BYTES: usize = 512 << 20;
const NATIVE_MEMORY_FLOOR_BYTES: usize = 128 << 20;
const CGROUP_MEMORY_LIMIT_PATHS: [&str; 2] = [
"/sys/fs/cgroup/memory.max",
"/sys/fs/cgroup/memory/memory.limit_in_bytes",
];
const CONTENT_TYPE_SNIFF_PREFIX_BYTES: usize = 8192;
pub(in crate::server) struct TransformRuntime {
@@ -89,6 +93,18 @@ impl TransformRuntime {
metrics: &Arc<metrics::Metrics>,
) -> anyhow::Result<Self> {
let limits = MediaLimits::default_from_config();
let memory_limit_bytes = process_memory_limit_bytes();
let memory_budget_bytes = native_memory_budget_bytes(
memory_limit_bytes,
cfg.media.transform_cache_capacity_bytes,
);
tracing::info!(
memory_limit_bytes,
transform_cache_bytes = cfg.media.transform_cache_capacity_bytes,
memory_budget_bytes,
max_native_transforms = cfg.media.max_native_transforms,
"native memory budget"
);
Ok(Self {
limits,
animation: AnimationLimits::new(
@@ -99,7 +115,7 @@ impl TransformRuntime {
tasks: NativeTaskExecutor::new(NativeTaskExecutorSettings {
max_native_transforms: cfg.media.max_native_transforms,
worker_queue_capacity: cfg.media.worker_queue_capacity,
decoded_bytes_per_transform: decoded_bytes_per_transform(&limits),
memory_budget_bytes,
native_metrics: metrics.native_transform(),
transform_metrics: metrics.transform(),
}),
@@ -144,12 +160,36 @@ impl TransformRuntime {
}
}
fn decoded_bytes_per_transform(limits: &MediaLimits) -> usize {
limits
.image_pixels()
.max(limits.animated_total_pixels())
.checked_mul(DECODED_BYTES_PER_PIXEL * DECODED_PIXEL_BUFFERS_PER_TRANSFORM)
.expect("the native decoded image budget must not overflow")
fn native_memory_budget_bytes(memory_limit_bytes: usize, transform_cache_bytes: usize) -> usize {
let cache = transform_cache_bytes.min(memory_limit_bytes / 2);
let headroom = NATIVE_MEMORY_HEADROOM_BYTES.min(memory_limit_bytes / 8);
memory_limit_bytes
.saturating_sub(cache)
.saturating_sub(headroom)
.max(NATIVE_MEMORY_FLOOR_BYTES)
}
fn process_memory_limit_bytes() -> usize {
let physical = physical_memory_bytes();
CGROUP_MEMORY_LIMIT_PATHS
.iter()
.find_map(|path| {
std::fs::read_to_string(path)
.ok()?
.trim()
.parse::<usize>()
.ok()
})
.filter(|limit| *limit < physical)
.unwrap_or(physical)
}
fn physical_memory_bytes() -> usize {
let pages = unsafe { libc::sysconf(libc::_SC_PHYS_PAGES) };
let page_size = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
usize::try_from(pages)
.unwrap_or(0)
.saturating_mul(usize::try_from(page_size).unwrap_or(0))
}
pub(in crate::server) async fn serve_bytes_or_transform(
@@ -199,3 +199,22 @@ async fn a_redundant_format_parameter_shares_the_implicit_transform_cache_entry(
);
}
}
#[test]
fn the_budget_leaves_room_for_the_cache_and_headroom() {
assert_eq!(
2560 << 20,
native_memory_budget_bytes(4096 << 20, 1024 << 20)
);
}
#[test]
fn a_cache_larger_than_the_pod_takes_at_most_half_of_it() {
assert_eq!(192 << 20, native_memory_budget_bytes(512 << 20, 1 << 30));
assert_eq!(384 << 20, native_memory_budget_bytes(1 << 30, 1 << 30));
}
#[test]
fn the_budget_never_drops_below_its_floor() {
assert_eq!(128 << 20, native_memory_budget_bytes(64 << 20, 0));
}
@@ -30,6 +30,9 @@ pub fn apng_header(frames: u32) -> Vec<u8> {
data
}
pub const JXL_4100_LOSSY_ALPHA: &[u8] = include_bytes!("jxl_4100_lossy_alpha.jxl");
pub const JXL_64_LOSSY_ALPHA: &[u8] = include_bytes!("jxl_64_lossy_alpha.jxl");
pub fn synthetic_png(width: u32, height: u32) -> Vec<u8> {
ensure_vips_init().unwrap();
let mut pixels = vec![0u8; width as usize * height as usize * 4];
+4 -3
View File
@@ -10,9 +10,10 @@ pub use adversarial::{
};
pub use ffmpeg_cli::{ffmpeg_gen_media, ffmpeg_gen_mp4, ffmpeg_gen_rotated_mp4, ffmpeg_mirror_mp4};
pub use images::{
animated_gif_fixture, animated_gif_frames, apng_header, first_webp_anim_frame_size,
gif_frame_delays_cs, gif_loop_count, minimal_gif, png_dimensions, synthetic_bmp, synthetic_png,
webp_animation_loop_count, webp_canvas_size, webp_chunk_payloads, webp_with_metadata_chunk,
JXL_64_LOSSY_ALPHA, JXL_4100_LOSSY_ALPHA, animated_gif_fixture, animated_gif_frames,
apng_header, first_webp_anim_frame_size, gif_frame_delays_cs, gif_loop_count, minimal_gif,
png_dimensions, synthetic_bmp, synthetic_png, webp_animation_loop_count, webp_canvas_size,
webp_chunk_payloads, webp_with_metadata_chunk,
};
pub use media::{
fixture_audio_mp3_with_png_cover_art, fixture_audio_mp4_with_attached_picture,