149 lines
6.9 KiB
Python
Executable File
149 lines
6.9 KiB
Python
Executable File
#!/usr/bin/env python3
|
||
# -*- coding: utf-8 -*-
|
||
"""S3 必做项4:collector 重放幂等取证(真库,复用 harness 的沙箱适配层,不改被测源码)。
|
||
|
||
对同一批 pbl_runtime_event 事件重放 collector 写入 N 次,验证:
|
||
- 主键 id 由应用层 pbl_evidence.pk.gen_pk 生成:VARCHAR(32)、长度 <=32、互不相同、
|
||
且**不依赖任何 DB trigger**(配合 s3_evidence_chain.sh 的 DROP TRIGGER 复验)
|
||
- upsert 幂等:count(*) 不增长、无 1062 Duplicate entry、skipped/updated 分支正常
|
||
|
||
路径全部由脚本自身位置推导(无绝对路径硬编码);期望行数默认由 harness 的种子常量
|
||
``EVENTS_SEED`` 派生(每个事件采出一条证据),也可用 --expect-rows 覆盖。
|
||
判定用 assert 语义落地到退出码:任一条不满足 -> 打印 REPLAY IDEMPOTENT FAIL 并 exit 1,
|
||
使上层 chain 脚本(set -euo pipefail)能真实拦截,不再出现「FAIL 也 RC=0」。
|
||
|
||
用法::
|
||
|
||
python3 tests/s3_replay_idempotency.py
|
||
python3 tests/s3_replay_idempotency.py --replays 3 --expect-rows 2
|
||
python3 tests/s3_replay_idempotency.py --help
|
||
"""
|
||
import argparse
|
||
import asyncio
|
||
import importlib.util
|
||
import pathlib
|
||
import sys
|
||
|
||
TESTS_DIR = pathlib.Path(__file__).resolve().parent
|
||
REPO_ROOT = TESTS_DIR.parent # modules/pbl_evidence
|
||
HARNESS = TESTS_DIR / "m5a_live_db_harness.py"
|
||
|
||
|
||
def load_harness():
|
||
spec = importlib.util.spec_from_file_location("m5a_h", str(HARNESS))
|
||
mod = importlib.util.module_from_spec(spec)
|
||
spec.loader.exec_module(mod)
|
||
return mod
|
||
|
||
|
||
def main(argv=None):
|
||
ap = argparse.ArgumentParser(
|
||
description="S3 collector 重放幂等取证(真库;判定失败 -> exit 1)")
|
||
ap.add_argument("--replays", type=int, default=3,
|
||
help="对同一批事件重放 collector 的次数(默认 3)")
|
||
ap.add_argument("--expect-rows", type=int, default=None,
|
||
help="期望 pbl_evidence 行数;缺省由 harness 种子常量 EVENTS_SEED 派生")
|
||
ap.add_argument("--tenant", default="t1", help="租户 ID(默认 t1,与 harness 种子一致)")
|
||
ap.add_argument("--id-max-len", type=int, default=32,
|
||
help="主键最大长度(models/pbl_evidence.json 定义 VARCHAR(32))")
|
||
ns = ap.parse_args(argv)
|
||
|
||
h = load_harness()
|
||
expect_rows = ns.expect_rows if ns.expect_rows is not None else len(h.EVENTS_SEED)
|
||
|
||
collector = h.load_collector()
|
||
box = h.Sandbox(h.db_kwargs())
|
||
box.open() # 复用已存在的沙箱库(不 reset,保留 harness 写入的数据)
|
||
h.install_stub(collector, box)
|
||
|
||
def count():
|
||
rows = box.raw("SELECT COUNT(*) AS c FROM `pbl_evidence`")
|
||
return int((rows or [{"c": -1}])[0]["c"])
|
||
|
||
def ids():
|
||
rows = box.raw("SELECT id, LENGTH(id) AS len FROM `pbl_evidence` ORDER BY id")
|
||
return [(r["id"], int(r["len"])) for r in rows]
|
||
|
||
problems = []
|
||
try:
|
||
print("沙箱 schema = %s | 凭据来源 = %s(口令不打印)"
|
||
% (h.SANDBOX_SCHEMA, h._CRED_SOURCE[0]))
|
||
print("期望行数(由 harness 种子 EVENTS_SEED 派生,%d 个事件 -> %d 条证据)= %d"
|
||
% (len(h.EVENTS_SEED), len(h.EVENTS_SEED), expect_rows))
|
||
|
||
before = count()
|
||
print("重放前 count(*) = %d" % before)
|
||
print("重放前主键回显 (id, LENGTH(id)) = %s" % ids())
|
||
print("事件表 count(*) = %s"
|
||
% box.raw("SELECT COUNT(*) AS c FROM `pbl_runtime_event`"))
|
||
print("当前 trigger 数(应为 0,证明 id 非 trigger 产物)= %s"
|
||
% len(box.raw("SHOW TRIGGERS") or []))
|
||
|
||
if before != expect_rows:
|
||
problems.append("重放前 count(*)=%d != 期望 %d" % (before, expect_rows))
|
||
|
||
for i in range(ns.replays):
|
||
loop = asyncio.get_event_loop_policy().new_event_loop()
|
||
try:
|
||
st = loop.run_until_complete(
|
||
collector.collect_evidence_from_events(tenant_id=ns.tenant))
|
||
finally:
|
||
loop.close()
|
||
slim = {k: st[k] for k in ("scanned", "created", "skipped", "updated", "failed")
|
||
if k in st}
|
||
print("第 %d 次重放 collector stats = %s | 重放后 count(*) = %d"
|
||
% (i + 1, slim, count()))
|
||
if count() != expect_rows:
|
||
problems.append("第 %d 次重放后 count(*)=%d != 期望 %d(行数增长=幂等破)"
|
||
% (i + 1, count(), expect_rows))
|
||
if int(slim.get("failed", 0) or 0) != 0:
|
||
problems.append("第 %d 次重放 failed=%s(出现 1062/写入异常)"
|
||
% (i + 1, slim.get("failed")))
|
||
|
||
after = ids()
|
||
uniq = {i for i, _ in after}
|
||
print("重放后主键回显 (id, LENGTH(id)) = %s" % after)
|
||
print("重复主键检测:ids 去重后数量 = %d / 总数 = %d" % (len(uniq), len(after)))
|
||
print("id 长度集合 = %s(models 定义 VARCHAR(%d),全 <=%d)"
|
||
% (sorted({l for _, l in after}), ns.id_max_len, ns.id_max_len))
|
||
print("全部 id 以 ev 前缀(gen_pk 产物)= %s"
|
||
% all(i.startswith("ev") for i, _ in after))
|
||
|
||
if len(uniq) != len(after):
|
||
problems.append("存在重复主键(去重后 %d < 总数 %d)" % (len(uniq), len(after)))
|
||
if any(l > ns.id_max_len for _, l in after):
|
||
problems.append("存在 id 长度 > %d" % ns.id_max_len)
|
||
if not after:
|
||
problems.append("pbl_evidence 为空,无证据可判")
|
||
|
||
# DB 层唯一键拦截探针:同三元组再插一条,必须被 uk_ev_dedup 拦下(证明不是应用层侥幸)
|
||
probe = box.raw("SELECT source_event_id, evidence_type FROM `pbl_evidence` LIMIT 1")
|
||
if probe:
|
||
sql = ("INSERT INTO `pbl_evidence` (id, tenant_id, source_event_id, evidence_type,"
|
||
" occurred_at, dedup_key) VALUES ('dup_probe','%s','%s','%s',"
|
||
"'2026-09-22 10:00:00','x')"
|
||
% (ns.tenant, probe[0]["source_event_id"], probe[0]["evidence_type"]))
|
||
try:
|
||
with box.conn.cursor() as cur:
|
||
cur.execute(sql)
|
||
problems.append("探针重复三元组 INSERT 未被 DB 拦截(uk_ev_dedup 失效)")
|
||
print("探针插入未被拦截(异常!)")
|
||
except Exception as exc: # noqa: BLE001
|
||
print("探针重复三元组 INSERT 的 DB 原始报错 = %s"
|
||
% h._mask("%s: %s" % (type(exc).__name__, exc)))
|
||
print("count(*) 最终 = %d(重放 %d 次未增长)" % (count(), ns.replays))
|
||
finally:
|
||
box.close()
|
||
|
||
if problems:
|
||
for p in problems:
|
||
print("FAIL: " + p)
|
||
print("REPLAY IDEMPOTENT FAIL")
|
||
return 1
|
||
print("REPLAY IDEMPOTENT OK")
|
||
return 0
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|