候选卡上线第一批:计划装配联入判决与证据线、当日 JSON 快照带版本戳、候选单与关注环节两节;接口顶层加 generated_at/plan_version/regime/snapshot 与访问日志(主榜观察档装配排序裁剪不动);取数层 sources.py;环境标签 regime.py 只展示分组;调度加 regime-append 步骤;入池加来源标记、切片与候选优先开关(默认关);方案文档同步

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
zlt 2026-09-02 16:39:52 +08:00
parent c540ea1b49
commit f09b087228
10 changed files with 715 additions and 20 deletions

24
api.py
View File

@ -13,15 +13,21 @@ cron 的 docker exec 构建/出计划照旧,互不影响。局域网内部服
统一任务调度平台XXL-JOB触发入口挂在 /api/v1/xxl/* xxl.py2026-08-03 统一任务调度平台XXL-JOB触发入口挂在 /api/v1/xxl/* xxl.py2026-08-03
盘前链build plan push-pool可由平台拉起并回调结案.env XXL_TRIGGER_KEY 才启用 盘前链build plan push-pool可由平台拉起并回调结案.env XXL_TRIGGER_KEY 才启用
""" """
import logging
import os
import pandas as pd import pandas as pd
from fastapi import FastAPI, HTTPException from fastapi import FastAPI, HTTPException, Request
from fastapi.responses import PlainTextResponse from fastapi.responses import PlainTextResponse
import db import db
import plan import plan
import plan_reconcile import plan_reconcile
import regime
from xxl import router as xxl_router from xxl import router as xxl_router
_access = logging.getLogger("plan.access")
app = FastAPI(title="akg-factor-bridge · 每日选股计划", version="0.1") app = FastAPI(title="akg-factor-bridge · 每日选股计划", version="0.1")
app.include_router(xxl_router) app.include_router(xxl_router)
@ -47,12 +53,26 @@ def plan_dates(limit: int = 30):
@app.get("/plan") @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): top: int = 20, obs_top: int = 10, theme_cap: int = 5):
"""向下兼容承诺2026-09-02 方案 2.7main / observe 的装配、排序、裁剪与既有字段
一字不动每行只多联入判决类字段顶层只新增 generated_atplan_versionregime
card_countscandidateswatchsegments_pointedsnapshotPMS 按字段名取值忽略未知键"""
try: try:
data = plan.collect(date, top, obs_top, theme_cap) data = plan.collect(date, top, obs_top, theme_cap)
except RuntimeError as e: except RuntimeError as e:
raise HTTPException(status_code=404, detail=str(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": if format == "md":
return PlainTextResponse(plan.render_md(data), return PlainTextResponse(plan.render_md(data),
media_type="text/markdown; charset=utf-8") media_type="text/markdown; charset=utf-8")

View File

@ -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 # 当晚 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_URL = os.environ.get("BIONIC_SCAN_URL", "").rstrip("/")
BIONIC_SCAN_KEY = os.environ.get("BIONIC_SCAN_KEY", "") 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"))

View File

@ -230,7 +230,8 @@ README"已知的上游约束"三条在基座侧已全部修掉:传导候选不
- 决策系统夜间扫描读 Mongo 股票池时取全部分组并集、丢掉分组编号,它自己分不清哪些票来自桥;盘中白名单另起任务重建。 - 决策系统夜间扫描读 Mongo 股票池时取全部分组并集、丢掉分组编号,它自己分不清哪些票来自桥;盘中白名单另起任务重建。
- 接口层缺口:/plan 无版本戳,同日重算多版时下游分不清拿的是哪一版;桥已实现升降档字段但 PMS 的解析代码不读它决策系统的资金强度流当前只留痕不触发任何动作夜间结果表会被盘中补扫就地改写参考位漂移PMS 设了 3% 漂移熔断)。 - 接口层缺口:/plan 无版本戳,同日重算多版时下游分不清拿的是哪一版;桥已实现升降档字段但 PMS 的解析代码不读它决策系统的资金强度流当前只留痕不触发任何动作夜间结果表会被盘中补扫就地改写参考位漂移PMS 设了 3% 漂移熔断)。
- PMS 不回写上游,是桥主动读 PMS 持仓账本 pms_position成交与收益没有回流通道。PMS 自己有动作账本 pms_action_ledger含"拒了的后来涨了多少"的判分事实源)与计划快照 pms_plan_snapshot是候选单复盘的天然终点。 - 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.pypool.py 从它引用)。
- PMS 运行参数实读09-02参数表 pms_runtime_paramPMS_PLAN_TOP_N 运行值为 **100**08-17 由用户改,代码默认 30 已被覆盖PMS_AUTONOMY 为 fullPMS_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 核实) ### 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-100schemas/evaluation.py:18-54。限制六因子的因子表历史只有 07-29 起约五周传导不可回填评价样本极短heat 与 event 可用 build history 回填更长历史upside 需先做 asof 重建。 平台 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-100schemas/evaluation.py:18-54。限制六因子的因子表历史只有 07-29 起约五周传导不可回填评价样本极短heat 与 event 可用 build history 回填更长历史upside 需先做 asof 重建。
@ -357,7 +358,9 @@ README"已知的上游约束"三条在基座侧已全部修掉:传导候选不
2. 环境闸桥出环境标签PMS 读字段定仓位。**拍板后实证修正1.8e**:用过去涨跌定义的闸没有预测力,标签当前只作展示与复盘分组,不作交易前置;"弱环境不加新票"挂起,是否正式退役见本节新增拍板点。 2. 环境闸桥出环境标签PMS 读字段定仓位。**拍板后实证修正1.8e**:用过去涨跌定义的闸没有预测力,标签当前只作展示与复盘分组,不作交易前置;"弱环境不加新票"挂起,是否正式退役见本节新增拍板点。
3. 新应用顺序:平台评价闭环作废;候选单复盘闭环、吸筹确认线排最前,公告进粮在回炉之后立项,行情侧两条展示线同批。 3. 新应用顺序:平台评价闭环作废;候选单复盘闭环、吸筹确认线排最前,公告进粮在回炉之后立项,行情侧两条展示线同批。
**本轮新增拍板点(方案批准后在对话里分组选项式提问):** **新增拍板点的拍板记录2026-09-02 方案批准后对话拍板,均选推荐):** 08:45 环境追加任务走统一调度平台新建任务;代码版本解析挂载进容器的 .git 文件冻结快照缺失的四个计划日07-29、08-07、08-10、08-20用当前投影反推并标注近似两个旋钮保持默认、赛道闸保持开启环境标签来源为请决策系统加只读日频接口就绪前标签为空弱势日定义为八个指数里弱势个数达到三个环境闸只展示与分组、不动 PMS 定时、"弱环境不加新票"退役缺评分的关注票走低优先入池切片。板块映射代定主板只出全市场标签科创板与创业板出板块标签依据是桥没有市值列、逐票主参考指数要等接口返回。补两条PMS 新建仓自主档现在就改为"只提议"作止血(运行时参数,零代码,一键回退,由用户在 PMS 侧操作);新应用顺序同意把"环节启动全景""环境标签验证"插在公告进粮之前。
**本轮新增拍板点(原清单存档):**
1. 快照与版本08:45 环境追加任务走统一调度平台还是宿主定时任务;代码版本读不到时用解析仓库文件还是宿主传环境变量。 1. 快照与版本08:45 环境追加任务走统一调度平台还是宿主定时任务;代码版本读不到时用解析仓库文件还是宿主传环境变量。
2. 视图与历史重建:冻结投影快照缺失的日子用当前投影反推并标注,是否接受。 2. 视图与历史重建:冻结投影快照缺失的日子用当前投影反推并标注,是否接受。
3. 候选卡两个旋钮(启动阈值 3%、评分日龄三十日)保持默认等样本外读数;赛道闸是否放开。 3. 候选卡两个旋钮(启动阈值 3%、评分日龄三十日)保持默认等样本外读数;赛道闸是否放开。

222
plan.py
View File

@ -11,14 +11,18 @@
升降档一节对比前一交易日的档位表数据到达本身是信号首次覆盖 / 升降档一节对比前一交易日的档位表数据到达本身是信号首次覆盖 /
新进传导链即升档 新进传导链即升档
""" """
import datetime as dt
import json import json
import os import os
import pandas as pd import pandas as pd
import card
import common import common
import config import config
import db import db
import sources
import version
# 分数编码(与 factors.build_score 一致):主榜 = 200 + 传导档位×20 + 组内分, # 分数编码(与 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) 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, def collect(date: str | None = None, top: int = 20, obs_top: int = 10,
theme_cap: int = 5) -> dict: theme_cap: int = 5) -> dict:
"""装配一天的计划为结构化字典。数据缺失抛 RuntimeErrorapi 侧转 404""" """装配一天的计划为结构化字典。数据缺失抛 RuntimeErrorapi 侧转 404
2026-09-02 起附带候选卡main / observe 的装配排序裁剪一字不动下游 PMS 只读
这两段的既有字段每行只是多联入判决类字段顶层新增 generated_atplan_version
card_countscandidates候选单全量不受裁剪watchsegments_pointed
`_full` 是全量主榜与观察档行只给 generate 落快照用api 返回前会去掉"""
ds = date or _latest_date("t_factor_akg_score") ds = date or _latest_date("t_factor_akg_score")
if not ds: if not ds:
raise RuntimeError("t_factor_akg_score 还没有数据——先 build akg_score。") 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) main = score[score >= _MAIN_MIN].sort_values(ascending=False)
obs = 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): def _pick(ranked: pd.Series, n: int):
"""分数从高到低取 n 条;每个传导主题最多 theme_cap 条0=不设限)—— """分数从高到低取 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: def _row(rank: int, k: str, s: float, with_tier: bool) -> dict:
e = ev.get(k) 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), r = {"rank": rank, "code": k, "name": names.get(k),
"score": round(float(s), 2), "score": round(float(s), 2),
"evidence": ({"theme": e[0], "n_sources": e[1], "evidence": evidence,
"moved_ratio": round(e[2], 4)} if e else None),
"heat": _val(heat, k), "upside": _val(upside, k)} "heat": _val(heat, k), "upside": _val(upside, k)}
if with_tier: if with_tier:
r["tier"] = _tier_label(s) 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 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 changes = None
prev_ds = _prev_date("t_factor_akg_gate", ds) prev_ds = _prev_date("t_factor_akg_gate", ds)
if prev_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)), "upgrades": _mv(up_df.head(15)),
"downgrades": _mv(down_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 { return {
"date": ds, "date": ds,
"generated_at": generated_at,
"plan_version": version.git_short_rev(),
"counts": {"main": int(len(main)), "observe": int(len(obs)), "counts": {"main": int(len(main)), "observe": int(len(obs)),
"gate_covered": int(len(gate))}, "gate_covered": int(len(gate))},
"market_snapshot_days": sorted(mkt_days), "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), "gate_on": bool(config.ENABLE_TRACK_GATE),
"encoding": "主榜分=200+传导档位×20+组内分(还没热、还便宜);" "encoding": "主榜分=200+传导档位×20+组内分(还没热、还便宜);"
"观察档分=100+0.6z(传导)+0.4z(−热度)", "观察档分=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: def _fmt_ev(e) -> str:
if not e: if not e:
return "" 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%}" 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: def render_md(d: dict) -> str:
L = [f"# 每日选股计划 · {d['date']}", ""] L = [f"# 每日选股计划 · {d['date']}", ""]
c = d["counts"] c = d["counts"]
@ -272,6 +415,59 @@ def render_md(d: dict) -> str:
f"(与计划日不同——历史降级日口径)。") f"(与计划日不同——历史降级日口径)。")
L.append("") 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 "" cap_txt = f",每主题限额 {d['theme_cap']}" if d["theme_cap"] else ""
gate_txt = "、在十五五赛道内" if d.get("gate_on") else "" gate_txt = "、在十五五赛道内" if d.get("gate_on") else ""
L.append(f"## 主榜 Top {len(d['main'])}" 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: except RuntimeError as e:
raise SystemExit(str(e)) raise SystemExit(str(e))
text = render_md(data) text = render_md(data)
os.makedirs("data/plan", exist_ok=True) os.makedirs(config.PLAN_SNAPSHOT_DIR, exist_ok=True)
out = f"data/plan/plan_{data['date']}.md" out = os.path.join(config.PLAN_SNAPSHOT_DIR, f"plan_{data['date']}.md")
with open(out, "w", encoding="utf-8") as f: with open(out, "w", encoding="utf-8") as f:
f.write(text + "\n") 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(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 return out

55
pool.py
View File

@ -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, data = plan.collect(date, top=max(top * 10, 200), obs_top=0,
theme_cap=config.POOL_THEME_CAP) theme_cap=config.POOL_THEME_CAP)
tiers = config.POOL_TIERS 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"] 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. 持仓(读不到直接抛,整轮不动池子)与旧池子 # 2. 持仓(读不到直接抛,整轮不动池子)与旧池子
holdings = _read_holdings() 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 col_name = config.POOL_COLLECTION
client = None if dry_run and not config_mongo_ready() else _mongo() client = None if dry_run and not config_mongo_ready() else _mongo()
old_doc, old_members, member_meta = None, set(), {} 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) config.POOL_MAX, today)
# 4. 打印明细(干跑到此为止) # 4. 打印明细(干跑到此为止)
print(f"计划日 {ds},档位白名单 {sorted(tiers) if tiers else '(不过滤)'}" src_cnt = {}
f"计划入选 {len(plan_codes)} 只;持仓 {len(holdings)} 只;" for r in plan_rows:
f"旧池 {len(old_members)} 只 → 新池 {len(d['pool'])} 只(上限 {config.POOL_MAX}") 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"]), for label, items in (("计划新进", d["new_entrants"]),
("持仓保留", d["retained_holdings"]), ("持仓保留", d["retained_holdings"]),
("留池观察", d["observers"]), ("留池观察", 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"), ctx[r["code"]] = {"factor_code": "akg_score", "score": r.get("score"),
"tier": r.get("tier"), "upside": r.get("upside"), "tier": r.get("tier"), "upside": r.get("upside"),
"theme": (r.get("evidence") or {}).get("theme"), "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"]: for c in d["retained_holdings"]:
ctx.setdefault(c, {"factor_code": "akg_score", "note": "持仓保留"}) ctx.setdefault(c, {"factor_code": "akg_score", "note": "持仓保留"})

125
regime.py Normal file
View File

@ -0,0 +1,125 @@
"""环境标签:只读决策系统的日频市场区制接口,写进当日计划快照;只展示与复盘分组,不作交易前置。
## 为什么要有它
实证docs/主观选股改进方案_2026-09-02.md 1.81.8e市场环境是第一解释变量
"看过去五日或十日涨跌"没有预测力相关 0.280.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

9
run.py
View File

@ -180,6 +180,8 @@ def main():
av = sub.add_parser("apply-views") av = sub.add_parser("apply-views")
av.add_argument("--file", default="sql/astock_kg_slot_views.sql") av.add_argument("--file", default="sql/astock_kg_slot_views.sql")
av.add_argument("--dry-run", action="store_true", help="只列语句不执行") 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 = sub.add_parser("freeze")
f.add_argument("--date", help="默认今天") f.add_argument("--date", help="默认今天")
b = sub.add_parser("build") b = sub.add_parser("build")
@ -210,6 +212,13 @@ def main():
elif a.cmd == "push-pool": elif a.cmd == "push-pool":
import pool import pool
pool.push(a.date, a.top, dry_run=a.dry_run, kick=not a.no_kick) 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": elif a.cmd == "tracks":
import tracks import tracks
tracks.coverage_report() tracks.coverage_report()

164
sources.py Normal file
View File

@ -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

91
test_regime.py Normal file
View File

@ -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()

13
xxl.py
View File

@ -46,12 +46,17 @@ HERE = os.path.dirname(os.path.abspath(__file__))
LOG_PATH = os.path.join(HERE, "data", "xxl_build.log") LOG_PATH = os.path.join(HERE, "data", "xxl_build.log")
# 步骤 → 命令(顺序即执行顺序;--date 由触发参数统一追加) # 步骤 → 命令(顺序即执行顺序;--date 由触发参数统一追加)
STEP_ORDER = ("build", "plan", "push-pool") STEP_ORDER = ("build", "plan", "push-pool", "regime-append")
STEP_CMDS = { STEP_CMDS = {
"build": ["run.py", "build", "all", "--mode", "daily"], "build": ["run.py", "build", "all", "--mode", "daily"],
"plan": ["run.py", "plan"], "plan": ["run.py", "plan"],
"push-pool": ["run.py", "push-pool"], "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 # 单步上限一小时,防呆死(正常盘前链全程分钟级) STEP_TIMEOUT_SEC = 3600 # 单步上限一小时,防呆死(正常盘前链全程分钟级)
_lock = threading.Lock() _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"]) @router.api_route("/daily-build", methods=["GET", "POST"])
def trigger_daily_build( def trigger_daily_build(
callbackUrl: str = Query(None, description="平台下发的结案回调地址(自动追加)"), callbackUrl: str = Query(None, description="平台下发的结案回调地址(自动追加)"),
steps: str = Query("build,plan,push-pool", steps: str = Query(",".join(DEFAULT_STEPS),
description="要跑哪几步,逗号分隔;顺序固定 build→plan→push-pool"), description="要跑哪几步,逗号分隔;默认三步 build→plan→push-pool。"
"环境追加 regime-append 不在默认里,平台 08:45 的任务单独传 "
"steps=regime-append决策系统区制快照 08:40 才预热)"),
date: str = Query(None, description="补跑指定数据日 YYYY-MM-DD空=最新数据日"), date: str = Query(None, description="补跑指定数据日 YYYY-MM-DD空=最新数据日"),
key: str = Query(None), key: str = Query(None),
x_job_key: str = Header(None), x_job_key: str = Header(None),