接着 /backtest 的那次改造,把其余会跑异步 Job 的页面也切到同一套
`useJobRunner` + `JobProgress`(用户原话:包括因子测试等所有测试都帮我完善用户反馈):
- **/factors**(一个因子一个 Job 的批量场景):删掉 `running/processed/current/jobId` 与
按「已处理数 ÷ 总数」自算的 `Progress` **假百分比**;每轮把 `第 i/共 n 个` 交给反馈条,
阶段/作业号/已用秒数全部来自后端。取消 = 用户明确意图 → 保留已出的报告、不再跑后续因子,
不报错;单个因子失败仍继续跑其余因子,并把后端原文累积成**批级清单**(多因子时反馈条
只能显示最后一个作业,前几个失败不能丢)。
- **/factors/compose**:删掉 `value={30}` 的假进度条;`archiveId` / 复用 `BacktestResultView` /
「新页面放大」全部保留;failed/cancelled 交给反馈条,页面 error 只留参数与目录错误
(同一失败不在两处各说一遍)。
- **/selection**:异步选股切反馈条;**同步**的「执行选股」保留原 loading(`POST /selections`
没有 job_id,套上会去 `GET /jobs/{signal}` 撞 404)。
- **/signals**:`POST /api/signals` 是同步接口,**不套**作业反馈(不编作业号、不编阶段),
改为点击即现的 `role="status"` 提示,如实写明「同步请求、请求期间不能关页、无阶段无取消、
出错显示后端原文」。
- **阶段圆点按作业类型区分**(修掉一个真实缺陷):原来全站共用一张含 `queued/done` 的
`STAGE_ORDER`,因子测试页实测出现过**裸英文** `factor_calculation` 且 4 个圆点全灰
(`indexOf` = -1),还画出了因子测试根本不存在的「逐择股日选股 / 撮合与净值结算」。
现在 `STAGE_PIPELINES = { factor_test: [加载→计算因子值→汇总], backtest: [加载→撮合→汇总],
selection: [逐择股日选股] }`(阶段序列**不含 queued/done**,那是作业状态不是阶段),
未知阶段显示「执行中(stage)」并保留已推进的圆点,不整排灰、不露裸枚举。
验证(真实浏览器 CDP,读数原文已记录):
- /factors:点击后 0.2s 内 `submitting`→`queued` + 作业号;+30s Pill「计算因子值」、
圆点 `✓ 加载行情与因子数据 / ● 计算因子值 / ○ 汇总指标与曲线`(无「选股/撮合」、无英文枚举);
成功态给「去对比 / 打开归档」。
- /selection:圆点只有 `逐择股日选股` 一段;取消 → 「已取消,没有归档」;
另实测撞并发上限时如实显示后端原文「系统繁忙:并发研究任务已达上限」。
- /factors/compose:`submitting→queued→撮合与净值结算→success(EXP-…)`,业务结果与 5 处放大入口照旧。
- /backtest 回归:圆点由 4 个变 3 个(去掉后端**从不上报**的 selection 阶段),取消仍「已取消 + 没归档」。
- 自检产生的 12 个实验归档已全部删除(bulk-delete count:12,复查无残留)。
345 lines
13 KiB
TypeScript
345 lines
13 KiB
TypeScript
/** 异步 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<string, unknown>;
|
||
started_at?: string | null;
|
||
finished_at?: string | null;
|
||
}
|
||
|
||
export function submitJob(spec: ResearchSpec): Promise<JobSubmit> {
|
||
return apiPost<JobSubmit>("/jobs", spec);
|
||
}
|
||
|
||
/** 任务终态(success / failed / cancelled / timeout)与结果 */
|
||
export interface JobOutcome<T> {
|
||
status: string;
|
||
result: T | null;
|
||
error?: string;
|
||
/** 归档实验 id:结果出来后前端可直接给出「去对比」入口 */
|
||
experimentId?: string | null;
|
||
}
|
||
|
||
/** 轮询直到 success / failed / cancelled,或超时(默认 10 分钟)。 */
|
||
export async function waitJob<T>(
|
||
jobId: string,
|
||
timeoutMs = 600_000,
|
||
onStage?: (info: { stage: string; elapsedMs: number }) => void,
|
||
): Promise<JobOutcome<T>> {
|
||
const deadline = Date.now() + timeoutMs;
|
||
const t0 = Date.now();
|
||
while (Date.now() < deadline) {
|
||
const job = await apiGet<JobStatusResp>(`/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 枚举 → 用户可读)。
|
||
*
|
||
* 必须覆盖每种作业**上报过的全部阶段**:漏一个,JobProgress 只能把裸枚举摆给用户
|
||
* —— 实测漏过 `factor_calculation`(/factors 反馈条上直接出现过英文 `factor_calculation`)。
|
||
* 新增阶段时请同时补这里的中文名与下方 `STAGE_PIPELINES` 里对应的序列。
|
||
*/
|
||
export const STAGE_LABEL: Record<string, string> = {
|
||
queued: "排队中",
|
||
data_loading: "加载行情与因子数据",
|
||
factor_calculation: "计算因子值",
|
||
selection: "逐择股日选股",
|
||
backtesting: "撮合与净值结算",
|
||
analysis: "汇总指标与曲线",
|
||
done: "完成",
|
||
};
|
||
|
||
/** 作业类型 —— **只列后端真的会异步执行的 kind**:
|
||
* - `signals` 不在其中:`POST /api/signals` 是同步接口(没有 /signals/jobs),没有作业阶段可言;
|
||
* - `combo`(回测组合)也不单列:它上报的阶段与 backtest 完全一致(证据见下),直接映射到 "backtest"。
|
||
*/
|
||
export type JobKind = "factor_test" | "backtest" | "selection";
|
||
|
||
/**
|
||
* 每种作业**真实会走**的阶段序列 —— 只放后端真的会上报的阶段,**不含 `queued` / `done`**:
|
||
* 「排队中」「已完成」是作业**状态**(由 `phase` 表达,见 JobRunState),不是后端上报的执行阶段,
|
||
* 混进圆点序列会让「尚未开始的第一个圆点」和「已结束」看起来像同一条进度。
|
||
*
|
||
* 为什么必须按作业类型分开:全站共用一条「阶段小圆点」之后,用一张全局表会把因子测试
|
||
* 根本不存在的阶段(逐择股日选股 / 撮合与净值结算)也画出来 —— 那正是要避免的
|
||
* 「看似有数据的假象」;反过来,后端新阶段如果不在表里,圆点会整排变灰。
|
||
*
|
||
* 证据(后端上报点,改动前请重新 grep 核对,勿凭印象):
|
||
* - factor_test:backend/app/quant/service.py:304 data_loading → :306 factor_calculation → :308 analysis
|
||
* - backtest :backend/app/quant/service.py:314 data_loading → :316 backtesting → :319 analysis
|
||
* - selection :backend/app/application/services/job_executor.py:156 selection(真的只有这一段)
|
||
* - combo(复用 backtest 表):backend/app/application/services/combo_service.py:68 data_loading
|
||
* → :70 backtesting → :90 analysis;组合的「逐择股日选股」是在 run_combo_backtest 内部完成的,
|
||
* 后端从未单独上报 selection 阶段,所以回测这条线不放 selection 圆点。
|
||
*/
|
||
export const STAGE_PIPELINES: Record<JobKind, string[]> = {
|
||
factor_test: ["data_loading", "factor_calculation", "analysis"],
|
||
backtest: ["data_loading", "backtesting", "analysis"],
|
||
selection: ["selection"],
|
||
};
|
||
|
||
/** @deprecated 已改为按作业类型取 `STAGE_PIPELINES[kind]`。保留导出仅为兼容旧引用,
|
||
* 其值等于回测的圆点序列(不再含 queued/done —— 它们是 phase 不是阶段)。 */
|
||
export const STAGE_ORDER = STAGE_PIPELINES.backtest;
|
||
|
||
/** 取消排队中/执行中的作业(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<JobSubmit> {
|
||
return apiPost<JobSubmit>("/combos/run", combo);
|
||
}
|
||
|
||
/**
|
||
* 运行**已保存**的回测组合(POST /api/combos/{id}/run)。
|
||
*
|
||
* 与 submitComboJob 的区别:这里用库里的那份参数,页面上未保存的改动不参与 ——
|
||
* 「从组合库直接运行」必须跑库里存的那套,否则用户改了一半的表单会污染既有组合的结果。
|
||
*/
|
||
export async function runSavedCombo(comboId: string): Promise<JobSubmit> {
|
||
return apiPost<JobSubmit>(`/combos/${encodeURIComponent(comboId)}/run`, {});
|
||
}
|
||
|
||
/* ------------------------------------------------------------------ *
|
||
* 统一的「提交 → 轮询 → 终态」状态机(useJobRunner)
|
||
*
|
||
* 为什么要有它:以前每个页面各写一遍 running/jobId/error,结果参差不齐 ——
|
||
* 有的页面点了「运行」只在按钮上转圈,用户不知道到底提交没提交、跑到哪一步;
|
||
* 有的页面「已用秒数」只在后端阶段变化时才更新,看起来像卡死了(0s 一动不动)。
|
||
* 这里把状态、秒表、取消、终态收敛到一处,页面只管渲染。
|
||
* ------------------------------------------------------------------ */
|
||
|
||
export type JobPhase =
|
||
| "idle"
|
||
| "submitting"
|
||
| "queued"
|
||
| "running"
|
||
| "success"
|
||
| "failed"
|
||
| "cancelled"
|
||
| "timeout";
|
||
|
||
export interface JobRunState<T> {
|
||
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<T>(label = ""): JobRunState<T> {
|
||
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<T> {
|
||
state: JobRunState<T>;
|
||
/** 提交并轮询到终态;返回终态结果(供页面继续处理,如渲染报告) */
|
||
run: (
|
||
submit: () => Promise<{ job_id: string }>,
|
||
opts?: JobRunOptions,
|
||
) => Promise<JobOutcome<T>>;
|
||
cancel: () => Promise<void>;
|
||
reset: () => void;
|
||
}
|
||
|
||
/**
|
||
* 作业运行状态机。
|
||
*
|
||
* 生命周期:idle → submitting(已点击,正在换 job_id)→ queued/running(轮询中,秒表走)
|
||
* → success | failed | cancelled | timeout。**submitting 也是状态**:点了按钮立刻就有
|
||
* 反馈文字,不会出现「点了没反应」的空白期。
|
||
*/
|
||
export function useJobRunner<T>(defaultLabel = ""): JobRunner<T> {
|
||
const [state, setState] = useState<JobRunState<T>>(() => emptyState<T>(defaultLabel));
|
||
const t0 = useRef(0);
|
||
const timer = useRef<number | null>(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<JobRunState<T>>) => {
|
||
if (!alive.current) return;
|
||
setState((prev) => ({ ...prev, ...next }));
|
||
}, []);
|
||
|
||
const reset = useCallback(() => {
|
||
stopTick();
|
||
currentJob.current = "";
|
||
if (alive.current) setState(emptyState<T>(defaultLabel));
|
||
}, [defaultLabel, stopTick]);
|
||
|
||
const run = useCallback(
|
||
async (
|
||
submit: () => Promise<{ job_id: string }>,
|
||
opts: JobRunOptions = {},
|
||
): Promise<JobOutcome<T>> => {
|
||
stopTick();
|
||
currentJob.current = "";
|
||
stageRef.current = "";
|
||
t0.current = Date.now();
|
||
setState({
|
||
...emptyState<T>(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<T>(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 };
|
||
}
|