"""统一任务调度平台(XXL-JOB)触发入口 —— 盘前链的外部拉起与结案回调。 参照决策系统 bionic_trader 的同款接入(app/api/xxl_jobs.py 的 trigger_daily_scan + tasks_periodic 的 notify_xxl),对平台表现完全一致: 触发: GET|POST /api/v1/xxl/daily-build?key=...&steps=...&date=... 立即返回 200 + task_id;平台会自动在 URL 末尾追加 &callbackUrl=... 状态: GET /api/v1/xxl/status/{task_id} 回调未送达时的兜底查询 结案: 任务真正跑完后 POST callbackUrl —— 成功 = 空包(对端默认按成功结案); 失败 = 表单 {"handleCode": 500, "handleMsg": 原因}。 有 HTTP 响应(含 4xx/5xx)即算送达、不重试;只有连接层失败(不可达/超时) 才隔 30 秒重试,最多 XXL_CALLBACK_MAX_RETRIES 次,放弃后靠状态查询兜底。 与决策系统那套只有一处实现差别:桥没有 celery,任务由 API 进程起一条后台线程、 逐步以**子进程**方式执行 `python run.py <步骤>`。用子进程而不是进程内调用,是为了 保住「每次执行都重读 .env、重载代码」的既有性质——常驻 API 自身不会自动重载, 但它派生的子进程会,和宿主 cron 的 docker exec 语义完全一致。 安全默认与决策系统相同:.env 不配 XXL_TRIGGER_KEY,这一组端点整体禁用。 同一时刻只允许一条链在跑:已有任务未结束时再触发回 409,由平台按失败重试策略处理。 步骤白名单(steps 参数,逗号分隔,默认三步全跑,执行顺序固定): build 因子构建(run.py build all --mode daily) plan 生成当日选股计划文件 push-pool 计划写入股票池 + 触发决策系统增量补扫 后面三步不在默认三步里,平台各建一个任务、显式传 steps 才跑: regime-append 把决策系统 08:40 预热的市场区制写进当天计划快照的环境段(平台 08:45 触发) judgement-snapshot 把数据基座当前这一版产业研判与环节评析抄一行,攒版本史(只写不判) plan-review 候选单复盘周报,滚动窗口(截止 7 天前、往前 21 天),只读(平台每周五收盘后触发) """ import datetime as dt import os import subprocess import sys import threading import time import urllib.error import urllib.parse import urllib.request import uuid from fastapi import APIRouter, Header, HTTPException, Query import config router = APIRouter(prefix="/api/v1/xxl", tags=["xxl"]) HERE = os.path.dirname(os.path.abspath(__file__)) LOG_PATH = os.path.join(HERE, "data", "xxl_build.log") # 步骤 → 命令(顺序即执行顺序;--date 由触发参数统一追加) STEP_ORDER = ("build", "plan", "push-pool", "regime-append", "judgement-snapshot", "plan-review") STEP_CMDS = { "build": ["run.py", "build", "all", "--mode", "daily"], "plan": ["run.py", "plan"], "push-pool": ["run.py", "push-pool"], # 环境追加(2026-09-02 方案第 2.2 节第一、四项):决策系统的市场区制快照 08:40 才预热, # 桥 07:10 出计划时拿不到;平台在 08:45 另建一个任务只跑这一步,把当日区制写进当天的 # 计划快照 regime 段。它不在默认三步里,必须显式 steps=regime-append 才跑。 "regime-append": ["run.py", "regime-append"], # 日频行业观点快照(2026-09-03 方案第四之五之三节,落地顺序第一步):把数据基座当前这一版 # 产业研判与环节评析抄进选股系统自己的表,攒版本史,只写不判。它与出计划互不依赖,同样 # 不在默认三步里,必须显式 steps=judgement-snapshot 才跑。 "judgement-snapshot": ["run.py", "judgement-snapshot"], # 复盘周报(2026-09-07 方案第五件之三):每周五收盘后平台单建一个任务只跑这一步,写 data/review。 # 只读,名单级观察收益,不是回测;--date 由平台追加时当作「今天」算滚动窗口。 "plan-review": ["run.py", "plan-review"], } DEFAULT_STEPS = ("build", "plan", "push-pool") STEP_TIMEOUT_SEC = 3600 # 单步上限一小时,防呆死(正常盘前链全程分钟级) _lock = threading.Lock() _jobs: dict = {} # task_id -> 状态字典(常驻进程内存,重启即清空) _running_id = None _KEEP_JOBS = 20 # 只留最近 20 条状态,防内存慢涨 def _require_key(x_job_key, key): cfg = config.XXL_TRIGGER_KEY if not cfg: raise HTTPException(status_code=403, detail="XXL trigger disabled: XXL_TRIGGER_KEY not set") if (x_job_key or key) != cfg: raise HTTPException(status_code=401, detail="invalid key") def _log(job, line: str): stamp = dt.datetime.now().strftime("%H:%M:%S") job["tail"].append(f"[{stamp}] {line}") job["tail"] = job["tail"][-60:] try: os.makedirs(os.path.dirname(LOG_PATH), exist_ok=True) with open(LOG_PATH, "a", encoding="utf-8") as f: f.write(f"{dt.datetime.now().isoformat(timespec='seconds')} " f"[{job['task_id'][:8]}] {line}\n") except Exception: # noqa: BLE001 —— 落盘日志失败不影响任务本体 pass def _run_step(cmd: list) -> tuple: """跑一步(子进程)。单测在这里打桩。返回 (退出码, 合并输出)。""" p = subprocess.run([sys.executable] + cmd, cwd=HERE, capture_output=True, text=True, timeout=STEP_TIMEOUT_SEC) return p.returncode, (p.stdout or "") + (p.stderr or "") def _post_callback(url: str, data: bytes) -> str: """回调一跳(单测在这里打桩)。返回送达说明;连接层失败抛异常。""" req = urllib.request.Request(url, data=data, method="POST") try: with urllib.request.urlopen(req, timeout=config.XXL_CALLBACK_READ_TIMEOUT) as resp: return f"delivered http={resp.status}" except urllib.error.HTTPError as e: return f"delivered http={e.code}" # 对端有响应即算送达,不重试 def _notify(job, callback_url: str, ok: bool, note: str): """结案回调,契约与决策系统 notify_xxl 一致。""" if not callback_url: _log(job, "无 callbackUrl,跳过结案回调(平台可查状态接口)") return data = (b"" if ok else urllib.parse.urlencode( {"handleCode": 500, "handleMsg": (note or "job failed")[:500]}).encode()) for attempt in range(config.XXL_CALLBACK_MAX_RETRIES + 1): try: _log(job, f"结案回调 ok={ok}: {_post_callback(callback_url, data)}") return except Exception as e: # noqa: BLE001 —— 连接层失败才会走到这 if attempt >= config.XXL_CALLBACK_MAX_RETRIES: _log(job, f"结案回调最终未达,放弃(靠状态查询兜底): {e!r}") return _log(job, f"结案回调连接失败,30 秒后第 {attempt + 1} 次重试: {e!r}") time.sleep(30) def _run_job(task_id: str, steps: list, date, callback_url): global _running_id job = _jobs[task_id] ok, note = True, "" try: for step in steps: job["step"] = step cmd = list(STEP_CMDS[step]) + (["--date", date] if date else []) _log(job, f"开始 {step}: python {' '.join(cmd)}") try: code, out = _run_step(cmd) except subprocess.TimeoutExpired: ok, note = False, f"{step} 超过 {STEP_TIMEOUT_SEC} 秒未结束,判失败" _log(job, note) break for line in (out or "").strip().splitlines()[-15:]: _log(job, f" {line}") if code != 0: ok, note = False, f"{step} 退出码 {code}" _log(job, note) break _log(job, f"完成 {step}") except Exception as e: # noqa: BLE001 —— 线程里绝不让异常无声消失 ok, note = False, f"意外异常: {type(e).__name__}: {e}" _log(job, note) job.update(state="success" if ok else "failed", note=note, ended_at=dt.datetime.now().isoformat(timespec="seconds"), step=None) with _lock: _running_id = None _notify(job, callback_url, ok, note) job["closed"] = True # 结案回调也处理完了(送达或放弃),状态接口可据此判断 @router.api_route("/daily-build", methods=["GET", "POST"]) def trigger_daily_build( callbackUrl: str = Query(None, description="平台下发的结案回调地址(自动追加)"), steps: str = Query(",".join(DEFAULT_STEPS), description="要跑哪几步,逗号分隔;默认三步 build→plan→push-pool。" "环境追加 regime-append 不在默认里,平台 08:45 的任务单独传 " "steps=regime-append(决策系统区制快照 08:40 才预热);" "行业观点快照 judgement-snapshot 同样不在默认里," "平台单独建任务传 steps=judgement-snapshot"), date: str = Query(None, description="补跑指定数据日 YYYY-MM-DD;空=最新数据日"), key: str = Query(None), x_job_key: str = Header(None), ): """触发盘前链。立即返回 task_id,跑完由结案回调通知平台。 平台执行参数示例(方式 GET,key 写进 URL,平台会自己追加 callbackUrl): http://<桥机IP>:8300/api/v1/xxl/daily-build?key= """ global _running_id _require_key(x_job_key, key) wanted = [s.strip() for s in (steps or "").split(",") if s.strip()] bad = [s for s in wanted if s not in STEP_CMDS] if bad or not wanted: raise HTTPException(status_code=400, detail=f"steps 只认 {list(STEP_ORDER)},收到 {wanted}") ordered = [s for s in STEP_ORDER if s in wanted] with _lock: if _running_id and _jobs.get(_running_id, {}).get("state") == "running": raise HTTPException(status_code=409, detail=f"已有任务在跑(task_id={_running_id}),本次拒绝") task_id = str(uuid.uuid4()) _jobs[task_id] = {"task_id": task_id, "state": "running", "steps": ordered, "step": None, "date": date, "note": "", "tail": [], "started_at": dt.datetime.now().isoformat(timespec="seconds"), "ended_at": None, "closed": False} while len(_jobs) > _KEEP_JOBS: # 只留最近 N 条 oldest = next(iter(_jobs)) if oldest == task_id: break _jobs.pop(oldest) _running_id = task_id threading.Thread(target=_run_job, args=(task_id, ordered, date, callbackUrl), daemon=True).start() return {"status": "accepted", "task_id": task_id, "steps": ordered, "date": date} @router.get("/status/{task_id}") def job_status(task_id: str, key: str = Query(None), x_job_key: str = Header(None)): """状态兜底:running / success / failed + 当前步骤 + 最近输出。""" _require_key(x_job_key, key) job = _jobs.get(task_id) if not job: raise HTTPException(status_code=404, detail="task_id 不存在(API 重启会清空历史状态)") return job