Skip to main content

tdm_server_rust/service/
episode_service.rs

1//! 话数业务服务 (Episode Service)
2//!
3//! 话数管理的核心业务逻辑:
4//! - 话数 CRUD(创建、批量创建、更新、删除)
5//! - 发布链接管理
6//! - 上传页数据查询
7//! - 统计信息(漫画统计、组员统计)
8//! - 最新话查询
9
10use crate::{
11    app::AppState,
12    cache::{
13        get_list_page_json_cached, get_member_statistics_json_cached,
14        get_statistics_list_json_cached, manga_episode_page_cache_key,
15        manga_episode_simple_cache_key, member_statistics_cache_key, on_episode_mutated,
16        set_list_page_json_cached, set_member_statistics_json_cached, set_result_json_cached,
17        set_statistics_list_json_cached, uploaded_submit_cache_key,
18    },
19    common::PageBean,
20    entity::enums::{MemberInternEnum, PostEnum},
21    entity::episode::{
22        EpisodeDetailVo, EpisodeEditDto, EpisodeListVo, EpisodeSimpleListVo, EpisodeType,
23        MemberStatistics, NewestEpisodeVo, PublishLinkRequest, Statistics, UploadPageVo,
24    },
25    entity::manga::EpisodeDownloadRequest,
26    error::{ApiResult, AppError},
27    repository::{
28        episode_repo::EpisodeRepository, member_repo::MemberRepository, oss_repo::OssRepository,
29    },
30    service::{
31        oss_service::OssService,
32        rss_service::{RssRefreshScope, RssService},
33    },
34    utils::page::slice_rows,
35};
36use axum::body::Bytes;
37use chrono::{DateTime, Utc};
38use std::collections::HashSet;
39use std::path::{Path, PathBuf};
40use std::sync::Arc;
41use std::time::Duration;
42
43/// 话数服务
44///
45/// 提供话数及关联数据(统计、发布)的完整业务逻辑。
46pub struct EpisodeService;
47
48impl EpisodeService {
49    /// 按漫画 ID 查询话数简单列表。
50    ///
51    /// 返回包含文件信息的精简话数列表,用于上传页展示。
52    /// 结果按 `manga_episode` 升序排列,仅包含未删除的话数记录。
53    ///
54    /// # 返回值
55    ///
56    /// 返回漫画下所有未删除的话数列表,按 `manga_episode` 排序。
57    ///
58    /// # Errors
59    ///
60    /// - `AppError::Database` — 数据库查询失败
61    ///
62    /// # Examples
63    ///
64    /// ```ignore
65    /// use tdm_server::service::episode_service::EpisodeService;
66    ///
67    /// let episodes = EpisodeService::select_list_by_manga_id(&state, manga_id).await?;
68    /// for ep in &episodes {
69    ///     println!("话数 {}: {}", ep.manga_episode, ep.manga_episode_name);
70    /// }
71    /// ```
72    #[tracing::instrument(skip_all, level = "debug")]
73    pub async fn select_list_by_manga_id(
74        state: &AppState,
75        manga_id: i32,
76    ) -> ApiResult<Vec<EpisodeSimpleListVo>> {
77        EpisodeRepository::new(state.db.clone())
78            .list_by_manga_id(manga_id)
79            .await
80    }
81
82    /// 漫画话数简单列表预序列化 JSON(短 TTL 缓存)
83    #[tracing::instrument(skip_all, level = "debug")]
84    pub async fn select_list_by_manga_id_json(
85        state: &AppState,
86        manga_id: i32,
87    ) -> ApiResult<Arc<Vec<u8>>> {
88        let key = manga_episode_simple_cache_key(state, manga_id);
89        if let Some(json) = get_list_page_json_cached(state, &key).await {
90            return Ok(json);
91        }
92        let list = Self::select_list_by_manga_id(state, manga_id).await?;
93        set_result_json_cached(state, key.clone(), list).await;
94        get_list_page_json_cached(state, &key)
95            .await
96            .ok_or_else(|| AppError::business("话数简单列表缓存写入失败"))
97    }
98
99    /// 批量更新话数发布链接。
100    ///
101    /// 接收一组发布链接请求,批量更新数据库中的链接字段。
102    /// 更新完成后会刷新任务追踪缓存,并为每个话数触发 RSS 订阅源刷新。
103    ///
104    /// # Errors
105    ///
106    /// - `AppError::Database` — 数据库更新失败
107    ///
108    /// # Examples
109    ///
110    /// ```ignore
111    /// let requests = vec![
112    ///     PublishLinkRequest { id: 1, publish_link: Some("https://example.com/1".into()) },
113    ///     PublishLinkRequest { id: 2, publish_link: Some("https://example.com/2".into()) },
114    /// ];
115    /// EpisodeService::update_manga_episodes(&state, requests).await?;
116    /// ```
117    #[tracing::instrument(skip_all, level = "debug")]
118    pub async fn update_manga_episodes(
119        state: &AppState,
120        requests: Vec<PublishLinkRequest>,
121    ) -> ApiResult<()> {
122        EpisodeRepository::new(state.db.clone())
123            .update_publish_links_batch(&requests)
124            .await?;
125        let mut manga_ids = HashSet::new();
126        let repo = EpisodeRepository::new(state.db.clone());
127        for id in requests.iter().map(|r| r.id) {
128            if let Ok(Some(manga_id)) = repo.get_manga_id_by_episode_id(id).await {
129                manga_ids.insert(manga_id);
130            }
131            RssService::refresh(state, RssRefreshScope::PublishLink { episode_id: id });
132        }
133        for manga_id in manga_ids {
134            on_episode_mutated(state, manga_id, &[]).await;
135        }
136        Ok(())
137    }
138
139    /// 更新单条话数的发布链接。
140    ///
141    /// 在写入前会检查新链接是否与其他话数冲突。更新成功后同样
142    /// 会刷新任务追踪缓存和 RSS 订阅源。
143    ///
144    /// # Errors
145    ///
146    /// - `AppError::business("发布链接已存在喵")` — 链接与其他话数冲突
147    /// - `AppError::Database` — 数据库操作失败
148    ///
149    /// # Examples
150    ///
151    /// ```ignore
152    /// EpisodeService::update_publish_link(
153    ///     &state,
154    ///     42,
155    ///     Some("https://example.com/ch42".to_string()),
156    /// ).await?;
157    /// ```
158    #[tracing::instrument(skip_all, level = "debug")]
159    pub async fn update_publish_link(
160        state: &AppState,
161        id: i32,
162        publish_link: Option<String>,
163    ) -> ApiResult<()> {
164        let repo = EpisodeRepository::new(state.db.clone());
165        if let Some(ref link) = publish_link {
166            if !link.is_empty() && repo.count_publish_link(link, Some(id)).await? > 0 {
167                return Err(AppError::business("发布链接已存在喵"));
168            }
169        }
170        let dto = EpisodeEditDto {
171            id: Some(id),
172            manga_id: None,
173            manga_episode: None,
174            manga_episode_end: None,
175            manga_episode_name: None,
176            episode_type: None,
177            provider_id: None,
178            translator_id: None,
179            proofreader_id: None,
180            letterer_id: None,
181            timer_id: None,
182            reviewer_id: None,
183            publish_link,
184        };
185        repo.update_publish_link_only(id, dto.publish_link).await?;
186        if let Some(manga_id) = repo.get_manga_id_by_episode_id(id).await? {
187            on_episode_mutated(state, manga_id, &[]).await;
188        }
189        RssService::refresh(state, RssRefreshScope::PublishLink { episode_id: id });
190        Ok(())
191    }
192
193    /// 分页查询漫画话数。
194    ///
195    /// 获取指定漫画下的全部话数列表,然后在内存中进行分页截取。
196    /// 返回标准分页对象 `PageBean`,包含当前页数据、总条数和总页数。
197    ///
198    /// # Errors
199    ///
200    /// - `AppError::Database` — 数据库查询失败
201    ///
202    /// # Examples
203    ///
204    /// ```ignore
205    /// // 查询第 1 页,每页 20 条
206    /// let page = EpisodeService::page_episode(&state, 1, 20, manga_id).await?;
207    /// println!("共 {} 条,当前页 {} 条", page.total, page.rows.len());
208    /// ```
209    #[tracing::instrument(skip_all, level = "debug")]
210    pub async fn page_episode(
211        state: &AppState,
212        page: i32,
213        page_size: i32,
214        manga_id: i32,
215        episode_type: EpisodeType,
216    ) -> ApiResult<PageBean<EpisodeListVo>> {
217        let (total, rows) = EpisodeRepository::new(state.db.clone())
218            .page_episodes_by_manga(manga_id, page, page_size, episode_type)
219            .await?;
220        Ok(slice_rows(rows, total))
221    }
222
223    /// 漫画话数分页预序列化 JSON(短 TTL 缓存)
224    #[tracing::instrument(skip_all, level = "debug")]
225    pub async fn page_episode_json(
226        state: &AppState,
227        page: i32,
228        page_size: i32,
229        manga_id: i32,
230        episode_type: EpisodeType,
231    ) -> ApiResult<Arc<Vec<u8>>> {
232        let key = manga_episode_page_cache_key(
233            state,
234            manga_id,
235            page,
236            page_size,
237            episode_type.as_str(),
238        );
239        if let Some(json) = get_list_page_json_cached(state, &key).await {
240            return Ok(json);
241        }
242        let page_data =
243            Self::page_episode(state, page, page_size, manga_id, episode_type).await?;
244        set_list_page_json_cached(state, key.clone(), page_data).await;
245        get_list_page_json_cached(state, &key)
246            .await
247            .ok_or_else(|| AppError::business("话数分页缓存写入失败"))
248    }
249
250    /// 删除指定话数(软删除)。
251    ///
252    /// 将话数标记为删除状态,同时刷新任务追踪缓存并触发 RSS 订阅源
253    /// 全量刷新。注意:此操作不会立即从数据库中物理删除数据。
254    ///
255    /// # Errors
256    ///
257    /// - `AppError::Database` — 数据库操作失败
258    ///
259    /// # Examples
260    ///
261    /// ```ignore
262    /// EpisodeService::delete_episode(&state, episode_id).await?;
263    /// ```
264    #[tracing::instrument(skip_all, level = "debug")]
265    pub async fn delete_episode(state: &AppState, id: i32) -> ApiResult<()> {
266        let repo = EpisodeRepository::new(state.db.clone());
267        let member_ids = repo.get_post_member_ids(id).await?;
268        let manga_id = repo.get_manga_id_by_episode_id(id).await?;
269        repo.delete_by_id(id).await?;
270        if let Some(manga_id) = manga_id {
271            on_episode_mutated(state, manga_id, &member_ids).await;
272        }
273        RssService::refresh(state, RssRefreshScope::EpisodePipeline);
274        Ok(())
275    }
276
277    /// 新增话数(对齐 Java `addEpisodes`)。
278    ///
279    /// 支持两种模式:
280    /// - **批量模式**:指定 `manga_episode` 起始值和 `manga_episode_end` 结束值,
281    ///   在区间内逐条创建话数。
282    /// - **单话模式**:仅指定 `manga_episode`,创建一条话数记录。
283    ///
284    /// 创建完成后会刷新任务追踪缓存并触发 RSS 全量刷新。
285    ///
286    /// # Errors
287    ///
288    /// - `AppError::business("缺少漫画 ID")` — DTO 中未提供 `manga_id`
289    /// - `AppError::business("起始话数格式不正确")` — 话数序号不是有效整数
290    /// - `AppError::business("结束话数格式不正确")` — 结束序号不是有效整数
291    /// - `AppError::business("该漫画单话已存在喵!")` — 话数序号重复
292    /// - `AppError::business("发布链接已存在喵")` — 链接与其他话数冲突
293    /// - `AppError::Database` — 数据库操作失败
294    ///
295    /// # Examples
296    ///
297    /// ```ignore
298    /// // 批量创建话数 1~10
299    /// let dto = EpisodeEditDto {
300    ///     manga_id: Some(100),
301    ///     manga_episode: Some("1".into()),
302    ///     manga_episode_end: Some("10".into()),
303    ///     ..Default::default()
304    /// };
305    /// EpisodeService::add_episodes(&state, dto).await?;
306    ///
307    /// // 创建单话
308    /// let dto = EpisodeEditDto {
309    ///     manga_id: Some(100),
310    ///     manga_episode: Some("11".into()),
311    ///     ..Default::default()
312    /// };
313    /// EpisodeService::add_episodes(&state, dto).await?;
314    /// ```
315    #[tracing::instrument(skip_all, level = "debug")]
316    pub async fn add_episodes(state: &AppState, dto: EpisodeEditDto) -> ApiResult<()> {
317        let manga_id = dto
318            .manga_id
319            .ok_or_else(|| AppError::business("缺少漫画 ID"))?;
320        let repo = EpisodeRepository::new(state.db.clone());
321
322        let end = dto
323            .manga_episode_end
324            .as_deref()
325            .map(str::trim)
326            .filter(|s| !s.is_empty());
327
328        if let Some(end_str) = end {
329            Self::validate_batch_options(&dto)?;
330            let start_str = dto.manga_episode.as_deref().unwrap_or("0");
331            let start: i32 = start_str
332                .parse()
333                .map_err(|_| AppError::business("起始话数格式不正确"))?;
334            let end_num: i32 = end_str
335                .parse()
336                .map_err(|_| AppError::business("结束话数格式不正确"))?;
337            if start > end_num {
338                return Err(AppError::business("结束话数不能小于起始话数"));
339            }
340
341            let episode_type = dto.episode_type.unwrap_or_default();
342            let mut batch = Vec::with_capacity((end_num - start + 1) as usize);
343            for i in start..=end_num {
344                let mut one = dto.clone();
345                one.manga_episode = Some(i.to_string());
346                one.manga_episode_end = None;
347                if repo
348                    .count_episode_by_number(manga_id, &i.to_string(), episode_type)
349                    .await?
350                    > 0
351                {
352                    return Err(AppError::business(format!(
353                        "该漫画第 {i} 话已存在喵!"
354                    )));
355                }
356                batch.push(one);
357            }
358            repo.insert_batch(&batch, manga_id).await?;
359        } else {
360            Self::add_single_episode(&repo, &dto, manga_id).await?;
361        }
362        on_episode_mutated(state, manga_id, &[]).await;
363        RssService::refresh(state, RssRefreshScope::EpisodePipeline);
364        Ok(())
365    }
366
367    fn validate_batch_options(dto: &EpisodeEditDto) -> ApiResult<()> {
368        if dto.episode_type.unwrap_or_default() == EpisodeType::Extra {
369            return Err(AppError::business("番外和特典不支持批量添加"));
370        }
371        if dto
372            .publish_link
373            .as_deref()
374            .is_some_and(|link| !link.is_empty())
375        {
376            return Err(AppError::business("批量添加不支持设置发布链接"));
377        }
378        Ok(())
379    }
380
381    /// 创建单条话数记录(内部方法)。
382    ///
383    /// 执行话数新增的核心逻辑:
384    /// 1. 校验话数序号是否与已有记录重复
385    /// 2. 校验发布链接是否与其他话数冲突
386    /// 3. 插入主表记录和详情记录
387    /// 4. 更新漫画的 `update_time` 时间戳
388    ///
389    /// # Errors
390    ///
391    /// - `AppError::business("缺少话数序号")` — `manga_episode` 为空
392    /// - `AppError::business("该漫画单话已存在喵!")` — 序号重复
393    /// - `AppError::business("发布链接已存在喵")` — 链接冲突
394    /// - `AppError::Database` — 数据库操作失败
395    #[tracing::instrument(skip_all, level = "debug")]
396    async fn add_single_episode(
397        repo: &EpisodeRepository,
398        dto: &EpisodeEditDto,
399        manga_id: i32,
400    ) -> ApiResult<()> {
401        let episode_label = dto
402            .manga_episode
403            .as_deref()
404            .ok_or_else(|| AppError::business("缺少话数序号"))?;
405        let episode_type = dto.episode_type.unwrap_or_default();
406        if repo
407            .count_episode_by_number(manga_id, episode_label, episode_type)
408            .await?
409            > 0
410        {
411            return Err(AppError::business("该漫画单话已存在喵!"));
412        }
413        if let Some(ref link) = dto.publish_link {
414            if !link.is_empty() && repo.count_publish_link(link, None).await? > 0 {
415                return Err(AppError::business("发布链接已存在喵"));
416            }
417        }
418        let episode_id = repo.insert(dto).await?;
419        repo.insert_detail(episode_id, dto).await?;
420        repo.touch_manga_update_time(manga_id).await?;
421        Ok(())
422    }
423
424    /// 按 ID 查询单条话数详情。
425    ///
426    /// 返回话数的完整信息,包括主表字段和详情字段。
427    /// 如果话数已被软删除或不存在,返回业务错误。
428    ///
429    /// # Errors
430    ///
431    /// - `AppError::business("话数不存在喵")` — 指定 ID 的话数不存在
432    /// - `AppError::Database` — 数据库查询失败
433    ///
434    /// # Examples
435    ///
436    /// ```ignore
437    /// let episode = EpisodeService::get_manga_episode_by_id(&state, 42).await?;
438    /// println!("话数名称: {}", episode.manga_episode_name);
439    /// ```
440    #[tracing::instrument(skip_all, level = "debug")]
441    pub async fn get_manga_episode_by_id(state: &AppState, id: i32) -> ApiResult<EpisodeDetailVo> {
442        EpisodeRepository::new(state.db.clone())
443            .get_by_id(id)
444            .await?
445            .ok_or_else(|| AppError::business("话数不存在喵"))
446    }
447
448    /// 查询漫画最新话。
449    ///
450    /// 获取指定漫画下序号最大的话数记录。常用于展示连载漫画的
451    /// 最近更新情况。
452    ///
453    /// # Errors
454    ///
455    /// - `AppError::business("该漫画尚无单话喵")` — 漫画下没有话数记录
456    /// - `AppError::Database` — 数据库查询失败
457    ///
458    /// # Examples
459    ///
460    /// ```ignore
461    /// let newest = EpisodeService::get_newest_manga_episode_by_id(&state, manga_id).await?;
462    /// println!("最新话: {}", newest.manga_episode);
463    /// ```
464    #[tracing::instrument(skip_all, level = "debug")]
465    pub async fn get_newest_manga_episode_by_id(
466        state: &AppState,
467        manga_id: i32,
468    ) -> ApiResult<NewestEpisodeVo> {
469        EpisodeRepository::new(state.db.clone())
470            .get_newest_by_manga_id(manga_id)
471            .await?
472            .ok_or_else(|| AppError::business("该漫画尚无单话喵"))
473    }
474
475    /// 更新话数信息。
476    ///
477    /// 根据 DTO 中提供的字段更新话数主表。更新完成后刷新任务追踪
478    /// 缓存并触发 RSS 全量刷新。
479    ///
480    /// # Errors
481    ///
482    /// - `AppError::Database` — 数据库更新失败
483    ///
484    /// # Examples
485    ///
486    /// ```ignore
487    /// let dto = EpisodeEditDto {
488    ///     id: Some(42),
489    ///     manga_episode_name: Some("新标题".into()),
490    ///     ..Default::default()
491    /// };
492    /// EpisodeService::update_manga_episode(&state, dto).await?;
493    /// ```
494    #[tracing::instrument(skip_all, level = "debug")]
495    pub async fn update_manga_episode(state: &AppState, dto: EpisodeEditDto) -> ApiResult<()> {
496        let should_refresh_rss = should_refresh_rss_for_episode_edit(&dto);
497        let repo = EpisodeRepository::new(state.db.clone());
498        let episode_id = dto.id;
499        let mut member_ids = HashSet::new();
500        let manga_id_before = if let Some(id) = episode_id {
501            for mid in repo.get_post_member_ids(id).await? {
502                member_ids.insert(mid);
503            }
504            repo.get_manga_id_by_episode_id(id).await?
505        } else {
506            None
507        };
508        repo.update(&dto).await?;
509        if let Some(id) = episode_id {
510            for mid in repo.get_post_member_ids(id).await? {
511                member_ids.insert(mid);
512            }
513        }
514        if let Some(manga_id) = dto.manga_id.or(manga_id_before) {
515            let ids: Vec<i32> = member_ids.into_iter().collect();
516            on_episode_mutated(state, manga_id, &ids).await;
517        }
518        if should_refresh_rss {
519            RssService::refresh(state, RssRefreshScope::EpisodePipeline);
520        }
521        Ok(())
522    }
523
524    /// 分页查询已上传稿件列表。
525    ///
526    /// 在数据库层面执行分页查询和计数,支持按漫画译名和用户名筛选。
527    /// 返回标准分页对象 `PageBean<UploadPageVo>`。
528    ///
529    /// # Errors
530    ///
531    /// - `AppError::Database` — 数据库查询失败
532    ///
533    /// # Examples
534    ///
535    /// ```ignore
536    /// // 按译名筛选,查询第 1 页
537    /// let page = EpisodeService::get_uploaded_submit(
538    ///     &state, 1, 20,
539    ///     Some("进击的巨人".into()),
540    ///     None,
541    /// ).await?;
542    /// ```
543    #[tracing::instrument(skip_all, level = "debug")]
544    pub async fn get_uploaded_submit(
545        state: &AppState,
546        page: i32,
547        page_size: i32,
548        manga_tran_name: Option<String>,
549        username: Option<String>,
550    ) -> ApiResult<PageBean<UploadPageVo>> {
551        let (total, rows) = EpisodeRepository::new(state.db.clone())
552            .page_uploaded_submit(
553                page,
554                page_size,
555                manga_tran_name.as_deref(),
556                username.as_deref(),
557            )
558            .await?;
559        Ok(slice_rows(rows, total))
560    }
561
562    /// 藏宝处分页预序列化 JSON(短 TTL 缓存)
563    #[tracing::instrument(skip_all, level = "debug")]
564    pub async fn get_uploaded_submit_json(
565        state: &AppState,
566        page: i32,
567        page_size: i32,
568        manga_tran_name: Option<String>,
569        username: Option<String>,
570    ) -> ApiResult<Arc<Vec<u8>>> {
571        let key = uploaded_submit_cache_key(
572            state,
573            page,
574            page_size,
575            manga_tran_name.as_deref(),
576            username.as_deref(),
577        );
578        if let Some(json) = get_list_page_json_cached(state, &key).await {
579            return Ok(json);
580        }
581        let page_data =
582            Self::get_uploaded_submit(state, page, page_size, manga_tran_name, username).await?;
583        set_list_page_json_cached(state, key.clone(), page_data).await;
584        get_list_page_json_cached(state, &key)
585            .await
586            .ok_or_else(|| AppError::business("藏宝处缓存写入失败"))
587    }
588    #[tracing::instrument(skip_all, level = "debug")]
589    pub async fn get_statistics(state: &AppState, start: &str, end: &str) -> ApiResult<Statistics> {
590        let start = DateTime::parse_from_rfc3339(start)
591            .map_err(|e| AppError::business(format!("时间格式错误: {e}")))?
592            .with_timezone(&Utc);
593        let end = DateTime::parse_from_rfc3339(end)
594            .map_err(|e| AppError::business(format!("时间格式错误: {e}")))?
595            .with_timezone(&Utc);
596        let repo = EpisodeRepository::new(state.db.clone());
597        let stats = repo.get_statistic_count(start, end).await?;
598        let result = EpisodeRepository::stats_from_post(&stats);
599        // #region agent log
600        crate::utils::agent_debug::log(
601            "H1",
602            "episode_service.rs:get_statistics",
603            "statistics_any",
604            serde_json::json!({
605                "translatorCount": result.translator_count,
606                "proofreaderCount": result.proofreader_count
607            }),
608        );
609        // #endregion
610        Ok(result)
611    }
612
613    /// 预设时间段统计列表(6 项:日/月/年 + 上期对比,UTC 日界)
614    #[tracing::instrument(skip_all, level = "debug")]
615    pub async fn get_statistics_list(state: &AppState) -> ApiResult<Vec<Statistics>> {
616        use chrono::{Datelike, TimeZone, Utc};
617        let now = Utc::now();
618        let today_start = Utc
619            .with_ymd_and_hms(now.year(), now.month(), now.day(), 0, 0, 0)
620            .single()
621            .ok_or_else(|| AppError::business("日期无效"))?;
622        let yesterday_start = today_start - chrono::Duration::days(1);
623        let first_day_of_month = Utc
624            .with_ymd_and_hms(now.year(), now.month(), 1, 0, 0, 0)
625            .single()
626            .ok_or_else(|| AppError::business("日期无效"))?;
627        let first_day_of_last_month = if now.month() == 1 {
628            Utc.with_ymd_and_hms(now.year() - 1, 12, 1, 0, 0, 0)
629                .single()
630                .ok_or_else(|| AppError::business("日期无效"))?
631        } else {
632            Utc.with_ymd_and_hms(now.year(), now.month() - 1, 1, 0, 0, 0)
633                .single()
634                .ok_or_else(|| AppError::business("日期无效"))?
635        };
636        let first_day_of_year = Utc
637            .with_ymd_and_hms(now.year(), 1, 1, 0, 0, 0)
638            .single()
639            .ok_or_else(|| AppError::business("日期无效"))?;
640        let first_day_of_last_year = Utc
641            .with_ymd_and_hms(now.year() - 1, 1, 1, 0, 0, 0)
642            .single()
643            .ok_or_else(|| AppError::business("日期无效"))?;
644
645        let ranges = [
646            (today_start, now),
647            (first_day_of_month, now),
648            (first_day_of_year, now),
649            (yesterday_start, today_start),
650            (first_day_of_last_month, first_day_of_month),
651            (first_day_of_last_year, first_day_of_year),
652        ];
653        let repo = EpisodeRepository::new(state.db.clone());
654        let stats = repo.get_statistic_counts_batch(&ranges).await?;
655        let list: Vec<Statistics> = stats
656            .iter()
657            .map(EpisodeRepository::stats_from_post)
658            .collect();
659        // #region agent log
660        crate::utils::agent_debug::log(
661            "H1",
662            "episode_service.rs:get_statistics_list",
663            "statistics_list",
664            serde_json::json!({
665                "len": list.len(),
666                "dayTranslatorCount": list.first().map(|s| s.translator_count)
667            }),
668        );
669        // #endregion
670        Ok(list)
671    }
672
673    /// 预设统计列表预序列化 JSON(短 TTL 缓存)
674    #[tracing::instrument(skip_all, level = "debug")]
675    pub async fn get_statistics_list_json(state: &AppState) -> ApiResult<std::sync::Arc<Vec<u8>>> {
676        if let Some(json) = get_statistics_list_json_cached(state).await {
677            return Ok(json);
678        }
679        let list = Self::get_statistics_list(state).await?;
680        set_statistics_list_json_cached(state, list).await;
681        get_statistics_list_json_cached(state)
682            .await
683            .ok_or_else(|| AppError::business("statistics 缓存写入失败"))
684    }
685
686    /// 后台预热 statistics 列表 JSON(启动立即执行,之后按间隔刷新)
687    pub fn spawn_statistics_warmup(state: AppState) {
688        tokio::spawn(async move {
689            let interval_secs = std::env::var("STATISTICS_WARMUP_SECS")
690                .ok()
691                .and_then(|v| v.parse().ok())
692                .unwrap_or(30);
693            loop {
694                if let Err(e) = Self::get_statistics_list_json(&state).await {
695                    tracing::warn!("statistics 预热失败: {e:?}");
696                }
697                tokio::time::sleep(Duration::from_secs(interval_secs)).await;
698            }
699        });
700    }
701
702    /// 组员完成统计
703    #[tracing::instrument(skip_all, level = "debug")]
704    pub async fn get_member_statistics(
705        state: &AppState,
706        start: DateTime<Utc>,
707        end: DateTime<Utc>,
708    ) -> ApiResult<Vec<MemberStatistics>> {
709        EpisodeRepository::new(state.db.clone())
710            .get_member_statistics(start, end, None)
711            .await
712    }
713
714    /// 组员完成统计预序列化 JSON(短 TTL 缓存)
715    #[tracing::instrument(skip_all, level = "debug")]
716    pub async fn get_member_statistics_json(
717        state: &AppState,
718        start: DateTime<Utc>,
719        end: DateTime<Utc>,
720    ) -> ApiResult<std::sync::Arc<Vec<u8>>> {
721        let key = member_statistics_cache_key(state, &start.to_rfc3339(), &end.to_rfc3339());
722        if let Some(json) = get_member_statistics_json_cached(state, &key).await {
723            return Ok(json);
724        }
725        let list = Self::get_member_statistics(state, start, end).await?;
726        set_member_statistics_json_cached(state, key.clone(), list).await;
727        get_member_statistics_json_cached(state, &key)
728            .await
729            .ok_or_else(|| AppError::business("memberStatistics 缓存写入失败"))
730    }
731
732    /// 回退流程
733    #[tracing::instrument(skip_all, level = "debug")]
734    pub async fn rollback_episode(
735        state: &AppState,
736        episode_id: i32,
737        workflow_type: &str,
738        member_id: Option<i32>,
739    ) -> ApiResult<()> {
740        let repo = EpisodeRepository::new(state.db.clone());
741        let mut member_ids: HashSet<i32> = repo
742            .get_post_member_ids(episode_id)
743            .await?
744            .into_iter()
745            .collect();
746        if let Some(mid) = member_id {
747            member_ids.insert(mid);
748        }
749        let manga_id = repo.get_manga_id_by_episode_id(episode_id).await?;
750        if let Some(oss_id) = repo
751            .rollback_episode(episode_id, workflow_type, member_id)
752            .await?
753        {
754            let cfg = (*state.config).clone();
755            let _ = OssRepository::new(state.db.clone(), cfg)
756                .delete_oss(oss_id)
757                .await;
758        }
759        for mid in repo.get_post_member_ids(episode_id).await? {
760            member_ids.insert(mid);
761        }
762        if let Some(manga_id) = manga_id {
763            let ids: Vec<i32> = member_ids.into_iter().collect();
764            on_episode_mutated(state, manga_id, &ids).await;
765        }
766        RssService::refresh(state, RssRefreshScope::EpisodePipeline);
767        Ok(())
768    }
769
770    /// 上传话数工作文件(对齐 Java MangaFileUtils + 更新 DB 路径)
771    #[tracing::instrument(skip_all, level = "debug")]
772    pub async fn upload_manga_file(
773        state: &AppState,
774        data: Bytes,
775        episode_id: i32,
776        my_name: &str,
777        manga_id: i32,
778        original_filename: &str,
779    ) -> ApiResult<()> {
780        let post_id = legacy_post_id(my_name)
781            .ok_or_else(|| AppError::business(format!("当前岗位文件不存在喵:{my_name}")))?;
782        if data.len() > 1024 * 1024 {
783            return Err(AppError::business("上传的文件不能大于1MB喵!"));
784        }
785        let filename = sanitize_filename(original_filename);
786        let repo = EpisodeRepository::new(state.db.clone());
787        if repo.get_by_id(episode_id).await?.is_none() {
788            return Err(AppError::business(format!("漫画单话 {episode_id} 不存在")));
789        }
790        let folder = PathBuf::from(&state.config.folder.base2)
791            .join(manga_id.to_string())
792            .join(episode_id.to_string())
793            .join(post_id.to_string());
794        tokio::fs::create_dir_all(&folder)
795            .await
796            .map_err(|e| AppError::Internal(e.to_string()))?;
797        if let Some(old) = repo.get_legacy_file_path(episode_id, my_name).await? {
798            let _ = tokio::fs::remove_file(old).await;
799        }
800        let file_path = folder.join(&filename);
801        tokio::fs::write(&file_path, data)
802            .await
803            .map_err(|e| AppError::business(format!("文件上传失败喵: {e}")))?;
804        let stored = file_path.to_string_lossy().replace('\\', "/");
805        repo.set_legacy_file_path(episode_id, my_name, &stored)
806            .await?;
807        on_episode_mutated(state, manga_id, &[]).await;
808        Ok(())
809    }
810
811    /// 下载话数工作文件(OSS 优先 → legacy 路径 → 规范目录)
812    #[tracing::instrument(skip_all, level = "debug")]
813    pub async fn download_episode_file(
814        state: &AppState,
815        req: EpisodeDownloadRequest,
816    ) -> ApiResult<(String, Bytes)> {
817        let member = MemberRepository::new(state.db.clone())
818            .get_by_id(req.member_id)
819            .await?;
820        if !MemberInternEnum::is_working_member(member.intern) {
821            return Err(AppError::download_unauth(
822                "目前的职阶无法下载喵!先联系管理员改职阶喵!",
823            ));
824        }
825        if req.my_name.eq_ignore_ascii_case("translator") {
826            let post_ids = &member.post_ids;
827            let is_letterer = post_ids.contains(&PostEnum::LETTERER);
828            let has_translation_access = post_ids.contains(&PostEnum::TRANSLATOR)
829                || post_ids.contains(&PostEnum::PROOFREADER)
830                || post_ids.contains(&PostEnum::REVIEWER);
831            if is_letterer && !has_translation_access {
832                return Err(AppError::download_unauth(
833                    "不!要!下载翻译稿!嵌字要下载校对稿!如果真的需要下载翻译稿请找Gum979",
834                ));
835            }
836        }
837
838        let repo = EpisodeRepository::new(state.db.clone());
839        if repo.get_by_id(req.id).await?.is_none() {
840            return Err(AppError::business(format!("漫画单话 {} 不存在", req.id)));
841        }
842
843        let oss_repo = OssRepository::new(state.db.clone(), (*state.config).clone());
844        if should_download_oss_first(&req.my_name)
845            && oss_repo
846                .episode_post_oss_id(req.id, &req.my_name)
847                .await?
848                .is_some()
849        {
850            return OssService::download_file_proxy(state, req.id, req.my_name.clone()).await;
851        }
852
853        let legacy_path = repo.get_legacy_file_path(req.id, &req.my_name).await?;
854        if let Some(stored_path) = legacy_path.as_deref() {
855            if let Some((filename, bytes)) =
856                read_stored_legacy(&state.config.folder.base2, stored_path).await
857            {
858                return Ok((filename, bytes));
859            }
860        }
861
862        if let Some((filename, bytes)) = read_canonical_legacy(
863            &state.config.folder.base2,
864            req.manga_id,
865            req.id,
866            &req.my_name,
867        )
868        .await
869        {
870            return Ok((filename, bytes));
871        }
872
873        if legacy_path.is_some() {
874            let name = legacy_path
875                .as_deref()
876                .and_then(|p| Path::new(p).file_name())
877                .and_then(|s| s.to_str())
878                .unwrap_or("download.dat");
879            return Err(AppError::download_failed(format!(
880                "文件不存在或无法读取: {name}"
881            )));
882        }
883
884        OssService::download_file_proxy(state, req.id, req.my_name.clone())
885            .await
886            .map_err(|e| match e {
887                AppError::DownloadUnAuth { .. } => e,
888                _ => AppError::download_unauth("文件不存在或无权限喵"),
889            })
890    }
891}
892
893/// 判断岗位下载是否需要优先使用 OSS/COS 文件。
894fn should_download_oss_first(post_name: &str) -> bool {
895    matches!(
896        post_name.to_lowercase().as_str(),
897        "provider" | "translator" | "proofreader" | "letterer" | "timer"
898    )
899}
900
901/// 判断话数编辑是否需要触发 RSS feed 刷新。
902///
903/// 话号、标题、发布链接、各岗位指派变更均刷新 XML;稳定 GUID 不会重复推送。
904fn should_refresh_rss_for_episode_edit(dto: &EpisodeEditDto) -> bool {
905    dto.manga_episode
906        .as_deref()
907        .is_some_and(|s| !s.trim().is_empty())
908        || dto
909            .manga_episode_name
910            .as_deref()
911            .is_some_and(|s| !s.trim().is_empty())
912        || dto
913            .publish_link
914            .as_deref()
915            .is_some_and(|s| !s.trim().is_empty())
916        || dto.provider_id.is_some()
917        || dto.translator_id.is_some()
918        || dto.proofreader_id.is_some()
919        || dto.letterer_id.is_some()
920        || dto.timer_id.is_some()
921        || dto.reviewer_id.is_some()
922}
923
924/// 岗位名 → legacy 目录 postId
925fn legacy_post_id(post_name: &str) -> Option<i32> {
926    match post_name.to_lowercase().as_str() {
927        "translator" => Some(PostEnum::TRANSLATOR),
928        "proofreader" => Some(PostEnum::PROOFREADER),
929        "timer" => Some(PostEnum::TIMER),
930        _ => None,
931    }
932}
933
934/// 清理上传文件名
935fn sanitize_filename(name: &str) -> String {
936    let base = Path::new(name)
937        .file_name()
938        .and_then(|s| s.to_str())
939        .unwrap_or("upload.dat");
940    if base.is_empty() || base.contains("..") {
941        "upload.dat".to_string()
942    } else {
943        base.to_string()
944    }
945}
946
947/// 生产环境 legacy 根目录前缀(DB 中存绝对路径,本地 dev 需映射到 folder.base2)
948const LEGACY_PROD_PREFIXES: &[&str] = &["/www/wwwroot/data", "/www/wwwroot/tdm/data"];
949
950/// 将 DB 中的 legacy 路径展开为候选本地路径(原路径 + base2 映射)
951fn legacy_path_candidates(base2: &str, stored: &str) -> Vec<PathBuf> {
952    let mut out = Vec::new();
953    let mut push = |p: PathBuf| {
954        if !out.iter().any(|x| x == &p) {
955            out.push(p);
956        }
957    };
958    push(PathBuf::from(stored));
959    for prefix in LEGACY_PROD_PREFIXES {
960        if let Some(suffix) = stored.strip_prefix(prefix) {
961            push(PathBuf::from(base2).join(suffix.trim_start_matches('/')));
962        }
963    }
964    out
965}
966
967/// 将 folder.base2 解析为可读绝对路径
968fn resolve_folder_base2(base2: &str) -> PathBuf {
969    let path = PathBuf::from(base2);
970    if path.is_absolute() {
971        path
972    } else {
973        PathBuf::from(env!("CARGO_MANIFEST_DIR")).join(path)
974    }
975}
976
977/// 按候选路径读取 DB 记录的 legacy 稿件
978#[tracing::instrument(skip_all, level = "debug")]
979async fn read_stored_legacy(base2: &str, stored: &str) -> Option<(String, Bytes)> {
980    let base2 = resolve_folder_base2(base2);
981    for path in legacy_path_candidates(&base2.to_string_lossy(), stored) {
982        let Ok(bytes) = tokio::fs::read(&path).await else {
983            continue;
984        };
985        let filename = path
986            .file_name()
987            .and_then(|s| s.to_str())
988            .unwrap_or("download.dat")
989            .to_string();
990        return Some((filename, Bytes::from(bytes)));
991    }
992    None
993}
994
995/// 读取 Java 规范目录下的首个稿件文件
996#[tracing::instrument(skip_all, level = "debug")]
997async fn read_canonical_legacy(
998    base2: &str,
999    manga_id: i32,
1000    episode_id: i32,
1001    post_name: &str,
1002) -> Option<(String, Bytes)> {
1003    let post_id = legacy_post_id(post_name)?;
1004    let dir = resolve_folder_base2(base2)
1005        .join(manga_id.to_string())
1006        .join(episode_id.to_string())
1007        .join(post_id.to_string());
1008    let mut entries = tokio::fs::read_dir(&dir).await.ok()?;
1009    while let Ok(Some(entry)) = entries.next_entry().await {
1010        let path = entry.path();
1011        if !path.is_file() {
1012            continue;
1013        }
1014        let bytes = tokio::fs::read(&path).await.ok()?;
1015        let filename = path.file_name()?.to_str()?.to_string();
1016        return Some((filename, Bytes::from(bytes)));
1017    }
1018    None
1019}
1020
1021#[cfg(test)]
1022mod rss_trigger_tests {
1023    use super::*;
1024
1025    /// 构造仅包含 ID 的话数编辑 DTO。
1026    fn edit_dto() -> EpisodeEditDto {
1027        EpisodeEditDto {
1028            id: Some(1),
1029            manga_id: Some(10),
1030            manga_episode: None,
1031            manga_episode_end: None,
1032            manga_episode_name: None,
1033            episode_type: None,
1034            provider_id: None,
1035            translator_id: None,
1036            proofreader_id: None,
1037            letterer_id: None,
1038            timer_id: None,
1039            reviewer_id: None,
1040            publish_link: None,
1041        }
1042    }
1043
1044    /// 岗位改派需刷新 RSS feed 内容(稳定 GUID 不重复推送)。
1045    #[test]
1046    fn post_assignment_edit_refreshes_rss() {
1047        let dto = EpisodeEditDto {
1048            provider_id: Some(1),
1049            translator_id: Some(2),
1050            proofreader_id: Some(3),
1051            letterer_id: Some(4),
1052            timer_id: Some(5),
1053            reviewer_id: Some(6),
1054            ..edit_dto()
1055        };
1056
1057        assert!(should_refresh_rss_for_episode_edit(&dto));
1058    }
1059
1060    /// 非岗位字段编辑需要保留 RSS 通知。
1061    #[test]
1062    fn non_post_edit_refreshes_rss() {
1063        let episode = EpisodeEditDto {
1064            manga_episode: Some("12".into()),
1065            ..edit_dto()
1066        };
1067        let name = EpisodeEditDto {
1068            manga_episode_name: Some("新标题".into()),
1069            ..edit_dto()
1070        };
1071        let link = EpisodeEditDto {
1072            publish_link: Some("https://example.com/1".into()),
1073            ..edit_dto()
1074        };
1075
1076        assert!(should_refresh_rss_for_episode_edit(&episode));
1077        assert!(should_refresh_rss_for_episode_edit(&name));
1078        assert!(should_refresh_rss_for_episode_edit(&link));
1079    }
1080
1081    #[test]
1082    fn batch_add_rejects_extra_type_and_publish_link() {
1083        let extra = EpisodeEditDto {
1084            episode_type: Some(EpisodeType::Extra),
1085            ..edit_dto()
1086        };
1087        let linked = EpisodeEditDto {
1088            episode_type: Some(EpisodeType::Main),
1089            publish_link: Some("https://example.com/batch".into()),
1090            ..edit_dto()
1091        };
1092        let empty_link = EpisodeEditDto {
1093            episode_type: Some(EpisodeType::Main),
1094            publish_link: Some("  ".into()),
1095            ..edit_dto()
1096        };
1097
1098        assert!(EpisodeService::validate_batch_options(&extra).is_err());
1099        assert!(EpisodeService::validate_batch_options(&linked).is_err());
1100        assert!(EpisodeService::validate_batch_options(&empty_link).is_err());
1101    }
1102}
1103
1104#[cfg(test)]
1105mod download_legacy_tests {
1106    use super::*;
1107
1108    /// 校验稿件下载岗位优先使用 OSS/COS。
1109    #[test]
1110    fn submitted_posts_prefer_oss_download() {
1111        assert!(should_download_oss_first("translator"));
1112        assert!(should_download_oss_first("proofreader"));
1113        assert!(should_download_oss_first("letterer"));
1114        assert!(should_download_oss_first("provider"));
1115        assert!(should_download_oss_first("timer"));
1116        assert!(!should_download_oss_first("reviewer"));
1117    }
1118
1119    /// 校验生产 legacy 路径可映射到 fixture 目录并读取
1120    #[tokio::test]
1121    async fn read_remapped_proofreader_file() {
1122        let base2 = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/legacy/data");
1123        let stored = "/www/wwwroot/data/33/肚子咕噜咕噜子她和肉肉在同居!04 翻译:Gum979 校对:萝莉控之魂.txt";
1124        assert!(
1125            read_stored_legacy(base2.to_string_lossy().as_ref(), stored)
1126                .await
1127                .is_some(),
1128            "remapped legacy file should be readable under tests/fixtures/legacy/data"
1129        );
1130    }
1131}