diff --git a/api.py b/api.py index e6b83f0..b32a49d 100644 --- a/api.py +++ b/api.py @@ -22,6 +22,7 @@ from fastapi.responses import PlainTextResponse import config import db +import logic_state_daily import plan import plan_reconcile import regime @@ -135,3 +136,53 @@ def plan_verdict(codes: str | None = None, code: str | None = None, verdicts.append(v) return {"date": d, "stale": L.get("stale", ""), "count": len(verdicts), "verdicts": verdicts} + + +@app.get("/logic_state") +def logic_state_lookup(codes: str | None = None, code: str | None = None, + date: str | None = None): + """逐票逻辑状态四态(只读)——给 PMS 早上拉完计划后查在持票用(2026-09-07 第三件桥侧前置)。 + + GET /logic_state?codes=300750,SH688041 多只(逗号分隔) + GET /logic_state?code=300750&date=2026-09-04 单只 + 指定数据日 + + 持仓票在候选筛选第一步就被整行剔掉,/plan 里读不到它,所以要有这个入口。 + 每只票优先回逐票日频表里当日那一行(早上生成计划时落定的,带代码版本,可回溯), + 标 source=daily;当日表里没有这只票(不在档位表、或当日还没生成)就按此刻的数据现算、 + 按表里的历史做抗抖动,标 source=computed——现算的是"此刻"不是"早上",复盘别把它当当日态。 + date 缺省=档位表最新日(与 /plan 的数据日同口径)。三种代码形态都收。纯 SELECT,不写任何库。 + """ + raw = (codes or code or "").strip() + want = [c.strip() for c in raw.split(",") if c.strip()] + if not want: + raise HTTPException(status_code=400, detail="缺少 code / codes 参数") + d = date or plan._latest_date("t_factor_akg_score") # noqa: SLF001 —— 桥内自用 + if not d: + raise HTTPException(status_code=404, detail="档位表为空——先跑当日构建") + keyed, bad = {}, [] + for c in want: + try: + keyed.setdefault(plan_reconcile._norm_code(c), c) # noqa: SLF001 + except SystemExit as e: # _norm_code 认不出的形态会 raise SystemExit + bad.append({"input": c, "error": str(e)}) + stored = logic_state_daily.lookup(list(keyed), d) + todo = [k for k in keyed if k not in stored] + computed = {} + if todo: + try: + computed = plan.logic_states_for(todo, d) + except Exception as e: # noqa: BLE001 —— 数据层异常统一收成 500 + raise HTTPException(status_code=500, detail=f"现算逻辑状态失败: {e!r}") + out = [] + for k, c in keyed.items(): + if k in stored: + row = stored[k] + out.append({"input": c, "code": k, "date": d, "source": "daily", + "verdict": row.get("verdict"), "card_rank": row.get("card_rank"), + "plan_version": row.get("plan_version"), + **logic_state_daily.row_to_out(row)}) + else: + out.append({"input": c, "code": k, "date": d, "source": "computed", + "verdict": None, "card_rank": None, "plan_version": None, + **(plan._state_out(computed.get(k)) or {})}) # noqa: SLF001 + return {"date": d, "count": len(out) + len(bad), "states": out + bad} diff --git a/config.py b/config.py index 7ad56bc..7e9a3cd 100644 --- a/config.py +++ b/config.py @@ -271,3 +271,16 @@ JUDGEMENT_SCOPES = {s.strip() for s in # 比对上一版时往前看几个自然日。取到窗口内最近一个计划日的行作为上一版;超过这个窗口没写过 # 快照的主题,本次按"第一次见到"处理,迁移记为空、陈旧天数从零重新起算。 JUDGEMENT_PREV_LOOKBACK_DAYS = int(os.environ.get("JUDGEMENT_PREV_LOOKBACK_DAYS", "60")) + +# --- 逐票日频逻辑状态表(2026-09-07 下一阶段方案第三件桥侧前置;模块见 logic_state_daily.py)------- +# 每个数据日为档位表里的每只票记一行原始态与落定态。只在早上 plan.generate 里写,/plan 实时重算只读。 +# 同库同理由:桥对数据基座只有只读账号,这张表又是选股系统自己的派生记录。键是(数据日,前缀码), +# 数据日与 t_factor_akg_* 的 trade_date 同口径,不是行业观点快照的 plan_date(那个是写入当天)。 +LOGIC_STATE_TABLE = os.environ.get("LOGIC_STATE_TABLE", "t_akg_logic_state_daily") +# 抗抖动的两个天数(计划日数),与 logic_state.settle 的默认值同源,一次定死、不按复盘读数回调: +# 其余方向的迁移要连续几个计划日同向;退出逻辑存疑要连续几个计划日不再存疑(进入存疑即刻成立)。 +LOGIC_SETTLE_CONFIRM_DAYS = int(os.environ.get("LOGIC_SETTLE_CONFIRM_DAYS", "2")) +LOGIC_SETTLE_EXIT_DAYS = int(os.environ.get("LOGIC_SETTLE_EXIT_DAYS", "3")) +# 读上一次落定态与近几日原始态时往前看几个自然日。要盖住国庆这种连休八天再加两个周末的空档; +# 超过这个窗口没记过行的票按第一次见到处理,落定态等于原始态。 +LOGIC_STATE_LOOKBACK_DAYS = int(os.environ.get("LOGIC_STATE_LOOKBACK_DAYS", "30")) diff --git a/logic_state_daily.py b/logic_state_daily.py new file mode 100644 index 0000000..21bfacc --- /dev/null +++ b/logic_state_daily.py @@ -0,0 +1,282 @@ +"""逐票日频逻辑状态表:每个数据日为档位表里的每只票记一行原始态与落定态,攒版本史并做抗抖动。 + +## 为什么要有它(2026-09-07,下一阶段方案第三件桥侧前置) + +逻辑状态四态(logic_state.py)在计划装配里已经随每张卡产出,但有两件事它自己做不到: + + 一,抗抖动。settle 那个函数写好了,可它要看"昨天落定的是什么、近几天原始态是什么", + 这些只有存下来才有。不存,四态一天一变,PMS 按它分流就会跟着一天一变。 + 二,给持仓票用。四态唯一真正有用的地方是持仓,而持仓票在候选筛选第一步就被整行剔掉, + /plan 里根本没有它。接口 /logic_state 要能回答"这只在持的票今天证据还在不在", + 就得有一张按票、按日可查的表。 + +## 落点、表名、键 + +写在平台因子库(写 t_factor_akg_* 的同一个 MySQL),与行业观点快照 t_akg_judgement_snapshot +同库同理由:桥对数据基座只有只读账号,这张表又是选股系统自己的派生记录。表名默认 +t_akg_logic_state_daily,不带 t_factor_ 前缀,免得平台的因子清单把它当因子表收进去。 + +键是(数据日,前缀码)。数据日就是计划文件名里那个日期(plan_),与 t_factor_akg_* 的 +trade_date 同口径——不是行业观点快照的 plan_date(那个是写入当天)。两张表的日期口径不同, +读的时候别混:计划在 D+1 凌晨构建、数据日是 D,这张表记的是 D。 + +## 只在早上那一次生成里写 + +写入点只有 plan.generate(run.py plan,调度中心的默认三步之一)。/plan 每次实时重算会读这张表 +做抗抖动,但不写——实时重算一天几十次,写进去会把"当天落定态"改成"最后一次有人查时的状态", +派生数据就不可回溯了。重跑同一天幂等:先删该日行再整批插。 + +## 落定的口径(settle_one) + + 上一次落定态 = 表里这只票严格早于数据日、回看窗口之内最近一行的 state。 + 近几日原始态 = 那几行的 raw_state 加上今天的原始态,最新的在最后。 + 落定规则在 logic_state.settle:进入逻辑存疑即刻成立;退出存疑要连续几个计划日不再存疑; + 其余迁移要连续几个计划日同向。天数由 config 给,一次定死,不按复盘读数回调。 + 表里没有这只票(第一次见到、或超出回看窗口)时,落定态就是原始态——没有历史就没有抖动可抗。 + +## 持仓票不在档位表里时的局限(写进设计不回避) + +这张表只记档位表里的票。一只在持的票掉出档位表之后,每天没有新行,接口现算时能用的历史只有 +它掉出去之前那几行;掉出去超过回看窗口,落定态就等于原始态。要让持仓票也逐日攒历史, +得让桥知道 PMS 的持仓,那是跨层的事,本阶段不做,先在接口返回里用 source=computed 标明。 + +离线单测见 test_logic_state_daily.py,不连库。 +""" +from __future__ import annotations + +import datetime as dt +import json + +import config +import db +import judgement +import logic_state +import sources + +# 落库列的顺序(插入语句按此顺序拼参数)。 +COLUMNS = ("trade_date", "code", "raw_state", "raw_why", "state", "settle_note", "prev_state", + "as_of", "usable", "missing", "paths", "reasons", + "verdict", "card_rank", "plan_version", "snapshot_at") + +# 主键(数据日,前缀码):一天一票最多一行,重跑覆盖;另建(前缀码,数据日)索引,按票回看走它。 +_CREATE_TABLE = """ +CREATE TABLE IF NOT EXISTS {t} ( + trade_date DATE NOT NULL, + code VARCHAR(16) NOT NULL, + raw_state VARCHAR(16), + raw_why VARCHAR(16), + state VARCHAR(16), + settle_note VARCHAR(128), + prev_state VARCHAR(16), + as_of DATE, + usable VARCHAR(64), + missing VARCHAR(64), + paths TEXT, + reasons TEXT, + verdict VARCHAR(16), + card_rank INT, + plan_version VARCHAR(32), + snapshot_at DATETIME, + PRIMARY KEY (trade_date, code), + KEY idx_code (code, trade_date) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 +""" + +# 四路各自的分隔:usable / missing 两列存"研报论断、券商行动"这种顿号串,读回来再拆。 +_SEP = "、" + + +# ============================================================================ +# 抗抖动:一只票的原始态配上历史,得到落定态 +# ============================================================================ + +def settle_one(code: str, raw: dict, hist: list | None, *, confirm_days: int | None = None, + exit_days: int | None = None) -> dict: + """纯函数:把一只票今天的原始合成结果,配上它在本表里的近日行,落定成今天的状态。 + + raw 是 logic_state.compose 的返回(state 是原始态)。hist 是这只票严格早于今天、 + 按数据日升序的近日行(history 的返回),每行至少有 raw_state 与 state 两列; + 没有历史传 None 或空列表,那时落定态就是原始态。 + + 返回一个新字典:state 改成落定态,原始态挪到 raw_state,另带 settle_note 与 prev_state; + 其余键(why、paths、usable、missing、reasons、as_of)原样保留。不改传入的 raw。 + """ + confirm = int(config.LOGIC_SETTLE_CONFIRM_DAYS if confirm_days is None else confirm_days) + exit_n = int(config.LOGIC_SETTLE_EXIT_DAYS if exit_days is None else exit_days) + rows = [r for r in (hist or []) if isinstance(r, dict)] + raw_state = raw.get("state") if isinstance(raw, dict) else None + prev_state = rows[-1].get("state") if rows else None + need = max(confirm, exit_n) - 1 + recent = [r.get("raw_state") for r in (rows[-need:] if need > 0 else [])] + [raw_state] + settled, note = logic_state.settle(prev_state, raw_state, recent, + confirm_days=confirm, exit_doubt_days=exit_n) + out = dict(raw or {}) + out.update(raw_state=raw_state, state=settled, settle_note=note, prev_state=prev_state) + return out + + +# ============================================================================ +# 读:历史行(做抗抖动)、当日行(给接口) +# ============================================================================ + +def _in_clause(codes) -> tuple[str, tuple]: + codes = [str(c).strip() for c in (codes or []) if str(c).strip()] + return ",".join(["%s"] * len(codes)), tuple(codes) + + +def history(ds: str, codes=None, read_mysql=None) -> dict: + """每只票严格早于数据日、回看窗口之内的行,按前缀码索引、行按数据日升序。 + + codes 不传就取全部票(计划装配一轮要看一千多只,一次查询取回,不逐票查);传了就只取 + 这几只(接口用)。只取抗抖动要读的几列。表还没建、或读不到,返回空字典并打印一行原因: + 那样全部票按第一次见到处理,落定态等于原始态,计划照出不断产。 + """ + table = config.LOGIC_STATE_TABLE + reader = read_mysql or db.read_mysql + try: + start = (dt.date.fromisoformat(ds) + - dt.timedelta(days=config.LOGIC_STATE_LOOKBACK_DAYS)).isoformat() + except ValueError: + print(f" (数据日 {ds!r} 不是合法日期,逻辑状态不做抗抖动)") + return {} + sql = (f"SELECT trade_date, code, raw_state, state FROM {table} " + f"WHERE trade_date >= %s AND trade_date < %s") + params: tuple = (start, ds) + if codes is not None: + marks, vals = _in_clause(codes) + if not vals: + return {} + sql += f" AND code IN ({marks})" + params = params + vals + try: + rows = sources._records(reader("factor", sql + " ORDER BY code, trade_date", params)) # noqa: SLF001 + except Exception as e: # noqa: BLE001 —— 头一次跑时表还不存在,属正常 + print(f" ({table} 读不到历史行,逻辑状态本次不做抗抖动,落定态等于原始态: {e!r})") + return {} + out: dict = {} + for r in rows: + k = str(r.get("code") or "").strip() + if k: + out.setdefault(k, []).append({ + "trade_date": sources._ymd(r.get("trade_date")), # noqa: SLF001 + "raw_state": judgement._blank_to_none(r.get("raw_state")), # noqa: SLF001 + "state": judgement._blank_to_none(r.get("state"))}) # noqa: SLF001 + return out + + +def lookup(codes, ds: str, read_mysql=None) -> dict: + """这几只票在数据日当天的行,按前缀码索引;没写过、表没建、读不到都返回空字典。""" + table = config.LOGIC_STATE_TABLE + reader = read_mysql or db.read_mysql + marks, vals = _in_clause(codes) + if not vals: + return {} + try: + rows = sources._records(reader( # noqa: SLF001 + "factor", f"SELECT * FROM {table} WHERE trade_date = %s AND code IN ({marks})", + (ds,) + vals)) + except Exception as e: # noqa: BLE001 + print(f" ({table} 当日行读取失败: {e!r})") + return {} + out = {} + for r in rows: + k = str(r.get("code") or "").strip() + if k: + out[k] = {c: judgement._blank_to_none(v) for c, v in dict(r).items()} # noqa: SLF001 + return out + + +def row_to_out(row: dict) -> dict: + """把表里的一行还原成发给下游的形状(与 plan._state_out 同形,多带 raw_state 等三键)。""" + def _split(v): + s = str(v or "").strip() + return [x for x in s.split(_SEP) if x] if s else [] + + def _loads(v, default): + if v is None or v == "": + return default + try: + return json.loads(v) if isinstance(v, str) else v + except (TypeError, ValueError): + return default + + return {"state": row.get("state"), "raw_state": row.get("raw_state"), "why": row.get("raw_why"), + "settle_note": row.get("settle_note"), "prev_state": row.get("prev_state"), + "as_of": sources._ymd(row.get("as_of")), # noqa: SLF001 + "usable": _split(row.get("usable")), "missing": _split(row.get("missing")), + "reasons": _loads(row.get("reasons"), []), "paths": _loads(row.get("paths"), [])} + + +# ============================================================================ +# 写:只在 plan.generate 里调 +# ============================================================================ + +def build_rows(ds: str, plan_rows: list, plan_version: str | None = None, + now: str | None = None) -> list: + """纯函数:把计划快照的行(plan.collect 里 _row 的输出,带 logic_state)拼成可落库的行。 + + 没有 logic_state 的行(档位表里有票但没装配出卡)跳过;同一只票出现多次只留最后一条。""" + stamp = now or dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S") + by_code = {} + for r in plan_rows or []: + st = r.get("logic_state") if isinstance(r, dict) else None + code = str((r or {}).get("code") or "").strip() + if not code or not isinstance(st, dict): + continue + by_code[code] = (r, st) + out = [] + for code in sorted(by_code): + r, st = by_code[code] + out.append({ + "trade_date": ds, "code": code[:16], + "raw_state": st.get("raw_state") or st.get("state"), "raw_why": st.get("why"), + "state": st.get("state"), "settle_note": (st.get("settle_note") or "")[:128] or None, + "prev_state": st.get("prev_state"), "as_of": st.get("as_of"), + "usable": _SEP.join(st.get("usable") or [])[:64] or None, + "missing": _SEP.join(st.get("missing") or [])[:64] or None, + "paths": json.dumps(st.get("paths") or [], ensure_ascii=False, default=str), + "reasons": json.dumps((st.get("reasons") or [])[:6], ensure_ascii=False, default=str), + "verdict": r.get("verdict"), + "card_rank": r.get("card_rank"), + "plan_version": (plan_version or "")[:32] or None, + "snapshot_at": stamp, + }) + return out + + +def save(ds: str, rows: list, conn_factory=None) -> None: + """幂等落库:建表(已存在就跳过)、删该日行、整批插,删与插同一事务。与 judgement.save 同范式。 + + rows 为空也照样删该日行:那表示这一天一张卡都没装配出来,不该留上一次重跑的残行。""" + table = config.LOGIC_STATE_TABLE + factory = conn_factory or db.factor_conn + cols = ",".join(COLUMNS) + marks = ",".join(["%s"] * len(COLUMNS)) + payload = [tuple(r.get(c) for c in COLUMNS) for r in rows] + with factory() as conn: + with conn.cursor() as cur: + cur.execute(_CREATE_TABLE.format(t=table)) + conn.commit() + with conn.cursor() as cur: + cur.execute(f"DELETE FROM {table} WHERE trade_date = %s", (ds,)) + if payload: + cur.executemany(f"INSERT INTO {table} ({cols}) VALUES ({marks})", payload) + conn.commit() + + +def persist(ds: str, plan_rows: list, plan_version: str | None = None, write=None) -> dict: + """plan.generate 的一步:从快照行拼表行、幂等写该日、打印读数。write 可注入(离线单测)。""" + rows = build_rows(ds, plan_rows, plan_version) + (write or save)(ds, rows) + by_state: dict = {} + for r in rows: + by_state[r["state"] or "空"] = by_state.get(r["state"] or "空", 0) + 1 + moved = [r for r in rows if r["prev_state"] and r["prev_state"] != r["state"]] + held = [r for r in rows if r["raw_state"] != r["state"]] + print(f"逐票逻辑状态 {ds}:写入 {len(rows)} 行" + f"({'、'.join(f'{k} {v}' for k, v in sorted(by_state.items())) or '无'})" + f",落定态迁移 {len(moved)} 只,被抗抖动按住 {len(held)} 只 -> {config.LOGIC_STATE_TABLE}") + for r in moved[:10]: + print(f" 迁移:{r['code']} {r['prev_state']} -> {r['state']}({r['settle_note']})") + return {"date": ds, "rows": len(rows), "by_state": by_state, "moved": len(moved), + "held": len(held), "table": config.LOGIC_STATE_TABLE} diff --git a/plan.py b/plan.py index 23c2cd3..5b68f1b 100644 --- a/plan.py +++ b/plan.py @@ -29,6 +29,7 @@ import common import config import judgement import logic_state +import logic_state_daily import db import sources import version @@ -144,6 +145,11 @@ def _state_out(st) -> dict | None: if not isinstance(st, dict): return None return {"state": st.get("state"), "why": st.get("why"), "as_of": st.get("as_of"), + # 落定态在 state 里(下游按它分流);原始态、落定说明、上一次落定态三键一并发出, + # 让人看得出"今天原始判的是什么、为什么被按住没迁"(2026-09-07 第三件桥侧前置)。 + # 没做过落定的(旧调用、单测)原始态就是 state 本身。 + "raw_state": st.get("raw_state", st.get("state")), + "settle_note": st.get("settle_note"), "prev_state": st.get("prev_state"), "usable": st.get("usable") or [], "missing": st.get("missing") or [], "reasons": (st.get("reasons") or [])[:6], "paths": [{"path": p.get("path"), "signal": p.get("signal"), @@ -200,6 +206,56 @@ def _logic_state_of(k: str, evd: dict, seg_of: dict, seg_view: dict, return logic_state.compose([a, b, c, logic_state.from_events()]) +def _logic_inputs(codes: list, ds: str) -> dict: + """逻辑状态四态的取数,一次取全(计划装配与接口现算共用这一处,口径改了两边一起变)。 + + 甲路(因果论断):2026-09-03 起挂在卡上作证据线,只展示不进判决;视图未建时为空。 + 2026-09-07 审查修:四态的甲路要吃**全量**论断,卡上展示只取最近几条。原先两者共用 + 一份截断到三条的列表,跨期翻转与同期分歧都只在三条上算,翻转两头都会判错—— + 历史八空两好加最新一空,全量判无法判断、三条判逻辑存疑;历史五好近三空则反过来。 + 全量取一次,展示从里面切前几条,取数层在截断前算的质量画像照旧挂在每条上。 + + 乙路(产业研判):行业观点快照在交易日 D 的晚上 20:40 写,plan_date 记的是 D;计划在 D+1 + 凌晨构建,数据日 ds 就是 D。所以要读的是 plan_date 不晚于 ds 的最近一版,不是 ds+1。 + 09-04 之前这里读的是 ds+1,能对上只因为那天上午手工跑过一次快照——把一次手工操作 + 的时序写进了代码,到周一就全部落空(台账 033 记的"快照断了"其实是这个错, + 调度中心一直在按时跑)。load_previous(D+1) 取的正是 plan_date <= D 的每簇最新一行, + 顺带也扛得住某个晚上没跑:会退回上一版,陈旧天数照常累加。seg_hist 是每簇近几个 + 计划日的快照,给乙路判"硬触发要不要维持"(2026-09-07 审查修)。 + + 丙路(券商行动):两个等长窗口的每股收益预测中位数与机构数。丁路无数据源,logic_state 恒出缺失。 + + hist 是逐票日频表里每只票的近日行,给抗抖动用(2026-09-07 第三件桥侧前置);表没建或 + 读不到为空字典,那时落定态等于原始态,计划照出。这张表只在早上 generate 里写,这里只读。 + """ + return { + "logic_full": sources.logic_claims(codes, ds, per_stock=config.LOGIC_CLAIMS_FULL), + "seg_view": judgement.by_segment_name(judgement.load_previous(_next_day(ds))), + "seg_hist": judgement.recent_rows(_next_day(ds), days=config.JUDGEMENT_HOLD_DAYS), + "broker": sources.broker_actions(codes, ds), + "seg_of": _segments_of(ds), + "hist": logic_state_daily.history(ds, codes=None if len(codes) > 50 else codes), + } + + +def logic_states_for(codes: list, ds: str) -> dict: + """给定几只票,按此刻的数据算它们的逻辑状态并按逐票日频表的历史落定,按前缀码索引。 + + 给接口 /logic_state 用:持仓票在候选筛选第一步就被整行剔掉,计划装配根本不会走到它。 + 取数与计划装配走同一个 _logic_inputs,口径逐字相同;落定与计划装配走同一个 settle_one。 + """ + codes = [str(c).strip() for c in (codes or []) if str(c).strip()] + if not codes: + return {} + inp = _logic_inputs(codes, ds) + out = {} + for k in codes: + raw = _logic_state_of(k, {"logic": []}, inp["seg_of"], inp["seg_view"], inp["broker"], ds, + full_logic=inp["logic_full"].get(k) or [], seg_hist=inp["seg_hist"]) + out[k] = logic_state_daily.settle_one(k, raw, inp["hist"].get(k)) + return out + + def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series, mkt_days: set, risk: set | None = None) -> tuple[dict, list]: """候选卡装配(2026-09-02 方案第 2.2 节第三项):对档位表里的全部票(主榜与观察档, @@ -214,28 +270,14 @@ def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series, moved = sources.moved_members(ds) daily = sources.stock_daily(ds) night = sources.night_conclusions(codes, ds) - # 因果论断(2026-09-03):数据基座抽取的论断挂在卡上作证据线,只展示不进判决;视图未建时为空。 - # 2026-09-07 审查修:四态的甲路要吃**全量**论断,卡上展示只取最近几条。原先两者共用 - # 一份截断到三条的列表,跨期翻转与同期分歧都只在三条上算,翻转两头都会判错—— - # 历史八空两好加最新一空,全量判无法判断、三条判逻辑存疑;历史五好近三空则反过来。 - # 全量取一次,展示从里面切前几条,取数层在截断前算的质量画像照旧挂在每条上。 - logic_full = sources.logic_claims(codes, ds, per_stock=config.LOGIC_CLAIMS_FULL) - logic = {k: v[:config.LOGIC_CLAIMS_PER_STOCK] for k, v in logic_full.items()} - # 逻辑状态四态的三路输入(丁路公司事件无数据源,logic_state 那边恒出缺失)。 + # 逻辑状态四态的取数(三路输入加逐票日频表的近日行;取数口径与理由见 _logic_inputs)。 # 这三路都是"研究证据还在不在"的跟踪,与候选卡的三门槛判决是正交的两维: # 判决回答今天要不要买,四态回答支撑它的研究证据还在不在。收敛规则在 # logic_state.apply_to_card,是单调的——强化只能提前卡内序、永远不升判决。 - # 行业观点快照在交易日 D 的晚上 20:40 写,plan_date 记的是 D;计划在 D+1 凌晨构建, - # 数据日 ds 就是 D。所以要读的是 plan_date 不晚于 ds 的最近一版,不是 ds+1。 - # 09-04 之前这里读的是 ds+1,能对上只因为那天上午手工跑过一次快照——把一次手工操作 - # 的时序写进了代码,到周一就全部落空(台账 033 记的"快照断了"其实是这个错, - # 调度中心一直在按时跑)。load_previous(D+1) 取的正是 plan_date <= D 的每簇最新一行, - # 顺带也扛得住某个晚上没跑:会退回上一版,陈旧天数照常累加。 - seg_view = judgement.by_segment_name(judgement.load_previous(_next_day(ds))) - # 每簇近几个计划日的快照,给乙路判"硬触发要不要维持"(2026-09-07 审查修)。 - seg_hist = judgement.recent_rows(_next_day(ds), days=config.JUDGEMENT_HOLD_DAYS) - broker = sources.broker_actions(codes, ds) - seg_of = _segments_of(ds) + inp = _logic_inputs(codes, ds) + logic_full, seg_view, seg_hist = inp["logic_full"], inp["seg_view"], inp["seg_hist"] + broker, seg_of, hist = inp["broker"], inp["seg_of"], inp["hist"] + logic = {k: v[:config.LOGIC_CLAIMS_PER_STOCK] for k, v in logic_full.items()} if risk is None: # collect 会传入读过一次的名单;单独调用时自己读 try: risk = factors._risk_set() or set() # noqa: SLF001 —— 同仓自用 @@ -256,8 +298,12 @@ def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series, "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, "logic": logic.get(k) or []} - state = _logic_state_of(k, evd, seg_of, seg_view, broker, ds, - full_logic=logic_full.get(k) or [], seg_hist=seg_hist) + raw_state = _logic_state_of(k, evd, seg_of, seg_view, broker, ds, + full_logic=logic_full.get(k) or [], seg_hist=seg_hist) + # 抗抖动(2026-09-07 第三件桥侧前置):配上这只票在逐票日频表里的近日行,落定今天的状态。 + # 进入逻辑存疑即刻成立,退出要连续几天不再存疑,其余迁移要连续几天同向;表里没有 + # 这只票时落定态就是原始态。落定态进 state,原始态另存 raw_state,两个都发给下游。 + state = logic_state_daily.settle_one(k, raw_state, hist.get(k)) j = card.judge(evd, start_pct=config.CARD_START_PCT, accum_max_age=config.CARD_ACCUM_MAX_AGE, neg_tol=config.UPSIDE_NEG_TOLERANCE, @@ -758,6 +804,13 @@ def generate(date: str | None = None, top: int = 20, obs_top: int = 10, with open(tmp, "w", encoding="utf-8") as f: json.dump(snap, f, ensure_ascii=False, indent=1, default=str) os.replace(tmp, jpath) + # ---- 逐票日频逻辑状态表(2026-09-07 第三件桥侧前置):只在这一次生成里写,/plan 的实时重算 + # 只读不写。写在快照落盘之后:表写失败只打印原因,计划文件与快照照出,不断产。---- + try: + logic_state_daily.persist(data["date"], snap["main"] + snap["observe"], + data.get("plan_version")) + except Exception as e: # noqa: BLE001 —— 表写失败不能拖垮出计划 + print(f" (逐票逻辑状态表写入失败,计划照出,抗抖动明天少一天历史: {e!r})") print(text) print(f"\n已写入 {out} 与快照 {jpath}" f"(主榜 {len(snap['main'])} 行、观察档 {len(snap['observe'])} 行、" diff --git a/test_logic_state_daily.py b/test_logic_state_daily.py new file mode 100644 index 0000000..ce780c0 --- /dev/null +++ b/test_logic_state_daily.py @@ -0,0 +1,239 @@ +"""逐票日频逻辑状态表的离线单测(不连库)。 + +钉住四件事: + 一,落定的口径:第一次见到落定态等于原始态;进入逻辑存疑即刻成立;退出存疑要连续三个计划日 + 不再存疑;其余迁移要连续两个计划日同向;落定不改传入的原始结果。 + 二,历史读不到不断产:表没建、日期不合法都返回空字典,落定态退回原始态。 + 三,落库幂等:同一数据日重跑先删该日行再整批插,列数与 COLUMNS 一致,不追加重复行。 + 四,表行与下游形状可往返:写进去的顿号串与 JSON 串读回来仍是列表。 + +落库那一路用一张内存里的假表接住真实的 SQL(与 test_judgement_snapshot.py 同一手法), +所以幂等测的是真语句不是桩。开发机没有 pandas 与数据库驱动时只给缺席的模块装最小桩。 + +跑法:python3 test_logic_state_daily.py 或 pytest test_logic_state_daily.py +""" +import sys +import types + +_STUBS = ("pandas", "psycopg", "pymysql", "dotenv") +for _n in _STUBS: + if _n not in sys.modules: + try: + __import__(_n) + except ImportError: + _m = types.ModuleType(_n) + if _n == "pandas": # db.py 的函数签名在定义时引用这两个名字 + _m.DataFrame = type("DataFrame", (), {}) + _m.Series = type("Series", (), {}) + sys.modules[_n] = _m + +import config # noqa: E402 +import logic_state as ls # noqa: E402 +import logic_state_daily as lsd # noqa: E402 + + +def t(name, cond): + assert cond, name + print(" ok", name) + + +DS = "2026-09-04" +TABLE = "t_akg_logic_state_daily_test" + + +def raw(state, why=None): + return {"state": state, "why": why, "as_of": "2026-09-03", + "usable": [ls.PATH_CLAIM, ls.PATH_BROKER], "missing": [ls.PATH_JUDGE, ls.PATH_EVENT], + "reasons": [f"{ls.PATH_CLAIM}(2026-08-25):最近一条利好"] * 8, + "paths": [{"path": ls.PATH_CLAIM, "signal": ls.SIG_UP, "as_of": "2026-08-25", + "why": "最近一条利好"}]} + + +def h(day, raw_state, state): + return {"trade_date": day, "raw_state": raw_state, "state": state} + + +def _boom(*a, **k): + raise OSError("connection refused") + + +# ---------------------------------------------------------------- 落定口径 +def test_settle_one(): + print("落定的口径") + r = raw(ls.STATE_HOLD) + s = lsd.settle_one("SH600000", r, None) + t("第一次见到:落定态等于原始态", s["state"] == ls.STATE_HOLD and s["raw_state"] == ls.STATE_HOLD) + t("上一次落定态为空、说明为维持", s["prev_state"] is None and s["settle_note"] == "维持") + t("不改传入的原始结果", "raw_state" not in r and r["state"] == ls.STATE_HOLD) + t("其余键原样保留", s["paths"] == r["paths"] and s["usable"] == r["usable"] and s["why"] is None) + + s = lsd.settle_one("SH600000", raw(ls.STATE_DOUBT), [h("2026-09-03", ls.STATE_HOLD, ls.STATE_HOLD)]) + t("进入逻辑存疑即刻成立", s["state"] == ls.STATE_DOUBT and s["prev_state"] == ls.STATE_HOLD) + + hist = [h("2026-09-01", ls.STATE_DOUBT, ls.STATE_DOUBT), h("2026-09-02", ls.STATE_HOLD, ls.STATE_DOUBT)] + s = lsd.settle_one("SH600000", raw(ls.STATE_HOLD), hist) + t("退出存疑:原始态只有两天不存疑,落定仍是存疑", + s["state"] == ls.STATE_DOUBT and s["raw_state"] == ls.STATE_HOLD and "现在 2 天" in s["settle_note"]) + hist.append(h("2026-09-03", ls.STATE_HOLD, ls.STATE_DOUBT)) + s = lsd.settle_one("SH600000", raw(ls.STATE_HOLD), hist) + t("退出存疑:连续三个计划日不再存疑才放行", s["state"] == ls.STATE_HOLD) + + s = lsd.settle_one("SH600000", raw(ls.STATE_STRONG), [h("2026-09-03", ls.STATE_HOLD, ls.STATE_HOLD)]) + t("其余迁移:只有一天同向,先维持上一次落定态", + s["state"] == ls.STATE_HOLD and s["raw_state"] == ls.STATE_STRONG) + s = lsd.settle_one("SH600000", raw(ls.STATE_STRONG), + [h("2026-09-02", ls.STATE_HOLD, ls.STATE_HOLD), h("2026-09-03", ls.STATE_STRONG, ls.STATE_HOLD)]) + t("其余迁移:连续两个计划日同向才迁移", s["state"] == ls.STATE_STRONG) + s = lsd.settle_one("SH600000", raw(ls.STATE_STRONG), [h("2026-09-03", ls.STATE_HOLD, ls.STATE_HOLD)], + confirm_days=1) + t("确认天数旋钮生效", s["state"] == ls.STATE_STRONG) + t("配置里的两个天数与 settle 的默认值同源", + config.LOGIC_SETTLE_CONFIRM_DAYS == 2 and config.LOGIC_SETTLE_EXIT_DAYS == 3) + + +# ---------------------------------------------------------------- 读 +def test_history(): + print("历史行的读取") + config.LOGIC_STATE_TABLE = TABLE + t("表还不存在时返回空字典、不抛错", lsd.history(DS, read_mysql=_boom) == {}) + t("数据日不合法时也只是没有历史", lsd.history("不是日期", read_mysql=_boom) == {}) + t("指定了票但一只都没有时不查库", lsd.history(DS, codes=[], read_mysql=_boom) == {}) + + seen = {} + + def _reader(which, sql, params): + seen["sql"], seen["params"] = " ".join(sql.split()), params + return [{"trade_date": "2026-09-02", "code": "SH600000", "raw_state": ls.STATE_HOLD, "state": ls.STATE_HOLD}, + {"trade_date": "2026-09-03", "code": "SH600000", "raw_state": ls.STATE_DOUBT, "state": ls.STATE_DOUBT}, + {"trade_date": "2026-09-03", "code": "SZ000001", "raw_state": None, "state": float("nan")}] + + hist = lsd.history(DS, codes=["SH600000", "SZ000001"], read_mysql=_reader) + t("只取严格早于数据日、回看窗口之内的行", + "trade_date >= %s AND trade_date < %s" in seen["sql"] and seen["params"][1] == DS + and seen["params"][0] == "2026-08-05") + t("指定票走 IN 子句", "code IN (%s,%s)" in seen["sql"] and seen["params"][2:] == ("SH600000", "SZ000001")) + t("按票索引、行按数据日升序", [r["trade_date"] for r in hist["SH600000"]] == ["2026-09-02", "2026-09-03"]) + t("空值归一成 None(NaN 不当成状态)", hist["SZ000001"][0]["state"] is None) + lsd.history(DS, read_mysql=_reader) + t("不指定票时不带 IN 子句", "IN (" not in seen["sql"]) + + t("当日行读失败返回空字典", lsd.lookup(["SH600000"], DS, read_mysql=_boom) == {}) + t("当日行一只票都不要时不查库", lsd.lookup([], DS, read_mysql=_boom) == {}) + + +# ---------------------------------------------------------------- 写 +class _FakeCursor: + def __init__(self, store): + self.store = store + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + def execute(self, sql, params=None): + s = " ".join(sql.split()) + if s.upper().startswith("CREATE TABLE"): + assert TABLE in s and "PRIMARY KEY (trade_date, code)" in s, s + self.store["created"] += 1 + return + if s.upper().startswith("DELETE"): + assert "WHERE trade_date = %s" in s, s + self.store["deleted"].append(params[0]) + self.store["rows"] = [r for r in self.store["rows"] if r[0] != params[0]] + return + raise AssertionError(f"意外的语句: {s}") + + def executemany(self, sql, rows): + s = " ".join(sql.split()) + assert s.upper().startswith("INSERT INTO") and TABLE in s, s + assert s.count("%s") == len(lsd.COLUMNS), s + for r in rows: + assert len(r) == len(lsd.COLUMNS) + self.store["rows"].extend(list(rows)) + + +class _FakeConn: + def __init__(self, store): + self.store = store + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + def cursor(self): + return _FakeCursor(self.store) + + def commit(self): + self.store["commits"] += 1 + + +def _plan_rows(): + """计划快照里的行(plan._row 的形状):两只有卡的、一只没卡的、一只重复的。""" + a = lsd.settle_one("SH600000", raw(ls.STATE_HOLD), None) + b = lsd.settle_one("SZ000001", raw(ls.STATE_STRONG), + [h("2026-09-03", ls.STATE_HOLD, ls.STATE_HOLD)]) # 被按住:原始强化、落定成立 + return [{"code": "SH600000", "verdict": "候选", "card_rank": 3, "logic_state": a}, + {"code": "SZ000001", "verdict": "仅展示", "card_rank": 40, "logic_state": b}, + {"code": "SZ000002", "verdict": "仅展示", "card_rank": 41}, + {"code": "SH600000", "verdict": "候选", "card_rank": 3, "logic_state": a}] + + +def test_build_and_save(): + print("拼行与幂等落库") + config.LOGIC_STATE_TABLE = TABLE + rows = lsd.build_rows(DS, _plan_rows(), plan_version="abc1234", now="2026-09-07 07:01:00") + t("没有逻辑状态的行跳过、重复的票只留一条", [r["code"] for r in rows] == ["SH600000", "SZ000001"]) + b = rows[1] + t("原始态与落定态分开存", b["raw_state"] == ls.STATE_STRONG and b["state"] == ls.STATE_HOLD + and b["prev_state"] == ls.STATE_HOLD and "先维持" in b["settle_note"]) + t("判决、卡内序、代码版本、写入时刻都带上", + b["verdict"] == "仅展示" and b["card_rank"] == 40 and b["plan_version"] == "abc1234" + and b["snapshot_at"] == "2026-09-07 07:01:00") + t("四路名单存成顿号串", b["usable"] == f"{ls.PATH_CLAIM}、{ls.PATH_BROKER}") + t("说明最多存六条", b["reasons"].count("最近一条利好") == 6) + + store = {"rows": [], "deleted": [], "created": 0, "commits": 0} + + def _write(day, rs): + lsd.save(day, rs, conn_factory=lambda: _FakeConn(store)) + + r1 = lsd.persist(DS, _plan_rows(), "abc1234", write=_write) + t("首次写两行", r1["rows"] == 2 and len(store["rows"]) == 2) + t("读数:按落定态计数、被按住一只、迁移零只", + r1["by_state"] == {ls.STATE_HOLD: 2} and r1["held"] == 1 and r1["moved"] == 0) + lsd.persist(DS, _plan_rows(), "abc1234", write=_write) + t("同一数据日重跑仍是两行,不追加重复行", len(store["rows"]) == 2) + t("重跑先删该日行(删的正是这个数据日)", store["deleted"] == [DS, DS]) + t("建表语句每次都发(已存在就跳过)、删与插同一事务提交", store["created"] == 2 and store["commits"] == 4) + lsd.persist(DS, [], "abc1234", write=_write) + t("一张卡都没有时该日行被清空,不留残行", store["rows"] == []) + + +def test_roundtrip(): + print("表行与下游形状往返") + rows = lsd.build_rows(DS, _plan_rows(), plan_version="abc1234", now="x") + out = lsd.row_to_out(rows[1]) + t("状态三键、子因、截止日都在", + out["state"] == ls.STATE_HOLD and out["raw_state"] == ls.STATE_STRONG + and out["prev_state"] == ls.STATE_HOLD and out["why"] is None and out["as_of"] == "2026-09-03") + t("顿号串读回来是列表", out["usable"] == [ls.PATH_CLAIM, ls.PATH_BROKER] + and out["missing"] == [ls.PATH_JUDGE, ls.PATH_EVENT]) + t("JSON 串读回来是列表", isinstance(out["paths"], list) and out["paths"][0]["path"] == ls.PATH_CLAIM + and len(out["reasons"]) == 6) + t("空行也能还原、不抛错", lsd.row_to_out({})["usable"] == [] and lsd.row_to_out({})["paths"] == []) + + +def main(): + test_settle_one() + test_history() + test_build_and_save() + test_roundtrip() + print("ALL OK — 落定口径 / 历史读取不断产 / 幂等落库 / 往返形状 全部通过") + + +if __name__ == "__main__": + main() diff --git a/test_plan_logic_state.py b/test_plan_logic_state.py index dbf5aec..de6175d 100644 --- a/test_plan_logic_state.py +++ b/test_plan_logic_state.py @@ -68,6 +68,18 @@ def main(): t("不发内部中间量(硬触发标记、出处原值不外泄)", all("hard" not in p and "refs" not in p for p in out["paths"])) t("没有状态时发 None,不硬拼一个空壳", plan._state_out(None) is None) # noqa: SLF001 + t("没做过落定的结果:原始态就是 state 本身、落定说明为空", + out["raw_state"] == out["state"] and out["settle_note"] is None and out["prev_state"] is None) + + print("落定态随卡发出(2026-09-07 第三件桥侧前置)") + import logic_state_daily as lsd + settled = lsd.settle_one("SH600000", st, [{"trade_date": "2026-09-02", "raw_state": ls.STATE_DOUBT, + "state": ls.STATE_DOUBT}]) + out2 = plan._state_out(settled) # noqa: SLF001 + t("原始态成立、上一次落定存疑、只有一天不存疑 -> 发出的 state 仍是存疑", + out2["raw_state"] == st["state"] and out2["state"] == ls.STATE_DOUBT + and out2["prev_state"] == ls.STATE_DOUBT and "现在 1 天" in out2["settle_note"]) + t("落定不改每路的截止日与说明", out2["paths"] == out["paths"]) print("四态不改判决(收敛规则单调)") for state in (ls.STATE_STRONG, ls.STATE_HOLD, ls.STATE_UNKNOWN, ls.STATE_DOUBT):