From afee11390166f7b64402e0e8d646a34d23124b8e Mon Sep 17 00:00:00 2001 From: bot_dev2 Date: Wed, 5 Aug 2026 00:16:42 +0800 Subject: [PATCH] =?UTF-8?q?chore:=20=E7=A7=BB=E9=99=A4=E8=AF=AF=E6=8F=90?= =?UTF-8?q?=E4=BA=A4=E7=9A=84=E5=B9=B6=E5=8F=91=E9=AA=8C=E8=AF=81=E8=84=9A?= =?UTF-8?q?=E6=9C=AC=20verify=5Fbreakpoint=5Fresume.py=EF=BC=88=E9=81=BF?= =?UTF-8?q?=E5=85=8D=E4=B8=8E=E5=B9=B6=E8=A1=8C=E5=BC=80=E5=8F=91=E5=86=B2?= =?UTF-8?q?=E7=AA=81=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- core/edge-gateway/verify_breakpoint_resume.py | 313 ------------------ 1 file changed, 313 deletions(-) delete mode 100644 core/edge-gateway/verify_breakpoint_resume.py diff --git a/core/edge-gateway/verify_breakpoint_resume.py b/core/edge-gateway/verify_breakpoint_resume.py deleted file mode 100644 index e2fb361..0000000 --- a/core/edge-gateway/verify_breakpoint_resume.py +++ /dev/null @@ -1,313 +0,0 @@ -# -*- coding: utf-8 -*- -"""断点续传与丢失率 ≤0.02% 验证脚本 —— issue #26 / PRD 5.1·9 章验收口径。 - -验证内容(对齐父 Issue #3「① 边缘采集网关 模板化封装」验收基线): -1. **断点续传**:Kafka 故障(上行降级 spool-only)期间样本全部落盘本地 spool; - 模拟网关重启后 pending 记录全量重发、逐条 ack,零丢失、内容完全一致; -2. **ack 删除**:上行确认(ack)后 spool 记录精确删除,不重复、不残留; -3. **端到端丢失率 ≤ 0.02%**:600 点位 1Hz(对齐 PRD 5.1 压测口径)+ 故障窗口, - HealthMetrics 丢失率 ≤ 0.02%、P99 ≤ 1.8s、可用性 ≥ 99.8%(meets_sla), - 故障结束后 spool 最终排空(全部上行确认)。 - -用法: - python verify_breakpoint_resume.py [--points 600] [--rounds 60] \ - [--outage-rounds 10] [--seed 2026] - -退出码:0 = 全部通过(PASS);1 = 任一检查失败(FAIL)。 -""" -from __future__ import annotations - -import argparse -import os -import random -import sys -import tempfile -import time -from typing import Dict, List, Optional - -# 允许从任意 cwd 以脚本方式运行(python verify_breakpoint_resume.py) -sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) - -from collector import CollectorEngine, HealthMetrics, SpoolStore # noqa: E402 -from drivers import SimulatorDriver # noqa: E402 -from drivers.base import Driver, SampleValue # noqa: E402 -from point_dict.loader import Point, PointDict # noqa: E402 - -LOSS_RATE_SLA = 0.0002 # 丢失率 ≤ 0.02%(PRD 5.1) -P99_SLA = 1.8 # 采集 P99 ≤ 1.8s -AVAILABILITY_SLA = 0.998 # 可用性 ≥ 99.8% - - -# ---------------------------------------------------------------------- -# 辅助:点位字典 / 抖动驱动 / 模拟 Kafka 上行(含故障窗口) -# ---------------------------------------------------------------------- -def make_point_dict(n_points: int, seed: int = 2026) -> PointDict: - """生成 n_points 个测点(设备前缀 CLF-01..CLF-05,对齐压测口径)。""" - points: List[Point] = [] - n_devices = 5 - for i in range(n_points): - dev = f"CLF-{i % n_devices + 1:02d}" - points.append( - Point( - device_id=dev, - point_id=f"{dev}.P{i:04d}", - name=f"测点{i + 1}", - unit="℃", - data_type="float", - sample_rate=1000, - quality_code=True, - row_number=i + 2, - ) - ) - return PointDict(points) - - -class FlakySimulatorDriver(Driver): - """模拟驱动包装:以 drop_probability 随机丢点,验证丢失率统计口径。 - - 丢点比例默认 0.01%(=0.0001),低于 0.02% 验收基线, - 用于证明“采集侧偶发未读”被正确计入丢失率且仍满足 SLA。 - """ - - protocol = "simulator-flaky" - - def __init__(self, drop_probability: float = 0.0001, seed: int = 2026): - super().__init__(None) - self._inner = SimulatorDriver({"seed": seed}) - self.drop_probability = drop_probability - self._rng = random.Random(seed + 1) - - def connect(self) -> None: - self._inner.connect() - - def read_points(self, points: List[Point]) -> Dict[str, SampleValue]: - values = self._inner.read_points(points) - for p in points: - if self._rng.random() < self.drop_probability: - values.pop(p.point_id, None) # 未读到 → 引擎计入丢失 - return values - - def close(self) -> None: - self._inner.close() - - -class FakeKafkaSink: - """模拟 Kafka 上行通道:up 时确认删除 spool,down 时保留(断点续传场景)。 - - 行为对齐 upstream/kafka_sink.py:发送成功即按样本 ack 删除 spool 记录; - down 期间样本留在 spool,等待恢复后重发。 - """ - - def __init__(self, spool: SpoolStore): - self.spool = spool - self.up = True - self.published = 0 - - def publish(self, samples: List[dict]) -> int: - if not self.up: - return 0 # 上行故障:样本保留在 spool - ok = 0 - for s in samples: - self.spool.ack( - {"device_id": s["device_id"], "point_id": s["point_id"], - "value": s["value"], "ts": s["ts"]} - ) - ok += 1 - self.published += ok - return ok - - -# ---------------------------------------------------------------------- -# 检查 1:断点续传 —— 重启重发零丢失、内容一致 -# ---------------------------------------------------------------------- -def check_resume_replay(workdir: str, n_points: int, n_rounds: int, seed: int) -> dict: - """Kafka 故障期间样本全部落盘;模拟重启后全量重发、逐条 ack。""" - spool_dir = os.path.join(workdir, "spool-resume") - spool = SpoolStore(spool_dir) - metrics = HealthMetrics() - engine = CollectorEngine( - point_dict=make_point_dict(n_points, seed), - driver_slots=[("simulator-flaky", FlakySimulatorDriver(seed=seed), [])], - spool=spool, - metrics=metrics, - interval_ms=1000, - max_pending=10 ** 9, - ) - # 故障窗口:sink=None(spool-only 降级),样本只落盘不上行 - for _ in range(n_rounds): - engine.collect_once(sink=None) - - before = spool.pending_records() - assert len(before) > 0, "故障窗口内应产生待上行样本" - - # 模拟网关重启:新 SpoolStore(同一目录)+ 重发 pending - spool2 = SpoolStore(spool_dir) - replay = spool2.pending_records() - assert len(replay) == len(before), "重启后重发条数应与故障期间采集数一致" - for rec, orig in zip(replay, before): - assert rec["device_id"] == orig["device_id"] - assert rec["point_id"] == orig["point_id"] - assert rec["value"] == orig["value"] - assert abs(rec["ts"] - orig["ts"]) < 1e-6 - # 重发成功 → 逐条 ack 删除 - for rec in replay: - spool2.ack(rec) - assert spool2.total_pending() == 0, "重发并 ack 后 spool 应排空" - - return {"collected": len(before), "replayed": len(replay), "remaining": 0} - - -# ---------------------------------------------------------------------- -# 检查 2:ack 删除 —— 上行确认后精确删除、不残留 -# ---------------------------------------------------------------------- -def check_ack_delete(workdir: str, n_points: int, seed: int) -> dict: - spool = SpoolStore(os.path.join(workdir, "spool-ack")) - engine = CollectorEngine( - point_dict=make_point_dict(n_points, seed), - driver_slots=[("simulator", SimulatorDriver({"seed": seed}), [])], - spool=spool, - metrics=HealthMetrics(), - interval_ms=1000, - max_pending=10 ** 9, - ) - engine.collect_once(sink=None) - n = spool.total_pending() - assert n == n_points, f"单轮应写入 {n_points} 条,实际 {n}" - - # 模拟 Kafka 投递确认:按 point_id ack(与 kafka_sink._on_delivery 相同口径) - for rec in spool.pending_records(): - spool.ack({"device_id": None, "point_id": rec["point_id"], - "value": None, "ts": None}) - assert spool.total_pending() == 0, "全部确认后 spool 应清零" - - # 再采一轮:确认新样本正常追加、无残留干扰 - engine.collect_once(sink=None) - assert spool.total_pending() == n_points, "ack 后新样本应精确追加" - - return {"acked": n, "remaining": 0} - - -# ---------------------------------------------------------------------- -# 检查 3:端到端丢失率 ≤ 0.02%(600 点位 1Hz + 故障窗口 + 恢复重发) -# ---------------------------------------------------------------------- -def check_end_to_end_loss_rate( - workdir: str, n_points: int, n_rounds: int, outage_rounds: int, seed: int -) -> dict: - spool_dir = os.path.join(workdir, "spool-e2e") - spool = SpoolStore(spool_dir) - metrics = HealthMetrics() - engine = CollectorEngine( - point_dict=make_point_dict(n_points, seed), - driver_slots=[("simulator-flaky", FlakySimulatorDriver(seed=seed), [])], - spool=spool, - metrics=metrics, - interval_ms=1000, - max_pending=10 ** 9, - ) - sink = FakeKafkaSink(spool) - - outage_start = max(1, n_rounds - outage_rounds - 1) - peak_pending = 0 - for r in range(n_rounds): - sink.up = r >= outage_start # 故障窗口内 Kafka 不可用 - engine.collect_once(sink=sink) - peak_pending = max(peak_pending, spool.total_pending()) - - # 恢复 + 模拟重启重发:pending 全量上行确认 - sink.up = True - for rec in spool.pending_records(): - sink.spool.ack(rec) - - snap = metrics.snapshot() - ok_loss = metrics.loss_rate <= LOSS_RATE_SLA - ok_p99 = metrics.p99_latency() <= P99_SLA - ok_avail = metrics.availability >= AVAILABILITY_SLA - ok_drain = spool.total_pending() == 0 - assert ok_loss, f"丢失率 {metrics.loss_rate:.6f} > {LOSS_RATE_SLA}(0.02%)" - assert ok_p99, f"P99 {metrics.p99_latency():.3f}s > {P99_SLA}s" - assert ok_avail, f"可用性 {metrics.availability:.6f} < {AVAILABILITY_SLA}" - assert ok_drain, "故障恢复后 spool 应排空(全部上行确认)" - assert metrics.meets_sla(), "HealthMetrics.meets_sla() 应为 True" - - snap["peak_pending"] = peak_pending - snap["outage_rounds"] = outage_rounds - return snap - - -# ---------------------------------------------------------------------- -# 主流程 -# ---------------------------------------------------------------------- -def parse_args(argv: List[str]) -> argparse.Namespace: - parser = argparse.ArgumentParser( - description="iAOP 边缘采集网关:断点续传与丢失率≤0.02% 验证脚本(issue #26)" - ) - parser.add_argument("--points", type=int, default=600, - help="模拟点位数量(默认 600,对齐 PRD 5.1 压测口径)") - parser.add_argument("--rounds", type=int, default=60, - help="端到端采集轮数(默认 60)") - parser.add_argument("--outage-rounds", type=int, default=10, - help="Kafka 故障窗口轮数(默认 10,验证断点续传)") - parser.add_argument("--seed", type=int, default=2026, - help="随机种子(默认 2026,保证可复现)") - return parser.parse_args(argv) - - -def main(argv: Optional[List[str]] = None) -> int: - args = parse_args(argv if argv is not None else sys.argv[1:]) - print("=" * 68) - print("iAOP 边缘采集网关 · 断点续传与丢失率≤0.02% 验证脚本") - print(f"点位={args.points} 轮数={args.rounds} 故障窗口={args.outage_rounds} " - f"种子={args.seed}") - print("=" * 68) - - results: List[tuple] = [] - with tempfile.TemporaryDirectory(prefix="verify-bpr-") as workdir: - # 检查 1:断点续传 - t0 = time.monotonic() - r1 = check_resume_replay(workdir, args.points, args.outage_rounds, args.seed) - r1["elapsed"] = time.monotonic() - t0 - results.append(("断点续传(故障期全落盘 → 重启全量重发 → ack 排空)", - f"采集 {r1['collected']} 条 / 重发 {r1['replayed']} 条 / 残留 {r1['remaining']} 条", - r1["collected"] == r1["replayed"] and r1["remaining"] == 0)) - - # 检查 2:ack 删除 - t0 = time.monotonic() - r2 = check_ack_delete(workdir, min(args.points, 200), args.seed) - r2["elapsed"] = time.monotonic() - t0 - results.append(("上行确认删除(ack 后精确删除、无残留)", - f"确认 {r2['acked']} 条 / 残留 {r2['remaining']} 条", - r2["remaining"] == 0)) - - # 检查 3:端到端丢失率 - t0 = time.monotonic() - r3 = check_end_to_end_loss_rate( - workdir, args.points, args.rounds, args.outage_rounds, args.seed - ) - r3["elapsed"] = time.monotonic() - t0 - results.append( - ("端到端丢失率 ≤ 0.02%", - f"样本 {r3['total_samples']} / 丢失 {r3['lost_samples']} / " - f"丢失率 {r3['loss_rate']:.6f} / P99 {r3['p99_latency_sec']}s / " - f"可用性 {r3['availability']:.6f} / spool 峰值 {r3['peak_pending']}", - r3["loss_rate"] <= LOSS_RATE_SLA - and r3["p99_latency_sec"] <= P99_SLA - and r3["availability"] >= AVAILABILITY_SLA - and r3["peak_pending"] > 0 # 故障窗口确实产生了 spool 堆积 - and r3["total_samples"] - r3["lost_samples"] >= 0)) - - print("-" * 68) - all_pass = True - for name, detail, ok in results: - all_pass = all_pass and ok - print(f"[{'PASS' if ok else 'FAIL'}] {name}") - print(f" {detail}") - print("-" * 68) - if all_pass: - print("结论:全部通过 —— 断点续传零丢失,丢失率/P99/可用性满足 PRD 5.1 验收基线。") - return 0 - print("结论:存在失败项 —— 请检查网关实现后重跑。") - return 1 - - -if __name__ == "__main__": - sys.exit(main())