Skip to main content

tdm_server_rust/service/
editor_import_service.rs

1//! 图源导入服务 (Editor Provider Import)
2//!
3//! 把话数的图源(`provider_file_oss_id` 指向的单图或压缩包)转换为编辑器页面:
4//!
5//! 1. 通过现有 OSS 通道下载图源字节
6//! 2. 解压/识别为按页序排序的图片序列
7//! 3. 逐张上传到图片 OSS(`tencent.image_bucket` / `image_cdn_domain`)
8//! 4. 写入 `editor_page` 记录并标记 `uploaded=true`
9//!
10//! 导入后编辑器只消费图片 OSS 地址,运行期不再读取压缩包。
11//!
12//! 幂等:已存在同话数页面时直接返回;`force=true` 时删除旧页面、旧标点后重建。
13
14use crate::app::AppState;
15use crate::cache::on_episode_mutated;
16use crate::entity::editor::{
17    EditorPageInfo, EditorUploadCredential, EditorUploadCredentialsResponse, EditorUploadFile,
18    ImportProviderResult, RegisterEditorPagesRequest, SourceArchiveUploadCredential,
19};
20use crate::error::{ApiResult, AppError};
21use crate::repository::editor_page_repo::{EditorPageRepository, NewEditorPage};
22use crate::repository::editor_unit_repo::EditorUnitRepository;
23use crate::repository::episode_repo::EpisodeRepository;
24use crate::repository::oss_repo::OssRepository;
25use crate::service::editor_source_archive_service::{archive_object_key, compute_archive_hash};
26use crate::service::oss_service::OssService;
27use crate::utils::archive_extract::{self, ExtractedImage};
28use reqwest::header::CONTENT_TYPE;
29use sea_orm::{ConnectionTrait, TransactionTrait};
30use std::collections::HashSet;
31use uuid::Uuid;
32
33/// 单次图源导入允许的最大页面数,避免异常请求耗尽凭证、内存和数据库资源。
34const MAX_EDITOR_PAGE_COUNT: usize = 500;
35/// 编辑器记录允许的最大图片边长。
36const MAX_EDITOR_IMAGE_DIMENSION: i32 = 100_000;
37
38/// 图源导入服务
39pub struct EditorImportService;
40
41impl EditorImportService {
42    /// 从话数图源导入并生成编辑器页面。
43    ///
44    /// ## 参数
45    /// - `state`: 应用状态
46    /// - `episode_id`: 话数 ID
47    /// - `force`: 是否强制重建(删除旧页面与旧标点后重新导入)
48    ///
49    /// ## 返回
50    /// - [`ImportProviderResult`]:是否新建、页面总数与页面列表
51    ///
52    /// ## 错误
53    /// - 话数不存在或无图源
54    /// - 图源下载、解压、上传失败
55    #[tracing::instrument(skip_all, level = "debug")]
56    pub async fn import_provider_pages(
57        state: &AppState,
58        episode_id: i32,
59        force: bool,
60    ) -> ApiResult<ImportProviderResult> {
61        let episode_repo = EpisodeRepository::new(state.db.clone());
62        // 校验话数存在并取得图源 OSS 记录
63        let manga_id = episode_repo
64            .get_manga_id_by_episode_id(episode_id)
65            .await?
66            .ok_or_else(|| AppError::business(format!("漫画单话 {episode_id} 不存在喵")))?;
67        let provider_oss_id = episode_repo
68            .get_provider_oss_id(episode_id)
69            .await?
70            .ok_or_else(|| AppError::business("当前话数尚未上传图源喵"))?;
71
72        let page_repo = EditorPageRepository::new(state.db.clone());
73        let existing_pages = page_repo.list_by_episode(episode_id).await?;
74
75        // 幂等:已有页面且非强制,直接返回
76        if !existing_pages.is_empty() && !force {
77            return Ok(ImportProviderResult {
78                episode_id,
79                imported: false,
80                page_count: existing_pages.len(),
81                pages: existing_pages
82                    .into_iter()
83                    .map(EditorPageInfo::from)
84                    .collect(),
85            });
86        }
87
88        if force && !existing_pages.is_empty() {
89            let db = crate::db::from_sqlx_pool(state.db.clone());
90            ensure_rebuild_allowed_with(&db, episode_id, &existing_pages).await?;
91        }
92        let expected_page_ids = sorted_page_ids(&existing_pages);
93
94        // 下载图源(复用现有 provider 下载通道)
95        let (filename, bytes) =
96            OssService::download_file_proxy(state, episode_id, "provider".to_string()).await?;
97
98        // 解压/识别为按页序排序的图片
99        let images = archive_extract::extract_images(&filename, &bytes)?;
100
101        // 逐张上传到图片 OSS
102        let oss_repo = OssRepository::new(state.db.clone(), (*state.config).clone());
103        let bucket = oss_repo.image_bucket_name();
104        let max_bytes = state.config.tencent.max_file_size * 1024 * 1024;
105        let mut new_pages = Vec::with_capacity(images.len());
106        for (idx, img) in images.iter().enumerate() {
107            let object_key = build_image_object_key(episode_id, idx, &img.name);
108            let put_url = oss_repo
109                .presigned_image_put_url(&object_key, max_bytes)
110                .await?;
111            upload_image(state, &put_url, img).await?;
112            let (width, height) = image_dimensions(&img.data);
113            new_pages.push(NewEditorPage {
114                page_index: idx as i32,
115                source_oss_id: Some(provider_oss_id),
116                source_file_name: Some(img.name.clone()),
117                image_object_key: object_key.clone(),
118                image_bucket: Some(bucket.clone()),
119                image_url: oss_repo.image_display_url(&object_key),
120                image_width: width,
121                image_height: height,
122            });
123        }
124
125        let (inserted, imported) =
126            replace_pages_atomically(state, episode_id, force, &expected_page_ids, new_pages)
127                .await?;
128        tracing::info!(
129            "图源导入完成 episode_id={episode_id} manga_id={manga_id} pages={}",
130            inserted.len()
131        );
132        on_episode_mutated(state, manga_id, &[]).await;
133        Ok(ImportProviderResult {
134            episode_id,
135            imported,
136            page_count: inserted.len(),
137            pages: inserted.into_iter().map(EditorPageInfo::from).collect(),
138        })
139    }
140
141    /// 为前端图源解压直传准备批量上传凭证。
142    ///
143    /// 前端把图源 zip 解压(或散图识别)后,调用本接口一次拿到全部页面的图片 OSS
144    /// 预签名 PUT URL(单次 STS),随后逐张直传图片字节,最后调用
145    /// [`register_pages`](Self::register_pages) 落库。
146    ///
147    /// ## 参数
148    /// - `episode_id`: 话数 ID
149    /// - `files`: 待上传文件(页序 + 原始文件名,用于推断对象键后缀)
150    #[tracing::instrument(skip_all, level = "debug")]
151    pub async fn prepare_upload_credentials(
152        state: &AppState,
153        episode_id: i32,
154        files: Vec<EditorUploadFile>,
155    ) -> ApiResult<EditorUploadCredentialsResponse> {
156        validate_upload_files(&files)?;
157        // 校验话数存在
158        EpisodeRepository::new(state.db.clone())
159            .get_manga_id_by_episode_id(episode_id)
160            .await?
161            .ok_or_else(|| AppError::business(format!("漫画单话 {episode_id} 不存在喵")))?;
162        let expected_page_ids = sorted_page_ids(
163            &EditorPageRepository::new(state.db.clone())
164                .list_by_episode(episode_id)
165                .await?,
166        );
167
168        let oss_repo = OssRepository::new(state.db.clone(), (*state.config).clone());
169        let bucket = oss_repo.image_bucket_name();
170        let prefix = format!("editor/episode_{episode_id}/");
171        let object_keys: Vec<String> = files
172            .iter()
173            .map(|f| build_image_object_key(episode_id, f.index as usize, &f.filename))
174            .collect();
175        // 编辑器页面图沿用图源导入上限,避免大于通用图片上限的页面被 STS 策略拒绝。
176        let max_bytes = state
177            .config
178            .tencent
179            .max_file_size
180            .max(state.config.tencent.image_max_file_size)
181            * 1024
182            * 1024;
183        let urls = oss_repo
184            .presigned_image_put_urls(&prefix, &object_keys, max_bytes)
185            .await?;
186
187        let archive_hash = compute_archive_hash(
188            episode_id,
189            object_keys
190                .iter()
191                .enumerate()
192                .map(|(idx, key)| (idx as i32, key.as_str())),
193        );
194        let archive_key = archive_object_key(episode_id, archive_hash);
195        let archive_max = max_bytes * files.len().max(1) as u64;
196        let archive_put_url = oss_repo
197            .presigned_image_put_url(&archive_key, archive_max)
198            .await?;
199
200        let file_creds = files
201            .into_iter()
202            .zip(object_keys)
203            .zip(urls)
204            .map(|((f, object_key), presigned_url)| EditorUploadCredential {
205                index: f.index,
206                image_url: oss_repo.image_display_url(&object_key),
207                object_key,
208                presigned_url,
209                bucket: bucket.clone(),
210            })
211            .collect();
212        Ok(EditorUploadCredentialsResponse {
213            files: file_creds,
214            source_archive: SourceArchiveUploadCredential {
215                object_key: archive_key.clone(),
216                presigned_url: archive_put_url,
217                download_url: oss_repo.image_display_url(&archive_key),
218            },
219            expected_page_ids,
220        })
221    }
222
223    /// 注册前端已上传到图片 OSS 的页面。
224    ///
225    /// 与服务端 `import_provider_pages` 等价的落库逻辑,但图片由前端直传,
226    /// 不再依赖 `provider_file_oss_id` 压缩包。
227    ///
228    /// 幂等:已有页面且非 `force` 时直接返回旧页面;强制重建会校验页面快照及下游工作,
229    /// 并在同一事务内替换旧页面与标点。
230    #[tracing::instrument(skip_all, level = "debug")]
231    pub async fn register_pages(
232        state: &AppState,
233        episode_id: i32,
234        req: RegisterEditorPagesRequest,
235    ) -> ApiResult<ImportProviderResult> {
236        let manga_id = EpisodeRepository::new(state.db.clone())
237            .get_manga_id_by_episode_id(episode_id)
238            .await?
239            .ok_or_else(|| AppError::business(format!("漫画单话 {episode_id} 不存在喵")))?;
240
241        let oss_repo = OssRepository::new(state.db.clone(), (*state.config).clone());
242        let bucket = oss_repo.image_bucket_name();
243
244        // 校验并按 page_index 排序,禁止伪造其他话数的对象键。
245        let pages = validate_registered_pages(episode_id, req.pages)?;
246        let new_pages: Vec<NewEditorPage> = pages
247            .into_iter()
248            .map(|p| NewEditorPage {
249                page_index: p.page_index,
250                source_oss_id: None,
251                source_file_name: p.source_file_name,
252                image_url: oss_repo.image_display_url(&p.object_key),
253                image_object_key: p.object_key,
254                image_bucket: Some(bucket.clone()),
255                image_width: p.image_width,
256                image_height: p.image_height,
257            })
258            .collect();
259
260        let (inserted, imported) = replace_pages_atomically(
261            state,
262            episode_id,
263            req.force,
264            &req.expected_page_ids,
265            new_pages,
266        )
267        .await?;
268        tracing::info!(
269            "前端直传页面注册完成 episode_id={episode_id} pages={}",
270            inserted.len()
271        );
272        on_episode_mutated(state, manga_id, &[]).await;
273        Ok(ImportProviderResult {
274            episode_id,
275            imported,
276            page_count: inserted.len(),
277            pages: inserted.into_iter().map(EditorPageInfo::from).collect(),
278        })
279    }
280}
281
282/// 在短事务内原子切换页面集,失败时完整保留旧页面和标点。
283async fn replace_pages_atomically(
284    state: &AppState,
285    episode_id: i32,
286    force: bool,
287    expected_page_ids: &[i64],
288    new_pages: Vec<NewEditorPage>,
289) -> ApiResult<(Vec<crate::sea_entity::editor_page::Model>, bool)> {
290    let db = crate::db::from_sqlx_pool(state.db.clone());
291    let txn = db.begin().await?;
292    let episode = EpisodeRepository::lock_by_id_with(&txn, episode_id)
293        .await?
294        .ok_or_else(|| AppError::business(format!("漫画单话 {episode_id} 不存在喵")))?;
295    let existing_pages = EditorPageRepository::list_by_episode_with(&txn, episode_id).await?;
296
297    // 网络超时后的同内容重试直接视为成功,避免再次删除已提交页面。
298    if page_sets_equal(&existing_pages, &new_pages) {
299        txn.commit().await?;
300        return Ok((existing_pages, false));
301    }
302
303    if !existing_pages.is_empty() && !force {
304        txn.commit().await?;
305        return Ok((existing_pages, false));
306    }
307
308    if force && sorted_page_ids(&existing_pages) != sorted_ids(expected_page_ids) {
309        return Err(AppError::business(
310            "图源页面已被其他人更新,请刷新后重新上传喵",
311        ));
312    }
313
314    if force && !existing_pages.is_empty() {
315        ensure_rebuild_allowed_with(&txn, episode_id, &existing_pages).await?;
316        if episode.translator_file.is_some()
317            || episode.proofreader_file.is_some()
318            || episode.timer_file.is_some()
319            || episode.translator_file_oss_id.is_some()
320            || episode.proofreader_file_oss_id.is_some()
321            || episode.letterer_file_oss_id.is_some()
322            || episode.timer_file_oss_id.is_some()
323            || episode.publish_link.is_some()
324        {
325            return Err(AppError::business(
326                "该话已有下游稿件或发布记录,禁止直接重建图源喵",
327            ));
328        }
329        EditorUnitRepository::delete_by_episode_with(&txn, episode_id).await?;
330        EditorPageRepository::delete_by_episode_with(&txn, episode_id).await?;
331    }
332
333    let inserted = EditorPageRepository::insert_pages_with(&txn, episode_id, new_pages).await?;
334    txn.commit().await?;
335    Ok((inserted, true))
336}
337
338/// 检查当前页面是否已有需要保护的在线编辑或岗位交稿。
339async fn ensure_rebuild_allowed_with<C: ConnectionTrait>(
340    db: &C,
341    episode_id: i32,
342    existing_pages: &[crate::sea_entity::editor_page::Model],
343) -> ApiResult<()> {
344    if existing_pages.iter().any(|page| page.translation_completed)
345        || EditorUnitRepository::count_by_episode_with(db, episode_id).await? > 0
346        || EpisodeRepository::has_downstream_submission_with(db, episode_id).await?
347    {
348        return Err(AppError::business(
349            "该话已有翻译或校对内容,禁止直接重建图源喵",
350        ));
351    }
352    Ok(())
353}
354
355/// 校验上传凭证请求,限制页面数量并要求页序连续。
356fn validate_upload_files(files: &[EditorUploadFile]) -> ApiResult<()> {
357    if files.is_empty() {
358        return Err(AppError::business("没有待上传的页面喵"));
359    }
360    if files.len() > MAX_EDITOR_PAGE_COUNT {
361        return Err(AppError::business(format!(
362            "单话图源不能超过 {MAX_EDITOR_PAGE_COUNT} 页喵"
363        )));
364    }
365    let mut indices: Vec<i32> = files.iter().map(|file| file.index).collect();
366    indices.sort_unstable();
367    if indices
368        .iter()
369        .enumerate()
370        .any(|(expected, actual)| *actual != expected as i32)
371    {
372        return Err(AppError::business("图源页序必须从 0 开始且连续喵"));
373    }
374    if files
375        .iter()
376        .any(|file| file.filename.trim().is_empty() || file.filename.len() > 255)
377    {
378        return Err(AppError::business("图源文件名为空或过长喵"));
379    }
380    Ok(())
381}
382
383/// 校验并排序前端已上传的页面清单。
384fn validate_registered_pages(
385    episode_id: i32,
386    mut pages: Vec<crate::entity::editor::RegisterEditorPage>,
387) -> ApiResult<Vec<crate::entity::editor::RegisterEditorPage>> {
388    if pages.is_empty() {
389        return Err(AppError::business("没有可注册的页面喵"));
390    }
391    if pages.len() > MAX_EDITOR_PAGE_COUNT {
392        return Err(AppError::business(format!(
393            "单话图源不能超过 {MAX_EDITOR_PAGE_COUNT} 页喵"
394        )));
395    }
396    pages.sort_by_key(|page| page.page_index);
397    let expected_prefix = format!("editor/episode_{episode_id}/");
398    let mut object_keys = HashSet::with_capacity(pages.len());
399    for (index, page) in pages.iter().enumerate() {
400        if page.page_index != index as i32 {
401            return Err(AppError::business("图源页序必须从 0 开始且连续喵"));
402        }
403        if !page.object_key.starts_with(&expected_prefix)
404            || !object_keys.insert(page.object_key.as_str())
405        {
406            return Err(AppError::business("图源对象键无效或重复喵"));
407        }
408        if page
409            .source_file_name
410            .as_ref()
411            .is_some_and(|name| name.len() > 255)
412        {
413            return Err(AppError::business("图源文件名过长喵"));
414        }
415        if [page.image_width, page.image_height]
416            .into_iter()
417            .flatten()
418            .any(|dimension| !(1..=MAX_EDITOR_IMAGE_DIMENSION).contains(&dimension))
419        {
420            return Err(AppError::business("图源图片尺寸无效喵"));
421        }
422    }
423    Ok(pages)
424}
425
426/// 判断当前页面集是否与本次注册内容完全相同。
427fn page_sets_equal(
428    existing: &[crate::sea_entity::editor_page::Model],
429    incoming: &[NewEditorPage],
430) -> bool {
431    existing.len() == incoming.len()
432        && existing.iter().zip(incoming).all(|(current, next)| {
433            current.page_index == next.page_index
434                && current.image_object_key == next.image_object_key
435        })
436}
437
438/// 对页面模型 ID 排序,生成并发更新快照。
439fn sorted_page_ids(pages: &[crate::sea_entity::editor_page::Model]) -> Vec<i64> {
440    sorted_ids(&pages.iter().map(|page| page.id).collect::<Vec<_>>())
441}
442
443/// 返回排序后的 ID 副本。
444fn sorted_ids(ids: &[i64]) -> Vec<i64> {
445    let mut sorted = ids.to_vec();
446    sorted.sort_unstable();
447    sorted
448}
449
450/// 生成页面图片在图片 OSS 中的对象键
451fn build_image_object_key(episode_id: i32, index: usize, source_name: &str) -> String {
452    let ext = source_name
453        .rsplit('.')
454        .next()
455        .filter(|e| e.len() <= 5 && !e.contains('/'))
456        .map(|e| e.to_ascii_lowercase())
457        .unwrap_or_else(|| "jpg".to_string());
458    format!(
459        "editor/episode_{episode_id}/{:04}_{}.{ext}",
460        index,
461        Uuid::new_v4()
462    )
463}
464
465/// 通过预签名 PUT URL 上传单张图片
466async fn upload_image(state: &AppState, put_url: &str, img: &ExtractedImage) -> ApiResult<()> {
467    let content_type = mime_guess::from_path(&img.name)
468        .first_or_octet_stream()
469        .to_string();
470    let resp = state
471        .http_client
472        .put(put_url)
473        .header(CONTENT_TYPE, content_type)
474        .body(img.data.clone())
475        .send()
476        .await
477        .map_err(|e| AppError::Oss {
478            code: None,
479            msg: format!("上传图片到 OSS 失败:{e}"),
480        })?;
481    if !resp.status().is_success() {
482        return Err(AppError::Oss {
483            code: None,
484            msg: format!("上传图片到 OSS 返回非 2xx:HTTP {}", resp.status()),
485        });
486    }
487    Ok(())
488}
489
490/// 解析图片像素宽高(失败返回 None,不阻断导入)
491fn image_dimensions(bytes: &[u8]) -> (Option<i32>, Option<i32>) {
492    match imagesize::blob_size(bytes) {
493        Ok(size) => (Some(size.width as i32), Some(size.height as i32)),
494        Err(_) => (None, None),
495    }
496}
497
498#[cfg(test)]
499mod tests {
500    use super::*;
501    use crate::entity::editor::RegisterEditorPage;
502
503    fn upload_file(index: i32, filename: &str) -> EditorUploadFile {
504        EditorUploadFile {
505            index,
506            filename: filename.to_string(),
507        }
508    }
509
510    fn registered_page(episode_id: i32, page_index: i32, suffix: &str) -> RegisterEditorPage {
511        RegisterEditorPage {
512            page_index,
513            object_key: format!("editor/episode_{episode_id}/{suffix}.jpg"),
514            source_file_name: Some(format!("{suffix}.jpg")),
515            image_width: Some(1200),
516            image_height: Some(1800),
517        }
518    }
519
520    #[test]
521    fn upload_files_require_contiguous_unique_indices() {
522        assert!(
523            validate_upload_files(&[upload_file(0, "001.jpg"), upload_file(1, "002.jpg"),]).is_ok()
524        );
525        assert!(
526            validate_upload_files(&[upload_file(0, "001.jpg"), upload_file(0, "002.jpg"),])
527                .is_err()
528        );
529        assert!(validate_upload_files(&[upload_file(1, "002.jpg")]).is_err());
530    }
531
532    #[test]
533    fn registered_pages_are_sorted_and_scoped_to_episode() {
534        let pages = validate_registered_pages(
535            42,
536            vec![
537                registered_page(42, 1, "0001_b"),
538                registered_page(42, 0, "0000_a"),
539            ],
540        )
541        .expect("valid pages");
542        assert_eq!(
543            pages.iter().map(|page| page.page_index).collect::<Vec<_>>(),
544            vec![0, 1]
545        );
546
547        assert!(validate_registered_pages(42, vec![registered_page(7, 0, "0000_a")]).is_err());
548    }
549
550    #[test]
551    fn registered_pages_reject_duplicate_object_keys_and_invalid_dimensions() {
552        let duplicate = registered_page(42, 0, "same");
553        let mut duplicate_second = duplicate.clone();
554        duplicate_second.page_index = 1;
555        assert!(validate_registered_pages(42, vec![duplicate, duplicate_second]).is_err());
556
557        let mut invalid_dimension = registered_page(42, 0, "0000_a");
558        invalid_dimension.image_width = Some(0);
559        assert!(validate_registered_pages(42, vec![invalid_dimension]).is_err());
560    }
561
562    #[test]
563    fn sorted_ids_returns_snapshot_in_stable_order() {
564        assert_eq!(sorted_ids(&[9, 2, 5]), vec![2, 5, 9]);
565    }
566}