Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 36 additions & 21 deletions apps/api/routers/pipelines.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
54 changes: 54 additions & 0 deletions frontend/src/hooks/paper-detail/pollTask.ts
Original file line number Diff line number Diff line change
@@ -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<T>(
taskId: string,
isCancelled: () => boolean,
): Promise<T> {
const start = Date.now();
let consecutiveErrors = 0;

// 递归 setTimeout 轮询(避免 setInterval 难以取消 + 退避)
const poll = async (): Promise<T> => {
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();
}
10 changes: 7 additions & 3 deletions frontend/src/hooks/paper-detail/useAutoAnalyze.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import type {
DeepDiveReport,
SkimReport,
} from "@/types";
import { pollTaskUntilDone } from "./pollTask";

type Toast = (type: ToastType, message: string) => void;

Expand Down Expand Up @@ -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);
Expand All @@ -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<SkimReport>(task_id, () => false);
setSkimReport(r);
} catch {}
setSkimLoading(false);
Expand All @@ -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<DeepDiveReport>(task_id, () => false);
setDeepReport(r);
} catch {}
setDeepLoading(false);
Expand Down
24 changes: 18 additions & 6 deletions frontend/src/hooks/paper-detail/useDeepRead.ts
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -20,21 +21,31 @@ export function useDeepRead({ id, toast, setReportTab }: UseDeepReadParams) {
} | null>(null);
const [deepLoading, setDeepLoading] = useState(false);
const deepAbort = useRef<AbortController | null>(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<DeepDiveReport>(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,
Expand All @@ -44,6 +55,7 @@ export function useDeepRead({ id, toast, setReportTab }: UseDeepReadParams) {
deepLoading,
setDeepLoading,
deepAbort,
cancelDeep,
handleDeep,
};
}
28 changes: 21 additions & 7 deletions frontend/src/hooks/paper-detail/useEmbed.ts
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -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 };
}
24 changes: 19 additions & 5 deletions frontend/src/hooks/paper-detail/useSkim.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -22,25 +23,37 @@ export function useSkim({ id, toast, setReportTab, setPaper }: UseSkimParams) {
} | null>(null);
const [skimLoading, setSkimLoading] = useState(false);
const skimAbort = useRef<AbortController | null>(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<SkimReport>(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,
Expand All @@ -50,6 +63,7 @@ export function useSkim({ id, toast, setReportTab, setPaper }: UseSkimParams) {
skimLoading,
setSkimLoading,
skimAbort,
cancelSkim,
handleSkim,
};
}
18 changes: 9 additions & 9 deletions frontend/src/pages/PaperDetail.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,7 @@ export default function PaperDetail() {
skimLoading,
setSkimLoading,
skimAbort,
cancelSkim,
handleSkim,
} = useSkim({ id, toast, setReportTab, setPaper });

Expand All @@ -136,6 +137,7 @@ export default function PaperDetail() {
deepLoading,
setDeepLoading,
deepAbort,
cancelDeep,
handleDeep,
} = useDeepRead({ id, toast, setReportTab });

Expand All @@ -160,7 +162,7 @@ export default function PaperDetail() {
});

// 嵌入 hook
const { embedLoading, setEmbedLoading, handleEmbed } = useEmbed({
const { embedLoading, setEmbedLoading, handleEmbed, cancelEmbed } = useEmbed({
id,
toast,
embedDone,
Expand Down Expand Up @@ -691,19 +693,17 @@ export default function PaperDetail() {
<PipelineProgress
type="skim"
onCancel={() => {
skimAbort.current?.abort();
setSkimLoading(false);
toast("info", "已取消(后台可能仍在处理,结果稍后可在详情页查看)");
cancelSkim();
toast("info", "已取消(后台任务可能仍在处理)");
}}
/>
)}
{deepLoading && (
<PipelineProgress
type="deep"
onCancel={() => {
deepAbort.current?.abort();
setDeepLoading(false);
toast("info", "已取消(后台可能仍在处理,结果稍后可在详情页查看)");
cancelDeep();
toast("info", "已取消(后台任务可能仍在处理)");
}}
/>
)}
Expand All @@ -730,8 +730,8 @@ export default function PaperDetail() {
<PipelineProgress
type="embed"
onCancel={() => {
setEmbedLoading(false);
toast("info", "已取消(后台可能仍在处理)");
cancelEmbed();
toast("info", "已取消(后台任务可能仍在处理)");
}}
/>
)}
Expand Down
Loading