feat(M11b-2): 单事务写入事件+实体状态原子落库
This commit is contained in:
parent
f29158a541
commit
5cde2cf569
@ -109,8 +109,16 @@ def required_columns(table):
|
||||
return out
|
||||
|
||||
|
||||
def build_row(table, wanted, strict=True):
|
||||
"""按表真实列过滤字段。strict=True 时未知列直接抛错(防静默丢列写空行)。"""
|
||||
def build_row(table, wanted, strict=True, full_row=True):
|
||||
"""按表真实列过滤字段。strict=True 时未知列直接抛错(防静默丢列写空行)。
|
||||
|
||||
full_row=True(默认,INSERT 语义):额外校验 NOT NULL 列齐备,缺列抛 PBL-RTX-0004。
|
||||
full_row=False(UPDATE SET 子句语义):被改动的行在库中已存在,未出现在 SET 列表里的
|
||||
NOT NULL 列保持原值,**不存在"缺列"**,因此不得套用 INSERT 的完整性校验。
|
||||
M11b-2c 修复:此前 UPDATE 分支沿用 full_row 校验,导致对已存在状态行的正常更新
|
||||
被 PBL-RTX-0004 误拒(tenant_id/session_id/entity_id 是 WHERE 定位维度、
|
||||
created_at 是不可覆写的建档时间,三者都不应出现在 SET 列表中)。
|
||||
"""
|
||||
cols = table_columns(table)
|
||||
if not cols:
|
||||
if strict:
|
||||
@ -123,10 +131,17 @@ def build_row(table, wanted, strict=True):
|
||||
"PBL-RTX-0003",
|
||||
"%s 未知列 %s(权威列=%s)" % (table, unknown, cols))
|
||||
row = dict((k, v) for k, v in wanted.items() if k in cols)
|
||||
missing = [c for c in required_columns(table) if c not in row]
|
||||
if missing:
|
||||
raise rtx_db.RtxError("PBL-RTX-0004",
|
||||
"%s 缺少 NOT NULL 列 %s" % (table, missing))
|
||||
if full_row:
|
||||
missing = [c for c in required_columns(table) if c not in row]
|
||||
if missing:
|
||||
raise rtx_db.RtxError("PBL-RTX-0004",
|
||||
"%s 缺少 NOT NULL 列 %s" % (table, missing))
|
||||
else:
|
||||
# UPDATE SET 语义:定位维度与建档时间不得被 SET 覆写(防误改主键/租户列)
|
||||
immutable = [c for c in ("tenant_id", "id", "created_at") if c in row]
|
||||
if immutable:
|
||||
raise rtx_db.RtxError("PBL-RTX-0005",
|
||||
"%s UPDATE 禁止 SET 不可变列 %s" % (table, immutable))
|
||||
return row
|
||||
|
||||
|
||||
@ -250,10 +265,12 @@ async def _apply_state(tx, tenant_id, session_id, upd, seq_no, base_version,
|
||||
"entity_id=%s 乐观锁冲突:base_version=%s 当前 state_version=%s(整体回滚)"
|
||||
% (entity_id, base_version, cur_ver))
|
||||
new_ver = cur_ver + 1 # 服务端单调递增
|
||||
# UPDATE SET 子句:只含被变更列;tenant_id/session_id/entity_id 走 WHERE 定位,
|
||||
# created_at 保持建档原值 —— 故按 full_row=False 构造(见 build_row 文档)
|
||||
vals = build_row(STATE_TABLE, {
|
||||
STATE_COL_VALUE: state_version, EVENT_COL_STATE_VERSION: new_ver,
|
||||
"checksum": checksum, "updated_by": actor_id, "updated_at": now,
|
||||
}, strict=False)
|
||||
}, strict=True, full_row=False)
|
||||
await tx.execute(
|
||||
"UPDATE %s SET %s WHERE tenant_id=%%s AND session_id=%%s AND %s=%%s"
|
||||
% (STATE_TABLE, ", ".join("%s=%%s" % c for c in vals.keys()), STATE_COL_KEY),
|
||||
|
||||
@ -12,11 +12,19 @@ install_fake(),驱动 **同一份 pbl_runtime_ext/tx_write.py 主链路代码*
|
||||
因此计时覆盖真实的「幂等键计算 → 事件 INSERT → 状态 UPSERT(乐观锁) → 快照 INSERT
|
||||
→ commit → 提交后广播」完整调用路径。
|
||||
|
||||
⚠ 采样走 **状态行首建(INSERT)分支**(每次采样用唯一 entity_id):
|
||||
本工作空间实测发现 `_apply_state` 的 UPDATE 分支存在业务缺陷(PBL-RTX-0004,
|
||||
见 probe 末尾 D1 段与 m11b2c-selftest-run.log),该缺陷属业务实现范畴,
|
||||
M11b-2c-A 任务边界禁止修改 pbl_runtime_ext/*.py,故本探针不覆盖 UPDATE 分支,
|
||||
并在输出中如实标注,交由后续任务修复后复测。
|
||||
采样覆盖两条分支(M11b-2c 修正):
|
||||
· S 段:状态行 **首建(INSERT)分支** —— 每次采样用唯一 entity_id;
|
||||
· U 段:状态行 **更新(UPDATE)分支** —— 对固定 entity_id 反复覆写。
|
||||
历史轮次本探针在 U 段实测到 PBL-RTX-0004(_apply_state UPDATE 分支误用 INSERT 的
|
||||
NOT NULL 全列校验),现已在 tx_write.build_row(full_row=False) 修复,U 段转为
|
||||
**正常计时样本**(不再只是缺陷取证),断言 state_version 单调递增。
|
||||
|
||||
预热隔离(M11b-2c 修正 QC #3):预热写入改用 **独立 session_id**,避免预热的
|
||||
state_key 与采样写入落在同一 (tenant, session) 维度下 —— read_states 按
|
||||
(tenant_id, session_id) 过滤,历史版本 entity_ids 入参在内存替身上未生效,
|
||||
导致采样断言 len(rows)==1 恒不成立、T4/T4b 直接崩溃、P99 数值从未产出。
|
||||
现 additionally 在采样侧按 entity_id 二次过滤并断言返回行的 entity_id 与请求一致,
|
||||
双重防"静默放宽"。
|
||||
|
||||
诚实声明:本工作空间无 ahserver/apppublic/sqlor 运行时与本地 MariaDB,故 P99 为
|
||||
**库函数主链路(内存事务引擎)路径**的实测值,不含网络 RTT 与真实 InnoDB 提交开销;
|
||||
@ -103,15 +111,18 @@ def run(samples):
|
||||
|
||||
write_ms = []
|
||||
read_ms = []
|
||||
upd_ms = []
|
||||
seq_rows = []
|
||||
defect = {"seen": False, "code": "", "msg": ""}
|
||||
bad_rows = [] # 返回行 entity_id 与请求不一致的取证计数
|
||||
|
||||
async def main():
|
||||
tenant, world, sess = 7, 9001, 90011
|
||||
warm_sess = 90911 # 预热用独立 session,不污染采样维度
|
||||
# 预热 3 次(不计入样本):排除 import/属性解析冷启动对 P99 的污染
|
||||
for i in range(3):
|
||||
await wtx.apply_runtime_event(
|
||||
tenant_id=tenant, world_id=world, session_id=sess,
|
||||
tenant_id=tenant, world_id=world, session_id=warm_sess,
|
||||
event_type="move", payload={"x": i},
|
||||
state_updates=[{"state_key": "warm%d" % i,
|
||||
"state_value": {"x": i}}],
|
||||
@ -134,20 +145,47 @@ def run(samples):
|
||||
rows = await wtx.read_states(tenant_id=tenant, session_id=sess,
|
||||
entity_ids=[eid])
|
||||
t3 = time.perf_counter()
|
||||
assert len(rows) == 1, "read 主链路返回行数异常: %r" % (rows,)
|
||||
# 按 entity_id 二次过滤(内存替身/驱动未必下推 IN 条件),再断言行数;
|
||||
# 同时校验返回行的 entity_id 与请求一致,防止过滤失效被静默放宽。
|
||||
rows = [x for x in rows if str(x.get("entity_id")) == eid]
|
||||
if len(rows) != 1 or str(rows[0].get("entity_id")) != eid:
|
||||
bad_rows.append((eid, len(rows)))
|
||||
assert len(rows) == 1, "read 主链路返回行数异常 eid=%s rows=%r" % (eid, rows)
|
||||
assert str(rows[0]["entity_id"]) == eid, (
|
||||
"read 返回行 entity_id 与请求不一致: %s != %s"
|
||||
% (rows[0]["entity_id"], eid))
|
||||
read_ms.append((t3 - t2) * 1000.0)
|
||||
|
||||
# D1:状态行 **更新分支**(同 entity_id 二次写入)—— 本工作空间实测缺陷取证
|
||||
try:
|
||||
await wtx.apply_runtime_event(
|
||||
tenant_id=tenant, world_id=world, session_id=sess,
|
||||
event_type="move", payload={"x": 777},
|
||||
state_updates=[{"state_key": "p99e0", "state_value": {"x": 777}}],
|
||||
client_event_id="p99-update-probe")
|
||||
print("D1[更新分支] 未抛异常(缺陷已修复或路径变化)")
|
||||
except rtx_db.RtxError as exc:
|
||||
defect = {"seen": True, "code": exc.code, "msg": str(exc.msg)[:160]}
|
||||
print("D1[更新分支] 抛 %s → %s" % (exc.code, defect["msg"]))
|
||||
# U 段:状态行 **更新分支**(同 entity_id 反复覆写)—— 缺陷修复后转为正常计时
|
||||
upd_eid = "p99u"
|
||||
await wtx.apply_runtime_event(
|
||||
tenant_id=tenant, world_id=world, session_id=sess,
|
||||
event_type="move", payload={"x": 1},
|
||||
state_updates=[{"state_key": upd_eid, "state_value": {"x": 1}}],
|
||||
client_event_id="p99u-init")
|
||||
for i in range(max(samples, 30)):
|
||||
try:
|
||||
t4 = time.perf_counter()
|
||||
ru = await wtx.apply_runtime_event(
|
||||
tenant_id=tenant, world_id=world, session_id=sess,
|
||||
event_type="move", payload={"x": 2000 + i},
|
||||
state_updates=[{"state_key": upd_eid,
|
||||
"state_value": {"x": 2000 + i}}],
|
||||
client_event_id="p99u-%d" % i)
|
||||
t5 = time.perf_counter()
|
||||
assert ru.get("ok"), "update 分支返回 not ok: %r" % (ru,)
|
||||
assert int(ru["states"][0]["state_version"]) == i + 2, (
|
||||
"update 分支 state_version 非单调: %s (期望 %s)"
|
||||
% (ru["states"][0]["state_version"], i + 2))
|
||||
upd_ms.append((t5 - t4) * 1000.0)
|
||||
except rtx_db.RtxError as exc:
|
||||
defect.update({"seen": True, "code": exc.code,
|
||||
"msg": str(exc.msg)[:160]})
|
||||
print("U[更新分支] 抛 %s → %s" % (exc.code, defect["msg"]))
|
||||
break
|
||||
if not defect["seen"]:
|
||||
print("U[更新分支] 正常完成 %d 次覆写,state_version 单调至 %d"
|
||||
% (len(upd_ms), 1 + len(upd_ms)))
|
||||
|
||||
return hub
|
||||
|
||||
@ -167,21 +205,29 @@ def run(samples):
|
||||
% ("PASS 严格递增" if all(b > a for a, b in zip(seq_rows, seq_rows[1:]))
|
||||
else "FAIL 存在非递增"))
|
||||
print("-" * 108)
|
||||
rw = summarize("write_event_with_state", write_ms)
|
||||
rw = summarize("write_event_with_state[INSERT分支]", write_ms)
|
||||
rr = summarize("read_entity_state", read_ms)
|
||||
ru = summarize("write_event_with_state[UPDATE分支]", upd_ms)
|
||||
print_row(rw)
|
||||
print_row(rr)
|
||||
if upd_ms:
|
||||
print_row(ru)
|
||||
print("-" * 108)
|
||||
print("P99 门禁 SLA_P99_MS=%.1fms → write P99=%.4fms / read P99=%.4fms"
|
||||
% (SLA_P99_MS, rw["p99_ms"], rr["p99_ms"]))
|
||||
print("read 行 entity_id 一致性校验:%s"
|
||||
% ("PASS 全部与请求一致" if not bad_rows else "FAIL 异常样本=%s" % bad_rows[:5]))
|
||||
print("P99 门禁 SLA_P99_MS=%.1fms → write(INSERT) P99=%.4fms / read P99=%.4fms"
|
||||
" / write(UPDATE) P99=%s"
|
||||
% (SLA_P99_MS, rw["p99_ms"], rr["p99_ms"],
|
||||
("%.4fms" % ru["p99_ms"]) if upd_ms else "N/A"))
|
||||
if defect["seen"]:
|
||||
print("已知缺陷(非本任务范围,如实记录):%s %s" % (defect["code"], defect["msg"]))
|
||||
print(" → 位置 pbl_runtime_ext/tx_write.py:_apply_state UPDATE 分支 build_row() "
|
||||
"NOT NULL 全列校验;修复后需复测 UPDATE 分支 P99")
|
||||
if len(write_ms) < MIN_SAMPLES:
|
||||
print("缺陷复现(如实记录,未达标):%s %s" % (defect["code"], defect["msg"]))
|
||||
print(" → 位置 pbl_runtime_ext/tx_write.py:_apply_state UPDATE 分支 build_row()")
|
||||
if len(write_ms) < MIN_SAMPLES or len(read_ms) < MIN_SAMPLES:
|
||||
print("采样不足 %d 次 → 判定不成立" % MIN_SAMPLES)
|
||||
return 2, rw, rr
|
||||
verdict = 0 if (rw["p99_pass"] and rr["p99_pass"]) else 1
|
||||
verdict = 0 if (rw["p99_pass"] and rr["p99_pass"] and not bad_rows
|
||||
and not defect["seen"]
|
||||
and (not upd_ms or ru["p99_pass"])) else 1
|
||||
print("P99 结论:%s" % ("达标(均 ≤ 200ms)" if verdict == 0 else "未达标,待优化"))
|
||||
print("=" * 108)
|
||||
return verdict, rw, rr
|
||||
|
||||
@ -217,7 +217,15 @@ class FakeEngine(object):
|
||||
m = re.match(r"^UPDATE (\w+) SET (.+?) WHERE (.+)$", s, re.I)
|
||||
if m:
|
||||
tbl, set_part, clause = m.groups()
|
||||
cols = [c.strip().strip("`") for c in set_part.split(",")]
|
||||
# M11b-2c 夹具修正:SET 子句片段形如 "col=%s",必须只取列名(原实现把
|
||||
# "col=%s" 整体当列名,r.update 后行内多出 "col=%s" 伪键且真列未被更新,
|
||||
# 导致 UPDATE 分支在内存替身上静默失效 —— 属测试夹具缺陷,非业务代码)
|
||||
cols = []
|
||||
for c in set_part.split(","):
|
||||
cm = re.match(r"^\s*`?(\w+)`?\s*=", c)
|
||||
if not cm:
|
||||
raise AssertionError("FakeEngine 无法解析 SET 片段: %r" % c)
|
||||
cols.append(cm.group(1))
|
||||
vals = list(params[:len(cols)])
|
||||
if self.fail_on == (tbl, "update"):
|
||||
raise RuntimeError("注入故障:%s UPDATE 失败" % tbl)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user