前置完美像素的端点并发闸
新增端点级并发许可并在首次 IO 之前取得 以取消安全的有界队列限制等待者数量 补齐队列边界、释放与闸位顺序的回归测试 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -2,14 +2,17 @@ use std::{
|
||||
borrow::Cow,
|
||||
collections::BTreeMap,
|
||||
io::Cursor,
|
||||
sync::{Arc, LazyLock},
|
||||
sync::{
|
||||
Arc, LazyLock,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
},
|
||||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
use axum::{
|
||||
Json,
|
||||
extract::{Extension, Path, Query, State, rejection::JsonRejection},
|
||||
http::StatusCode,
|
||||
http::{HeaderValue, StatusCode},
|
||||
};
|
||||
use module_assets::{
|
||||
AssetObjectAccessPolicy, AssetObjectFieldError, build_asset_object_upsert_input,
|
||||
@@ -133,6 +136,21 @@ const EDITOR_GENERATION_UNSUPPORTED_STYLE_WARNING_CODE: &str = "unsupported-imag
|
||||
pub(crate) const EDITOR_GENERATION_MULTIPLE_WARNINGS_CODE: &str = "multiple-generation-warnings";
|
||||
const EDITOR_GENERATION_MAX_ASPECT_RATIO_DRIFT: f64 = 0.05;
|
||||
const EDITOR_PIXEL_ART_CPU_MAX_CONCURRENCY: usize = 2;
|
||||
/// 中文注释:inline 完美像素的端点级并发闸。它是仓库里唯一「用户主动触发 + 同步执行 +
|
||||
/// 下载大图 + 吃 CPU」且不经生成队列的路径——其余像素规整入口都由 job + worker 承担准入。
|
||||
/// 闸此前只有 CPU 许可那一道,而它设在下载之后:请求先把最多 32 MiB 读进内存、先打完
|
||||
/// 几轮全账号 SpacetimeDB 扫描,才被拦下排队,等于闸在资源已被消耗之后才检查。
|
||||
/// 客户端只发几百字节 JSON 就能让服务端拉取 32 MiB,放大比约 64000 : 1,且本操作免费
|
||||
/// (generation_cost_mud_points = 0)、无 per-user 配额。
|
||||
///
|
||||
/// 这道闸覆盖从第一次 IO 到 handler 结束的全过程,一个数字同时封住并发 SpacetimeDB
|
||||
/// 扫描数、并发源图缓冲数与并发 OSS PUT 数。取 4 而不是等于 CPU 槽的 2:2 会让下载完全
|
||||
/// 串在 snap 后面,4 才能让两个请求下载的同时另两个在算,把下载延迟藏进 CPU 时间里。
|
||||
const EDITOR_PIXEL_ART_SNAP_MAX_CONCURRENCY: usize = 4;
|
||||
/// 中文注释:等待队列保险丝,对齐 BgFilter 的 Q。部署配置里全局准入是 512,小于这个值,
|
||||
/// 所以正常部署下它打不到;它兜的是 max_concurrent_requests 未配置(代码默认 None,
|
||||
/// 即无全局上限)时的连接风暴。它只防雪崩,不做流量整形。
|
||||
const EDITOR_PIXEL_ART_SNAP_MAX_QUEUE_DEPTH: usize = 2048;
|
||||
const EDITOR_PIXEL_ART_MAX_PROCESSING_DURATION: Duration = Duration::from_secs(30);
|
||||
const EDITOR_PIXEL_ART_SNAP_ASSET_KIND: &str = "editor_pixel_art_snap";
|
||||
const EDITOR_PIXEL_ART_SNAP_MODEL: &str = "Perfect Pixel";
|
||||
@@ -142,6 +160,12 @@ static EDITOR_PIXEL_ART_CPU_LIMITER: LazyLock<Arc<tokio::sync::Semaphore>> = Laz
|
||||
EDITOR_PIXEL_ART_CPU_MAX_CONCURRENCY,
|
||||
))
|
||||
});
|
||||
static EDITOR_PIXEL_ART_SNAP_LIMITER: LazyLock<Arc<tokio::sync::Semaphore>> = LazyLock::new(|| {
|
||||
Arc::new(tokio::sync::Semaphore::new(
|
||||
EDITOR_PIXEL_ART_SNAP_MAX_CONCURRENCY,
|
||||
))
|
||||
});
|
||||
static EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH: AtomicUsize = AtomicUsize::new(0);
|
||||
static EDITOR_ICON_SPRITESHEET_CPU_LIMITER: LazyLock<Arc<tokio::sync::Semaphore>> =
|
||||
LazyLock::new(|| {
|
||||
Arc::new(tokio::sync::Semaphore::new(
|
||||
@@ -3215,6 +3239,78 @@ fn take_arc_downloaded_image(image: Arc<DownloadedOpenAiImage>) -> DownloadedOpe
|
||||
Arc::try_unwrap(image).unwrap_or_else(|image| image.as_ref().clone())
|
||||
}
|
||||
|
||||
// 中文注释:只在 depth < max_depth 时递增,用 CAS 而不是「先读后加」——两个线程同时读到
|
||||
// max_depth - 1 各自加一就会越界。抽成自由函数是为了能直接单测边界与并发行为。
|
||||
fn try_enter_bounded_queue(depth: &AtomicUsize, max_depth: usize) -> bool {
|
||||
depth
|
||||
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
|
||||
(current < max_depth).then_some(current + 1)
|
||||
})
|
||||
.is_ok()
|
||||
}
|
||||
|
||||
/// 中文注释:递减必须放在 Drop 里。等待中的 future 随时可能被丢弃(客户端断连、超时触发、
|
||||
/// 上层取消),若把递减写在正常返回路径上,计数就会只增不减,最终队列永久「满」、接口
|
||||
/// 彻底不可用——这是本改动里唯一一处写错会造成永久性故障的地方。
|
||||
struct EditorPixelArtSnapQueueGuard;
|
||||
|
||||
impl EditorPixelArtSnapQueueGuard {
|
||||
fn try_enter() -> Option<Self> {
|
||||
try_enter_bounded_queue(
|
||||
&EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH,
|
||||
EDITOR_PIXEL_ART_SNAP_MAX_QUEUE_DEPTH,
|
||||
)
|
||||
.then_some(Self)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for EditorPixelArtSnapQueueGuard {
|
||||
fn drop(&mut self) {
|
||||
EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH.fetch_sub(1, Ordering::AcqRel);
|
||||
}
|
||||
}
|
||||
|
||||
// 中文注释:端点级并发闸,必须在第一次 IO 之前取得,许可持有到 handler 结束。等待时间与
|
||||
// 后续处理共用同一份 processing_deadline,所以排队不会让请求重新获得完整预算。
|
||||
async fn acquire_editor_pixel_art_snap_permit(
|
||||
processing_deadline: Instant,
|
||||
) -> Result<tokio::sync::OwnedSemaphorePermit, AppError> {
|
||||
// 中文注释:这句预检不能省。timeout_at 会先 poll 一次内层 future,许可空闲时
|
||||
// acquire_owned 立刻就绪,于是预算已耗尽的请求照样拿到许可,白占一个名额再去打
|
||||
// 几轮全账号扫描,直到下载那步才失败。既有 CPU 许可的同名预检也是为此存在。
|
||||
if Instant::now() >= processing_deadline {
|
||||
return Err(editor_pixel_art_snap_failure(
|
||||
StatusCode::GATEWAY_TIMEOUT,
|
||||
"完美像素处理预算已耗尽。",
|
||||
));
|
||||
}
|
||||
let queue_guard = EditorPixelArtSnapQueueGuard::try_enter().ok_or_else(|| {
|
||||
editor_pixel_art_snap_failure(
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
"完美像素排队已满,请稍后重试。",
|
||||
)
|
||||
.with_header("retry-after", HeaderValue::from_static("1"))
|
||||
})?;
|
||||
let acquired = tokio::time::timeout_at(
|
||||
tokio::time::Instant::from_std(processing_deadline),
|
||||
Arc::clone(&*EDITOR_PIXEL_ART_SNAP_LIMITER).acquire_owned(),
|
||||
)
|
||||
.await;
|
||||
// 中文注释:排队阶段到此结束,先让出队列位再判定结果,避免持有许可期间还占着队列名额。
|
||||
drop(queue_guard);
|
||||
match acquired {
|
||||
Ok(Ok(permit)) => Ok(permit),
|
||||
Ok(Err(error)) => Err(editor_pixel_art_snap_failure(
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
format!("完美像素并发门限不可用:{error}"),
|
||||
)),
|
||||
Err(_) => Err(editor_pixel_art_snap_failure(
|
||||
StatusCode::GATEWAY_TIMEOUT,
|
||||
"完美像素等待处理槽位时预算已耗尽。",
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn acquire_editor_pixel_art_cpu_permit(
|
||||
processing_deadline: Instant,
|
||||
) -> Result<tokio::sync::OwnedSemaphorePermit, String> {
|
||||
@@ -4518,6 +4614,10 @@ pub async fn snap_editor_image_to_pixel_art(
|
||||
// 中文注释:在 CPU 处理和 OSS PUT 前完成所有无副作用校验;主动操作的像素规整
|
||||
// 失败必须直接返回错误,不能先创建与原图相同的派生资源。
|
||||
serialize_editor_asset_metadata(payload.generation_inputs.clone())?;
|
||||
// 中文注释:端点级并发闸设在第一次 IO 之前——上面几步都是纯内存校验,让畸形请求也去
|
||||
// 排队既浪费名额又让 400 拖到 30 秒。闸之后的全部 IO(全账号扫描、下载、规整、持久化)
|
||||
// 都在许可覆盖范围内,许可随 handler 返回自动释放。
|
||||
let _snap_permit = acquire_editor_pixel_art_snap_permit(processing_deadline).await?;
|
||||
let owner_user_id = current_owner_user_id(&authenticated);
|
||||
let project = state
|
||||
.spacetime_client()
|
||||
@@ -10385,6 +10485,55 @@ mod tests {
|
||||
assert_eq!(warning.reason, "像素规整未完成,已保留原始生成结果。");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pixel_art_snap_queue_guard_is_bounded_and_released_on_drop() {
|
||||
assert_eq!(EDITOR_PIXEL_ART_SNAP_MAX_CONCURRENCY, 4);
|
||||
assert_eq!(EDITOR_PIXEL_ART_SNAP_MAX_QUEUE_DEPTH, 2048);
|
||||
|
||||
let depth = AtomicUsize::new(0);
|
||||
assert!(try_enter_bounded_queue(&depth, 2));
|
||||
assert!(try_enter_bounded_queue(&depth, 2));
|
||||
// 中文注释:满了必须拒绝,且拒绝时不得把计数推过上限。
|
||||
assert!(!try_enter_bounded_queue(&depth, 2));
|
||||
assert_eq!(depth.load(Ordering::Acquire), 2);
|
||||
|
||||
// 上限为 0 时任何进入都必须失败。
|
||||
let closed = AtomicUsize::new(0);
|
||||
assert!(!try_enter_bounded_queue(&closed, 0));
|
||||
assert_eq!(closed.load(Ordering::Acquire), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pixel_art_snap_queue_depth_returns_to_zero_after_guards_drop() {
|
||||
// 中文注释:递减写在 Drop 里,被取消的等待者也必须归还名额;否则计数只增不减,
|
||||
// 队列会永久「满」,接口彻底不可用。
|
||||
let before = EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH.load(Ordering::Acquire);
|
||||
{
|
||||
let _first = EditorPixelArtSnapQueueGuard::try_enter().expect("queue should accept");
|
||||
let _second = EditorPixelArtSnapQueueGuard::try_enter().expect("queue should accept");
|
||||
assert_eq!(
|
||||
EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH.load(Ordering::Acquire),
|
||||
before + 2
|
||||
);
|
||||
}
|
||||
assert_eq!(
|
||||
EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH.load(Ordering::Acquire),
|
||||
before
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn pixel_art_snap_permit_reports_exhausted_budget_without_waiting() {
|
||||
let expired = Instant::now() - Duration::from_secs(1);
|
||||
let error = acquire_editor_pixel_art_snap_permit(expired)
|
||||
.await
|
||||
.expect_err("expired budget should not acquire a permit");
|
||||
|
||||
assert_eq!(error.status_code(), StatusCode::GATEWAY_TIMEOUT);
|
||||
// 中文注释:失败路径也必须归还队列名额。
|
||||
assert_eq!(EDITOR_PIXEL_ART_SNAP_QUEUE_DEPTH.load(Ordering::Acquire), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pixel_art_processing_deadline_uses_earlier_local_or_request_budget() {
|
||||
assert_eq!(EDITOR_PIXEL_ART_CPU_MAX_CONCURRENCY, 2);
|
||||
@@ -10680,6 +10829,9 @@ mod tests {
|
||||
"ensure_editor_reference_image_source_is_stable",
|
||||
"validate_editor_pixel_art_snap_canvas_completion",
|
||||
"serialize_editor_asset_metadata",
|
||||
// 中文注释:并发闸必须排在第一次 IO 之前。挪到 .get_editor_project 之后,
|
||||
// 全账号扫描与下载就重新回到闸外,等于闸在资源被消耗之后才检查。
|
||||
"acquire_editor_pixel_art_snap_permit",
|
||||
".get_editor_project",
|
||||
"validate_editor_pixel_art_snap_placeholder_exists",
|
||||
"resolve_editor_pixel_art_source_for_owner",
|
||||
|
||||
Reference in New Issue
Block a user