Skip to main content

yggdrasil/components/assets/
upload_pool.rs

1//! 上传引擎内部状态机:与 UI 完全解耦的纯逻辑(无 `rsx!` / 无 Dioxus 组件),
2//! 供 [`crate::components::assets::AssetUploadModal`] 组合使用。
3//!
4//! 三条入口(点击选择 / 拖拽 / 粘贴)收敛到同一个 `enqueue_files`:拿到的
5//! `web_sys::File` 先过 `validate_file`(镜像服务端 5MiB / 四种 MIME 的硬限制,
6//! 不合格立即记失败行、不发请求),合格项入共享队列后由 **worker 池**并发上传:
7//! 在跑 worker 数上限 = 并发配置(「站点配置」面板 / `UPLOAD_CONCURRENCY` env 播种,
8//! 挂载时经 `get_upload_settings` 拉取,失败回退默认 3)。每个 worker 张间停顿
9//! `500ms × 当前并发数`——N 路并行时聚合速率恒 ≤ 2/s,与默认上传限流桶
10//! (`RATE_LIMIT_UPLOAD_PER_SEC=2` / `BURST=15`)对齐:N=1 时与旧顺序版逐张 500ms
11//! 完全一致,N 调大也不触发 429。上传管线与文章编辑器完全同源
12//! (`upload_image_file` → `POST /api/upload`)。
13//!
14//! 注:worker 池共享状态用 `Rc<UploadPool>` 而非裸 Cell/RefCell——`use_hook` 每次渲染
15//! 都会 clone 存储值(dioxus-core `use_hook_inner` 走 `.cloned()`),裸结构会被按值
16//! 拷贝,各入口闭包拿到互不同步的副本导致 id 冲突与队列分裂;Rc 克隆共享同一实例。
17
18#[cfg(target_arch = "wasm32")]
19use dioxus::prelude::*;
20
21#[cfg(target_arch = "wasm32")]
22use crate::bridges::tiptap::upload_image_file;
23#[cfg(target_arch = "wasm32")]
24use crate::utils::format_bytes;
25
26/// 单文件大小硬上限(5MiB),镜像服务端 `crate::utils::server::MAX_FILE_SIZE`。
27#[cfg(any(test, target_arch = "wasm32"))]
28const MAX_UPLOAD_BYTES: u64 = 5 * 1024 * 1024;
29/// 允许的 MIME 白名单,镜像服务端 `api/upload.rs` 的 `ALLOWED_MIME_TYPES`。
30#[cfg(any(test, target_arch = "wasm32"))]
31const ALLOWED_MIME: &[&str] = &["image/jpeg", "image/png", "image/gif", "image/webp"];
32
33/// 单条上传状态机:Queued → Uploading → Done / Failed(原因)。
34// 变体仅在 WASM 端构造,server SSR 只匹配渲染,非 wasm 构建放行 dead_code。
35#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
36#[derive(Clone, PartialEq)]
37pub(crate) enum UploadStatus {
38    Queued,
39    Uploading,
40    Done,
41    Failed(String),
42}
43
44/// 列表行数据(纯数据,两端都可编译;`web_sys::File` 句柄不进 signal,
45/// 由 WASM 端的 files 表以 id 关联保存,供重试取回)。
46// 同 UploadStatus:实例仅在 WASM 端构造。
47#[cfg_attr(not(target_arch = "wasm32"), allow(dead_code))]
48#[derive(Clone, PartialEq)]
49pub(crate) struct UploadItem {
50    pub(crate) id: u64,
51    pub(crate) name: String,
52    /// 入队时已用 `format_bytes` 格式化好的可读大小。
53    pub(crate) size: String,
54    pub(crate) status: UploadStatus,
55    /// 用户点了 ×:先播退出动画(animate-row-leave),EXIT_ANIM_MS 后真正摘除。
56    pub(crate) removing: bool,
57}
58
59/// 一批入队(一次拖拽/选择/粘贴)的完成追踪:remaining 归零且 ≥1 成功时回调一次。
60#[cfg(target_arch = "wasm32")]
61struct BatchCtx {
62    remaining: std::cell::Cell<usize>,
63    any_done: std::cell::Cell<bool>,
64}
65
66/// Worker 池共享状态:id 分配、重试句柄表、待传队列、在跑 worker 数。
67/// `use_hook` 持 `Rc<UploadPool>`:渲染期 clone 的是 Rc(共享同一实例),各入口闭包
68/// 与 worker 看到的永远是同一份队列/句柄表。字段私有——UI 层只经下方方法访问,
69/// 不直接触碰 files/queue 等内部表示。
70#[cfg(target_arch = "wasm32")]
71pub(crate) struct UploadPool {
72    next_id: std::cell::Cell<u64>,
73    files: std::cell::RefCell<Vec<(u64, web_sys::File)>>,
74    queue: std::cell::RefCell<std::collections::VecDeque<(u64, std::rc::Rc<BatchCtx>)>>,
75    active_workers: std::cell::Cell<u32>,
76}
77
78#[cfg(target_arch = "wasm32")]
79impl UploadPool {
80    pub(crate) fn new() -> Self {
81        Self {
82            next_id: std::cell::Cell::new(0),
83            files: std::cell::RefCell::new(Vec::new()),
84            queue: std::cell::RefCell::new(std::collections::VecDeque::new()),
85            active_workers: std::cell::Cell::new(0),
86        }
87    }
88
89    /// 取回指定 id 的文件句柄克隆(供 UI 层「重试」重新读取原始文件;不移除句柄表条目)。
90    pub(crate) fn find_file(&self, id: u64) -> Option<web_sys::File> {
91        self.files
92            .borrow()
93            .iter()
94            .find(|(fid, _)| *fid == id)
95            .map(|(_, f)| f.clone())
96    }
97
98    /// 移除指定 id 的文件句柄(UI 层用户点「移除」摘除该条目时调用,释放持有)。
99    pub(crate) fn remove_file(&self, id: u64) {
100        self.files.borrow_mut().retain(|(fid, _)| *fid != id);
101    }
102}
103
104/// 预校验:MIME 白名单 + 5MiB 上限。失败返回可读原因(直接展示在行内,不发请求)。
105///
106/// `pub(crate)`:`asset_picker.rs` 的内嵌上传入口复用同一份校验规则,与本模块
107/// worker 池入队时执行的规则保持单一实现来源。
108#[cfg(any(test, target_arch = "wasm32"))]
109pub(crate) fn validate_file(mime: &str, size: u64) -> Result<(), String> {
110    if !ALLOWED_MIME.contains(&mime) {
111        return Err("不支持的文件类型(仅 JPEG / PNG / GIF / WebP)".into());
112    }
113    if size > MAX_UPLOAD_BYTES {
114        return Err("大小超过 5MB 限制".into());
115    }
116    Ok(())
117}
118
119/// 更新单条状态:write guard 在函数内立即释放,绝不跨 await 持有。
120#[cfg(target_arch = "wasm32")]
121pub(crate) fn set_status(items: &mut Signal<Vec<UploadItem>>, id: u64, status: UploadStatus) {
122    let mut guard = items.write();
123    if let Some(it) = guard.iter_mut().find(|it| it.id == id) {
124        it.status = status;
125    }
126}
127
128/// 单个 worker 的消费循环:队列空即退出并归还名额。
129///
130/// 张间停顿 `500ms × 当前并发数`:N 个 worker 并行时聚合速率恒 ≤ 2/s(停顿随 N
131/// 线性放大),与默认上传限流桶(`RATE_LIMIT_UPLOAD_PER_SEC=2` / `BURST=15`)对齐——
132/// N=1 时与旧顺序版逐张 500ms 完全一致,N 调大也不会触发 429。
133#[cfg(target_arch = "wasm32")]
134async fn worker_loop(
135    mut items: Signal<Vec<UploadItem>>,
136    pool: std::rc::Rc<UploadPool>,
137    concurrency: Signal<i32>,
138    on_uploaded: EventHandler<()>,
139) {
140    loop {
141        // borrow 不出块,guard 不跨 await。
142        let next = pool.queue.borrow_mut().pop_front();
143        let Some((id, batch)) = next else { break };
144        // 句柄可能已被「×」移除:跳过上传但仍计入批次完成度(否则该批永不合拢,
145        // 其他成功项的 on_uploaded 无法触发)。
146        let file = pool.find_file(id);
147        if let Some(file) = file {
148            set_status(&mut items, id, UploadStatus::Uploading);
149            match upload_image_file(file).await {
150                Ok(_) => {
151                    set_status(&mut items, id, UploadStatus::Done);
152                    batch.any_done.set(true);
153                }
154                Err(msg) => set_status(&mut items, id, UploadStatus::Failed(msg)),
155            }
156        }
157        // 本批最后一条收尾且 ≥1 成功 → 回调一次(父组件刷新网格)。
158        let remaining = batch.remaining.get() - 1;
159        batch.remaining.set(remaining);
160        if remaining == 0 && batch.any_done.get() {
161            on_uploaded.call(());
162        }
163        // 队列未空则停顿压速率;停顿随并发数线性放大(见 fn doc),live 读取
164        // 让面板改动即时生效。clamp(1, 32) 纯防御:服务端已钳到 1–8,
165        // 此处只防异常值导致 500*n 溢出或路径级长停。
166        if !pool.queue.borrow().is_empty() {
167            let n = (*concurrency.peek()).clamp(1, 32) as u32;
168            crate::utils::time::sleep_ms(500 * n).await;
169        }
170    }
171    pool.active_workers.set(pool.active_workers.get() - 1);
172}
173
174/// 三入口收敛点:校验入队 + 按需补足 worker。仅 WASM 端存在(`web_sys::File` /
175/// `spawn` / `upload_image_file` 都是 WASM-only 符号)。
176#[cfg(target_arch = "wasm32")]
177pub(crate) fn enqueue_files(
178    mut items: Signal<Vec<UploadItem>>,
179    pool: std::rc::Rc<UploadPool>,
180    concurrency: Signal<i32>,
181    on_uploaded: EventHandler<()>,
182    new_files: Vec<web_sys::File>,
183) {
184    // 1) 校验入队:不合格直接记 Failed(不发请求);合格记 Queued 并留存句柄供重试。
185    let mut valid_ids = Vec::new();
186    for file in new_files {
187        let id = pool.next_id.get() + 1;
188        pool.next_id.set(id);
189        let item = UploadItem {
190            id,
191            name: file.name(),
192            size: format_bytes(file.size() as i64),
193            removing: false,
194            status: match validate_file(&file.type_(), file.size() as u64) {
195                Ok(()) => {
196                    pool.files.borrow_mut().push((id, file));
197                    valid_ids.push(id);
198                    UploadStatus::Queued
199                }
200                Err(msg) => UploadStatus::Failed(msg),
201            },
202        };
203        items.write().push(item);
204    }
205    if valid_ids.is_empty() {
206        return;
207    }
208
209    // 2) 一批一个 BatchCtx 追踪完成度;id 入共享队列后补足 worker:在跑数低于
210    //    并发上限且队列非空时逐个 spawn(worker 队列空自退,不会超额驻留)。
211    let batch = std::rc::Rc::new(BatchCtx {
212        remaining: std::cell::Cell::new(valid_ids.len()),
213        any_done: std::cell::Cell::new(false),
214    });
215    {
216        let mut q = pool.queue.borrow_mut();
217        for id in valid_ids {
218            q.push_back((id, batch.clone()));
219        }
220    }
221    // clamp 上界与 worker_loop 的停顿同款防御(正常值 1–8)。
222    let target = (*concurrency.peek()).clamp(1, 32) as u32;
223    while pool.active_workers.get() < target && !pool.queue.borrow().is_empty() {
224        pool.active_workers.set(pool.active_workers.get() + 1);
225        spawn(worker_loop(items, pool.clone(), concurrency, on_uploaded));
226    }
227}
228
229#[cfg(test)]
230mod tests {
231    use super::{validate_file, ALLOWED_MIME, MAX_UPLOAD_BYTES};
232
233    /// 四种支持的类型在上限边缘全部接受。
234    #[test]
235    fn validate_file_accepts_supported_types_at_limit() {
236        for mime in ALLOWED_MIME {
237            assert!(
238                validate_file(mime, MAX_UPLOAD_BYTES).is_ok(),
239                "{mime} 应被接受"
240            );
241        }
242    }
243
244    /// svg 不在白名单(服务端同样拒绝)。
245    #[test]
246    fn validate_file_rejects_svg() {
247        assert!(validate_file("image/svg+xml", 1024).is_err());
248    }
249
250    /// 空 MIME(浏览器给不出类型)也拒绝。
251    #[test]
252    fn validate_file_rejects_empty_mime() {
253        assert!(validate_file("", 1024).is_err());
254    }
255
256    /// 上限 +1 字节拒绝。
257    #[test]
258    fn validate_file_rejects_oversize() {
259        assert!(validate_file("image/png", MAX_UPLOAD_BYTES + 1).is_err());
260    }
261}