diff --git a/apps/api/routers/pipelines.py b/apps/api/routers/pipelines.py index 421883c..e405d2f 100644 --- a/apps/api/routers/pipelines.py +++ b/apps/api/routers/pipelines.py @@ -24,44 +24,59 @@ @router.post("/pipelines/skim/{paper_id}") def run_skim(paper_id: UUID) -> dict: - tid = f"skim_{paper_id.hex[:8]}" + """粗读 — 后台任务化,立即返回 task_id(此前同步阻塞 5-30s 占请求线程)""" title = get_paper_title(paper_id) or str(paper_id)[:8] - global_tracker.start(tid, "skim", f"粗读:{title[:30]}", total=1, category="analysis") - try: + + def _fn(progress_callback=None): + if progress_callback: + progress_callback("正在粗读...", 30, 100) skim = pipelines.skim(paper_id) - global_tracker.finish(tid, success=True) + if progress_callback: + progress_callback("完成", 100, 100) return skim.model_dump() - except Exception as exc: - global_tracker.finish(tid, success=False, error=str(exc)[:100]) - raise + + task_id = global_tracker.submit( + "skim", f"粗读:{title[:30]}", _fn, total=100, category="analysis" + ) + return {"task_id": task_id, "status": "running"} @router.post("/pipelines/deep/{paper_id}") def run_deep(paper_id: UUID) -> dict: - tid = f"deep_{paper_id.hex[:8]}" + """精读 — 后台任务化,立即返回 task_id(此前同步阻塞 30s-2min 占请求线程)""" title = get_paper_title(paper_id) or str(paper_id)[:8] - global_tracker.start(tid, "deep_read", f"精读:{title[:30]}", total=1, category="analysis") - try: + + def _fn(progress_callback=None): + if progress_callback: + progress_callback("正在精读...", 20, 100) deep = pipelines.deep_dive(paper_id) - global_tracker.finish(tid, success=True) + if progress_callback: + progress_callback("完成", 100, 100) return deep.model_dump() - except Exception as exc: - global_tracker.finish(tid, success=False, error=str(exc)[:100]) - raise + + task_id = global_tracker.submit( + "deep_read", f"精读:{title[:30]}", _fn, total=100, category="analysis" + ) + return {"task_id": task_id, "status": "running"} @router.post("/pipelines/embed/{paper_id}") def run_embed(paper_id: UUID) -> dict: - tid = f"embed_{paper_id.hex[:8]}" + """嵌入 — 后台任务化,立即返回 task_id(此前同步阻塞 0.5-3s 占请求线程)""" title = get_paper_title(paper_id) or str(paper_id)[:8] - global_tracker.start(tid, "embed", f"嵌入:{title[:30]}", total=1, category="analysis") - try: + + def _fn(progress_callback=None): + if progress_callback: + progress_callback("正在计算向量嵌入...", 50, 100) pipelines.embed_paper(paper_id) - global_tracker.finish(tid, success=True) + if progress_callback: + progress_callback("完成", 100, 100) return {"status": "embedded", "paper_id": str(paper_id)} - except Exception as exc: - global_tracker.finish(tid, success=False, error=str(exc)[:100]) - raise + + task_id = global_tracker.submit( + "embed", f"嵌入:{title[:30]}", _fn, total=100, category="analysis" + ) + return {"task_id": task_id, "status": "running"} @router.get("/pipelines/runs") diff --git a/frontend/src/hooks/paper-detail/pollTask.ts b/frontend/src/hooks/paper-detail/pollTask.ts new file mode 100644 index 0000000..43cdbdd --- /dev/null +++ b/frontend/src/hooks/paper-detail/pollTask.ts @@ -0,0 +1,54 @@ +/** + * 后台任务轮询 helper — 配合 PR5 后台化的 skim/deep/embed 端点 + * 端点现在返回 {task_id, status},前端轮询 tasksApi.getStatus 直到 finished,再 getResult 取结果 + * @author Color2333 + */ +import { tasksApi } from "@/services/api"; +import type { TaskStatus } from "@/types"; + +const POLL_INTERVAL_MS = 2000; +const MAX_WAIT_MS = 5 * 60 * 1000; // 5 分钟超时 +const MAX_ERRORS = 10; // 连续查询失败上限 + +/** + * 轮询任务直到完成/失败/超时。 + * @returns 完成时返回 getResult 的结果;失败/超时抛错。 + */ +export async function pollTaskUntilDone( + taskId: string, + isCancelled: () => boolean, +): Promise { + const start = Date.now(); + let consecutiveErrors = 0; + + // 递归 setTimeout 轮询(避免 setInterval 难以取消 + 退避) + const poll = async (): Promise => { + if (isCancelled()) throw new Error("已取消"); + if (Date.now() - start > MAX_WAIT_MS) { + throw new Error("任务超时,请稍后在详情页查看结果"); + } + let status: TaskStatus; + try { + status = await tasksApi.getStatus(taskId); + consecutiveErrors = 0; + } catch { + consecutiveErrors += 1; + if (consecutiveErrors >= MAX_ERRORS) { + throw new Error("任务状态查询持续失败"); + } + await new Promise((r) => setTimeout(r, POLL_INTERVAL_MS)); + return poll(); + } + if (!status.finished) { + await new Promise((r) => setTimeout(r, POLL_INTERVAL_MS)); + return poll(); + } + if (!status.success) { + throw new Error(status.error || "任务失败"); + } + // 完成:取结果 + return (await tasksApi.getResult(taskId)) as T; + }; + + return poll(); +} diff --git a/frontend/src/hooks/paper-detail/useAutoAnalyze.ts b/frontend/src/hooks/paper-detail/useAutoAnalyze.ts index 1f2f2b7..bc828e8 100644 --- a/frontend/src/hooks/paper-detail/useAutoAnalyze.ts +++ b/frontend/src/hooks/paper-detail/useAutoAnalyze.ts @@ -9,6 +9,7 @@ import type { DeepDiveReport, SkimReport, } from "@/types"; +import { pollTaskUntilDone } from "./pollTask"; type Toast = (type: ToastType, message: string) => void; @@ -65,7 +66,8 @@ export function useAutoAnalyze({ setAutoStage("向量嵌入中..."); setEmbedLoading(true); try { - await pipelineApi.embed(id); + const { task_id } = await pipelineApi.embed(id); + await pollTaskUntilDone<{ status: string; paper_id: string }>(task_id, () => false); setEmbedDone(true); } catch {} setEmbedLoading(false); @@ -77,7 +79,8 @@ export function useAutoAnalyze({ setSkimLoading(true); setReportTab("skim"); try { - const r = await pipelineApi.skim(id); + const { task_id } = await pipelineApi.skim(id); + const r = await pollTaskUntilDone(task_id, () => false); setSkimReport(r); } catch {} setSkimLoading(false); @@ -90,7 +93,8 @@ export function useAutoAnalyze({ setDeepLoading(true); setReportTab("deep"); try { - const r = await pipelineApi.deep(id); + const { task_id } = await pipelineApi.deep(id); + const r = await pollTaskUntilDone(task_id, () => false); setDeepReport(r); } catch {} setDeepLoading(false); diff --git a/frontend/src/hooks/paper-detail/useDeepRead.ts b/frontend/src/hooks/paper-detail/useDeepRead.ts index 468837d..90743ef 100644 --- a/frontend/src/hooks/paper-detail/useDeepRead.ts +++ b/frontend/src/hooks/paper-detail/useDeepRead.ts @@ -1,8 +1,9 @@ -import { useState, useRef } from "react"; +import { useState, useRef, useCallback } from "react"; import type { Dispatch, SetStateAction } from "react"; import type { ToastType } from "@/contexts/ToastContext"; import { pipelineApi } from "@/services/api"; import type { DeepDiveReport } from "@/types"; +import { pollTaskUntilDone } from "./pollTask"; type Toast = (type: ToastType, message: string) => void; @@ -20,21 +21,31 @@ export function useDeepRead({ id, toast, setReportTab }: UseDeepReadParams) { } | null>(null); const [deepLoading, setDeepLoading] = useState(false); const deepAbort = useRef(null); + const cancelledRef = useRef(false); - const handleDeep = async () => { + const handleDeep = useCallback(async () => { if (!id) return; setDeepLoading(true); setReportTab("deep"); + cancelledRef.current = false; try { - const report = await pipelineApi.deep(id); + const { task_id } = await pipelineApi.deep(id); + const report = await pollTaskUntilDone(task_id, () => cancelledRef.current); setDeepReport(report); toast("success", "精读完成"); - } catch { - toast("error", "精读失败"); + } catch (e) { + if (!cancelledRef.current) { + toast("error", e instanceof Error ? e.message : "精读失败"); + } } finally { setDeepLoading(false); } - }; + }, [id, setReportTab, toast]); + + const cancelDeep = useCallback(() => { + cancelledRef.current = true; + setDeepLoading(false); + }, []); return { deepReport, @@ -44,6 +55,7 @@ export function useDeepRead({ id, toast, setReportTab }: UseDeepReadParams) { deepLoading, setDeepLoading, deepAbort, + cancelDeep, handleDeep, }; } diff --git a/frontend/src/hooks/paper-detail/useEmbed.ts b/frontend/src/hooks/paper-detail/useEmbed.ts index a975541..a194bb7 100644 --- a/frontend/src/hooks/paper-detail/useEmbed.ts +++ b/frontend/src/hooks/paper-detail/useEmbed.ts @@ -1,7 +1,8 @@ -import { useState } from "react"; +import { useState, useRef, useCallback } from "react"; import type { Dispatch, SetStateAction } from "react"; import type { ToastType } from "@/contexts/ToastContext"; import { pipelineApi } from "@/services/api"; +import { pollTaskUntilDone } from "./pollTask"; type Toast = (type: ToastType, message: string) => void; @@ -14,20 +15,33 @@ interface UseEmbedParams { export function useEmbed({ id, toast, setEmbedDone }: UseEmbedParams) { const [embedLoading, setEmbedLoading] = useState(false); + const cancelledRef = useRef(false); - const handleEmbed = async () => { + const handleEmbed = useCallback(async () => { if (!id) return; setEmbedLoading(true); + cancelledRef.current = false; try { - await pipelineApi.embed(id); + const { task_id } = await pipelineApi.embed(id); + await pollTaskUntilDone<{ status: string; paper_id: string }>( + task_id, + () => cancelledRef.current, + ); setEmbedDone(true); toast("success", "嵌入完成"); - } catch { - toast("error", "嵌入失败"); + } catch (e) { + if (!cancelledRef.current) { + toast("error", e instanceof Error ? e.message : "嵌入失败"); + } } finally { setEmbedLoading(false); } - }; + }, [id, toast, setEmbedDone]); - return { embedLoading, setEmbedLoading, handleEmbed }; + const cancelEmbed = useCallback(() => { + cancelledRef.current = true; + setEmbedLoading(false); + }, []); + + return { embedLoading, setEmbedLoading, handleEmbed, cancelEmbed }; } diff --git a/frontend/src/hooks/paper-detail/useSkim.ts b/frontend/src/hooks/paper-detail/useSkim.ts index f0edaa3..a07bd36 100644 --- a/frontend/src/hooks/paper-detail/useSkim.ts +++ b/frontend/src/hooks/paper-detail/useSkim.ts @@ -3,6 +3,7 @@ import type { Dispatch, SetStateAction } from "react"; import type { ToastType } from "@/contexts/ToastContext"; import { paperApi, pipelineApi } from "@/services/api"; import type { Paper, SkimReport } from "@/types"; +import { pollTaskUntilDone } from "./pollTask"; type Toast = (type: ToastType, message: string) => void; @@ -22,25 +23,37 @@ export function useSkim({ id, toast, setReportTab, setPaper }: UseSkimParams) { } | null>(null); const [skimLoading, setSkimLoading] = useState(false); const skimAbort = useRef(null); + const cancelledRef = useRef(false); - const handleSkim = async () => { + const handleSkim = useCallback(async () => { if (!id) return; setSkimLoading(true); setReportTab("skim"); + cancelledRef.current = false; try { - const report = await pipelineApi.skim(id); + // 后台任务化:端点返回 task_id,轮询直到完成再 getResult 取结果 + const { task_id } = await pipelineApi.skim(id); + const report = await pollTaskUntilDone(task_id, () => cancelledRef.current); setSkimReport(report); // 刷新论文信息,更新粗读报告 const updated = await paperApi.detail(id); setPaper(updated); if (updated.skim_report) setSavedSkim(updated.skim_report); toast("success", "粗读完成"); - } catch { - toast("error", "粗读失败"); + } catch (e) { + if (!cancelledRef.current) { + toast("error", e instanceof Error ? e.message : "粗读失败"); + } } finally { setSkimLoading(false); } - }; + }, [id, setReportTab, toast, setPaper]); + + // onCancel:标记取消,pollTaskUntilDone 下次检查时抛错退出 + const cancelSkim = useCallback(() => { + cancelledRef.current = true; + setSkimLoading(false); + }, []); return { skimReport, @@ -50,6 +63,7 @@ export function useSkim({ id, toast, setReportTab, setPaper }: UseSkimParams) { skimLoading, setSkimLoading, skimAbort, + cancelSkim, handleSkim, }; } diff --git a/frontend/src/pages/PaperDetail.tsx b/frontend/src/pages/PaperDetail.tsx index 8bb81fb..0a1261a 100644 --- a/frontend/src/pages/PaperDetail.tsx +++ b/frontend/src/pages/PaperDetail.tsx @@ -124,6 +124,7 @@ export default function PaperDetail() { skimLoading, setSkimLoading, skimAbort, + cancelSkim, handleSkim, } = useSkim({ id, toast, setReportTab, setPaper }); @@ -136,6 +137,7 @@ export default function PaperDetail() { deepLoading, setDeepLoading, deepAbort, + cancelDeep, handleDeep, } = useDeepRead({ id, toast, setReportTab }); @@ -160,7 +162,7 @@ export default function PaperDetail() { }); // 嵌入 hook - const { embedLoading, setEmbedLoading, handleEmbed } = useEmbed({ + const { embedLoading, setEmbedLoading, handleEmbed, cancelEmbed } = useEmbed({ id, toast, embedDone, @@ -691,9 +693,8 @@ export default function PaperDetail() { { - skimAbort.current?.abort(); - setSkimLoading(false); - toast("info", "已取消(后台可能仍在处理,结果稍后可在详情页查看)"); + cancelSkim(); + toast("info", "已取消(后台任务可能仍在处理)"); }} /> )} @@ -701,9 +702,8 @@ export default function PaperDetail() { { - deepAbort.current?.abort(); - setDeepLoading(false); - toast("info", "已取消(后台可能仍在处理,结果稍后可在详情页查看)"); + cancelDeep(); + toast("info", "已取消(后台任务可能仍在处理)"); }} /> )} @@ -730,8 +730,8 @@ export default function PaperDetail() { { - setEmbedLoading(false); - toast("info", "已取消(后台可能仍在处理)"); + cancelEmbed(); + toast("info", "已取消(后台任务可能仍在处理)"); }} /> )} diff --git a/frontend/src/services/api.ts b/frontend/src/services/api.ts index 8805f8b..95c2038 100644 --- a/frontend/src/services/api.ts +++ b/frontend/src/services/api.ts @@ -439,10 +439,11 @@ export const translateApi = { /* ========== Pipeline ========== */ export const pipelineApi = { - skim: (paperId: string) => post(`/pipelines/skim/${paperId}`), - deep: (paperId: string) => post(`/pipelines/deep/${paperId}`), + // 后台任务化:返回 task_id,前端轮询 tasksApi.getStatus + getResult 取结果 + skim: (paperId: string) => post<{ task_id: string; status: string }>(`/pipelines/skim/${paperId}`), + deep: (paperId: string) => post<{ task_id: string; status: string }>(`/pipelines/deep/${paperId}`), embed: (paperId: string) => - post<{ status: string; paper_id: string }>(`/pipelines/embed/${paperId}`), + post<{ task_id: string; status: string }>(`/pipelines/embed/${paperId}`), runs: (limit = 30) => get<{ items: PipelineRun[] }>(`/pipelines/runs?limit=${limit}`), }; diff --git a/frontend/src/types/index.ts b/frontend/src/types/index.ts index df74a55..3501746 100644 --- a/frontend/src/types/index.ts +++ b/frontend/src/types/index.ts @@ -1010,11 +1010,14 @@ export interface TaskStatus { task_type: string; title: string; status: "pending" | "running" | "completed" | "failed"; - progress: number; + progress: number; // 0-1 小数 + progress_pct?: number; // 0-100 百分比 message: string; error: string | null; created_at: number; - updated_at: number; + finished: boolean; // 后端 to_dict 真实字段 + success: boolean; + elapsed_seconds?: number; has_result: boolean; } diff --git a/packages/domain/task_tracker.py b/packages/domain/task_tracker.py index 2d6e049..05c4372 100644 --- a/packages/domain/task_tracker.py +++ b/packages/domain/task_tracker.py @@ -45,6 +45,12 @@ class TaskInfo: def to_dict(self) -> dict: elapsed = time.time() - self.started_at + # status 便利字段:前端 TaskStatus.status 期望 "pending"|"running"|"completed"|"failed" + # 此前前端读 status.status 拿 undefined,靠 finished/success 判断的轮询点失配 + if not self.finished: + status = "running" if self.current > 0 else "pending" + else: + status = "completed" if self.success else "failed" return { "task_id": self.task_id, "task_type": self.task_type, @@ -56,8 +62,12 @@ def to_dict(self) -> dict: "created_at": self.created_at, "elapsed_seconds": round(elapsed, 1), "progress_pct": round(self.current / self.total * 100) if self.total > 0 else 0, + "progress": round(self.current / self.total, 4) + if self.total > 0 + else 0, # 0-1 小数(前端期望) "finished": self.finished, "success": self.success, + "status": status, # 便利字段,由 finished/success 派生 "error": self.error, "has_result": self.result is not None, }