pbl_evidence/tests/m5a_live_db_harness.py
2026-09-22 17:20:20 +08:00

449 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""M5a 幂等 / C-3 fail-closed **真实库**验收 harness(任务 u7RBcE_UjcEOkfhw6WRlV 产出)。
与既有 ``tests/test_m5a_idempotency.py``(离线、无 DB、只断 fail-closed 与纯函数)
的分工:本 harness 是**增量证据**——它在一次性沙箱库里建出
``pbl_runtime_event`` + ``pbl_evidence`` 两张真实表,跑**同一条采集路径两次**,
用 MariaDB 的 ``uk_ev_dedup`` 唯一索引与 ``INSERT ... ON DUPLICATE KEY UPDATE``
实测幂等,并把 DB 原样回显(SHOW INDEX / 1062 Duplicate entry)打进日志。
被测代码是**真源码**(``pbl_evidence/collector.py``),不是复刻逻辑:
只把最底层数据访问适配层 ``db.q_all/q_one/q_exec/table_exists/table_columns``
替换成 pymysql 直连沙箱库的等价实现(sqlor 的 ``${name}$`` 占位符 → pymysql
``%(name)s``),``collector`` 的扫描/映射/去重/幂等写入逻辑一行未改。
宿主包(sqlor/appPublic/ahserver)在本机未安装,故用 importlib 造假包绕过
``pbl_evidence/__init__.py``,与离线测试同一手法。
用法::
python3 tests/m5a_live_db_harness.py # 跑完自动清理沙箱库
python3 tests/m5a_live_db_harness.py --keep # 保留沙箱库供人工核对
PBL_M5A_DB=host:port:user:pwd 覆盖默认 test/test123@127.0.0.1:3306
退出码:0 = 全部断言通过(stdout 末行 ``LIVE DB HARNESS ALL PASS``);1 = 失败。
"""
import argparse
import contextlib
import importlib
import io
import os
import pathlib
import re
import sys
import traceback
import types
REPO_ROOT = pathlib.Path(__file__).resolve().parent.parent # modules/pbl_evidence
PKG_DIR = REPO_ROOT / 'pbl_evidence'
SANDBOX_SCHEMA = os.environ.get('PBL_M5A_SANDBOX', 'pbl_m5a_u7rb')
_TENANT = 't1'
_EVENT_ID_A = 'evt_u7rb_0001'
_EVENT_ID_B = 'evt_u7rb_0002'
# 沙箱 DDL:列集合与 models/pbl_evidence.json 一致(MariaDB 方言,本机实测可建)
DDL_PBL_EVIDENCE = """
CREATE TABLE `pbl_evidence` (
`id` VARCHAR(32) NOT NULL,
`tenant_id` VARCHAR(32) NOT NULL DEFAULT '0',
`artifact_id` VARCHAR(32) NOT NULL DEFAULT '0',
`evidence_type` VARCHAR(64) NOT NULL DEFAULT 'observation',
`source_event_id` VARCHAR(32) NOT NULL,
`session_id` VARCHAR(32) NOT NULL DEFAULT '0',
`learner_id` VARCHAR(32) NOT NULL DEFAULT '0',
`blueprint_id` VARCHAR(32) NOT NULL DEFAULT '0',
`payload_json` LONGTEXT NULL,
`occurred_at` DATETIME NOT NULL,
`dedup_key` VARCHAR(32) NOT NULL DEFAULT '',
`collect_batch_id` VARCHAR(32) NOT NULL DEFAULT '',
`created_at` DATETIME NOT NULL DEFAULT '1970-01-01 00:00:00',
`updated_at` DATETIME NOT NULL DEFAULT '1970-01-01 00:00:00',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_ev_dedup` (`tenant_id`, `source_event_id`, `evidence_type`),
KEY `idx_ev_learner` (`tenant_id`, `learner_id`, `occurred_at`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci
"""
DDL_PBL_RUNTIME_EVENT = """
CREATE TABLE `pbl_runtime_event` (
`event_id` VARCHAR(64) NOT NULL,
`tenant_id` VARCHAR(32) NOT NULL DEFAULT '0',
`event_type` VARCHAR(64) NOT NULL,
`payload_json` LONGTEXT NULL,
`occurred_at` DATETIME NOT NULL,
`learner_id` VARCHAR(32) NOT NULL DEFAULT '0',
`session_id` VARCHAR(32) NOT NULL DEFAULT '0',
`blueprint_id` VARCHAR(32) NOT NULL DEFAULT '0',
`artifact_id` VARCHAR(32) NOT NULL DEFAULT '0',
`created_at` DATETIME NOT NULL DEFAULT '1970-01-01 00:00:00',
PRIMARY KEY (`event_id`),
KEY `idx_re_time` (`tenant_id`, `occurred_at`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci
"""
EVENTS_SEED = (
(_EVENT_ID_A, 'artifact_submitted',
'{"title":"bridge-design","score":88}', '2026-09-22 10:00:00', 'stu_01', 'sess_01', 'bp_01', 'art_01'),
(_EVENT_ID_B, 'team_message_sent',
'{"channel":"team","chars":120}', '2026-09-22 10:05:00', 'stu_01', 'sess_01', 'bp_01', '0'),
)
_failures = []
def check(name, ok, detail=''):
"""记录一条断言并原样打印(PASS/FAIL),失败不中断,最后统一 exit 1。"""
print('[%s] %s%s' % ('PASS' if ok else 'FAIL', name,
(' -> ' + detail) if detail else ''))
if not ok:
_failures.append(name)
return ok
# ── 被测代码隔离加载(不触发宿主依赖)────────────────────────────────
def load_collector():
if 'pbl_evidence' not in sys.modules:
pkg = types.ModuleType('pbl_evidence')
pkg.__path__ = [str(PKG_DIR)]
pkg.__package__ = 'pbl_evidence'
sys.modules['pbl_evidence'] = pkg
return importlib.import_module('pbl_evidence.collector')
# ── 沙箱库 + db 适配层替身 ───────────────────────────────────────────
class Sandbox:
"""一次性沙箱库;把 collector 的数据访问函数换成 pymysql 直连实现。"""
def __init__(self, conn_kwargs):
self.conn_kwargs = dict(conn_kwargs)
self.conn = None
def connect_root(self):
import pymysql
kw = dict(self.conn_kwargs)
kw.pop('database', None)
return pymysql.connect(**kw)
def reset_schema(self):
conn = self.connect_root()
try:
with conn.cursor() as cur:
cur.execute('DROP DATABASE IF EXISTS `%s`' % SANDBOX_SCHEMA)
cur.execute('CREATE DATABASE `%s` CHARACTER SET utf8mb4 '
'COLLATE utf8mb4_general_ci' % SANDBOX_SCHEMA)
conn.commit()
finally:
conn.close()
def drop_schema(self):
conn = self.connect_root()
try:
with conn.cursor() as cur:
cur.execute('DROP DATABASE IF EXISTS `%s`' % SANDBOX_SCHEMA)
conn.commit()
finally:
conn.close()
def open(self):
import pymysql
kw = dict(self.conn_kwargs)
kw['database'] = SANDBOX_SCHEMA
kw['cursorclass'] = pymysql.cursors.DictCursor
kw['autocommit'] = True
self.conn = pymysql.connect(**kw)
def close(self):
if self.conn is not None:
try:
self.conn.close()
except Exception:
pass
self.conn = None
def ddl(self, stmt):
with self.conn.cursor() as cur:
cur.execute(stmt)
def raw(self, sql, args=None):
"""执行并返回 list[dict](原样回显用)。"""
with self.conn.cursor() as cur:
cur.execute(sql, args or {})
try:
return list(cur.fetchall() or [])
except Exception:
return []
# ── collector.db 替身:${name}$ → %(name)s ──
@staticmethod
def to_pymysql(sql):
return re.sub(r'\$\{(\w+)\}\$', r'%(\1)s', sql)
def q_all(self, sql, params=None):
with self.conn.cursor() as cur:
cur.execute(self.to_pymysql(sql), dict(params or {}))
return [dict(r) for r in (cur.fetchall() or [])]
def q_one(self, sql, params=None):
rows = self.q_all(sql, params)
return rows[0] if rows else None
def q_exec(self, sql, params=None):
with self.conn.cursor() as cur:
n = cur.execute(self.to_pymysql(sql), dict(params or {}))
return {'affected_rows': n}
def table_exists(self, table):
row = self.q_one(
'SELECT COUNT(*) AS c FROM information_schema.TABLES '
'WHERE TABLE_SCHEMA=${db}$ AND TABLE_NAME=${t}$',
{'db': SANDBOX_SCHEMA, 't': table})
return bool(row) and int(list(row.values())[0] or 0) > 0
def table_columns(self, table):
rows = self.q_all(
'SELECT COLUMN_NAME AS col FROM information_schema.COLUMNS '
'WHERE TABLE_SCHEMA=${db}$ AND TABLE_NAME=${t}$',
{'db': SANDBOX_SCHEMA, 't': table})
return {str(r.get('col')) for r in rows if r.get('col')}
def get_dbname(self):
return SANDBOX_SCHEMA
def install_stub(collector, box):
"""把 collector 模块命名空间里的数据访问函数指向沙箱实现(真源码其余不动)。
collector 里这些调用都是 ``await``,故同步的 pymysql 实现要包成协程;
``get_dbname`` 在真源码里是同步函数,保持同步。
"""
def _async(fn):
async def wrapper(*a, **kw):
return fn(*a, **kw)
wrapper.__name__ = getattr(fn, '__name__', 'stub')
return wrapper
for attr, fn in (('q_all', box.q_all), ('q_one', box.q_one), ('q_exec', box.q_exec),
('table_exists', box.table_exists), ('table_columns', box.table_columns),
('invalidate_cache', lambda t=None: None)):
setattr(collector, attr, _async(fn))
setattr(collector, 'get_dbname', box.get_dbname)
# ── 阶段 1:C-3 fail-closed(事件表缺失)────────────────────────────
def phase_fail_closed(collector, box):
print('\n--- 阶段 1:C-3 fail-closed(pbl_runtime_event 缺失)---')
print('沙箱库 %s 内表清单(原样回显): %s'
% (SANDBOX_SCHEMA, box.raw('SHOW TABLES') or '[] (空库)'))
resolved = None
import asyncio
resolved = asyncio.get_event_loop_policy().new_event_loop().run_until_complete(
collector.resolve_event_table())
print('resolve_event_table() 实测返回: %r' % (resolved,))
check('C-3.1 resolve_event_table() 在事件表缺失时返回 None', resolved is None,
repr(resolved))
err = None
tb_text = ''
import asyncio
try:
asyncio.get_event_loop_policy().new_event_loop().run_until_complete(
collector.collect_evidence_from_events(tenant_id=_TENANT))
except collector.CollectError as exc: # 真源码抛出的类型
err = exc
tb_text = ''.join(traceback.format_exc())
except Exception as exc: # noqa: BLE001
err = exc
tb_text = ''.join(traceback.format_exc())
print('抛出的异常类型: %s' % (type(err).__name__ if err else 'None(未抛错=静默空成功)'))
print('异常消息原文: %s' % (str(err) if err else '<无>'))
if tb_text:
print('traceback 原文:')
for line in tb_text.rstrip().splitlines():
print(' | ' + line)
check('C-3.2 collect_evidence_from_events(tenant_id="t1") 抛 CollectError',
isinstance(err, collector.CollectError), type(err).__name__ if err else 'no raise')
msg = str(err or '')
check("C-3.3 异常消息含 'M11b'(硬依赖声明,非空成功)", 'M11b' in msg,
'消息=%s' % msg[:120])
check('C-3.4 未产生任何证据行(fail-closed 不写脏数据)',
box.table_exists('pbl_evidence') is False or
int((box.raw('SELECT COUNT(*) AS c FROM `pbl_evidence`') or [{'c': 0}])[0]['c']) == 0)
# ── 阶段 2:真实双次采集幂等 ────────────────────────────────────────
def phase_idempotent(collector, box):
print('\n--- 阶段 2:同 source_event_id + evidence_type 连续采集两次(真库)---')
box.ddl(DDL_PBL_RUNTIME_EVENT)
box.ddl(DDL_PBL_EVIDENCE)
print('建表后 SHOW TABLES: %s'
% [list(r.values())[0] for r in box.raw('SHOW TABLES')])
print('uk_ev_dedup 索引定义(SHOW INDEX 原样回显):')
for r in box.raw("SHOW INDEX FROM `pbl_evidence` WHERE Key_name='uk_ev_dedup'"):
print(' Seq_in_index={0} Column_name={1} Non_unique={2} Index_name={3}'
.format(r.get('Seq_in_index'), r.get('Column_name'),
r.get('Non_unique'), r.get('Key_name')))
with box.conn.cursor() as cur:
for ev in EVENTS_SEED:
cur.execute(
'INSERT INTO `pbl_runtime_event` (event_id, tenant_id, event_type, '
'payload_json, occurred_at, learner_id, session_id, blueprint_id, '
'artifact_id) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)', ev)
print('事件表写入 %d 行(pbl_runtime_event)' % len(EVENTS_SEED))
import asyncio
def run_collect(label):
loop = asyncio.get_event_loop_policy().new_event_loop()
try:
stats = loop.run_until_complete(
collector.collect_evidence_from_events(tenant_id=_TENANT))
finally:
loop.close()
slim = {k: stats[k] for k in ('scanned', 'created', 'skipped', 'updated',
'ignored', 'failed', 'event_table', 'watermark')
if k in stats}
print('%s stats = %s' % (label, slim))
return stats
s1 = run_collect('第 1 次采集')
rows1 = box.raw('SELECT tenant_id, source_event_id, evidence_type, dedup_key '
'FROM `pbl_evidence` ORDER BY source_event_id, evidence_type')
print('第 1 次后表内行(原样回显,%d 行):' % len(rows1))
for r in rows1:
print(' %s' % r)
s2 = run_collect('第 2 次采集(同一批事件重放)')
rows2 = box.raw('SELECT tenant_id, source_event_id, evidence_type, dedup_key '
'FROM `pbl_evidence` ORDER BY source_event_id, evidence_type')
print('第 2 次后表内行(原样回显,%d 行):' % len(rows2))
for r in rows2:
print(' %s' % r)
check('IDM.1 第 1 次采集 created == 2(两条事件各落一条证据)',
s1.get('created') == 2, 'created=%s' % s1.get('created'))
check('IDM.2 第 2 次采集 skipped >= 1(命中唯一键走跳过分支)',
s2.get('skipped', 0) >= 1, 'skipped=%s' % s2.get('skipped'))
check('IDM.3 第 2 次采集 created == 0(重放不新增)',
s2.get('created') == 0, 'created=%s' % s2.get('created'))
check('IDM.4 两次采集后表内仍只有 2 行(每三元组唯一一行)',
len(rows2) == 2 and len(rows1) == 2, 'rows=%d' % len(rows2))
# 唯一索引的 DB 级行为:绕过预查,直接重复 INSERT,必须 1062 Duplicate entry
dup_err = ''
try:
with box.conn.cursor() as cur:
cur.execute(
'INSERT INTO `pbl_evidence` (id, tenant_id, source_event_id, '
'evidence_type, occurred_at, dedup_key) '
'VALUES (%s,%s,%s,%s,%s,%s)',
('manual_dup_1', _TENANT, _EVENT_ID_A, 'artifact',
'2026-09-22 10:00:00', 'x' * 32))
except Exception as exc: # noqa: BLE001
dup_err = '%s: %s' % (type(exc).__name__, exc)
print('直接重复 INSERT 的 DB 原始报错(uk_ev_dedup 真防重): %s'
% (dup_err or '<未报错>'))
check('IDM.6 唯一键 uk_ev_dedup 在 DB 层拦截重复三元组(1062 Duplicate entry)',
'1062' in dup_err or 'Duplicate entry' in dup_err, dup_err[:160])
# collector 实际写入语句形态确认(ON DUPLICATE KEY UPDATE 双保险)
src = (PKG_DIR / 'collector.py').read_text(encoding='utf-8')
check("IDM.7 写入路径含 'ON DUPLICATE KEY UPDATE'(并发/重放第二层防重)",
'ON DUPLICATE KEY UPDATE' in src)
check('IDM.8 dedup_key 与 uk 三元组同源(md5(tenant|event|type))',
_dedup_matches(rows2))
return rows2
def _dedup_matches(rows):
from hashlib import md5
for r in rows:
want = md5('%s|%s|%s' % (r['tenant_id'], r['source_event_id'],
r['evidence_type']).encode('utf-8')).hexdigest()
if str(r.get('dedup_key')) != want:
return False
return True
# ── 阶段 3:第三次采集(幂等可重放,补跑一次并对比)─────────────────
def phase_third_run(collector, box, s2_stats):
import asyncio
loop = asyncio.get_event_loop_policy().new_event_loop()
try:
s3 = loop.run_until_complete(
collector.collect_evidence_from_events(tenant_id=_TENANT))
finally:
loop.close()
slim3 = {k: s3[k] for k in ('scanned', 'created', 'skipped', 'updated', 'ignored', 'failed')
if k in s3}
slim2 = {k: s2_stats[k] for k in ('scanned', 'created', 'skipped', 'updated', 'ignored', 'failed')
if k in s2_stats}
print('第 3 次采集 stats = %s(与第 2 次对比 %s)'
% (slim3, '一致' if slim2 == slim3 else '不一致'))
rows = box.raw('SELECT COUNT(*) AS c FROM `pbl_evidence`')
cnt = int((rows or [{'c': -1}])[0]['c'])
print('第 3 次后表内总行数 = %d' % cnt)
check('IDM.5 第三次采集统计与第二次逐字段一致', slim2 == slim3, '%s vs %s' % (slim2, slim3))
check('IDM.9 三次采集后表内行数仍为 2(无重复行累积)', cnt == 2, 'count=%d' % cnt)
def db_kwargs():
env = os.environ.get('PBL_M5A_DB')
if env:
host, port, user, pwd = env.split(':')
return dict(host=host, port=int(port), user=user, password=pwd)
# 本机测试库凭据来自 projects/pbls/env/test.json 的 db 段(唯一事实源,未新增秘密)
return dict(host='127.0.0.1', port=3306, user='test', password='test123',
charset='utf8mb4', connect_timeout=5)
def main():
ap = argparse.ArgumentParser()
ap.add_argument('--keep', action='store_true', help='跑完保留沙箱库')
args = ap.parse_args()
print('M5a live-DB harness | sandbox schema = %s | target = modules/pbl_evidence'
% SANDBOX_SCHEMA)
try:
import pymysql # noqa: F401
except Exception as exc: # noqa: BLE001
print('FAIL 环境缺少 pymysql,无法跑真实库幂等: %r' % (exc,))
return 1
collector = load_collector()
box = Sandbox(db_kwargs())
box.reset_schema()
box.open()
try:
install_stub(collector, box)
phase_fail_closed(collector, box)
_rows = phase_idempotent(collector, box)
# 重跑一次统计对比(放在 phase 2 之后,避免影响其内部断言顺序)
import asyncio
loop = asyncio.get_event_loop_policy().new_event_loop()
try:
s2 = loop.run_until_complete(
collector.collect_evidence_from_events(tenant_id=_TENANT))
finally:
loop.close()
phase_third_run(collector, box, s2)
finally:
box.close()
if not args.keep:
box.drop_schema()
print('沙箱库已清理: DROP DATABASE %s' % SANDBOX_SCHEMA)
else:
print('--keep:沙箱库 %s 保留' % SANDBOX_SCHEMA)
if _failures:
print('LIVE DB HARNESS FAILED (%d): %s' % (len(_failures), ' | '.join(_failures)))
return 1
print('LIVE DB HARNESS ALL PASS')
return 0
if __name__ == '__main__':
sys.exit(main())