第三件桥侧前置:逐票日频逻辑状态表(抗抖动)与按代码查询接口

新模块 logic_state_daily.py:每个数据日为档位表里的每只票记一行原始态与落定态,
落定走 logic_state.settle(进入存疑即刻成立、退出要连续三天、其余迁移要连续两天),
只在 plan.generate 里写、/plan 实时重算只读;同日重跑幂等。
plan.py:四路取数抽成 _logic_inputs(计划装配与接口现算共用一处);卡上的 state 改为
落定态,另发 raw_state、settle_note、prev_state;新增 logic_states_for 给接口现算。
api.py:GET /logic_state?codes=…,优先回当日表行(source=daily),不在表里的按此刻
现算并按历史落定(source=computed)。config 加表名与三个天数。
单测:新建 test_logic_state_daily.py(落定口径、读不到不断产、幂等落库、往返形状),
test_plan_logic_state.py 加落定态发出两例;十个测试文件在 155 容器隔离副本全绿。

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
zlt 2026-09-07 14:21:42 +08:00
parent b2ebf6d368
commit faede9ab68
6 changed files with 671 additions and 21 deletions

51
api.py
View File

@ -22,6 +22,7 @@ from fastapi.responses import PlainTextResponse
import config import config
import db import db
import logic_state_daily
import plan import plan
import plan_reconcile import plan_reconcile
import regime import regime
@ -135,3 +136,53 @@ def plan_verdict(codes: str | None = None, code: str | None = None,
verdicts.append(v) verdicts.append(v)
return {"date": d, "stale": L.get("stale", ""), return {"date": d, "stale": L.get("stale", ""),
"count": len(verdicts), "verdicts": verdicts} "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}

View File

@ -271,3 +271,16 @@ JUDGEMENT_SCOPES = {s.strip() for s in
# 比对上一版时往前看几个自然日。取到窗口内最近一个计划日的行作为上一版;超过这个窗口没写过 # 比对上一版时往前看几个自然日。取到窗口内最近一个计划日的行作为上一版;超过这个窗口没写过
# 快照的主题,本次按"第一次见到"处理,迁移记为空、陈旧天数从零重新起算。 # 快照的主题,本次按"第一次见到"处理,迁移记为空、陈旧天数从零重新起算。
JUDGEMENT_PREV_LOOKBACK_DAYS = int(os.environ.get("JUDGEMENT_PREV_LOOKBACK_DAYS", "60")) 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"))

282
logic_state_daily.py Normal file
View File

@ -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_<ds> t_factor_akg_*
trade_date 同口径不是行业观点快照的 plan_date那个是写入当天两张表的日期口径不同
读的时候别混计划在 D+1 凌晨构建数据日是 D这张表记的是 D
## 只在早上那一次生成里写
写入点只有 plan.generaterun.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
其余键whypathsusablemissingreasonsas_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}

95
plan.py
View File

@ -29,6 +29,7 @@ import common
import config import config
import judgement import judgement
import logic_state import logic_state
import logic_state_daily
import db import db
import sources import sources
import version import version
@ -144,6 +145,11 @@ def _state_out(st) -> dict | None:
if not isinstance(st, dict): if not isinstance(st, dict):
return None return None
return {"state": st.get("state"), "why": st.get("why"), "as_of": st.get("as_of"), 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 [], "usable": st.get("usable") or [], "missing": st.get("missing") or [],
"reasons": (st.get("reasons") or [])[:6], "reasons": (st.get("reasons") or [])[:6],
"paths": [{"path": p.get("path"), "signal": p.get("signal"), "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()]) 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, def _assemble_cards(ds: str, codes: list, ev: dict, upside: pd.Series,
mkt_days: set, risk: set | None = None) -> tuple[dict, list]: mkt_days: set, risk: set | None = None) -> tuple[dict, list]:
"""候选卡装配2026-09-02 方案第 2.2 节第三项):对档位表里的全部票(主榜与观察档, """候选卡装配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) moved = sources.moved_members(ds)
daily = sources.stock_daily(ds) daily = sources.stock_daily(ds)
night = sources.night_conclusions(codes, ds) night = sources.night_conclusions(codes, ds)
# 因果论断2026-09-03数据基座抽取的论断挂在卡上作证据线只展示不进判决视图未建时为空。 # 逻辑状态四态的取数(三路输入加逐票日频表的近日行;取数口径与理由见 _logic_inputs
# 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_state.apply_to_card是单调的——强化只能提前卡内序、永远不升判决。 # logic_state.apply_to_card是单调的——强化只能提前卡内序、永远不升判决。
# 行业观点快照在交易日 D 的晚上 20:40 写plan_date 记的是 D计划在 D+1 凌晨构建, inp = _logic_inputs(codes, ds)
# 数据日 ds 就是 D。所以要读的是 plan_date 不晚于 ds 的最近一版,不是 ds+1。 logic_full, seg_view, seg_hist = inp["logic_full"], inp["seg_view"], inp["seg_hist"]
# 09-04 之前这里读的是 ds+1能对上只因为那天上午手工跑过一次快照——把一次手工操作 broker, seg_of, hist = inp["broker"], inp["seg_of"], inp["hist"]
# 的时序写进了代码,到周一就全部落空(台账 033 记的"快照断了"其实是这个错, logic = {k: v[:config.LOGIC_CLAIMS_PER_STOCK] for k, v in logic_full.items()}
# 调度中心一直在按时跑。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)
if risk is None: # collect 会传入读过一次的名单;单独调用时自己读 if risk is None: # collect 会传入读过一次的名单;单独调用时自己读
try: try:
risk = factors._risk_set() or set() # noqa: SLF001 —— 同仓自用 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_state": n.get("accum_state"), "accum_score": n.get("accum_score"),
"accum_age": n.get("accum_age"), "y_signal": n.get("signal"), "accum_age": n.get("accum_age"), "y_signal": n.get("signal"),
"stale_snapshot": stale, "logic": logic.get(k) or []} "stale_snapshot": stale, "logic": logic.get(k) or []}
state = _logic_state_of(k, evd, seg_of, seg_view, broker, ds, raw_state = _logic_state_of(k, evd, seg_of, seg_view, broker, ds,
full_logic=logic_full.get(k) or [], seg_hist=seg_hist) 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, j = card.judge(evd, start_pct=config.CARD_START_PCT,
accum_max_age=config.CARD_ACCUM_MAX_AGE, accum_max_age=config.CARD_ACCUM_MAX_AGE,
neg_tol=config.UPSIDE_NEG_TOLERANCE, 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: with open(tmp, "w", encoding="utf-8") as f:
json.dump(snap, f, ensure_ascii=False, indent=1, default=str) json.dump(snap, f, ensure_ascii=False, indent=1, default=str)
os.replace(tmp, jpath) 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(text)
print(f"\n已写入 {out} 与快照 {jpath}" print(f"\n已写入 {out} 与快照 {jpath}"
f"(主榜 {len(snap['main'])} 行、观察档 {len(snap['observe'])} 行、" f"(主榜 {len(snap['main'])} 行、观察档 {len(snap['observe'])} 行、"

239
test_logic_state_daily.py Normal file
View File

@ -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("空值归一成 NoneNaN 不当成状态)", 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()

View File

@ -68,6 +68,18 @@ def main():
t("不发内部中间量(硬触发标记、出处原值不外泄)", t("不发内部中间量(硬触发标记、出处原值不外泄)",
all("hard" not in p and "refs" not in p for p in out["paths"])) 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("没有状态时发 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("四态不改判决(收敛规则单调)") print("四态不改判决(收敛规则单调)")
for state in (ls.STATE_STRONG, ls.STATE_HOLD, ls.STATE_UNKNOWN, ls.STATE_DOUBT): for state in (ls.STATE_STRONG, ls.STATE_HOLD, ls.STATE_UNKNOWN, ls.STATE_DOUBT):