Files
lingniu-vehicle-ingest/tools/vehicle_ingest_live_verify.py
lingniu 2a4150adc9
Some checks failed
ci/woodpecker/push/woodpecker Pipeline was canceled
fix: verify jt808 locations from current tdengine table
2026-06-29 21:38:36 +08:00

393 lines
14 KiB
Python
Executable File
Raw Permalink 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
"""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("vehicle_locations", ["protocol = 'JT808'"], 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_pattern(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_pattern(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_pattern(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 gb32960_command_count_sql(date_from: str, date_to: str, peer_like: str = "") -> str:
filters = ["protocol = 'GB32960'"]
if peer_like:
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like_pattern(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 message_id,COUNT(*) FROM raw_frames WHERE "
+ " AND ".join(where)
+ " GROUP BY message_id ORDER BY message_id"
)
def gb32960_recent_frame_sql(date_from: str, date_to: str, peer_like: str = "", limit: int = 10) -> str:
filters = ["protocol = 'GB32960'"]
if peer_like:
filters.append("peer LIKE " + tdengine_smoke.sql_literal(peer_like_pattern(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))
safe_limit = max(1, min(limit, 100))
return (
"SELECT ts,message_id,vin,vehicle_key,peer,raw_uri FROM raw_frames WHERE "
+ " AND ".join(where)
+ " ORDER BY ts DESC LIMIT "
+ str(safe_limit)
)
def peer_like_pattern(value: str) -> str:
value = value.strip()
if "%" in value or "_" in value:
return value
return "%" + value + "%"
def command_name(message_id: int) -> str:
return {
1: "VEHICLE_LOGIN",
2: "REALTIME_REPORT",
3: "RESEND_REPORT",
4: "VEHICLE_LOGOUT",
5: "PLATFORM_LOGIN",
6: "PLATFORM_LOGOUT",
7: "HEARTBEAT",
8: "TIME_CALIBRATION",
}.get(message_id, "UNKNOWN_0x" + format(message_id, "02X"))
def response_rows(response: dict[str, Any]) -> list[list[Any]]:
if response.get("code") != 0:
raise RuntimeError(f"tdengine query failed: {response}")
return response.get("data") or []
def command_count_rows(response: dict[str, Any]) -> list[dict[str, Any]]:
rows = []
for item in response_rows(response):
message_id = int(item[0])
rows.append({
"messageId": message_id,
"command": command_name(message_id),
"count": int(item[1]),
})
return rows
def recent_frame_rows(response: dict[str, Any]) -> list[dict[str, Any]]:
rows = []
for item in response_rows(response):
message_id = int(item[1])
rows.append({
"ts": item[0],
"messageId": message_id,
"command": command_name(message_id),
"vin": item[2],
"vehicleKey": item[3],
"peer": item[4],
"rawUri": item[5],
})
return rows
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 query_gb32960_command_counts(args: argparse.Namespace) -> list[dict[str, Any]]:
response = query_json(
args.tdengine_rest_url,
args.tdengine_username,
args.tdengine_password,
gb32960_command_count_sql(args.date_from, args.date_to, args.gb32960_peer_like),
args.http_timeout,
)
return command_count_rows(response)
def query_gb32960_recent_frames(args: argparse.Namespace) -> list[dict[str, Any]]:
response = query_json(
args.tdengine_rest_url,
args.tdengine_username,
args.tdengine_password,
gb32960_recent_frame_sql(
args.date_from,
args.date_to,
args.gb32960_peer_like,
args.gb32960_recent_limit),
args.http_timeout,
)
return recent_frame_rows(response)
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] = {}
command_counts = query_gb32960_command_counts(args)
recent_frames = query_gb32960_recent_frames(args)
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],
"gb32960CommandCounts": command_counts,
"gb32960RecentFrames": recent_frames,
"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("--gb32960-recent-limit", type=int, default=10)
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()