diff --git a/README.md b/README.md index 6310fec..8b3d99c 100644 --- a/README.md +++ b/README.md @@ -57,7 +57,7 @@ scripts/ test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 19 例 test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例 test_batch5_units.py 决策系统信号流解析与消化口径 8 例 - test_batch6_units.py ws 通道: 协议测试向量/签名/水位/单表守卫 44 例 + test_batch6_units.py ws 通道: 测试向量/签名/公钥形态/水位/DDL体检 51 例 test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 32 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练) check_db.py 实机连通性与表结构自检 (需真实 .env) diff --git a/app/repo/qmt_repo.py b/app/repo/qmt_repo.py index d438996..df63b5d 100644 --- a/app/repo/qmt_repo.py +++ b/app/repo/qmt_repo.py @@ -149,36 +149,45 @@ def inbox_put(*, seq: int, msg_id: str, msg_type: str, payload: dict, msg_ts: in corr_id=None, dedup_key=None, processed: int = 2, note=None) -> str: """落一条上行消息。返回 NEW / DUP_SEQ / DUP_KEY。 - 双层去重 (§6.1 + §5.5) 在这里靠两个唯一键实现: 主键 seq 是第一层, 唯一索引 - dedup_key (trade_no) 是第二层。 + 双层去重 (§6.1 + §5.5): 主键 seq 是第一层, 唯一索引 dedup_key (trade_no) 是第二层。 - **DUP_KEY 这条分支容易漏, 单独说明**: trade_no 撞了但 seq 是新的 —— 说明 QMT 用 - 新序号重推了一条我们已入过账的成交。这一笔不能再入账, 但**这个 seq 仍必须占住一行**, - 否则连续水位永远卡在它前面, ack_seq 再也推不动, QMT 那边的消息也就永远清理不掉。 - 所以这里补插一行 dedup_key=NULL、processed=2 的存档行。 + **判重必须先 SELECT, 不能拿 rowcount 当判据。** 这条踩过: + `INSERT ... ON DUPLICATE KEY UPDATE seq = seq` 在"行已存在、值没变"时, MySQL 手册说 + affected-rows 是 0 —— 但那是**没开 CLIENT_FOUND_ROWS** 的前提。SQLAlchemy 的 MySQL + 方言默认就开着这个标志 (它要让 rowcount 反映"匹配到几行"而不是"改了几行"), 于是重复 + 插入照样回 1。实测: 表里始终只有 1 行, rowcount 却三次都是 1。 + 后果不是报错而是**静默双记**: 断线重连后 QMT 按 §6.1 重发的成交会被当成新成交, + 再入账一次 —— 持仓和摊薄成本直接算错。多一次 SELECT 换这个确定性, 非常值。 + + **DUP_KEY 这条分支容易漏**: trade_no 撞了但 seq 是新的 —— QMT 用新序号重推了一条我们 + 已入过账的成交。这一笔不能再入账, 但**这个 seq 仍必须占住一行**, 否则连续水位永远卡在 + 它前面, ack_seq 再也推不动, 对端的消息也永远清理不掉。故补插一行 dedup_key=NULL、 + processed=2 的存档行。 """ now = _NOW() - p = {"s": int(seq), "mi": str(msg_id)[:64], "mt": str(msg_type)[:24], + seq = int(seq) + if fetch_one("SELECT seq FROM pms_qmt_inbox WHERE seq = :s", {"s": seq}): + return PUT_DUP_SEQ + + verdict = PUT_NEW + if dedup_key and fetch_one("SELECT seq FROM pms_qmt_inbox WHERE dedup_key = :k", + {"k": str(dedup_key)[:80]}): + verdict = PUT_DUP_KEY + note = f"重复成交 (dedup_key={dedup_key}), 不入账, 仅占位以推进 seq 水位" + dedup_key, processed = None, 2 + + p = {"s": seq, "mi": str(msg_id)[:64], "mt": str(msg_type)[:24], "ci": (str(corr_id)[:64] if corr_id else None), "dk": (str(dedup_key)[:80] if dedup_key else None), "pj": _dumps(payload or {}), "mts": int(msg_ts or 0), "ts": now, "pc": int(processed), "nt": (str(note)[:300] if note else None), "pa": now if int(processed) != 0 else None} - sql = ("INSERT INTO pms_qmt_inbox (seq, msg_id, msg_type, corr_id, dedup_key, " - "payload_json, msg_ts, received_at, processed, processed_at, process_note) " - "VALUES (:s, :mi, :mt, :ci, :dk, :pj, :mts, :ts, :pc, :pa, :nt) " - "ON DUPLICATE KEY UPDATE seq = seq") - if execute(sql, p): - return PUT_NEW - if fetch_one("SELECT seq FROM pms_qmt_inbox WHERE seq = :s", {"s": int(seq)}): - return PUT_DUP_SEQ - # 见 docstring: trade_no 重复但 seq 是新的 —— 占位存档, 让水位能继续往前推 - p["dk"] = None - p["pc"] = 2 - p["pa"] = now - p["nt"] = f"重复成交 (dedup_key={dedup_key}), 不入账, 仅占位以推进 seq 水位" - execute(sql, p) - return PUT_DUP_KEY + # ON DUPLICATE KEY UPDATE 只作兜底 (万一并发/时序意外), 判重结论以上面的 SELECT 为准 + execute("INSERT INTO pms_qmt_inbox (seq, msg_id, msg_type, corr_id, dedup_key, " + "payload_json, msg_ts, received_at, processed, processed_at, process_note) " + "VALUES (:s, :mi, :mt, :ci, :dk, :pj, :mts, :ts, :pc, :pa, :nt) " + "ON DUPLICATE KEY UPDATE seq = seq", p) + return verdict def inbox_recover_watermark(stored_last_seq: int, limit: int = 20000) -> int: diff --git a/ddl_pms_v1.sql b/ddl_pms_v1.sql index e793ab6..ba78aa0 100644 --- a/ddl_pms_v1.sql +++ b/ddl_pms_v1.sql @@ -216,8 +216,7 @@ CREATE TABLE IF NOT EXISTS pms_qmt_order ( intent VARCHAR(8) NOT NULL DEFAULT 'OPEN' COMMENT 'OPEN/FILL/ADD/DCA/TRIM/EXIT/T0', note VARCHAR(200) NULL, status VARCHAR(16) NOT NULL DEFAULT 'QUEUED' - COMMENT '本地: QUEUED待发/SENDING已认领/SENT已发出/SEND_FAILED/ABORTED未发即作废; ' - '协议: ACCEPTED/SUBMITTED/PARTIAL/FILLED/CANCELLED/EXPIRED/REJECTED', + COMMENT '本地: QUEUED待发/SENDING已认领/SENT已发出/SEND_FAILED/ABORTED未发即作废; 协议: ACCEPTED/SUBMITTED/PARTIAL/FILLED/CANCELLED/EXPIRED/REJECTED', broker_order_id VARCHAR(64) NULL COMMENT 'QMT 回的委托号', cum_qty INT NOT NULL DEFAULT 0 COMMENT '本委托累计成交股数 (§5.4 口径)', cum_avg_price DECIMAL(10,3) NULL COMMENT '本委托累计成交均价', diff --git a/scripts/init_db.py b/scripts/init_db.py index d7fd4e2..27db4d0 100644 --- a/scripts/init_db.py +++ b/scripts/init_db.py @@ -59,6 +59,41 @@ def split_sql(text: str) -> list: return out +def find_adjacent_literals(stmt: str) -> list: + """找相邻字符串字面量 `'a' 'b'` —— MySQL 的 COMMENT 子句不接受这种写法。 + + Python 里 `'a' 'b'` 会自动拼成 `'ab'`, 手写长注释时很容易顺手断行写成两段。SQL 标准 + 确实允许相邻字面量拼接, MySQL 在**表达式**里也认, 但 `COMMENT` 是表/列选项, 语法上只 + 接受**一个** string_literal —— 于是报一个 1064, 而且错误信息把整条 CREATE TABLE 原样 + 吐出来, 光看那坨东西根本定位不到是哪一行。2026-07-28 pms_qmt_order 就栽在这上面。 + + 放在演练阶段拦, 是因为这类错误不连库也能发现, 没必要等到执行时才炸。 + """ + hits, i, n = [], 0, len(stmt) + while i < n: + if stmt[i] != "'": + i += 1 + continue + j = i + 1 # 进入字符串, 找闭合引号 + while j < n: + if stmt[j] == "\\": + j += 2 + continue + if stmt[j] == "'": + if j + 1 < n and stmt[j + 1] == "'": # '' 转义 + j += 2 + continue + break + j += 1 + k = j + 1 # 闭合引号之后, 跳过空白 + while k < n and stmt[k] in " \t\r\n": + k += 1 + if k < n and stmt[k] == "'" and k > j + 1: # 中间隔了空白又来一个引号 + hits.append(stmt[max(0, i):min(n, k + 40)].strip()) + i = j + 1 + return hits + + def parse_statements(sql_text: str) -> list: """去掉整行 `--` 注释后切语句。返回 [(种类, 表名, 语句)]。 @@ -113,6 +148,14 @@ def main(): f"说明 DDL 切分有误, 已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}") sys.exit(1) + lint = [(t, h) for _, t, s in stmts for h in find_adjacent_literals(s)] + if lint: + print(f"FAIL {len(lint)} 处相邻字符串字面量 —— MySQL 的 COMMENT 只接受单个字面量, " + f"会报 1064。已中止 (不执行任何语句)。把它们合成一个字符串即可:") + for tbl, snippet in lint: + print(f"\n {tbl}:\n {snippet}") + sys.exit(1) + tables = [t for k, t, _ in stmts if k == "table"] seeds = [t for k, t, _ in stmts if k == "seed"] print(f"DDL 文件: {DDL_FILE}") diff --git a/scripts/run_tests.py b/scripts/run_tests.py index bcb5673..835173f 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -10,7 +10,7 @@ test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 (19 例) test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例) test_batch5_units.py 决策系统信号流解析与消化口径 (8 例) - test_batch6_units.py ws 通道: 协议测试向量/签名/公钥形态/水位/单表守卫 (49 例) + test_batch6_units.py ws 通道: 协议测试向量/签名/公钥形态/水位/DDL 体检 (51 例) test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (32 例) 任一子集失败即整体失败 (退出码 1)。 """ diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py index b34deb0..a6b8dc1 100644 --- a/scripts/test_batch6_units.py +++ b/scripts/test_batch6_units.py @@ -398,6 +398,35 @@ def run(): except MultiTableSQL: pass + # --- DDL 体检。2026-07-28 建表时 pms_qmt_order 报 1064: COMMENT 写成了 Python 风格的 + # 相邻字符串字面量 `'前半' '后半'`。MySQL 在表达式里认这种拼接, 但 COMMENT 是列选项, + # 语法上只接受一个 string_literal —— 而它报错时会把整条 CREATE TABLE 原样吐出来, + # 光看那坨东西定位不到是哪一行。不连库也能拦, 就该在这里拦。 + @case("DDL 体检: 揪出相邻字符串字面量") + def _(): + sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + from init_db import find_adjacent_literals + bad_ = ("status VARCHAR(16) NOT NULL DEFAULT 'QUEUED'\n" + " COMMENT '本地: QUEUED待发; '\n '协议: ACCEPTED/FILLED',") + assert find_adjacent_literals(bad_), "这正是 1064 的成因, 必须揪出来" + for good in ("COMMENT '单个字符串, 里面有 '' 转义也不算'", + "COMMENT 'a', other VARCHAR(8) COMMENT 'b'", # 逗号隔开, 不是相邻 + "DEFAULT 'NONE' COMMENT 'NONE/REQUESTED/SENT'"): + eq(find_adjacent_literals(good), [], f"误报: {good[:40]}") + + @case("DDL 文件本身体检通过 (13 张表 + 1 条初始行)") + def _(): + sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + from init_db import DDL_FILE, find_adjacent_literals, parse_statements + with open(DDL_FILE, encoding="utf-8") as f: + stmts = parse_statements(f.read()) + broken = [s for k, _, s in stmts if k == "?"] + eq(broken, [], "有认不出的残句, 说明 DDL 切分被破坏") + for _, tbl, s in stmts: + eq(find_adjacent_literals(s), [], f"{tbl} 有相邻字面量") + eq(len([1 for k, _, _ in stmts if k == "table"]), 13) + eq(len([1 for k, _, _ in stmts if k == "seed"]), 1) + @case("通道三表的 SQL 全部单表合规") def _(): from app.db.session import assert_single_table