Skip to main content

tdm_server_rust/service/
task_tracking_service.rs

1//! 稿件监控服务 (Task Tracking Service)
2//!
3//! 组员任务看板的业务逻辑:
4//! - 组员任务数量统计(按岗位分类)
5//! - 全部待做稿件查询(带缓存)
6//! - 待发布/待审稿漫画分页
7
8use crate::{
9    app::AppState,
10    cache::{
11        get_episode_tasks_cached, get_episode_tasks_json_cached, get_member_task_counts_cached,
12        get_member_task_counts_json_cached, get_pending_manga_tasks_cached,
13        pending_manga_cache_key, set_episode_tasks_cached, set_member_task_counts_cached,
14        set_pending_manga_tasks_cached,
15    },
16    common::PageBean,
17    entity::episode::{
18        MangaEpisodeTb, MemberTaskCount, PendingEpisode, PendingMangaTask, TaskTrackingResponse,
19        WorkflowEpisodeItem,
20    },
21    error::ApiResult,
22    repository::{
23        task_tracking_repo::{PendingEpisodeRow, TaskTrackingRepository},
24        workflow_filter::WorkflowEpisodeQuery,
25    },
26    utils::{agent_debug, page::slice_rows},
27};
28use std::collections::HashMap;
29
30/// 稿件监控服务
31pub struct TaskTrackingService;
32
33impl TaskTrackingService {
34    /// 查询所有组员的任务计数(带缓存)。
35    ///
36    /// 优先返回缓存数据,缓存未命中时查库并写入缓存。
37    ///
38    /// # 返回值
39    ///
40    /// 返回所有组员的各岗位任务计数列表。
41    #[tracing::instrument(skip_all, level = "debug")]
42    pub async fn get_member_task_count_list(state: &AppState) -> ApiResult<Vec<MemberTaskCount>> {
43        if let Some(cached) = get_member_task_counts_cached(state).await {
44            return Ok(cached);
45        }
46        let repo = TaskTrackingRepository::new(state.db.clone());
47        let rows = repo.list_member_task_counts().await?;
48        let data = TaskTrackingRepository::to_member_task_counts(&rows);
49        set_member_task_counts_cached(state, data.clone()).await;
50        Ok(data)
51    }
52
53    /// 查询组员任务计数预序列化 JSON(缓存命中零拷贝)
54    #[tracing::instrument(skip_all, level = "debug")]
55    pub async fn get_member_task_counts_json(
56        state: &AppState,
57    ) -> ApiResult<std::sync::Arc<Vec<u8>>> {
58        if let Some(json) = get_member_task_counts_json_cached(state).await {
59            return Ok(json);
60        }
61        Self::get_member_task_count_list(state).await?;
62        get_member_task_counts_json_cached(state)
63            .await
64            .ok_or_else(|| crate::error::AppError::business("memberTaskCounts 缓存写入失败"))
65    }
66
67    /// 查询全部待做稿件(带缓存)。
68    ///
69    /// 返回按岗位分组的待处理话数,用于任务看板主视图。
70    #[tracing::instrument(skip_all, level = "debug")]
71    pub async fn get_all_task(state: &AppState) -> ApiResult<TaskTrackingResponse> {
72        if let Some(cached) = get_episode_tasks_cached(state).await {
73            return Ok(cached);
74        }
75        let repo = TaskTrackingRepository::new(state.db.clone());
76        let data = repo.list_episode_tasks_response().await?;
77        set_episode_tasks_cached(state, data.clone()).await;
78        Ok(data)
79    }
80
81    /// 查询待做稿件并返回预序列化 JSON(缓存命中零拷贝)
82    #[tracing::instrument(skip_all, level = "debug")]
83    pub async fn get_episode_tasks_json(state: &AppState) -> ApiResult<std::sync::Arc<Vec<u8>>> {
84        if let Some(json) = get_episode_tasks_json_cached(state).await {
85            return Ok(json);
86        }
87        Self::get_all_task(state).await?;
88        get_episode_tasks_json_cached(state)
89            .await
90            .ok_or_else(|| crate::error::AppError::business("episodeTasks 缓存写入失败"))
91    }
92
93    /// 待发布/待审稿漫画分页(对齐 Java TaskTrackingServiceImpl.getPendingMangaTasks)
94    #[tracing::instrument(skip_all, level = "debug")]
95    pub async fn get_pending_manga_tasks(
96        state: &AppState,
97        page: i32,
98        page_size: i32,
99        manga_tran_name: Option<String>,
100    ) -> ApiResult<PageBean<PendingMangaTask>> {
101        let cache_key = pending_manga_cache_key(page, page_size, manga_tran_name.as_deref());
102        if let Some(cached) = get_pending_manga_tasks_cached(state, &cache_key).await {
103            return Ok(cached);
104        }
105
106        let repo = TaskTrackingRepository::new(state.db.clone());
107        let name_ref = manga_tran_name.as_deref();
108        let (total, episodes) = tokio::try_join!(
109            repo.count_pending_publish_episodes(name_ref),
110            repo.list_pending_publish_episodes(name_ref, page, page_size),
111        )?;
112
113        // #region agent log
114        agent_debug::log(
115            "A",
116            "task_tracking_service.rs:get_pending_manga_tasks",
117            "pending episodes fetched",
118            serde_json::json!({
119                "total": total,
120                "pageCount": episodes.len(),
121                "sampleMangaIds": episodes.iter().take(3).map(|e| e.manga_id).collect::<Vec<_>>()
122            }),
123        );
124        // #endregion
125
126        let manga_ids: Vec<i32> = episodes
127            .iter()
128            .map(|e| e.manga_id)
129            .collect::<std::collections::HashSet<_>>()
130            .into_iter()
131            .collect();
132
133        let (latest, next_publish) = repo.map_publish_episode_context(&manga_ids).await?;
134
135        let mut grouped: HashMap<i32, Vec<&PendingEpisodeRow>> = HashMap::new();
136        let mut manga_order: Vec<i32> = Vec::new();
137        for ep in &episodes {
138            if !grouped.contains_key(&ep.manga_id) {
139                manga_order.push(ep.manga_id);
140            }
141            grouped.entry(ep.manga_id).or_default().push(ep);
142        }
143
144        let mut tasks: Vec<PendingMangaTask> = Vec::new();
145        for manga_id in manga_order {
146            let Some(rows) = grouped.get(&manga_id) else {
147                continue;
148            };
149            let first = rows[0];
150            let mut pending_list: Vec<PendingEpisode> = rows
151                .iter()
152                .map(|row| {
153                    let status = if row.reviewer_update_time.is_some() {
154                        "publisher".to_string()
155                    } else {
156                        "reviewer".to_string()
157                    };
158                    let next = next_publish
159                        .get(&row.manga_id)
160                        .map(|id| *id == row.episode_id)
161                        .unwrap_or(false);
162                    PendingEpisode {
163                        mangaepisodetb: row_to_episode_tb(row),
164                        status,
165                        next_publish: next,
166                    }
167                })
168                .collect();
169            pending_list.sort_by(|a, b| {
170                a.mangaepisodetb
171                    .manga_episode
172                    .cmp(&b.mangaepisodetb.manga_episode)
173            });
174
175            tasks.push(PendingMangaTask {
176                mangatb: TaskTrackingRepository::row_mangatb(first),
177                newest_manga_episode: latest.get(&manga_id).cloned(),
178                pending_episode_list: pending_list,
179            });
180        }
181
182        tasks.sort_by(|a, b| a.mangatb.manga_tran_name.cmp(&b.mangatb.manga_tran_name));
183
184        // #region agent log
185        agent_debug::log(
186            "A",
187            "task_tracking_service.rs:get_pending_manga_tasks",
188            "pending manga tasks built",
189            serde_json::json!({
190                "taskCount": tasks.len(),
191                "firstTaskEpisodes": tasks.first().map(|t| t.pending_episode_list.len())
192            }),
193        );
194        // #endregion
195
196        let result = slice_rows(tasks, total);
197        set_pending_manga_tasks_cached(state, cache_key, result.clone()).await;
198        Ok(result)
199    }
200
201    /// 工序汇总筛选分页(稿件监控 workflowEpisodes)
202    #[tracing::instrument(skip_all, level = "debug")]
203    pub async fn get_workflow_episodes(
204        state: &AppState,
205        query: WorkflowEpisodeQuery,
206    ) -> ApiResult<PageBean<WorkflowEpisodeItem>> {
207        let repo = TaskTrackingRepository::new(state.db.clone());
208        let (total, rows) = repo.page_workflow_episodes(&query).await?;
209        Ok(PageBean::new(total, rows))
210    }
211}
212
213/// PendingEpisodeRow 转 API 话数对象
214fn row_to_episode_tb(row: &PendingEpisodeRow) -> MangaEpisodeTb {
215    MangaEpisodeTb {
216        id: row.episode_id,
217        manga_id: row.manga_id,
218        manga_episode: row.manga_episode.clone(),
219        manga_episode_name: row.manga_episode_name.clone(),
220        episode_type: row.episode_type,
221        provider_id: row.provider_id,
222        translator_id: row.translator_id,
223        proofreader_id: row.proofreader_id,
224        letterer_id: row.letterer_id,
225        timer_id: row.timer_id,
226        reviewer_id: row.reviewer_id,
227        setup_time: row.setup_time,
228        update_time: row.update_time,
229        translator_file: row.translator_file.clone(),
230        proofreader_file: row.proofreader_file.clone(),
231        timer_file: row.timer_file.clone(),
232        publish_link: row.publish_link.clone(),
233        provider_file_oss_id: row.provider_file_oss_id,
234        translator_file_oss_id: row.translator_file_oss_id,
235        proofreader_file_oss_id: row.proofreader_file_oss_id,
236        letterer_file_oss_id: row.letterer_file_oss_id,
237        timer_file_oss_id: row.timer_file_oss_id,
238    }
239}
240
241#[cfg(test)]
242mod tests {
243    use super::*;
244    use crate::entity::episode::EpisodeType;
245
246    fn pending_row(episode_type: EpisodeType) -> PendingEpisodeRow {
247        PendingEpisodeRow {
248            episode_id: 1,
249            manga_id: 2,
250            manga_episode: Some("1".into()),
251            manga_episode_name: None,
252            episode_type,
253            provider_id: None,
254            translator_id: None,
255            proofreader_id: None,
256            letterer_id: None,
257            timer_id: None,
258            reviewer_id: None,
259            setup_time: None,
260            update_time: None,
261            translator_file: None,
262            proofreader_file: None,
263            timer_file: None,
264            publish_link: None,
265            provider_file_oss_id: None,
266            translator_file_oss_id: None,
267            proofreader_file_oss_id: None,
268            letterer_file_oss_id: None,
269            timer_file_oss_id: None,
270            reviewer_update_time: None,
271            manga_tran_name: None,
272            manga_ori_name: None,
273            category: None,
274            manga_status: None,
275            image: None,
276            manga_setup_time: None,
277            manga_update_time: None,
278            link: None,
279            introduction: None,
280        }
281    }
282
283    #[test]
284    fn pending_episode_mapping_preserves_episode_type() {
285        let main = row_to_episode_tb(&pending_row(EpisodeType::Main));
286        let extra = row_to_episode_tb(&pending_row(EpisodeType::Extra));
287
288        assert_eq!(main.episode_type, EpisodeType::Main);
289        assert_eq!(extra.episode_type, EpisodeType::Extra);
290        assert_eq!(serde_json::to_value(extra).unwrap()["episodeType"], "EXTRA");
291    }
292}