From c0d3b4155058787b99c2d93a7a39e3f8928e5169 Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 14:49:53 +0800 Subject: [PATCH] feat(go): filter raw parsed fields --- go/vehicle-gateway/internal/history/query.go | 53 +++++++++++++++++++ .../internal/history/query_test.go | 40 ++++++++++++++ 2 files changed, 93 insertions(+) diff --git a/go/vehicle-gateway/internal/history/query.go b/go/vehicle-gateway/internal/history/query.go index c30c851d..9f2faabf 100644 --- a/go/vehicle-gateway/internal/history/query.go +++ b/go/vehicle-gateway/internal/history/query.go @@ -28,6 +28,7 @@ type RawFrameQuery struct { IncludeFields bool IncludePayload bool IncludeTotal bool + ParsedFields []string DateFrom string DateTo string Limit int @@ -165,6 +166,7 @@ func (r *RawFrameRepository) Query(ctx context.Context, query RawFrameQuery) ([] if err := r.hydratePayloadChunks(ctx, out); err != nil { return nil, err } + filterParsedFields(out, query.ParsedFields) return out, nil } @@ -297,6 +299,10 @@ func normalizeRawFrameQuery(query RawFrameQuery) RawFrameQuery { query.Phone = strings.TrimSpace(query.Phone) query.DeviceID = strings.TrimSpace(query.DeviceID) query.MessageID = strings.TrimSpace(query.MessageID) + query.ParsedFields = normalizeParsedFieldNames(query.ParsedFields) + if len(query.ParsedFields) > 0 { + query.IncludeFields = true + } query.OrderBy = normalizeRawFrameOrderBy(query.OrderBy) query.DateFrom = normalizeDateTimeLiteral(query.DateFrom) query.DateTo = normalizeDateTimeLiteral(query.DateTo) @@ -457,6 +463,52 @@ func rawJSONMessage(value string) json.RawMessage { return json.RawMessage(jsonString(value)) } +func filterParsedFields(rows []RawFrameRow, fieldNames []string) { + fieldNames = normalizeParsedFieldNames(fieldNames) + if len(fieldNames) == 0 { + return + } + for index := range rows { + if len(rows[index].ParsedFields) == 0 { + continue + } + var fields map[string]any + if err := json.Unmarshal(rows[index].ParsedFields, &fields); err != nil { + continue + } + filtered := make(map[string]any, len(fieldNames)) + for _, name := range fieldNames { + if value, ok := fields[name]; ok { + filtered[name] = value + } + } + if len(filtered) == 0 { + rows[index].ParsedFields = json.RawMessage(`{}`) + continue + } + rows[index].ParsedFields = rawJSONMessage(jsonString(filtered)) + } +} + +func normalizeParsedFieldNames(values []string) []string { + out := make([]string, 0, len(values)) + seen := map[string]struct{}{} + for _, value := range values { + for _, part := range strings.Split(value, ",") { + name := strings.TrimSpace(part) + if name == "" { + continue + } + if _, ok := seen[name]; ok { + continue + } + seen[name] = struct{}{} + out = append(out, name) + } + } + return out +} + func payloadKindMatches(candidate string, manifest string) bool { if candidate == manifest { return true @@ -704,6 +756,7 @@ func parseRawFrameQuery(r *http.Request) (RawFrameQuery, error) { IncludeFields: strings.EqualFold(strings.TrimSpace(values.Get("includeFields")), "true"), IncludePayload: strings.EqualFold(strings.TrimSpace(values.Get("includePayload")), "true"), IncludeTotal: strings.EqualFold(strings.TrimSpace(values.Get("includeTotal")), "true"), + ParsedFields: append(values["fields"], values["parsedFields"]...), 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 21e1c895..0cb8bc0d 100644 --- a/go/vehicle-gateway/internal/history/query_test.go +++ b/go/vehicle-gateway/internal/history/query_test.go @@ -146,6 +146,46 @@ func TestRawFrameHandlerReturnsRawFrames(t *testing.T) { } } +func TestRawFrameHandlerFiltersParsedFields(t *testing.T) { + db, mock, err := sqlmock.New() + if err != nil { + t.Fatalf("sqlmock.New() error = %v", err) + } + defer db.Close() + mock.ExpectQuery("SELECT ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, .* AS parsed_fields, parse_status, parse_error, source_endpoint, protocol, vehicle_key, vin, phone, device_id FROM lingniu_vehicle_ts.raw_gb32960_"). + WillReturnRows(sqlmock.NewRows([]string{ + "ts", "frame_id", "event_id", "message_id", "event_time", "received_at", "raw_size_bytes", + "raw_hex", "raw_text", "parsed_fields", "parse_status", "parse_error", "source_endpoint", + "protocol", "vehicle_key", "vin", "phone", "device_id", + }).AddRow( + "2026-07-01 22:28:25", "go_frame", "event-2", 2, "2026-07-01 22:28:25", "2026-07-01 22:28:25", + 128, "2323", "", `{"gb32960.vehicle.soc_percent":"85","gb32960.vehicle.speed_kmh":"30","gb32960.vehicle.total_mileage_km":"10000"}`, + "OK", "", "8.134.95.166:53702", "GB32960", "LB9A32A21R0LS1707", "LB9A32A21R0LS1707", "", "", + )) + + handler := NewRawFrameHandler(NewRawFrameRepository(db, "lingniu_vehicle_ts")) + request := httptest.NewRequest(http.MethodGet, "/api/history/raw-frames?protocol=GB32960&vin=LB9A32A21R0LS1707&limit=1&fields=gb32960.vehicle.soc_percent,gb32960.vehicle.speed_kmh", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", response.Code, response.Body.String()) + } + body := response.Body.String() + for _, want := range []string{`"gb32960.vehicle.soc_percent":"85"`, `"gb32960.vehicle.speed_kmh":"30"`} { + if !strings.Contains(body, want) { + t.Fatalf("response missing %s: %s", want, body) + } + } + if strings.Contains(body, "gb32960.vehicle.total_mileage_km") { + t.Fatalf("response should filter unrequested parsed field: %s", body) + } + if err := mock.ExpectationsWereMet(); err != nil { + t.Fatalf("sql expectations: %v", err) + } +} + func TestRawFrameHandlerFiltersByVehicleKey(t *testing.T) { db, mock, err := sqlmock.New() if err != nil {