修复初始化建表语句
This commit is contained in:
parent
de6e85cb0b
commit
ceaf4eef12
|
|
@ -57,7 +57,7 @@ scripts/
|
||||||
test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 19 例
|
test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 19 例
|
||||||
test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例
|
test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例
|
||||||
test_batch5_units.py 决策系统信号流解析与消化口径 8 例
|
test_batch5_units.py 决策系统信号流解析与消化口径 8 例
|
||||||
test_batch6_units.py ws 通道: 协议测试向量/签名/水位/单表守卫 44 例
|
test_batch6_units.py ws 通道: 测试向量/签名/公钥形态/水位/DDL体检 51 例
|
||||||
test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 32 例
|
test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 32 例
|
||||||
init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练)
|
init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练)
|
||||||
check_db.py 实机连通性与表结构自检 (需真实 .env)
|
check_db.py 实机连通性与表结构自检 (需真实 .env)
|
||||||
|
|
|
||||||
|
|
@ -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:
|
corr_id=None, dedup_key=None, processed: int = 2, note=None) -> str:
|
||||||
"""落一条上行消息。返回 NEW / DUP_SEQ / DUP_KEY。
|
"""落一条上行消息。返回 NEW / DUP_SEQ / DUP_KEY。
|
||||||
|
|
||||||
双层去重 (§6.1 + §5.5) 在这里靠两个唯一键实现: 主键 seq 是第一层, 唯一索引
|
双层去重 (§6.1 + §5.5): 主键 seq 是第一层, 唯一索引 dedup_key (trade_no) 是第二层。
|
||||||
dedup_key (trade_no) 是第二层。
|
|
||||||
|
|
||||||
**DUP_KEY 这条分支容易漏, 单独说明**: trade_no 撞了但 seq 是新的 —— 说明 QMT 用
|
**判重必须先 SELECT, 不能拿 rowcount 当判据。** 这条踩过:
|
||||||
新序号重推了一条我们已入过账的成交。这一笔不能再入账, 但**这个 seq 仍必须占住一行**,
|
`INSERT ... ON DUPLICATE KEY UPDATE seq = seq` 在"行已存在、值没变"时, MySQL 手册说
|
||||||
否则连续水位永远卡在它前面, ack_seq 再也推不动, QMT 那边的消息也就永远清理不掉。
|
affected-rows 是 0 —— 但那是**没开 CLIENT_FOUND_ROWS** 的前提。SQLAlchemy 的 MySQL
|
||||||
所以这里补插一行 dedup_key=NULL、processed=2 的存档行。
|
方言默认就开着这个标志 (它要让 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()
|
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),
|
"ci": (str(corr_id)[:64] if corr_id else None),
|
||||||
"dk": (str(dedup_key)[:80] if dedup_key else None),
|
"dk": (str(dedup_key)[:80] if dedup_key else None),
|
||||||
"pj": _dumps(payload or {}), "mts": int(msg_ts or 0), "ts": now,
|
"pj": _dumps(payload or {}), "mts": int(msg_ts or 0), "ts": now,
|
||||||
"pc": int(processed), "nt": (str(note)[:300] if note else None),
|
"pc": int(processed), "nt": (str(note)[:300] if note else None),
|
||||||
"pa": now if int(processed) != 0 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, "
|
# 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) "
|
"payload_json, msg_ts, received_at, processed, processed_at, process_note) "
|
||||||
"VALUES (:s, :mi, :mt, :ci, :dk, :pj, :mts, :ts, :pc, :pa, :nt) "
|
"VALUES (:s, :mi, :mt, :ci, :dk, :pj, :mts, :ts, :pc, :pa, :nt) "
|
||||||
"ON DUPLICATE KEY UPDATE seq = seq")
|
"ON DUPLICATE KEY UPDATE seq = seq", p)
|
||||||
if execute(sql, p):
|
return verdict
|
||||||
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
|
|
||||||
|
|
||||||
|
|
||||||
def inbox_recover_watermark(stored_last_seq: int, limit: int = 20000) -> int:
|
def inbox_recover_watermark(stored_last_seq: int, limit: int = 20000) -> int:
|
||||||
|
|
|
||||||
|
|
@ -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',
|
intent VARCHAR(8) NOT NULL DEFAULT 'OPEN' COMMENT 'OPEN/FILL/ADD/DCA/TRIM/EXIT/T0',
|
||||||
note VARCHAR(200) NULL,
|
note VARCHAR(200) NULL,
|
||||||
status VARCHAR(16) NOT NULL DEFAULT 'QUEUED'
|
status VARCHAR(16) NOT NULL DEFAULT 'QUEUED'
|
||||||
COMMENT '本地: QUEUED待发/SENDING已认领/SENT已发出/SEND_FAILED/ABORTED未发即作废; '
|
COMMENT '本地: QUEUED待发/SENDING已认领/SENT已发出/SEND_FAILED/ABORTED未发即作废; 协议: ACCEPTED/SUBMITTED/PARTIAL/FILLED/CANCELLED/EXPIRED/REJECTED',
|
||||||
'协议: ACCEPTED/SUBMITTED/PARTIAL/FILLED/CANCELLED/EXPIRED/REJECTED',
|
|
||||||
broker_order_id VARCHAR(64) NULL COMMENT 'QMT 回的委托号',
|
broker_order_id VARCHAR(64) NULL COMMENT 'QMT 回的委托号',
|
||||||
cum_qty INT NOT NULL DEFAULT 0 COMMENT '本委托累计成交股数 (§5.4 口径)',
|
cum_qty INT NOT NULL DEFAULT 0 COMMENT '本委托累计成交股数 (§5.4 口径)',
|
||||||
cum_avg_price DECIMAL(10,3) NULL COMMENT '本委托累计成交均价',
|
cum_avg_price DECIMAL(10,3) NULL COMMENT '本委托累计成交均价',
|
||||||
|
|
|
||||||
|
|
@ -59,6 +59,41 @@ def split_sql(text: str) -> list:
|
||||||
return out
|
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:
|
def parse_statements(sql_text: str) -> list:
|
||||||
"""去掉整行 `--` 注释后切语句。返回 [(种类, 表名, 语句)]。
|
"""去掉整行 `--` 注释后切语句。返回 [(种类, 表名, 语句)]。
|
||||||
|
|
||||||
|
|
@ -113,6 +148,14 @@ def main():
|
||||||
f"说明 DDL 切分有误, 已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}")
|
f"说明 DDL 切分有误, 已中止 (不执行任何语句)。首条残句:\n{broken[0][:200]}")
|
||||||
sys.exit(1)
|
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"]
|
tables = [t for k, t, _ in stmts if k == "table"]
|
||||||
seeds = [t for k, t, _ in stmts if k == "seed"]
|
seeds = [t for k, t, _ in stmts if k == "seed"]
|
||||||
print(f"DDL 文件: {DDL_FILE}")
|
print(f"DDL 文件: {DDL_FILE}")
|
||||||
|
|
|
||||||
|
|
@ -10,7 +10,7 @@
|
||||||
test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 (19 例)
|
test_batch3_units.py 规则闸 / 择时执行器实现B 纯逻辑 (19 例)
|
||||||
test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例)
|
test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例)
|
||||||
test_batch5_units.py 决策系统信号流解析与消化口径 (8 例)
|
test_batch5_units.py 决策系统信号流解析与消化口径 (8 例)
|
||||||
test_batch6_units.py ws 通道: 协议测试向量/签名/公钥形态/水位/单表守卫 (49 例)
|
test_batch6_units.py ws 通道: 协议测试向量/签名/公钥形态/水位/DDL 体检 (51 例)
|
||||||
test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (32 例)
|
test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (32 例)
|
||||||
任一子集失败即整体失败 (退出码 1)。
|
任一子集失败即整体失败 (退出码 1)。
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
|
|
@ -398,6 +398,35 @@ def run():
|
||||||
except MultiTableSQL:
|
except MultiTableSQL:
|
||||||
pass
|
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 全部单表合规")
|
@case("通道三表的 SQL 全部单表合规")
|
||||||
def _():
|
def _():
|
||||||
from app.db.session import assert_single_table
|
from app.db.session import assert_single_table
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue