From 318486b1e98d59dc471c0f8519673ab73f443f08 Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 09:16:45 +0800 Subject: [PATCH] fix: verify raw freshness by received time --- go/vehicle-gateway/internal/history/query.go | 22 ++++++- .../internal/history/query_test.go | 14 +++++ tools/go_native_prod_smoke.py | 57 ++++++++++++------- tools/test_go_native_prod_smoke.py | 57 +++++++++++++++++++ 4 files changed, 130 insertions(+), 20 deletions(-) diff --git a/go/vehicle-gateway/internal/history/query.go b/go/vehicle-gateway/internal/history/query.go index e68e75e2..d3a0173c 100644 --- a/go/vehicle-gateway/internal/history/query.go +++ b/go/vehicle-gateway/internal/history/query.go @@ -23,6 +23,7 @@ type RawFrameQuery struct { Phone string DeviceID string MessageID string + OrderBy string DateFrom string DateTo string Limit int @@ -366,6 +367,7 @@ func normalizeRawFrameQuery(query RawFrameQuery) RawFrameQuery { query.Phone = strings.TrimSpace(query.Phone) query.DeviceID = strings.TrimSpace(query.DeviceID) query.MessageID = strings.TrimSpace(query.MessageID) + query.OrderBy = normalizeRawFrameOrderBy(query.OrderBy) query.DateFrom = normalizeDateTimeLiteral(query.DateFrom) query.DateTo = normalizeDateTimeLiteral(query.DateTo) if query.Limit <= 0 { @@ -408,7 +410,7 @@ func buildRawFrameSQL(table string, query RawFrameQuery) (string, []any) { if len(where) > 0 { sqlText += " WHERE " + strings.Join(where, " AND ") } - sqlText += " ORDER BY ts DESC LIMIT " + strconv.Itoa(query.Limit) + " OFFSET " + strconv.Itoa(query.Offset) + sqlText += " ORDER BY " + rawFrameOrderColumn(query.OrderBy) + " DESC LIMIT " + strconv.Itoa(query.Limit) + " OFFSET " + strconv.Itoa(query.Offset) return sqlText, nil } @@ -490,6 +492,23 @@ func rawFrameWhere(query RawFrameQuery) []string { return where } +func normalizeRawFrameOrderBy(value string) string { + value = strings.ToLower(strings.TrimSpace(value)) + switch value { + case "receivedat", "received_at": + return "receivedAt" + default: + return "ts" + } +} + +func rawFrameOrderColumn(value string) string { + if normalizeRawFrameOrderBy(value) == "receivedAt" { + return "received_at" + } + return "ts" +} + func locationWhere(query LocationQuery) []string { var where []string add := func(clause string) { @@ -697,6 +716,7 @@ func parseRawFrameQuery(r *http.Request) (RawFrameQuery, error) { Phone: values.Get("phone"), DeviceID: values.Get("deviceId"), MessageID: values.Get("messageId"), + OrderBy: values.Get("orderBy"), DateFrom: values.Get("dateFrom"), DateTo: values.Get("dateTo"), Limit: limit, diff --git a/go/vehicle-gateway/internal/history/query_test.go b/go/vehicle-gateway/internal/history/query_test.go index e59b98ff..f4edcab6 100644 --- a/go/vehicle-gateway/internal/history/query_test.go +++ b/go/vehicle-gateway/internal/history/query_test.go @@ -370,3 +370,17 @@ func TestBuildRawFrameSQLUsesLiteralsForTDengine(t *testing.T) { } } } + +func TestBuildRawFrameSQLCanOrderByReceivedAt(t *testing.T) { + sqlText, args := buildRawFrameSQL("lingniu_vehicle_ts.raw_frames", RawFrameQuery{ + Protocol: "JT808", + OrderBy: "receivedAt", + Limit: 1, + }) + if len(args) != 0 { + t.Fatalf("expected no query args for TDengine, got %#v", args) + } + if !strings.Contains(sqlText, "ORDER BY received_at DESC LIMIT 1 OFFSET 0") { + t.Fatalf("sql should order by received_at: %s", sqlText) + } +} diff --git a/tools/go_native_prod_smoke.py b/tools/go_native_prod_smoke.py index b6675cfa..1bfb3e66 100755 --- a/tools/go_native_prod_smoke.py +++ b/tools/go_native_prod_smoke.py @@ -40,6 +40,14 @@ class Check: message: str +@dataclass(frozen=True) +class CheckWindow: + date_from: str + date_to: str + stat_date: str + max_raw_age_minutes: float | None + + def shanghai_day_window(now: dt.datetime | None = None) -> tuple[str, str]: now = now or dt.datetime.now(tz=SHANGHAI) local = now.astimezone(SHANGHAI) @@ -48,6 +56,23 @@ def shanghai_day_window(now: dt.datetime | None = None) -> tuple[str, str]: return start.isoformat(timespec="seconds"), end.isoformat(timespec="seconds") +def resolve_check_window( + raw_date: str | None, + max_raw_age_minutes: float, + now: dt.datetime | None = None, +) -> CheckWindow: + current = (now or dt.datetime.now(tz=SHANGHAI)).astimezone(SHANGHAI) + today = current.date() + if raw_date: + day = dt.date.fromisoformat(raw_date) + start = dt.datetime(day.year, day.month, day.day, tzinfo=SHANGHAI) + date_from, date_to = shanghai_day_window(start) + freshness = max_raw_age_minutes if day == today else None + return CheckWindow(date_from, date_to, raw_date, freshness) + date_from, date_to = shanghai_day_window(current) + return CheckWindow(date_from, date_to, date_from[:10], max_raw_age_minutes) + + def api_url(base_url: str, path: str, params: dict[str, Any]) -> str: url = base_url.rstrip("/") + path if not params: @@ -129,12 +154,15 @@ def parse_tdengine_utc_timestamp(value: str) -> dt.datetime: def latest_age_minutes(items: list[Any], now: dt.datetime | None = None) -> float | None: - if not items or not isinstance(items[0], dict) or not items[0].get("ts"): + if not items or not isinstance(items[0], dict): + return None + timestamp = items[0].get("received_at") or items[0].get("ts") + if not timestamp: return None current = now or dt.datetime.now(tz=dt.timezone.utc) if current.tzinfo is None: current = current.replace(tzinfo=dt.timezone.utc) - latest = parse_tdengine_utc_timestamp(str(items[0]["ts"])) + latest = parse_tdengine_utc_timestamp(str(timestamp)) return max(0.0, (current.astimezone(dt.timezone.utc) - latest).total_seconds() / 60) @@ -291,7 +319,7 @@ def build_check_specs( min_stat: int, max_raw_age_minutes: float | None = None, ) -> list[CheckSpec]: - raw_params = {"dateFrom": date_from, "dateTo": date_to, "limit": 1} + raw_params = {"dateFrom": date_from, "dateTo": date_to, "limit": 1, "orderBy": "receivedAt"} gb32960_history_params = {"protocol": "GB32960", "dateFrom": date_from, "dateTo": date_to, "limit": 1} jt808_history_params = {"protocol": "JT808", "dateFrom": date_from, "dateTo": date_to, "limit": 1} yutong_mqtt_history_params = {"protocol": "YUTONG_MQTT", "dateFrom": date_from, "dateTo": date_to, "limit": 1} @@ -489,24 +517,15 @@ def parse_args(argv: list[str]) -> argparse.Namespace: def main(argv: list[str]) -> int: args = parse_args(argv) - if args.date: - day = dt.date.fromisoformat(args.date) - start = dt.datetime(day.year, day.month, day.day, tzinfo=SHANGHAI) - date_from, date_to = shanghai_day_window(start) - stat_date = args.date - max_raw_age_minutes = None - else: - date_from, date_to = shanghai_day_window() - stat_date = date_from[:10] - max_raw_age_minutes = args.max_raw_age_minutes + window = resolve_check_window(args.date, args.max_raw_age_minutes) specs = build_check_specs( - date_from=date_from, - date_to=date_to, - stat_date=stat_date, + date_from=window.date_from, + date_to=window.date_to, + stat_date=window.stat_date, min_raw=args.min_raw, min_history=args.min_history, min_stat=args.min_stat, - max_raw_age_minutes=max_raw_age_minutes, + max_raw_age_minutes=window.max_raw_age_minutes, ) checks, payloads = query_payloads(args.base_url, specs, args.timeout) checks.extend(run_checks(args.base_url, build_realtime_specs(payloads), args.timeout)) @@ -514,8 +533,8 @@ def main(argv: list[str]) -> int: print(json.dumps({ "status": status, "baseUrl": args.base_url, - "dateFrom": date_from, - "dateTo": date_to, + "dateFrom": window.date_from, + "dateTo": window.date_to, "checks": [asdict(check) for check in checks], }, ensure_ascii=False, indent=2)) return 0 if status == "pass" else 1 diff --git a/tools/test_go_native_prod_smoke.py b/tools/test_go_native_prod_smoke.py index ee65fdd9..e633763c 100644 --- a/tools/test_go_native_prod_smoke.py +++ b/tools/test_go_native_prod_smoke.py @@ -32,6 +32,26 @@ class GoNativeProdSmokeTest(unittest.TestCase): "&dateTo=2026-07-03T00%3A00%3A00%2B08%3A00&limit=1", ) + def test_explicit_today_date_still_enforces_raw_freshness(self): + now = dt.datetime(2026, 7, 2, 9, 30, tzinfo=smoke.SHANGHAI) + + window = smoke.resolve_check_window("2026-07-02", 15.0, now) + + self.assertEqual(window.date_from, "2026-07-02T00:00:00+08:00") + self.assertEqual(window.date_to, "2026-07-03T00:00:00+08:00") + self.assertEqual(window.stat_date, "2026-07-02") + self.assertEqual(window.max_raw_age_minutes, 15.0) + + def test_explicit_historical_date_disables_raw_freshness(self): + now = dt.datetime(2026, 7, 2, 9, 30, tzinfo=smoke.SHANGHAI) + + window = smoke.resolve_check_window("2026-07-01", 15.0, now) + + self.assertEqual(window.date_from, "2026-07-01T00:00:00+08:00") + self.assertEqual(window.date_to, "2026-07-02T00:00:00+08:00") + self.assertEqual(window.stat_date, "2026-07-01") + self.assertIsNone(window.max_raw_age_minutes) + def test_total_check_passes_when_total_reaches_minimum(self): check = smoke.check_total( "jt808.raw", @@ -82,6 +102,23 @@ class GoNativeProdSmokeTest(unittest.TestCase): self.assertIn("gb32960.daily_mileage", names) self.assertIn("jt808.daily_total_mileage", names) + def test_raw_checks_order_by_received_at_for_ingest_freshness(self): + specs = smoke.build_check_specs( + date_from="2026-07-02T00:00:00+08:00", + date_to="2026-07-03T00:00:00+08:00", + stat_date="2026-07-02", + min_raw=1, + min_history=1, + min_stat=1, + max_raw_age_minutes=15, + ) + + raw_specs = [spec for spec in specs if spec.name.endswith(".raw")] + + self.assertTrue(raw_specs) + for spec in raw_specs: + self.assertEqual(spec.params.get("orderBy"), "receivedAt") + def test_vehicle_identifier_prefers_vin_over_vehicle_key(self): identifier = smoke.vehicle_identifier({ "vin": "LB9A32A20R0LS1343", @@ -185,6 +222,26 @@ class GoNativeProdSmokeTest(unittest.TestCase): self.assertEqual(check.status, "fail") self.assertIn("latest_age_minutes=60.0", check.message) + def test_total_check_uses_received_at_for_raw_freshness_when_present(self): + now = dt.datetime(2026, 7, 1, 18, 0, 0, tzinfo=dt.timezone.utc) + + check = smoke.check_total( + "gb32960.raw", + { + "total": 1, + "items": [{ + "ts": "2026-07-01 17:00:00", + "received_at": "2026-07-01 17:58:00", + }], + }, + minimum=1, + max_age_minutes=15, + now=now, + ) + + self.assertEqual(check.status, "pass") + self.assertIn("latest_age_minutes=2.0", check.message) + def test_total_check_passes_when_latest_sample_is_fresh(self): now = dt.datetime(2026, 7, 1, 18, 0, 0, tzinfo=dt.timezone.utc)