593 lines
26 KiB
Python
593 lines
26 KiB
Python
#!/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 # 保留沙箱库供人工核对
|
||
|
||
沙箱库凭据(**代码内零明文口令**,唯一事实源二选一,取不到即 fail-fast)::
|
||
|
||
1) 环境变量 PBL_M5A_DB=host:port:user:pwd
|
||
2) 环境文件 projects/pbls/env/test.json 的 db.sandbox 段(scope=sandbox_only)
|
||
|
||
该沙箱账号**只用于一次性沙箱 schema**(本脚本自建库、跑完 DROP DATABASE),
|
||
不接触业务库、不承载任何应用运行时连接。日志对口令做脱敏,绝不打印。
|
||
|
||
退出码:0 = 全部断言通过(stdout 末行 ``LIVE DB HARNESS ALL PASS``);1 = 失败
|
||
(含未捕获异常——统一收敛为显式 FAIL 项,不靠 traceback 截断证据)。
|
||
"""
|
||
|
||
import argparse
|
||
import importlib
|
||
import io
|
||
import json
|
||
import os
|
||
import pathlib
|
||
import re
|
||
import sys
|
||
import traceback
|
||
import types
|
||
|
||
REPO_ROOT = pathlib.Path(__file__).resolve().parent.parent # modules/pbl_evidence
|
||
WORKSPACE_ROOT = REPO_ROOT.parent.parent # 机构工作空间
|
||
PKG_DIR = REPO_ROOT / 'pbl_evidence'
|
||
ENV_FILE_DEFAULT = str(WORKSPACE_ROOT / 'projects' / 'pbls' / 'env' / 'test.json')
|
||
|
||
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 缺口,collector 源码一行未改):
|
||
# models/pbl_evidence.json 的 id 是 str(32) NOT NULL 且无 default,而 collector
|
||
# _upsert_evidence 的 INSERT 列集合来自 record(不含 id),写完才 SELECT 回读 id
|
||
# ——即它假定 id 由 DB 侧生成。MariaDB STRICT_TRANS_TABLES 下这会直接 1364
|
||
# "Field 'id' doesn't have a default value'". 真实部署 DDL 必须由 build.sh/xls2ddl 生成带默认值
|
||
# 的 id 列或应用层 getID 赋值;本 harness 用 BEFORE INSERT trigger 在**沙箱库**内
|
||
# 等价补上该生成能力,并在日志显式声明为发现项(见 DISC.1 / IDM.10)。
|
||
TRG_PBL_EVIDENCE_ID = """
|
||
CREATE TRIGGER `trg_pbl_evidence_id` BEFORE INSERT ON `pbl_evidence`
|
||
FOR EACH ROW
|
||
BEGIN
|
||
IF NEW.`id` IS NULL OR NEW.`id` = '' THEN
|
||
SET NEW.`id` = REPLACE(UUID(), '-', '');
|
||
END IF;
|
||
END
|
||
"""
|
||
|
||
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
|
||
"""
|
||
|
||
# 种子事件:具名字段 + 字典(pymysql 具名参数 %(field)s),避免位置元组与
|
||
# 列数不一致时触发 "not enough arguments for format string" 类崩溃。
|
||
SEED_FIELDS = ('event_id', 'tenant_id', 'event_type', 'payload_json', 'occurred_at',
|
||
'learner_id', 'session_id', 'blueprint_id', 'artifact_id')
|
||
|
||
EVENTS_SEED = (
|
||
{'event_id': _EVENT_ID_A, 'tenant_id': _TENANT, 'event_type': 'artifact_submitted',
|
||
'payload_json': '{"title":"bridge-design","score":88}',
|
||
'occurred_at': '2026-09-22 10:00:00', 'learner_id': 'stu_01',
|
||
'session_id': 'sess_01', 'blueprint_id': 'bp_01', 'artifact_id': 'art_01'},
|
||
{'event_id': _EVENT_ID_B, 'tenant_id': _TENANT, 'event_type': 'team_message_sent',
|
||
'payload_json': '{"channel":"team","chars":120}',
|
||
'occurred_at': '2026-09-22 10:05:00', 'learner_id': 'stu_01',
|
||
'session_id': 'sess_01', 'blueprint_id': 'bp_01', 'artifact_id': '0'},
|
||
)
|
||
|
||
_failures = []
|
||
_secret_values = set() # 运行期已解析出的口令,日志统一脱敏
|
||
_CRED_SOURCE = ['(未解析)'] # 凭据来源标签(只记来源,不记口令)
|
||
|
||
|
||
def _mask(text):
|
||
"""把任何可能混入日志/异常消息的口令替换成 ******。"""
|
||
out = str(text)
|
||
for secret in _secret_values:
|
||
if secret:
|
||
out = out.replace(secret, '******')
|
||
return out
|
||
|
||
|
||
def check(name, ok, detail=''):
|
||
"""记录一条断言并原样打印(PASS/FAIL),失败不中断,最后统一 exit 1。"""
|
||
detail = _mask(detail)
|
||
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](原样回显用)。
|
||
|
||
注意:无参 SQL 必须传 None 而不是空 dict —— DictCursor 下传 ``{}``
|
||
会让 pymysql 走 ``query % {}`` 格式化,SQL 里任何字面 ``%``(LIKE/DATE_FORMAT)
|
||
都会抛 ``TypeError: not enough arguments for format string``。
|
||
"""
|
||
with self.conn.cursor() as cur:
|
||
cur.execute(sql, args if args else None)
|
||
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) if params else None)
|
||
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) if params else None)
|
||
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' % (_mask(str(err)) if err else '<无>'))
|
||
if tb_text:
|
||
print('traceback 原文:')
|
||
for line in tb_text.rstrip().splitlines():
|
||
print(' | ' + _mask(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 = _mask(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)
|
||
box.ddl(TRG_PBL_EVIDENCE_ID)
|
||
print('[DISC.1] 沙箱适配声明:pbl_evidence.id 在 models 定义里 NOT NULL 无 default,'
|
||
'而 collector 的 INSERT 不带 id(写完才回读)——沙箱库用 BEFORE INSERT '
|
||
'trigger 等价补 id 生成能力,被测源码零改动;真实部署需由 DDL 默认值或'
|
||
'应用层 getID 落定(已作为发现项上报,不在本 harness 修复范围)。')
|
||
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')))
|
||
|
||
# 具名参数插入(列名与占位符同源 SEED_FIELDS,列/值不可能错位)
|
||
seed_cols = ', '.join('`' + f + '`' for f in SEED_FIELDS)
|
||
seed_vals = ', '.join('%(' + f + ')s' for f in SEED_FIELDS)
|
||
seed_sql = ('INSERT INTO `pbl_runtime_event` (' + seed_cols + ') VALUES ('
|
||
+ seed_vals + ')')
|
||
print('事件种子 INSERT(具名参数): %s' % seed_sql)
|
||
with box.conn.cursor() as cur:
|
||
for ev in EVENTS_SEED:
|
||
cur.execute(seed_sql, dict(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 (%(id)s,%(tenant_id)s,%(source_event_id)s,'
|
||
'%(evidence_type)s,%(occurred_at)s,%(dedup_key)s)',
|
||
{'id': 'manual_dup_1', 'tenant_id': _TENANT,
|
||
'source_event_id': _EVENT_ID_A, 'evidence_type': 'artifact',
|
||
'occurred_at': '2026-09-22 10:00:00', 'dedup_key': 'x' * 32})
|
||
except Exception as exc: # noqa: BLE001
|
||
dup_err = '%s: %s' % (type(exc).__name__, _mask(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))
|
||
id_rows = box.raw('SELECT id, tenant_id, source_event_id, evidence_type '
|
||
'FROM `pbl_evidence` ORDER BY source_event_id, evidence_type')
|
||
ids = [str(r.get('id') or '') for r in id_rows]
|
||
print('证据行主键回显(id 由沙箱 trigger 生成,非空且互异): %s' % ids)
|
||
check('IDM.10 落库行 id 非空且互不相同(沙箱 id 生成适配生效,非静默吞错)',
|
||
len(ids) == 2 and all(ids) and len(set(ids)) == 2, 'ids=%s' % ids)
|
||
return s2, rows2
|
||
|
||
|
||
def _dedup_matches(rows):
|
||
from hashlib import md5
|
||
for r in rows:
|
||
key = '%s|%s|%s' % (r['tenant_id'], r['source_event_id'],
|
||
r['evidence_type'])
|
||
want = md5(key.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():
|
||
"""沙箱库连接参数。凭据唯一事实源(**代码内零明文口令**,取不到即 fail-fast):
|
||
|
||
1. 环境变量 ``PBL_M5A_DB=host:port:user:pwd``(口令允许含冒号,按前 3 段切分)
|
||
2. 环境文件 ``projects/pbls/env/test.json`` 的 ``db.sandbox`` 段
|
||
(``scope=sandbox_only``,仅授权一次性沙箱 schema 的建/删)
|
||
|
||
两者都拿不到 → 抛 ``RuntimeError``(不回退任何默认账号、不猜凭据)。
|
||
"""
|
||
raw_env = (os.environ.get('PBL_M5A_DB') or '').strip()
|
||
if raw_env:
|
||
parts = raw_env.split(':', 3)
|
||
if len(parts) != 4:
|
||
raise RuntimeError(
|
||
'PBL_M5A_DB 格式必须为 host:port:user:pwd(4 段冒号分隔),'
|
||
'实际段数=%d' % len(parts))
|
||
host, port, user, pwd = parts
|
||
if not host or not user:
|
||
raise RuntimeError('PBL_M5A_DB 的 host/user 段为空,拒绝连接')
|
||
try:
|
||
port_i = int(port)
|
||
except ValueError:
|
||
raise RuntimeError('PBL_M5A_DB 的 port 段不是整数: %r' % (port,))
|
||
_secret_values.add(pwd)
|
||
_CRED_SOURCE[0] = '环境变量 PBL_M5A_DB'
|
||
return dict(host=host, port=port_i, user=user, password=pwd,
|
||
charset='utf8mb4', connect_timeout=5)
|
||
|
||
env_file = (os.environ.get('PBL_M5A_ENV_FILE') or '').strip() or ENV_FILE_DEFAULT
|
||
if not os.path.isfile(env_file):
|
||
raise RuntimeError(
|
||
'拿不到沙箱库凭据:环境变量 PBL_M5A_DB 未设置,且环境文件不存在: '
|
||
+ env_file
|
||
+ '。请设置 PBL_M5A_DB=host:port:user:pwd,或在该文件的 db.sandbox 段'
|
||
'(engine/host/port/user/password/charset/scope=sandbox_only)补齐沙箱账号。')
|
||
try:
|
||
with io.open(env_file, 'r', encoding='utf-8') as fh:
|
||
cfg = json.load(fh)
|
||
except ValueError as exc:
|
||
raise RuntimeError('环境文件不是合法 JSON: %s (%s)' % (env_file, _mask(exc)))
|
||
sandbox = ((cfg.get('db') or {}).get('sandbox')) or {}
|
||
if not sandbox:
|
||
raise RuntimeError(
|
||
'环境文件缺少 db.sandbox 段(沙箱库凭据唯一事实源): ' + env_file
|
||
+ ';或改用环境变量 PBL_M5A_DB=host:port:user:pwd')
|
||
scope = str(sandbox.get('scope') or '')
|
||
if scope and scope != 'sandbox_only':
|
||
raise RuntimeError(
|
||
'db.sandbox.scope=' + scope + ' 非 sandbox_only,拒绝用非沙箱账号跑 harness')
|
||
missing = [k for k in ('host', 'port', 'user') if sandbox.get(k) in (None, '')]
|
||
if 'password' not in sandbox:
|
||
missing.append('password')
|
||
if missing:
|
||
raise RuntimeError('db.sandbox 缺少必填字段: ' + ', '.join(missing)
|
||
+ '(文件 ' + env_file + ')')
|
||
pwd = str(sandbox['password'])
|
||
_secret_values.add(pwd)
|
||
_CRED_SOURCE[0] = '环境文件 ' + env_file + ' 的 db.sandbox 段'
|
||
return dict(host=str(sandbox['host']), port=int(sandbox['port']),
|
||
user=str(sandbox['user']), password=pwd,
|
||
charset=str(sandbox.get('charset') or 'utf8mb4'),
|
||
connect_timeout=5)
|
||
|
||
|
||
def run(args):
|
||
"""harness 主体(异常由 main() 统一收敛为显式 FAIL)。"""
|
||
try:
|
||
import pymysql # noqa: F401
|
||
except Exception as exc: # noqa: BLE001
|
||
check('ENV.1 运行环境具备 pymysql', False, repr(exc))
|
||
return 1
|
||
|
||
conn_kwargs = db_kwargs()
|
||
print('沙箱库凭据来源: %s | 目标 = %s@%s:%s(口令不打印,日志已脱敏)'
|
||
% (_CRED_SOURCE[0], conn_kwargs['user'], conn_kwargs['host'],
|
||
conn_kwargs['port']))
|
||
|
||
collector = load_collector()
|
||
box = Sandbox(conn_kwargs)
|
||
box.reset_schema()
|
||
box.open()
|
||
try:
|
||
install_stub(collector, box)
|
||
phase_fail_closed(collector, box)
|
||
_s2, _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)
|
||
return 0
|
||
|
||
|
||
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:
|
||
run(args)
|
||
except SystemExit:
|
||
raise
|
||
except Exception as exc: # noqa: BLE001
|
||
# 未捕获异常(凭据缺失/连不上库/被测代码抛错)统一收敛为显式 FAIL 项,
|
||
# 打印异常类型 + 脱敏消息,绝不靠 traceback 截断证据。
|
||
check('HARNESS.1 harness 无未捕获异常(凭据可用 + 沙箱库可连 + 各阶段跑完)',
|
||
False, '%s: %s' % (type(exc).__name__, _mask(exc)))
|
||
print('未捕获异常汇总(脱敏): %s' % _mask(traceback.format_exc()))
|
||
if _failures:
|
||
print('LIVE DB HARNESS FAILED (%d): %s'
|
||
% (len(_failures), ' | '.join(_mask(f) for f in _failures)))
|
||
return 1
|
||
print('LIVE DB HARNESS ALL PASS')
|
||
return 0
|
||
|
||
|
||
if __name__ == '__main__':
|
||
sys.exit(main())
|