diff --git a/api.py b/api.py index 65e97ea..53f5d60 100644 --- a/api.py +++ b/api.py @@ -13,15 +13,21 @@ cron 的 docker exec 构建/出计划照旧,互不影响。局域网内部服 统一任务调度平台(XXL-JOB)触发入口挂在 /api/v1/xxl/*(见 xxl.py,2026-08-03): 盘前链(build → plan → push-pool)可由平台拉起并回调结案,.env 配 XXL_TRIGGER_KEY 才启用。 """ +import logging +import os + import pandas as pd -from fastapi import FastAPI, HTTPException +from fastapi import FastAPI, HTTPException, Request from fastapi.responses import PlainTextResponse import db import plan import plan_reconcile +import regime from xxl import router as xxl_router +_access = logging.getLogger("plan.access") + app = FastAPI(title="akg-factor-bridge · 每日选股计划", version="0.1") app.include_router(xxl_router) @@ -47,12 +53,26 @@ def plan_dates(limit: int = 30): @app.get("/plan") -def get_plan(date: str | None = None, format: str = "json", +def get_plan(request: Request, date: str | None = None, format: str = "json", top: int = 20, obs_top: int = 10, theme_cap: int = 5): + """向下兼容承诺(2026-09-02 方案 2.7):main / observe 的装配、排序、裁剪与既有字段 + 一字不动,每行只多联入判决类字段;顶层只新增 generated_at、plan_version、regime、 + card_counts、candidates、watch、segments_pointed、snapshot。PMS 按字段名取值、忽略未知键。""" try: data = plan.collect(date, top, obs_top, theme_cap) except RuntimeError as e: raise HTTPException(status_code=404, detail=str(e)) + data.pop("_full", None) + ds = data["date"] + reg = regime.read_from_snapshot(ds) + data["snapshot"] = "present" if os.path.exists(regime.snapshot_path(ds)) else "missing" + data["regime"] = reg or {"status": regime.UNKNOWN, "weak_day": None, + "source": "当日快照无环境段(08:45 追加未跑或快照缺失)"} + _access.info("plan client=%s date=%s regime=%s generated_at=%s version=%s " + "top=%s obs_top=%s theme_cap=%s", + request.client.host if request.client else "-", ds, + data["regime"].get("status"), data.get("generated_at"), + data.get("plan_version"), top, obs_top, theme_cap) if format == "md": return PlainTextResponse(plan.render_md(data), media_type="text/markdown; charset=utf-8") diff --git a/config.py b/config.py index 3d39547..32dac24 100644 --- a/config.py +++ b/config.py @@ -185,3 +185,28 @@ POOL_MAX = int(os.environ.get("POOL_MAX", "60")) # 当晚 22:30 全量扫兜底。例:http://192.168.16.188:38000/api/v1/xxl/daily-scan BIONIC_SCAN_URL = os.environ.get("BIONIC_SCAN_URL", "").rstrip("/") BIONIC_SCAN_KEY = os.environ.get("BIONIC_SCAN_KEY", "") + +# --- 候选卡(2026-09-02 主观选股改进方案,docs/主观选股改进方案_2026-09-02.md)------ +# 桥从"打分排序器"改成"候选卡装配器":分数与档位不动,另出带理由的候选单。规则在 card.py。 +# 两个旋钮保持默认、不在历史样本上挑参数,等样本外复盘读数再拍(09-02 拍板)。 +CARD_START_PCT = float(os.environ.get("CARD_START_PCT", "3")) # 门槛二:数据日涨幅达到几个百分点算已启动(与基座热点扫描同口径) +CARD_ACCUM_MAX_AGE = int(os.environ.get("CARD_ACCUM_MAX_AGE", "30")) # 确认线:吸筹评分日龄上限(交易日) +# 计划快照落点(容器内路径;仓库根挂在 /app,故默认落在仓库 data/plan/)。 +# 每日 JSON 含主榜与观察档全部行与全部证据线,是复盘与对账的唯一底本。 +PLAN_SNAPSHOT_DIR = os.environ.get("PLAN_SNAPSHOT_DIR", "data/plan") + +# --- 环境标签(只展示与复盘分组,不作交易前置;09-02 拍板)-------------------------- +# 来源是决策系统的只读日频区制接口(请它加,桥零依赖);空串 = 未接入,标签一律 UNKNOWN, +# 不拦任何票。快照 08:40 预热,桥 07:10 出计划时拿不到,由 08:45 的追加步骤写进当日快照。 +REGIME_API_URL = os.environ.get("REGIME_API_URL", "").rstrip("/") +REGIME_API_TIMEOUT = float(os.environ.get("REGIME_API_TIMEOUT", "6")) +# 弱势日定义(09-02 拍板,首份周报前锁定、之后不改;改它算新一轮验证):八个指数里弱势个数达到几个。 +REGIME_WEAK_COUNT = int(os.environ.get("REGIME_WEAK_COUNT", "3")) + +# --- 入池与候选卡的联动(09-02;第一步只加字段与切片,PMS 行为不变)------------------ +# POOL_SOURCE:现状 = 强传导档前 POOL_TOP(默认,不变);candidate = 候选优先、再按强传导档补足到 POOL_TOP +# (第二步,复盘读数齐后拍板再切)。 +POOL_SOURCE = os.environ.get("POOL_SOURCE", "tier") +# 低优先入池切片:让"门槛全过但缺吸筹评分"的关注票入池、当晚获得评分,否则复盘缺数据是环状依赖。 +# 0 = 不开(默认);受 POOL_MAX 约束,只改 Mongo 池成分,PMS 不读 Mongo 池。 +POOL_WATCH_SLICE = int(os.environ.get("POOL_WATCH_SLICE", "0")) diff --git a/docs/主观选股改进方案_2026-09-02.md b/docs/主观选股改进方案_2026-09-02.md index 3fd4563..6282da6 100644 --- a/docs/主观选股改进方案_2026-09-02.md +++ b/docs/主观选股改进方案_2026-09-02.md @@ -230,7 +230,8 @@ README"已知的上游约束"三条在基座侧已全部修掉:传导候选不 - 决策系统夜间扫描读 Mongo 股票池时取全部分组并集、丢掉分组编号,它自己分不清哪些票来自桥;盘中白名单另起任务重建。 - 接口层缺口:/plan 无版本戳,同日重算多版时下游分不清拿的是哪一版;桥已实现升降档字段但 PMS 的解析代码不读它;决策系统的资金强度流当前只留痕不触发任何动作;夜间结果表会被盘中补扫就地改写(参考位漂移,PMS 设了 3% 漂移熔断)。 - PMS 不回写上游,是桥主动读 PMS 持仓账本 pms_position;成交与收益没有回流通道。PMS 自己有动作账本 pms_action_ledger(含"拒了的后来涨了多少"的判分事实源)与计划快照 pms_plan_snapshot,是候选单复盘的天然终点。 -- 三处同源纪律:坏信号集合 {SELL, AVOID, DROPPED} 在桥 pool.py、决策系统 pms_advisor.py、PMS rule_gate.py 三处写死,改一处必须三处同改。 +- 三处同源纪律:坏信号集合 {SELL, AVOID, DROPPED} 在桥 pool.py、决策系统 pms_advisor.py、PMS rule_gate.py 三处写死,改一处必须三处同改(09-02 起桥内归一到 card.py,pool.py 从它引用)。 +- PMS 运行参数实读(09-02,参数表 pms_runtime_param):PMS_PLAN_TOP_N 运行值为 **100**(08-17 由用户改,代码默认 30 已被覆盖),PMS_AUTONOMY 为 full;PMS_OPEN_AUTONOMY 不在表里、走代码默认 full(新建仓全自动),已拍板改为 propose_only。桥生产旋钮 POOL_TOP=50、POOL_MAX=100、赛道闸开启。**池深不变式"POOL_TOP 不低于 PMS_PLAN_TOP_N"现在不成立**(50 对 100):主榜第 51 到 100 名的强传导票 PMS 能买、却不在夜间分析池,这正是对接说明第 7 节讲过的缺口,参数被改大后又出现了。处置列入第五项入池联动:要么 POOL_TOP 提到 100(夜扫时长翻倍),要么 PMS 候选深度回到 50,由用户拍板。 ### 1.12 平台评价接口的可行性(09-02 核实) 平台 POST /api/v1/evaluation/tasks 一次调用即可建评价任务:参数 task_name、factor_codes(列表)、start_date 与 end_date(或给 market_statuses 自动取近一年)、exclude_st、可选指数与行业范围;auto_run 默认真,自动走 B 到 D 环节,结果含 factor_code、forward_period、ic_mean 与 IC 序列(backend/app/api/v1/endpoints/evaluation.py:70-100,schemas/evaluation.py:18-54)。限制:六因子的因子表历史只有 07-29 起约五周(传导不可回填),评价样本极短;heat 与 event 可用 build history 回填更长历史,upside 需先做 asof 重建。 @@ -357,7 +358,9 @@ README"已知的上游约束"三条在基座侧已全部修掉:传导候选不 2. 环境闸:桥出环境标签,PMS 读字段定仓位。**拍板后实证修正(1.8e)**:用过去涨跌定义的闸没有预测力,标签当前只作展示与复盘分组,不作交易前置;"弱环境不加新票"挂起,是否正式退役见本节新增拍板点。 3. 新应用顺序:平台评价闭环作废;候选单复盘闭环、吸筹确认线排最前,公告进粮在回炉之后立项,行情侧两条展示线同批。 -**本轮新增拍板点(方案批准后在对话里分组选项式提问):** +**新增拍板点的拍板记录(2026-09-02 方案批准后对话拍板,均选推荐):** 08:45 环境追加任务走统一调度平台新建任务;代码版本解析挂载进容器的 .git 文件;冻结快照缺失的四个计划日(07-29、08-07、08-10、08-20)用当前投影反推并标注近似;两个旋钮保持默认、赛道闸保持开启;环境标签来源为请决策系统加只读日频接口(就绪前标签为空);弱势日定义为八个指数里弱势个数达到三个;环境闸只展示与分组、不动 PMS 定时、"弱环境不加新票"退役;缺评分的关注票走低优先入池切片。板块映射代定:主板只出全市场标签,科创板与创业板出板块标签,依据是桥没有市值列、逐票主参考指数要等接口返回。补两条:PMS 新建仓自主档现在就改为"只提议"作止血(运行时参数,零代码,一键回退,由用户在 PMS 侧操作);新应用顺序同意把"环节启动全景""环境标签验证"插在公告进粮之前。 + +**本轮新增拍板点(原清单存档):** 1. 快照与版本:08:45 环境追加任务走统一调度平台还是宿主定时任务;代码版本读不到时用解析仓库文件还是宿主传环境变量。 2. 视图与历史重建:冻结投影快照缺失的日子用当前投影反推并标注,是否接受。 3. 候选卡两个旋钮(启动阈值 3%、评分日龄三十日)保持默认等样本外读数;赛道闸是否放开。 diff --git a/plan.py b/plan.py index f55ef6a..adc33cd 100644 --- a/plan.py +++ b/plan.py @@ -11,14 +11,18 @@ 升降档一节对比前一交易日的档位表——数据到达本身是信号(首次覆盖 / 新进传导链即升档)。 """ +import datetime as dt import json import os import pandas as pd +import card import common import config import db +import sources +import version # 分数编码(与 factors.build_score 一致):主榜 = 200 + 传导档位×20 + 组内分, @@ -96,9 +100,91 @@ def _val(series: pd.Series, k: str): return None if v is None or pd.isna(v) else float(v) +def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series, + mkt_days: set) -> tuple[dict, list]: + """候选卡装配(2026-09-02 方案第 2.2 节第三项):对档位表里的全部票(主榜与观察档, + 裁剪之前)读三路证据、逐票判决。规则在 card.py,取数在 sources.py,这里只做对齐。 + + 返回 (cards, segments_pointed):cards 按前缀码索引,含判决、理由、缺失、风险、卡内序 + 与证据线原值;segments_pointed 是"关注环节"聚合——今日被传导指向的每个环节的源数、 + 链符、成员数、已启动成员、领涨者、候选数。这是拍板记录第一项"强传导档降为关注环节" + 在计划文本层的落地。""" + import factors + + moved = sources.moved_members(ds) + daily = sources.stock_daily(ds) + night = sources.night_conclusions(codes, ds) + try: + risk = factors._risk_set() or set() # noqa: SLF001 —— 同仓自用 + except Exception: # noqa: BLE001 —— 风险名单拿不到时不当 ST 处理,与 gate 的宽容一致 + risk = set() + stale = bool(mkt_days) and any(x != ds for x in mkt_days) + + cards: dict = {} + for k in codes: + mv, e, d, n = moved.get(k), ev.get(k), daily.get(k, {}), night.get(k, {}) + theme = (mv or {}).get("theme") or (e[0] if e else None) + n_sources = (mv or {}).get("n_sources") or (e[1] if e else None) + up = _val(upside, k) + evd = {"pointed": bool(mv or e), "theme": theme, "n_sources": n_sources, + "chain_fit": (mv or {}).get("chain_fit"), + "pct0": d.get("pct0"), "covered": up is not None, "upside": up, + "risk_name": k in risk, + "accum_state": n.get("accum_state"), "accum_score": n.get("accum_score"), + "accum_age": n.get("accum_age"), "y_signal": n.get("signal"), + "stale_snapshot": stale} + j = card.judge(evd, start_pct=config.CARD_START_PCT, + accum_max_age=config.CARD_ACCUM_MAX_AGE, + neg_tol=config.UPSIDE_NEG_TOLERANCE) + cards[k] = { + **j, + "theme": theme, "n_sources": n_sources, "chain_fit": evd["chain_fit"], + "started_source": "moved_view" if mv else None, + "pct0": d.get("pct0"), "net_z": d.get("net_z"), "heat_chg": d.get("heat_chg"), + "accum": ({"state": n.get("accum_state"), "score": n.get("accum_score"), + "age": n.get("accum_age"), "pos_tag": n.get("accum_pos_tag"), + "date": n.get("conclusion_date")} if n else None), + "night": ({"signal": n.get("signal"), "support": n.get("support"), + "pressure": n.get("pressure"), "date": n.get("conclusion_date")} + if n else None), + } + for i, (k, c) in enumerate(sorted(cards.items(), key=lambda kv: card.sort_key(kv[1])), 1): + c["card_rank"] = i + + segs: dict = {} + for k, mv in moved.items(): + s = segs.setdefault(mv["theme"], { + "segment": mv["theme"], "n_sources": mv["n_sources"], "chain_fit": mv["chain_fit"], + "members_total": mv.get("members_total"), "moved": mv.get("moved"), + "started": [], "candidates": 0, "leader": None}) + s["started"].append(k) + if cards.get(k, {}).get("verdict") == card.VERDICT_CANDIDATE: + s["candidates"] += 1 + for e in ev.values(): # 只有未动名单的环节也是"被指向",列入但无已启动 + segs.setdefault(e[0], {"segment": e[0], "n_sources": e[1], "chain_fit": None, + "members_total": None, "moved": None, + "started": [], "candidates": 0, "leader": None}) + for s in segs.values(): + best, best_pct = None, None + for k in s["started"]: + p = cards.get(k, {}).get("pct0") + if p is not None and (best_pct is None or p > best_pct): + best, best_pct = k, p + s["leader"] = {"code": best, "pct0": best_pct} if best else None + s["started_count"] = len(s["started"]) + ordered = sorted(segs.values(), key=lambda s: (-s["candidates"], -(s["n_sources"] or 0), + -(s["chain_fit"] or 0), s["segment"])) + return cards, ordered + + def collect(date: str | None = None, top: int = 20, obs_top: int = 10, theme_cap: int = 5) -> dict: - """装配一天的计划为结构化字典。数据缺失抛 RuntimeError(api 侧转 404)。""" + """装配一天的计划为结构化字典。数据缺失抛 RuntimeError(api 侧转 404)。 + + 2026-09-02 起附带候选卡:main / observe 的装配、排序、裁剪一字不动(下游 PMS 只读 + 这两段的既有字段),每行只是多联入判决类字段;顶层新增 generated_at、plan_version、 + card_counts、candidates(候选单全量,不受裁剪)、watch、segments_pointed。 + `_full` 是全量主榜与观察档行,只给 generate 落快照用,api 返回前会去掉。""" ds = date or _latest_date("t_factor_akg_score") if not ds: raise RuntimeError("t_factor_akg_score 还没有数据——先 build akg_score。") @@ -122,6 +208,9 @@ def collect(date: str | None = None, top: int = 20, obs_top: int = 10, main = score[score >= _MAIN_MIN].sort_values(ascending=False) obs = score[score < _MAIN_MIN].sort_values(ascending=False) + generated_at = dt.datetime.now().isoformat(timespec="seconds") + cards, segments_pointed = _assemble_cards( + ds, list(main.index) + list(obs.index), ev, upside, mkt_days) def _pick(ranked: pd.Series, n: int): """分数从高到低取 n 条;每个传导主题最多 theme_cap 条(0=不设限)—— @@ -140,15 +229,40 @@ def collect(date: str | None = None, top: int = 20, obs_top: int = 10, def _row(rank: int, k: str, s: float, with_tier: bool) -> dict: e = ev.get(k) + c = cards.get(k) or {} + evidence = ({"theme": e[0], "n_sources": e[1], "moved_ratio": round(e[2], 4)} + if e else None) + # 已启动成员在传导视图里没有证据行(视图只摊平未动名单):只在原本为空时用 + # 已动成员视图的目标环节与源数补上,不覆盖已有值——否则到 PMS 会全落进"无主题"桶。 + if evidence is None and c.get("started_source") == "moved_view" and c.get("theme"): + evidence = {"theme": c["theme"], "n_sources": c.get("n_sources"), + "moved_ratio": None, "source": "moved_view"} r = {"rank": rank, "code": k, "name": names.get(k), "score": round(float(s), 2), - "evidence": ({"theme": e[0], "n_sources": e[1], - "moved_ratio": round(e[2], 4)} if e else None), + "evidence": evidence, "heat": _val(heat, k), "upside": _val(upside, k)} if with_tier: r["tier"] = _tier_label(s) + if c: + r.update(verdict=c["verdict"], reasons=c["reasons"], missing=c["missing"], + risk=c["risk"], card_rank=c["card_rank"], + card={"pct0": c.get("pct0"), "net_z": c.get("net_z"), + "heat_chg": c.get("heat_chg"), "accum": c.get("accum"), + "night": c.get("night"), "gates": c.get("gates"), + "confirm": c.get("confirm")}) return r + def _full_rows(ranked: pd.Series, with_tier: bool) -> list: + return [_row(i, k, s, with_tier) for i, (k, s) in enumerate(ranked.items(), 1)] + + def _by_verdict(v: str) -> list: + rows = [_row(0, k, score[k], k in main.index) for k, c in cards.items() + if c.get("verdict") == v] + rows.sort(key=lambda r: r["card_rank"]) + for i, r in enumerate(rows, 1): + r["rank"] = i + return rows + changes = None prev_ds = _prev_date("t_factor_akg_gate", ds) if prev_ds: @@ -229,8 +343,12 @@ def collect(date: str | None = None, top: int = 20, obs_top: int = 10, "upgrades": _mv(up_df.head(15)), "downgrades": _mv(down_df.head(15))} + card_counts = {v: sum(1 for c in cards.values() if c.get("verdict") == v) + for v in (card.VERDICT_CANDIDATE, card.VERDICT_WATCH, card.VERDICT_SHOW)} return { "date": ds, + "generated_at": generated_at, + "plan_version": version.git_short_rev(), "counts": {"main": int(len(main)), "observe": int(len(obs)), "gate_covered": int(len(gate))}, "market_snapshot_days": sorted(mkt_days), @@ -244,6 +362,17 @@ def collect(date: str | None = None, top: int = 20, obs_top: int = 10, "gate_on": bool(config.ENABLE_TRACK_GATE), "encoding": "主榜分=200+传导档位×20+组内分(还没热、还便宜);" "观察档分=100+0.6z(传导)+0.4z(−热度)", + # ---- 候选卡(2026-09-02):分数与档位之外的另一份产物,不受 top / theme_cap 裁剪 ---- + "card_params": {"start_pct": config.CARD_START_PCT, + "accum_max_age": config.CARD_ACCUM_MAX_AGE, + "neg_tol": config.UPSIDE_NEG_TOLERANCE, + "rules": "候选=环节被指向∧当日涨幅达标∧券商覆盖且非ST∧明确吸筹∧无硬风险;" + "关注=无硬风险且(只差覆盖 或 门槛全过无确认);其余仅展示"}, + "card_counts": card_counts, + "candidates": _by_verdict(card.VERDICT_CANDIDATE), + "watch": _by_verdict(card.VERDICT_WATCH), + "segments_pointed": segments_pointed, + "_full": {"main": _full_rows(main, True), "observe": _full_rows(obs, False)}, } @@ -258,9 +387,23 @@ def _fmt_num(v) -> str: def _fmt_ev(e) -> str: if not e: return "—" + if e.get("moved_ratio") is None: # 已动成员视图补的证据行,没有已动比例 + return f"{e['theme']}({e.get('n_sources') or '—'} 源,本票已启动)" return f"{e['theme']}({e['n_sources']} 源,已动 {e['moved_ratio']:.0%})" +def _fmt_pct0(v) -> str: + return "—" if v is None else f"{v:+.1f}%" + + +def _fmt_accum(ac: dict) -> str: + if not ac or not ac.get("state"): + return "无评分" + st = str(ac["state"]).split("·")[0] + age = ac.get("age") + return f"{st}({age} 日前)" if isinstance(age, int) else st + + def render_md(d: dict) -> str: L = [f"# 每日选股计划 · {d['date']}", ""] c = d["counts"] @@ -272,6 +415,59 @@ def render_md(d: dict) -> str: f"(与计划日不同——历史降级日口径)。") L.append("") + # ---- 候选单与关注环节(2026-09-02):放在主榜之前,这是新的主产物 ---- + cc = d.get("card_counts") or {} + cands = d.get("candidates") or [] + L.append(f"## 候选单(环节被指向、当日已启动、券商覆盖且非 ST、明确吸筹、无硬风险;" + f"共 {cc.get('候选', len(cands))} 只,全量列出不受裁剪)") + L.append("") + if not cands: + L.append("(今日无候选——候选为空不是故障:环节没被指向、成员没启动或没有明确吸筹,都会为空。)") + else: + L.append("| # | 代码 | 名称 | 环节 | 源数 | 当日涨幅 | 吸筹 | 预期空间 | 理由 |") + L.append("|---|------|------|------|------|----------|------|----------|------|") + for r in cands: + c = r.get("card") or {} + ac = c.get("accum") or {} + ev_ = r.get("evidence") or {} + L.append(f"| {r['rank']} | {r['code']} | {r['name'] or '—'} | {ev_.get('theme') or '—'} " + f"| {ev_.get('n_sources') or '—'} | {_fmt_pct0(c.get('pct0'))} " + f"| {_fmt_accum(ac)} | {_fmt_pct(r.get('upside'))} " + f"| {';'.join(r.get('reasons') or [])} |") + L.append("") + segs = d.get("segments_pointed") or [] + L.append(f"## 关注环节(今日被传导指向的 {len(segs)} 个环节:定位对不对看这里,挑票看候选单)") + L.append("") + if segs: + L.append("| 环节 | 源数 | 链符 | 成员 | 已启动 | 领涨 | 候选 |") + L.append("|------|------|------|------|--------|------|------|") + for s in segs: + ld = s.get("leader") or {} + L.append(f"| {s['segment']} | {s.get('n_sources') or '—'} | {_fmt_num(s.get('chain_fit'))} " + f"| {s.get('members_total') if s.get('members_total') is not None else '—'} " + f"| {s.get('started_count', 0)} " + f"| {(ld.get('code') or '—') + (' ' + _fmt_pct0(ld.get('pct0')) if ld.get('code') else '')} " + f"| {s.get('candidates', 0)} |") + L.append("") + watch = d.get("watch") or [] + L.append(f"## 关注单(无硬风险,只差券商覆盖或缺明确吸筹;共 {cc.get('关注', len(watch))} 只,列前 20)") + L.append("") + if watch: + L.append("| # | 代码 | 名称 | 环节 | 当日涨幅 | 吸筹 | 缺什么 |") + L.append("|---|------|------|------|----------|------|--------|") + for r in watch[:20]: + c = r.get("card") or {} + ev_ = r.get("evidence") or {} + L.append(f"| {r['rank']} | {r['code']} | {r['name'] or '—'} | {ev_.get('theme') or '—'} " + f"| {_fmt_pct0(c.get('pct0'))} | {_fmt_accum(c.get('accum') or {})} " + f"| {';'.join(r.get('missing') or [])} |") + L.append("") + reg = d.get("regime") + if reg: + L.append(f"环境标签:{reg.get('status')},弱势指数 {reg.get('weak_count')}/8" + f"{',弱势日' if reg.get('weak_day') else ''}(只展示与复盘分组,不作交易前置)。") + L.append("") + cap_txt = f",每主题限额 {d['theme_cap']}" if d["theme_cap"] else "" gate_txt = "、在十五五赛道内" if d.get("gate_on") else "" L.append(f"## 主榜 Top {len(d['main'])}" @@ -334,10 +530,24 @@ def generate(date: str | None = None, top: int = 20, obs_top: int = 10, except RuntimeError as e: raise SystemExit(str(e)) text = render_md(data) - os.makedirs("data/plan", exist_ok=True) - out = f"data/plan/plan_{data['date']}.md" + os.makedirs(config.PLAN_SNAPSHOT_DIR, exist_ok=True) + out = os.path.join(config.PLAN_SNAPSHOT_DIR, f"plan_{data['date']}.md") with open(out, "w", encoding="utf-8") as f: f.write(text + "\n") + # ---- 当日 JSON 快照(2026-09-02 方案第 2.2 节第一项):主榜与观察档全部行、全部证据线, + # 不裁剪、不设主题限额;带生成时刻与代码版本。它是复盘与对账的唯一底本; + # 08:45 的 regime-append 步骤会往里追加 regime 段,/plan 的 regime 段只读这份。---- + full = data.pop("_full", None) or {} + snap = {**data, "main_shown": data["main"], "observe_shown": data["observe"], + "main": full.get("main", []), "observe": full.get("observe", []), + "shown_params": {"top": top, "obs_top": obs_top, "theme_cap": theme_cap}} + jpath = os.path.join(config.PLAN_SNAPSHOT_DIR, f"plan_{data['date']}.json") + tmp = jpath + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(snap, f, ensure_ascii=False, indent=1, default=str) + os.replace(tmp, jpath) print(text) - print(f"\n已写入 {out}") + print(f"\n已写入 {out} 与快照 {jpath}" + f"(主榜 {len(snap['main'])} 行、观察档 {len(snap['observe'])} 行、" + f"候选 {len(data.get('candidates') or [])} 只,版本 {data.get('plan_version')})") return out diff --git a/pool.py b/pool.py index 26afbda..9b983bd 100644 --- a/pool.py +++ b/pool.py @@ -244,13 +244,43 @@ def push(date: str | None = None, top: int | None = None, data = plan.collect(date, top=max(top * 10, 200), obs_top=0, theme_cap=config.POOL_THEME_CAP) tiers = config.POOL_TIERS - plan_rows = [r for r in data["main"] - if not tiers or r.get("tier") in tiers][:top] - plan_codes = [r["code"] for r in plan_rows] ds = data["date"] + tier_rows = [r for r in data["main"] if not tiers or r.get("tier") in tiers] + if config.POOL_SOURCE == "candidate": + # 第二步(2026-09-02 方案第 2.2 节第五项,复盘读数齐后拍板再切):候选优先, + # 再按强传导档补足到 top;默认 POOL_SOURCE=tier 不走这里,行为与现状一致。 + cand = list(data.get("candidates") or []) + seen = {r["code"] for r in cand} + plan_rows = (cand + [r for r in tier_rows if r["code"] not in seen])[:top] + else: + plan_rows = tier_rows[:top] + for r in plan_rows: + r["_pool_source"] = ("candidate" if config.POOL_SOURCE == "candidate" + and r.get("verdict") == "候选" else "tier_top") + plan_codes = [r["code"] for r in plan_rows] # 2. 持仓(读不到直接抛,整轮不动池子)与旧池子 holdings = _read_holdings() + + # 2.5 低优先入池切片(09-02 拍板):门槛全过但缺吸筹评分的关注票,让它们当晚获得评分, + # 否则复盘缺评分是环状依赖(评分只覆盖池内票)。只改 Mongo 池成分,PMS 不读 Mongo 池。 + # 受 POOL_MAX 约束:额度 = 上限 − 计划入选 − 持仓,切片只填得下的部分。默认 0 = 不开。 + slice_rows = [] + if config.POOL_WATCH_SLICE > 0: + room = max(0, config.POOL_MAX - len(plan_codes) - len(holdings - set(plan_codes))) + have = set(plan_codes) | holdings + for r in data.get("watch") or []: + g = (r.get("card") or {}).get("gates") or {} + if (all(g.get(x) for x in ("pointed", "started", "covered", "clean_name")) + and any("无吸筹评分" in m for m in (r.get("missing") or [])) + and r["code"] not in have): + r["_pool_source"] = "watch_slice" + slice_rows.append(r) + have.add(r["code"]) + if len(slice_rows) >= min(config.POOL_WATCH_SLICE, room): + break + plan_rows = plan_rows + slice_rows + plan_codes = [r["code"] for r in plan_rows] col_name = config.POOL_COLLECTION client = None if dry_run and not config_mongo_ready() else _mongo() old_doc, old_members, member_meta = None, set(), {} @@ -272,9 +302,14 @@ def push(date: str | None = None, top: int | None = None, config.POOL_MAX, today) # 4. 打印明细(干跑到此为止) - print(f"计划日 {ds},档位白名单 {sorted(tiers) if tiers else '(不过滤)'}," - f"计划入选 {len(plan_codes)} 只;持仓 {len(holdings)} 只;" - f"旧池 {len(old_members)} 只 → 新池 {len(d['pool'])} 只(上限 {config.POOL_MAX})") + src_cnt = {} + for r in plan_rows: + src_cnt[r.get("_pool_source")] = src_cnt.get(r.get("_pool_source"), 0) + 1 + print(f"计划日 {ds},入池来源 {config.POOL_SOURCE},档位白名单 {sorted(tiers) if tiers else '(不过滤)'}," + f"计划入选 {len(plan_codes)} 只({src_cnt},其中候选卡判为候选 " + f"{sum(1 for r in plan_rows if r.get('verdict') == '候选')} 只);持仓 {len(holdings)} 只;" + f"旧池 {len(old_members)} 只 → 新池 {len(d['pool'])} 只(上限 {config.POOL_MAX});" + f"计划版本 {data.get('plan_version')}") for label, items in (("计划新进", d["new_entrants"]), ("持仓保留", d["retained_holdings"]), ("留池观察", d["observers"]), @@ -299,7 +334,13 @@ def push(date: str | None = None, top: int | None = None, ctx[r["code"]] = {"factor_code": "akg_score", "score": r.get("score"), "tier": r.get("tier"), "upside": r.get("upside"), "theme": (r.get("evidence") or {}).get("theme"), - "plan_date": ds} + "plan_date": ds, + # 候选卡摘要与来源标记(09-02 方案第五项第一步:只加字段) + "source": r.get("_pool_source"), + "verdict": r.get("verdict"), + "reasons": (r.get("reasons") or [])[:4], + "card_rank": r.get("card_rank"), + "plan_version": data.get("plan_version")} for c in d["retained_holdings"]: ctx.setdefault(c, {"factor_code": "akg_score", "note": "持仓保留"}) diff --git a/regime.py b/regime.py new file mode 100644 index 0000000..a70b45f --- /dev/null +++ b/regime.py @@ -0,0 +1,125 @@ +"""环境标签:只读决策系统的日频市场区制接口,写进当日计划快照;只展示与复盘分组,不作交易前置。 + +## 为什么要有它 + +实证(docs/主观选股改进方案_2026-09-02.md 1.8、1.8e):市场环境是第一解释变量,但 +"看过去五日或十日涨跌"没有预测力(相关 −0.28、−0.55)。所以桥不自算环境,只读 +决策系统已有的、带迟滞的区制判断(主参考指数跌破二十日线且缩量或连续破位判弱势), +把它当标签挂在每张卡上,供复盘按"事前标签"分组。标签有没有预测力由复盘过线条件验证; +过线前不拦任何票,"弱环境不加新票"已退役(09-02 拍板)。 + +## 时序 + +决策系统 08:40 预热当日快照;桥 07:10 出计划时拿不到,标 UNKNOWN。统一调度平台在 08:45 +再触发一次追加步骤(xxl.py 的 regime-append),把当日快照写进当天 JSON 的 regime 段; +/plan 应答的 regime 段只读那份落盘的快照,盘中不再向来源发请求。 + +## 接口契约(桥按此消费;决策系统侧按此实现只读端点) + +GET {REGIME_API_URL}?date=YYYY-MM-DD -> + {"status": "OK" | "UNKNOWN" | "DISABLED", "data_date": "YYYY-MM-DD", + "indices": [{"code": "000300.SH", "name": "沪深300", "weak": true, "reason": "..."}], + "weak_count": 3, "degraded": false, "computed_at": "..."} +DISABLED 归 UNKNOWN 并在 source 里注明总开关关闭。任何失败 -> UNKNOWN,绝不抛错。 +""" +from __future__ import annotations + +import datetime as dt +import json +import os +import urllib.parse +import urllib.request + +import config + +UNKNOWN = "UNKNOWN" + + +def fetch(day: str) -> dict: + """读一次当日区制。返回归一后的 regime 段。""" + base = {"status": UNKNOWN, "data_date": None, "weak_count": None, "weak_day": None, + "indices": [], "degraded": None, "source": None, + "fetched_at": dt.datetime.now().isoformat(timespec="seconds"), + "weak_threshold": config.REGIME_WEAK_COUNT} + if not config.REGIME_API_URL: + base["source"] = "unconfigured(决策系统只读区制接口未接入)" + return base + url = f"{config.REGIME_API_URL}?{urllib.parse.urlencode({'date': day})}" + base["source"] = url + try: + with urllib.request.urlopen(url, timeout=config.REGIME_API_TIMEOUT) as resp: + payload = json.loads(resp.read().decode("utf-8", "replace")) + except Exception as e: # noqa: BLE001 —— 来源不可达就是 UNKNOWN,不拦票 + base["source"] = f"{url}(不可达: {type(e).__name__})" + return base + if not isinstance(payload, dict): + return base + status = str(payload.get("status") or UNKNOWN).upper() + if status == "DISABLED": + base["source"] += "(总开关关闭)" + status = UNKNOWN + # 两种形状都收:契约形状 indices 为列表;决策系统缓存快照的原始形状 indices 为 + # {指数代码: {name, weak, weak_via, qrs_stale, ...}}、广度字段叫 breadth_weak。 + raw_idx = payload.get("indices") + if isinstance(raw_idx, dict): + indices = [{"code": c, **(v if isinstance(v, dict) else {})} for c, v in raw_idx.items()] + elif isinstance(raw_idx, list): + indices = [x for x in raw_idx if isinstance(x, dict)] + else: + indices = [] + weak_count = payload.get("weak_count") + if weak_count is None: + weak_count = payload.get("breadth_weak") + if weak_count is None and indices: + weak_count = sum(1 for x in indices if x.get("weak")) + degraded = payload.get("degraded") + if degraded is None and indices: + degraded = any(x.get("qrs_stale") for x in indices) + base.update({ + "status": "OK" if status == "OK" else UNKNOWN, + "data_date": payload.get("data_date"), + "weak_count": weak_count, + "weak_total": payload.get("breadth_total") or (len(indices) or None), + "indices": [{"code": x.get("code"), "name": x.get("name"), "weak": bool(x.get("weak")), + "weak_via": x.get("weak_via"), "qrs_stale": x.get("qrs_stale"), + "dev_pct": x.get("dev_pct")} for x in indices], + "degraded": bool(degraded) if degraded is not None else None, + "computed_at": payload.get("computed_at"), + }) + if base["status"] == "OK" and isinstance(weak_count, int): + base["weak_day"] = weak_count >= config.REGIME_WEAK_COUNT + return base + + +def snapshot_path(day: str) -> str: + return os.path.join(config.PLAN_SNAPSHOT_DIR, f"plan_{day}.json") + + +def append_to_snapshot(day: str) -> dict: + """08:45 追加步骤:把当日区制写进当天 JSON 快照的 regime 段。 + 快照不存在(当日构建失败或 07:10 之前)就只打印,不创建空快照。""" + reg = fetch(day) + path = snapshot_path(day) + if not os.path.exists(path): + print(f"快照 {path} 不存在,环境段未写入(状态 {reg['status']})") + return reg + with open(path, "r", encoding="utf-8") as f: + doc = json.load(f) + doc["regime"] = reg + tmp = path + ".tmp" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(doc, f, ensure_ascii=False, indent=1) + os.replace(tmp, path) + print(f"环境段已写入 {path}:status={reg['status']} weak_count={reg['weak_count']} " + f"weak_day={reg['weak_day']}") + return reg + + +def read_from_snapshot(day: str) -> dict | None: + """/plan 用:只读当日快照里的 regime 段,没有就 None(应答里给 UNKNOWN)。""" + path = snapshot_path(day) + try: + with open(path, "r", encoding="utf-8") as f: + return json.load(f).get("regime") + except (OSError, ValueError): + return None diff --git a/run.py b/run.py index 149bd99..ac3caf2 100644 --- a/run.py +++ b/run.py @@ -180,6 +180,8 @@ def main(): av = sub.add_parser("apply-views") av.add_argument("--file", default="sql/astock_kg_slot_views.sql") av.add_argument("--dry-run", action="store_true", help="只列语句不执行") + ra = sub.add_parser("regime-append") # 08:45 环境追加(写当日计划快照的 regime 段) + ra.add_argument("--date", help="默认取 score 表最新日(与 plan 同口径)") f = sub.add_parser("freeze") f.add_argument("--date", help="默认今天") b = sub.add_parser("build") @@ -210,6 +212,13 @@ def main(): elif a.cmd == "push-pool": import pool pool.push(a.date, a.top, dry_run=a.dry_run, kick=not a.no_kick) + elif a.cmd == "regime-append": + import plan + import regime + day = a.date or plan._latest_date("t_factor_akg_score") # noqa: SLF001 —— 同仓自用 + if not day: + raise SystemExit("t_factor_akg_score 还没有数据,没有当日快照可追加。") + regime.append_to_snapshot(day) elif a.cmd == "tracks": import tracks tracks.coverage_report() diff --git a/sources.py b/sources.py new file mode 100644 index 0000000..1be146d --- /dev/null +++ b/sources.py @@ -0,0 +1,164 @@ +"""候选卡的取数层:把各条证据线从三处库读成"按前缀码索引的字典"(全部只读)。 + +## 为什么要有它 + +候选卡的规则在 card.py,是纯函数;证据从哪来、怎么对齐代码格式、缺了怎么办, +全部收在这里,plan.py 只做装配。三处来源: + + 基座 PG v_factor_transmission_moved 已启动成员及其所在环节的传导证据(只认 Segment 目标) + v_factor_stock_daily 数据日涨幅、主力净额异常值、热度变化 + 153 代理 strategy_daily_results 决策系统昨夜结论:信号、支撑压力位、吸筹块 + (吸筹的评分、状态、评分日三项一次读齐;基座落库的吸筹版本没有评分日,所以不从基座取) + 平台 MySQL gp_day_data 交易日历(算评分日龄用;取一只长期存在的票的日期序列) + +代码格式:基座是点后缀式 600000.SH,决策系统与桥是前缀式 SH600000,进出都过 common.to_prefix。 +读失败的语义:每一路读不到都返回空字典并打印一行原因,候选卡按"缺失"处理(进关注或仅展示), +不让计划断产——与 pool.py 的安全边界一致。 +""" +from __future__ import annotations + +import datetime as dt +import json + +import pandas as pd + +import common +import db + + +def moved_members(ds: str) -> dict[str, dict]: + """数据日 ds 被传导指向的环节里,已启动(不在未动名单)的成员。 + 同一票挂在多个被指向环节上时,取源数最多、其次链符最高的那条作卡上的证据。""" + try: + df = db.read_pg( + "SELECT ts_code, target, n_sources, chain_fit, members_total, moved, " + "moved_ratio, mkt_trade_date FROM v_factor_transmission_moved " + "WHERE scan_date = %s", (ds,)) + except Exception as e: # noqa: BLE001 + print(f" (已动成员视图读取失败,候选卡的传导门槛整体缺席: {e!r})") + return {} + if df.empty: + return {} + df["k"] = df["ts_code"].map(lambda s: common.to_prefix(str(s).strip())) + df["n_sources"] = pd.to_numeric(df["n_sources"], errors="coerce").fillna(0) + df["chain_fit"] = pd.to_numeric(df["chain_fit"], errors="coerce").fillna(0) + df = df.sort_values(["n_sources", "chain_fit"], ascending=False).drop_duplicates("k") + out = {} + for r in df.itertuples(): + out[r.k] = {"theme": str(r.target), "n_sources": int(r.n_sources), + "chain_fit": float(r.chain_fit), + "members_total": None if pd.isna(r.members_total) else int(r.members_total), + "moved": None if pd.isna(r.moved) else int(r.moved), + "mkt_trade_date": None if pd.isna(r.mkt_trade_date) else str(r.mkt_trade_date)} + return out + + +def stock_daily(ds: str) -> dict[str, dict]: + """数据日 ds 的个股行情三列:涨幅(百分数)、主力净额异常值、热度五日变化。""" + try: + df = db.read_pg( + "SELECT ts_code, pct_change, net_z, heat_chg FROM v_factor_stock_daily " + "WHERE trade_date = %s", (ds,)) + except Exception as e: # noqa: BLE001 + print(f" (个股日行情视图读取失败,候选卡的已启动门槛整体缺席: {e!r})") + return {} + out = {} + for r in df.itertuples(): + k = common.to_prefix(str(r.ts_code).strip()) + out[k] = {"pct0": _f(r.pct_change), "net_z": _f(r.net_z), "heat_chg": _f(r.heat_chg)} + return out + + +def trading_days(upto: str, back_days: int = 90) -> list[str]: + """交易日历:取一只长期存在的票在 gp_day_data 里的日期序列(表有一千四百万行, + 全表 DISTINCT 太慢;按 symbol 走索引)。返回升序 ISO 日期串。""" + start = (dt.date.fromisoformat(upto) - dt.timedelta(days=back_days)).isoformat() + try: + df = db.read_mysql( + "factor", "SELECT DISTINCT DATE(`timestamp`) d FROM gp_day_data " + "WHERE symbol = %s AND `timestamp` >= %s AND `timestamp` <= %s " + "ORDER BY d", ("SH600519", start, upto)) + return [pd.Timestamp(x).date().isoformat() for x in df["d"].tolist()] + except Exception as e: # noqa: BLE001 + print(f" (交易日历读取失败,评分日龄按自然日×5/7 近似: {e!r})") + return [] + + +def night_conclusions(codes, ds: str) -> dict[str, dict]: + """决策系统昨夜结论(每票最新一行):信号、支撑压力位、吸筹块与评分日龄。 + + 评分日龄 = 结论行的 trade_date 到数据日 ds 之间的交易日数(含头不含尾)。 + 该表会被盘中补扫就地改写、没有落库时刻列,日龄以行的 trade_date 为准, + 这是已知局限(方案 2.2 第三项)。""" + codes = sorted({common.to_prefix(str(c).strip()) for c in codes if c}) + if not codes: + return {} + try: + marks = ",".join(["%s"] * len(codes)) + df = db.read_mysql( + "pms", + f"SELECT stock_code, trade_date, signal_type, support_level, pressure_level, " + f"raw_logic_json FROM strategy_daily_results WHERE stock_code IN ({marks})", + tuple(codes)) + except Exception as e: # noqa: BLE001 + print(f" (决策系统结论表读取失败,吸筹确认线与坏信号风险整体缺席: {e!r})") + return {} + if df.empty: + return {} + df = df.sort_values("trade_date").drop_duplicates("stock_code", keep="last") + cal = trading_days(ds) + cal_index = {d: i for i, d in enumerate(cal)} + out = {} + for r in df.itertuples(): + k = common.to_prefix(str(r.stock_code).strip()) + tdate = _ymd(r.trade_date) + age = _age(tdate, ds, cal_index) + ff = {} + try: + raw = json.loads(r.raw_logic_json) if isinstance(r.raw_logic_json, str) else (r.raw_logic_json or {}) + ff = (raw or {}).get("fund_flow") or {} + except Exception: # noqa: BLE001 —— 坏 JSON 当无吸筹块 + ff = {} + out[k] = { + "signal": str(r.signal_type or "").strip().upper() or None, + "support": _f(r.support_level), "pressure": _f(r.pressure_level), + "conclusion_date": tdate, + "accum_state": str(ff.get("state") or "") or None, + "accum_score": _f(ff.get("score")), + "accum_structure": ff.get("structure"), "accum_pos_tag": ff.get("pos_tag"), + "accum_age": age, + } + return out + + +def _ymd(v) -> str | None: + """strategy_daily_results.trade_date 是整数 YYYYMMDD(也可能是日期),统一成 ISO 串。""" + if v is None or (isinstance(v, float) and pd.isna(v)): + return None + s = str(v).strip() + if len(s) == 8 and s.isdigit(): + return f"{s[:4]}-{s[4:6]}-{s[6:]}" + try: + return pd.Timestamp(s).date().isoformat() + except Exception: # noqa: BLE001 + return None + + +def _age(tdate: str | None, ds: str, cal_index: dict) -> int | None: + if not tdate: + return None + if tdate in cal_index and ds in cal_index: + return cal_index[ds] - cal_index[tdate] + try: # 日历缺失或日期在日历之外:自然日 × 5/7 近似 + nat = (dt.date.fromisoformat(ds) - dt.date.fromisoformat(tdate)).days + return max(0, round(nat * 5 / 7)) + except ValueError: + return None + + +def _f(v): + try: + x = float(v) + except (TypeError, ValueError): + return None + return None if x != x else x diff --git a/test_regime.py b/test_regime.py new file mode 100644 index 0000000..0760e27 --- /dev/null +++ b/test_regime.py @@ -0,0 +1,91 @@ +"""regime.py 离线单测(无需网络、无需库):未配置 -> UNKNOWN;两种快照形状都能归一; +弱势日按旋钮判;追加进临时快照并能读回。 + +跑法:python3 test_regime.py 或 pytest test_regime.py +""" +import io +import json +import os +import tempfile +import urllib.request + +import config +import regime + + +def t(name, cond): + assert cond, name + print(" ok", name) + + +class _Resp(io.BytesIO): + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + +def main(): + config.REGIME_API_URL = "" + r = regime.fetch("2026-09-01") + t("未配置来源 -> UNKNOWN 且不拦票(weak_day 为空)", + r["status"] == "UNKNOWN" and r["weak_day"] is None and "unconfigured" in r["source"]) + + config.REGIME_API_URL = "http://127.0.0.1:1/api/v1/market/regime" + config.REGIME_WEAK_COUNT = 3 + # 决策系统缓存快照的原始形状:indices 是字典、广度叫 breadth_weak + raw = {"status": "OK", "data_date": "2026-09-01", "computed_at": "x", + "breadth_weak": 3, "breadth_total": 8, + "indices": {"000300.SH": {"name": "沪深300", "ok": True, "weak": True, "weak_via": "ma20", "qrs_stale": False}, + "399006.SZ": {"name": "创业板指", "ok": True, "weak": False, "qrs_stale": True}, + "000688.SH": {"name": "科创50", "ok": True, "weak": True}, + "000001.SH": {"name": "上证指数", "ok": True, "weak": True}}} + orig = urllib.request.urlopen + urllib.request.urlopen = lambda url, timeout=0: _Resp(json.dumps(raw).encode("utf-8")) + try: + r = regime.fetch("2026-09-01") + t("原始形状归一:OK、弱势数 3、指数转列表、降级按 qrs_stale 判", + r["status"] == "OK" and r["weak_count"] == 3 and len(r["indices"]) == 4 + and r["degraded"] is True and r["weak_total"] == 8) + t("弱势日:弱势数 3 达到旋钮 3", r["weak_day"] is True) + config.REGIME_WEAK_COUNT = 4 + t("旋钮改 4 -> 非弱势日", regime.fetch("2026-09-01")["weak_day"] is False) + config.REGIME_WEAK_COUNT = 3 + + # 契约形状:indices 为列表、weak_count 直给 + contract = {"status": "OK", "data_date": "2026-09-01", "weak_count": 1, + "indices": [{"code": "000300.SH", "name": "沪深300", "weak": True}]} + urllib.request.urlopen = lambda url, timeout=0: _Resp(json.dumps(contract).encode("utf-8")) + r = regime.fetch("2026-09-01") + t("契约形状:弱势数 1、非弱势日", r["weak_count"] == 1 and r["weak_day"] is False) + + urllib.request.urlopen = lambda url, timeout=0: _Resp(json.dumps({"status": "DISABLED"}).encode()) + r = regime.fetch("2026-09-01") + t("DISABLED 归 UNKNOWN 并注明", r["status"] == "UNKNOWN" and "总开关" in r["source"]) + + def boom(url, timeout=0): + raise OSError("refused") + urllib.request.urlopen = boom + r = regime.fetch("2026-09-01") + t("来源不可达 -> UNKNOWN,不抛错", r["status"] == "UNKNOWN" and "不可达" in r["source"]) + + # 追加进快照 + urllib.request.urlopen = lambda url, timeout=0: _Resp(json.dumps(raw).encode("utf-8")) + with tempfile.TemporaryDirectory() as d: + config.PLAN_SNAPSHOT_DIR = d + t("快照不存在时不创建空快照", regime.append_to_snapshot("2026-09-01")["status"] == "OK" + and not os.path.exists(regime.snapshot_path("2026-09-01"))) + with open(regime.snapshot_path("2026-09-01"), "w", encoding="utf-8") as f: + json.dump({"date": "2026-09-01", "main": []}, f) + regime.append_to_snapshot("2026-09-01") + back = regime.read_from_snapshot("2026-09-01") + t("追加后能从快照读回 regime 段", back and back["status"] == "OK" and back["weak_count"] == 3) + t("快照其余内容不丢", json.load(open(regime.snapshot_path("2026-09-01"), encoding="utf-8"))["main"] == []) + finally: + urllib.request.urlopen = orig + print("ALL OK — 环境标签:未配置 / 两种形状 / 弱势日旋钮 / 关闭 / 不可达 / 快照追加 全部通过") + + +if __name__ == "__main__": + main() diff --git a/xxl.py b/xxl.py index ee77376..38b47a9 100644 --- a/xxl.py +++ b/xxl.py @@ -46,12 +46,17 @@ 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") +STEP_ORDER = ("build", "plan", "push-pool", "regime-append") 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"], } +DEFAULT_STEPS = ("build", "plan", "push-pool") STEP_TIMEOUT_SEC = 3600 # 单步上限一小时,防呆死(正常盘前链全程分钟级) _lock = threading.Lock() @@ -154,8 +159,10 @@ def _run_job(task_id: str, steps: list, date, callback_url): @router.api_route("/daily-build", methods=["GET", "POST"]) def trigger_daily_build( callbackUrl: str = Query(None, description="平台下发的结案回调地址(自动追加)"), - steps: str = Query("build,plan,push-pool", - description="要跑哪几步,逗号分隔;顺序固定 build→plan→push-pool"), + steps: str = Query(",".join(DEFAULT_STEPS), + description="要跑哪几步,逗号分隔;默认三步 build→plan→push-pool。" + "环境追加 regime-append 不在默认里,平台 08:45 的任务单独传 " + "steps=regime-append(决策系统区制快照 08:40 才预热)"), date: str = Query(None, description="补跑指定数据日 YYYY-MM-DD;空=最新数据日"), key: str = Query(None), x_job_key: str = Header(None),