test: add live ingest verification gate
Some checks failed
ci/woodpecker/push/woodpecker Pipeline was canceled
Some checks failed
ci/woodpecker/push/woodpecker Pipeline was canceled
This commit is contained in:
72
tools/test_vehicle_ingest_live_verify.py
Normal file
72
tools/test_vehicle_ingest_live_verify.py
Normal file
@@ -0,0 +1,72 @@
|
||||
import unittest
|
||||
|
||||
from tools import vehicle_ingest_live_verify as live_verify
|
||||
|
||||
|
||||
class VehicleIngestLiveVerifyTest(unittest.TestCase):
|
||||
def test_warns_when_formal_gb32960_platform_login_exists_without_vehicle_frames(self):
|
||||
checks = live_verify.evaluate_counts(
|
||||
{
|
||||
"jt808RawFrames": 36090,
|
||||
"jt808Locations": 8783,
|
||||
"gb32960PlatformLogins": 3,
|
||||
"gb32960VehicleRealtimeFrames": 0,
|
||||
},
|
||||
live_verify.Thresholds(
|
||||
min_jt808_raw_frames=1,
|
||||
min_jt808_locations=1,
|
||||
min_gb32960_platform_logins=1,
|
||||
min_gb32960_vehicle_realtime_frames=1,
|
||||
require_gb32960_vehicle_realtime=False,
|
||||
),
|
||||
)
|
||||
|
||||
by_name = {check.name: check for check in checks}
|
||||
|
||||
self.assertEqual(live_verify.overall_status(checks), "warn")
|
||||
self.assertEqual(by_name["gb32960.vehicle_realtime"].status, "warn")
|
||||
self.assertIn("尚未收到正式车辆 0x02", by_name["gb32960.vehicle_realtime"].message)
|
||||
|
||||
def test_fails_when_vehicle_frames_are_required_but_missing(self):
|
||||
checks = live_verify.evaluate_counts(
|
||||
{
|
||||
"jt808RawFrames": 10,
|
||||
"jt808Locations": 2,
|
||||
"gb32960PlatformLogins": 1,
|
||||
"gb32960VehicleRealtimeFrames": 0,
|
||||
},
|
||||
live_verify.Thresholds(require_gb32960_vehicle_realtime=True),
|
||||
)
|
||||
|
||||
by_name = {check.name: check for check in checks}
|
||||
|
||||
self.assertEqual(live_verify.overall_status(checks), "fail")
|
||||
self.assertEqual(by_name["gb32960.vehicle_realtime"].status, "fail")
|
||||
|
||||
def test_gb32960_vehicle_realtime_sql_excludes_test_vins(self):
|
||||
sql = live_verify.gb32960_vehicle_realtime_count_sql(
|
||||
date_from="2026-06-29 00:00:00",
|
||||
date_to="2026-06-29 23:59:59",
|
||||
)
|
||||
|
||||
self.assertIn("message_id = 2", sql)
|
||||
self.assertIn("vin <> ''", sql)
|
||||
self.assertIn("vin NOT LIKE 'LTEST%'", sql)
|
||||
self.assertIn("ts >= '2026-06-29 00:00:00'", sql)
|
||||
self.assertIn("ts <= '2026-06-29 23:59:59'", sql)
|
||||
|
||||
def test_jt808_location_sql_uses_dedicated_location_stable_without_protocol_filter(self):
|
||||
sql = live_verify.jt808_location_count_sql(
|
||||
date_from="2026-06-29 00:00:00",
|
||||
date_to="2026-06-29 23:59:59",
|
||||
)
|
||||
|
||||
self.assertEqual(
|
||||
sql,
|
||||
"SELECT COUNT(*) FROM jt808_locations WHERE "
|
||||
"ts >= '2026-06-29 00:00:00' AND ts <= '2026-06-29 23:59:59'",
|
||||
)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
273
tools/vehicle_ingest_live_verify.py
Executable file
273
tools/vehicle_ingest_live_verify.py
Executable file
@@ -0,0 +1,273 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Verify the live 32960/JT808 ingest-to-history chain with production traffic."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import base64
|
||||
import json
|
||||
import os
|
||||
import urllib.parse
|
||||
import urllib.request
|
||||
from dataclasses import asdict, dataclass
|
||||
from typing import Any
|
||||
|
||||
try:
|
||||
from tools import tdengine_smoke
|
||||
except ModuleNotFoundError:
|
||||
import tdengine_smoke
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Thresholds:
|
||||
min_jt808_raw_frames: int = 1
|
||||
min_jt808_locations: int = 1
|
||||
min_gb32960_platform_logins: int = 1
|
||||
min_gb32960_vehicle_realtime_frames: int = 1
|
||||
require_gb32960_vehicle_realtime: bool = False
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Check:
|
||||
name: str
|
||||
status: str
|
||||
count: int
|
||||
minimum: int
|
||||
message: str
|
||||
|
||||
|
||||
def count_sql(table: str, filters: list[str], date_from: str, date_to: str) -> str:
|
||||
where = list(filters)
|
||||
if date_from:
|
||||
where.append("ts >= " + tdengine_smoke.sql_literal(date_from))
|
||||
if date_to:
|
||||
where.append("ts <= " + tdengine_smoke.sql_literal(date_to))
|
||||
return "SELECT COUNT(*) FROM " + table + " WHERE " + " AND ".join(where)
|
||||
|
||||
|
||||
def jt808_raw_count_sql(date_from: str, date_to: str, peer_like: str = "") -> str:
|
||||
filters = ["protocol = 'JT808'"]
|
||||
if peer_like:
|
||||
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like))
|
||||
return count_sql("raw_frames", filters, date_from, date_to)
|
||||
|
||||
|
||||
def jt808_location_count_sql(date_from: str, date_to: str) -> str:
|
||||
return count_sql("jt808_locations", [], date_from, date_to)
|
||||
|
||||
|
||||
def gb32960_platform_login_count_sql(date_from: str, date_to: str, peer_like: str = "") -> str:
|
||||
filters = ["protocol = 'GB32960'", "message_id = 5"]
|
||||
if peer_like:
|
||||
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like))
|
||||
return count_sql("raw_frames", filters, date_from, date_to)
|
||||
|
||||
|
||||
def gb32960_vehicle_realtime_count_sql(date_from: str, date_to: str, peer_like: str = "") -> str:
|
||||
filters = [
|
||||
"protocol = 'GB32960'",
|
||||
"message_id = 2",
|
||||
"vin <> ''",
|
||||
"vin NOT LIKE 'LTEST%'",
|
||||
]
|
||||
if peer_like:
|
||||
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like))
|
||||
return count_sql("raw_frames", filters, date_from, date_to)
|
||||
|
||||
|
||||
def latest_gb32960_platform_raw_uri_sql(date_from: str, date_to: str, peer_like: str = "") -> str:
|
||||
filters = ["protocol = 'GB32960'", "message_id = 5", "raw_uri <> ''"]
|
||||
if peer_like:
|
||||
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like))
|
||||
where = list(filters)
|
||||
if date_from:
|
||||
where.append("ts >= " + tdengine_smoke.sql_literal(date_from))
|
||||
if date_to:
|
||||
where.append("ts <= " + tdengine_smoke.sql_literal(date_to))
|
||||
return "SELECT raw_uri FROM raw_frames WHERE " + " AND ".join(where) + " ORDER BY ts DESC LIMIT 1"
|
||||
|
||||
|
||||
def evaluate_counts(counts: dict[str, int], thresholds: Thresholds) -> list[Check]:
|
||||
checks = [
|
||||
threshold_check(
|
||||
"jt808.raw_frames",
|
||||
counts.get("jt808RawFrames", 0),
|
||||
thresholds.min_jt808_raw_frames,
|
||||
"JT808 raw_frames 已有正式报文",
|
||||
"JT808 raw_frames 未达到最低验收数量",
|
||||
),
|
||||
threshold_check(
|
||||
"jt808.locations",
|
||||
counts.get("jt808Locations", 0),
|
||||
thresholds.min_jt808_locations,
|
||||
"JT808 位置历史已入 TDengine",
|
||||
"JT808 位置历史未达到最低验收数量",
|
||||
),
|
||||
threshold_check(
|
||||
"gb32960.platform_login",
|
||||
counts.get("gb32960PlatformLogins", 0),
|
||||
thresholds.min_gb32960_platform_logins,
|
||||
"GB32960 正式平台登录已入 TDengine",
|
||||
"GB32960 正式平台登录未达到最低验收数量",
|
||||
),
|
||||
]
|
||||
vehicle_count = counts.get("gb32960VehicleRealtimeFrames", 0)
|
||||
vehicle_minimum = thresholds.min_gb32960_vehicle_realtime_frames
|
||||
if vehicle_count >= vehicle_minimum:
|
||||
checks.append(Check(
|
||||
"gb32960.vehicle_realtime",
|
||||
"pass",
|
||||
vehicle_count,
|
||||
vehicle_minimum,
|
||||
"GB32960 正式车辆 0x02 实时上报已入 TDengine",
|
||||
))
|
||||
else:
|
||||
status = "fail" if thresholds.require_gb32960_vehicle_realtime else "warn"
|
||||
checks.append(Check(
|
||||
"gb32960.vehicle_realtime",
|
||||
status,
|
||||
vehicle_count,
|
||||
vehicle_minimum,
|
||||
"尚未收到正式车辆 0x02;平台登录在线不等于车辆数据在线",
|
||||
))
|
||||
return checks
|
||||
|
||||
|
||||
def threshold_check(name: str, count: int, minimum: int, ok_message: str, bad_message: str) -> Check:
|
||||
if count >= minimum:
|
||||
return Check(name, "pass", count, minimum, ok_message)
|
||||
return Check(name, "fail", count, minimum, bad_message)
|
||||
|
||||
|
||||
def overall_status(checks: list[Check]) -> str:
|
||||
statuses = {check.status for check in checks}
|
||||
if "fail" in statuses:
|
||||
return "fail"
|
||||
if "warn" in statuses:
|
||||
return "warn"
|
||||
return "pass"
|
||||
|
||||
|
||||
def query_json(rest_url: str, username: str, password: str, sql: str, timeout: float) -> dict[str, Any]:
|
||||
request = urllib.request.Request(rest_url, data=sql.encode("utf-8"), method="POST")
|
||||
if username:
|
||||
token = base64.b64encode(f"{username}:{password}".encode("utf-8")).decode("ascii")
|
||||
request.add_header("Authorization", "Basic " + token)
|
||||
request.add_header("Content-Type", "text/plain; charset=utf-8")
|
||||
request.add_header("Accept", "application/json")
|
||||
with urllib.request.urlopen(request, timeout=timeout) as response:
|
||||
return json.loads(response.read().decode("utf-8"))
|
||||
|
||||
|
||||
def first_cell(response: dict[str, Any]) -> Any:
|
||||
if response.get("code") != 0:
|
||||
raise RuntimeError(f"tdengine query failed: {response}")
|
||||
data = response.get("data") or []
|
||||
if not data or not data[0]:
|
||||
return None
|
||||
return data[0][0]
|
||||
|
||||
|
||||
def query_counts(args: argparse.Namespace) -> dict[str, int]:
|
||||
queries = {
|
||||
"jt808RawFrames": jt808_raw_count_sql(args.date_from, args.date_to, args.jt808_peer_like),
|
||||
"jt808Locations": jt808_location_count_sql(args.date_from, args.date_to),
|
||||
"gb32960PlatformLogins": gb32960_platform_login_count_sql(
|
||||
args.date_from, args.date_to, args.gb32960_peer_like),
|
||||
"gb32960VehicleRealtimeFrames": gb32960_vehicle_realtime_count_sql(
|
||||
args.date_from, args.date_to, args.gb32960_peer_like),
|
||||
}
|
||||
out: dict[str, int] = {}
|
||||
for key, sql in queries.items():
|
||||
out[key] = tdengine_smoke.query_count(
|
||||
args.tdengine_rest_url,
|
||||
args.tdengine_username,
|
||||
args.tdengine_password,
|
||||
sql,
|
||||
args.http_timeout,
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
def query_latest_platform_raw_uri(args: argparse.Namespace) -> str:
|
||||
sql = latest_gb32960_platform_raw_uri_sql(args.date_from, args.date_to, args.gb32960_peer_like)
|
||||
response = query_json(
|
||||
args.tdengine_rest_url,
|
||||
args.tdengine_username,
|
||||
args.tdengine_password,
|
||||
sql,
|
||||
args.http_timeout,
|
||||
)
|
||||
value = first_cell(response)
|
||||
return value if isinstance(value, str) else ""
|
||||
|
||||
|
||||
def decode_frame(history_base_url: str, raw_uri: str, timeout: float) -> dict[str, Any]:
|
||||
if not raw_uri:
|
||||
return {}
|
||||
params = urllib.parse.urlencode({"rawArchiveUri": raw_uri})
|
||||
request = urllib.request.Request(
|
||||
history_base_url.rstrip("/") + "/api/event-history/gb32960/frame?" + params,
|
||||
headers={"Accept": "application/json"},
|
||||
)
|
||||
with urllib.request.urlopen(request, timeout=timeout) as response:
|
||||
return json.loads(response.read().decode("utf-8"))
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> dict[str, Any]:
|
||||
counts = query_counts(args)
|
||||
thresholds = Thresholds(
|
||||
min_jt808_raw_frames=args.min_jt808_raw_frames,
|
||||
min_jt808_locations=args.min_jt808_locations,
|
||||
min_gb32960_platform_logins=args.min_gb32960_platform_logins,
|
||||
min_gb32960_vehicle_realtime_frames=args.min_gb32960_vehicle_realtime_frames,
|
||||
require_gb32960_vehicle_realtime=args.require_gb32960_vehicle_realtime,
|
||||
)
|
||||
checks = evaluate_counts(counts, thresholds)
|
||||
latest_raw_uri = ""
|
||||
decoded_platform_frame: dict[str, Any] = {}
|
||||
if args.history_base_url:
|
||||
latest_raw_uri = query_latest_platform_raw_uri(args)
|
||||
decoded_platform_frame = decode_frame(args.history_base_url, latest_raw_uri, args.http_timeout)
|
||||
return {
|
||||
"status": overall_status(checks),
|
||||
"dateFrom": args.date_from,
|
||||
"dateTo": args.date_to,
|
||||
"counts": counts,
|
||||
"checks": [asdict(check) for check in checks],
|
||||
"gb32960LatestPlatformRawUri": latest_raw_uri,
|
||||
"gb32960LatestPlatformCommand": decoded_platform_frame.get("command"),
|
||||
}
|
||||
|
||||
|
||||
def parser() -> argparse.ArgumentParser:
|
||||
p = argparse.ArgumentParser(description=__doc__)
|
||||
p.add_argument("--tdengine-rest-url", default=os.environ.get("TDENGINE_REST_URL", ""))
|
||||
p.add_argument("--tdengine-username", default=os.environ.get("TDENGINE_USERNAME", "root"))
|
||||
p.add_argument("--tdengine-password", default=os.environ.get("TDENGINE_PASSWORD", "taosdata"))
|
||||
p.add_argument("--history-base-url", default=os.environ.get("HISTORY_BASE_URL", "http://127.0.0.1:20200"))
|
||||
p.add_argument("--date-from", required=True)
|
||||
p.add_argument("--date-to", required=True)
|
||||
p.add_argument("--jt808-peer-like", default=os.environ.get("JT808_PEER_LIKE", ""))
|
||||
p.add_argument("--gb32960-peer-like", default=os.environ.get("GB32960_PEER_LIKE", ""))
|
||||
p.add_argument("--min-jt808-raw-frames", type=int, default=1)
|
||||
p.add_argument("--min-jt808-locations", type=int, default=1)
|
||||
p.add_argument("--min-gb32960-platform-logins", type=int, default=1)
|
||||
p.add_argument("--min-gb32960-vehicle-realtime-frames", type=int, default=1)
|
||||
p.add_argument("--require-gb32960-vehicle-realtime", action="store_true")
|
||||
p.add_argument("--http-timeout", type=float, default=5)
|
||||
return p
|
||||
|
||||
|
||||
def main() -> None:
|
||||
args = parser().parse_args()
|
||||
if not args.tdengine_rest_url:
|
||||
raise SystemExit("--tdengine-rest-url is required")
|
||||
result = run(args)
|
||||
print(json.dumps(result, ensure_ascii=False, indent=2))
|
||||
if result["status"] == "fail":
|
||||
raise SystemExit(1)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user