/** 异步 Job 研究执行(Phase 4):POST /api/jobs 提交 → 轮询 GET /api/jobs/{id}。 * * 大样本研究(全市场)可能耗时数十秒到分钟级,经异步 Job 后台执行, * 避免 HTTP 长阻塞(AGENT §19)。页面提交后即时返回 job_id,再轮询到终态。 */ import { useCallback, useEffect, useRef, useState } from "react"; import { apiGet, apiPost } from "./api"; import type { ResearchSpec } from "./types"; export interface JobSubmit { job_id: string; status: string; } export interface JobStatusResp { job_id: string; status: string; /** 执行阶段(data_loading / backtesting / analysis…)—— 用于给用户真实进度而非假进度条 */ stage?: string | null; /** 归档后的实验 id(后端在 Job 完成时写入) */ experiment_id?: string | null; result?: unknown; error?: string | null; spec?: Record; started_at?: string | null; finished_at?: string | null; } export function submitJob(spec: ResearchSpec): Promise { return apiPost("/jobs", spec); } /** 任务终态(success / failed / cancelled / timeout)与结果 */ export interface JobOutcome { status: string; result: T | null; error?: string; /** 归档实验 id:结果出来后前端可直接给出「去对比」入口 */ experimentId?: string | null; } /** 轮询直到 success / failed / cancelled,或超时(默认 10 分钟)。 */ export async function waitJob( jobId: string, timeoutMs = 600_000, onStage?: (info: { stage: string; elapsedMs: number }) => void, ): Promise> { const deadline = Date.now() + timeoutMs; const t0 = Date.now(); while (Date.now() < deadline) { const job = await apiGet(`/jobs/${jobId}`); // 阶段来自后端真实执行状态,避免前端用假进度条假装在跑 if (onStage && job.stage) onStage({ stage: job.stage, elapsedMs: Date.now() - t0 }); if (job.status === "success") { return { status: job.status, result: (job.result as T) ?? null, experimentId: job.experiment_id ?? null, }; } if (job.status === "failed" || job.status === "cancelled") { return { status: job.status, result: null, error: job.error ?? undefined }; } await new Promise((resolve) => setTimeout(resolve, 1500)); } return { status: "timeout", result: null, error: "等待结果超时,请稍后在「实验」页查看归档" }; } /** 阶段中文名(后端 stage 枚举 → 用户可读) */ export const STAGE_LABEL: Record = { queued: "排队中", data_loading: "加载行情与因子数据", selection: "逐择股日选股", backtesting: "撮合与净值结算", analysis: "汇总指标与曲线", done: "完成", }; /** 执行阶段时间线(顺序即推进顺序,用于「阶段小圆点」) */ export const STAGE_ORDER = ["queued", "data_loading", "selection", "backtesting", "analysis", "done"]; /** 取消排队中/执行中的作业(POST /api/jobs/{id}/cancel;已结束的会返回 cancelled=false)。 */ export function cancelJob( jobId: string, ): Promise<{ job_id: string; status: string; cancelled: boolean }> { return apiPost<{ job_id: string; status: string; cancelled: boolean }>( `/jobs/${encodeURIComponent(jobId)}/cancel`, {}, ); } /** 提交一个回测组合为异步 Job(POST /api/combos/run,不保存组合)。 */ export async function submitComboJob(combo: unknown): Promise { return apiPost("/combos/run", combo); } /** * 运行**已保存**的回测组合(POST /api/combos/{id}/run)。 * * 与 submitComboJob 的区别:这里用库里的那份参数,页面上未保存的改动不参与 —— * 「从组合库直接运行」必须跑库里存的那套,否则用户改了一半的表单会污染既有组合的结果。 */ export async function runSavedCombo(comboId: string): Promise { return apiPost(`/combos/${encodeURIComponent(comboId)}/run`, {}); } /* ------------------------------------------------------------------ * * 统一的「提交 → 轮询 → 终态」状态机(useJobRunner) * * 为什么要有它:以前每个页面各写一遍 running/jobId/error,结果参差不齐 —— * 有的页面点了「运行」只在按钮上转圈,用户不知道到底提交没提交、跑到哪一步; * 有的页面「已用秒数」只在后端阶段变化时才更新,看起来像卡死了(0s 一动不动)。 * 这里把状态、秒表、取消、终态收敛到一处,页面只管渲染。 * ------------------------------------------------------------------ */ export type JobPhase = | "idle" | "submitting" | "queued" | "running" | "success" | "failed" | "cancelled" | "timeout"; export interface JobRunState { phase: JobPhase; /** 作业号(提交成功后立即拿到;这是「到底提交没提交」的凭据) */ jobId: string; /** 后端上报的执行阶段(queued/data_loading/selection/backtesting/analysis/done) */ stage: string; /** 已用毫秒:**每秒自增**,不依赖后端阶段变化,所以界面不会看起来卡住 */ elapsedMs: number; /** 这次跑的是什么(如「因子测试 · momentum_60」),显示在反馈条上 */ label: string; /** 批量场景的进度(第 index / total 个),单跑时为 null */ progress: { index: number; total: number } | null; error: string; result: T | null; experimentId: string | null; /** 提交时刻(本地时间,用于「发起于 hh:mm:ss」) */ startedAt: Date | null; /** 正在请求取消 */ cancelling: boolean; } function emptyState(label = ""): JobRunState { return { phase: "idle", jobId: "", stage: "", elapsedMs: 0, label, progress: null, error: "", result: null, experimentId: null, startedAt: null, cancelling: false, }; } export interface JobRunOptions { /** 反馈条上显示的任务名(如「因子测试 · momentum_60」) */ label?: string; /** 批量进度(第 index / total 个) */ progress?: { index: number; total: number } | null; timeoutMs?: number; } export interface JobRunner { state: JobRunState; /** 提交并轮询到终态;返回终态结果(供页面继续处理,如渲染报告) */ run: ( submit: () => Promise<{ job_id: string }>, opts?: JobRunOptions, ) => Promise>; cancel: () => Promise; reset: () => void; } /** * 作业运行状态机。 * * 生命周期:idle → submitting(已点击,正在换 job_id)→ queued/running(轮询中,秒表走) * → success | failed | cancelled | timeout。**submitting 也是状态**:点了按钮立刻就有 * 反馈文字,不会出现「点了没反应」的空白期。 */ export function useJobRunner(defaultLabel = ""): JobRunner { const [state, setState] = useState>(() => emptyState(defaultLabel)); const t0 = useRef(0); const timer = useRef(null); const alive = useRef(true); const currentJob = useRef(""); // 后端最近上报的阶段放在 ref 里:终态那一帧要写「最后一个阶段」,但不想让 // `run` 依赖 state(否则每次 state 变化都换一个新函数,页面 effect 会连环触发)。 const stageRef = useRef(""); const stopTick = useCallback(() => { if (timer.current !== null) { window.clearInterval(timer.current); timer.current = null; } }, []); useEffect( () => () => { alive.current = false; if (timer.current !== null) window.clearInterval(timer.current); }, [], ); const patch = useCallback((next: Partial>) => { if (!alive.current) return; setState((prev) => ({ ...prev, ...next })); }, []); const reset = useCallback(() => { stopTick(); currentJob.current = ""; if (alive.current) setState(emptyState(defaultLabel)); }, [defaultLabel, stopTick]); const run = useCallback( async ( submit: () => Promise<{ job_id: string }>, opts: JobRunOptions = {}, ): Promise> => { stopTick(); currentJob.current = ""; stageRef.current = ""; t0.current = Date.now(); setState({ ...emptyState(opts.label ?? defaultLabel), phase: "submitting", label: opts.label ?? defaultLabel, progress: opts.progress ?? null, startedAt: new Date(), }); // 秒表独立于后端阶段:即使后端几秒内不报新阶段,「已用 Xs」也在动 timer.current = window.setInterval(() => { if (!alive.current) return; setState((prev) => prev.phase === "submitting" || prev.phase === "queued" || prev.phase === "running" ? { ...prev, elapsedMs: Date.now() - t0.current } : prev, ); }, 500); try { const { job_id } = await submit(); currentJob.current = job_id; stageRef.current = "queued"; patch({ jobId: job_id, phase: "queued", stage: "queued" }); const out = await waitJob(job_id, opts.timeoutMs ?? 900_000, (info) => { stageRef.current = info.stage; patch({ stage: info.stage, elapsedMs: info.elapsedMs, phase: info.stage && info.stage !== "queued" ? "running" : "queued", }); }); const finalPhase: JobPhase = out.status === "success" ? "success" : out.status === "cancelled" ? "cancelled" : out.status === "timeout" ? "timeout" : "failed"; patch({ phase: finalPhase, stage: out.status === "success" ? "done" : stageRef.current, elapsedMs: Date.now() - t0.current, result: out.result, experimentId: out.experimentId ?? null, error: out.error ?? "", }); return out; } catch (e) { const msg = (e as Error).message; patch({ phase: "failed", error: msg, elapsedMs: Date.now() - t0.current }); return { status: "failed", result: null, error: msg }; } finally { stopTick(); } }, [defaultLabel, patch, stopTick], ); const cancel = useCallback(async () => { const jobId = currentJob.current; if (!jobId) return; patch({ cancelling: true }); try { await cancelJob(jobId); // 取消后轮询会在 1.5s 内看到 cancelled 终态;这里先把阶段文字改掉, // 让用户马上知道「取消请求已发出」而不是等下一个轮询周期。 patch({ stage: "cancelling" }); } finally { patch({ cancelling: false }); } }, [patch]); return { state, run, cancel, reset }; }