1use 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
33const MAX_EDITOR_PAGE_COUNT: usize = 500;
35const MAX_EDITOR_IMAGE_DIMENSION: i32 = 100_000;
37
38pub struct EditorImportService;
40
41impl EditorImportService {
42 #[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 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 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 let (filename, bytes) =
96 OssService::download_file_proxy(state, episode_id, "provider".to_string()).await?;
97
98 let images = archive_extract::extract_images(&filename, &bytes)?;
100
101 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 #[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 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 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 #[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 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
282async 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 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
338async 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
355fn 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
383fn 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
426fn 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
438fn 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
443fn sorted_ids(ids: &[i64]) -> Vec<i64> {
445 let mut sorted = ids.to_vec();
446 sorted.sort_unstable();
447 sorted
448}
449
450fn 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
465async 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
490fn 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}