1use 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
30pub struct TaskTrackingService;
32
33impl TaskTrackingService {
34 #[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 #[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 #[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 #[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 #[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 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 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 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 let result = slice_rows(tasks, total);
197 set_pending_manga_tasks_cached(state, cache_key, result.clone()).await;
198 Ok(result)
199 }
200
201 #[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
213fn 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}