Files
lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime/kv.go

484 lines
12 KiB
Go

package realtime
import (
"encoding/json"
"fmt"
"sort"
"strconv"
"strings"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
type RealtimeKVField struct {
Protocol envelope.Protocol
VIN string
Domain string
Field string
Value string
ValueType string
EventTimeMS int64
ReceivedAtMS int64
EventID string
}
func BuildFieldsEnvelope(env envelope.FrameEnvelope) (envelope.FrameEnvelope, bool) {
if !envelope.IsRealtimeTelemetryFrame(env) {
return envelope.FrameEnvelope{}, false
}
fields, fieldTypes, ok := ParsedFieldsForEnvelope(env)
if !ok {
return envelope.FrameEnvelope{}, false
}
statFields := fieldsWithCanonicalStats(fields, env.Fields)
out := envelope.FrameEnvelope{
EventID: env.StableEventID() + ":fields",
TraceID: env.TraceID,
Protocol: env.Protocol,
MessageID: env.MessageID,
Sequence: env.Sequence,
VIN: env.VIN,
VehicleKeyHint: env.VehicleKeyHint,
Phone: env.Phone,
DeviceID: env.DeviceID,
Plate: env.Plate,
SourceEndpoint: env.SourceEndpoint,
EventTimeMS: env.EventTimeMS,
ReceivedAtMS: env.ReceivedAtMS,
Fields: statFields,
ParsedFields: fields,
ParsedFieldTypes: fieldTypes,
Parsed: map[string]any{
"source_event_id": env.StableEventID(),
"field_mapping": realtimeFieldMappingVersion,
},
ParseStatus: envelope.ParseOK,
}
return out, true
}
func fieldsWithCanonicalStats(mapped map[string]any, canonical map[string]any) map[string]any {
out := cloneAnyMap(mapped)
for _, key := range []string{
envelope.FieldTotalMileageKM,
envelope.FieldSpeedKMH,
envelope.FieldSOCPercent,
envelope.FieldLongitude,
envelope.FieldLatitude,
} {
if value, ok := canonical[key]; ok && value != nil {
out[key] = value
}
}
return out
}
func EnsureParsedFields(env *envelope.FrameEnvelope) bool {
if env == nil {
return false
}
if len(env.ParsedFields) > 0 {
if env.ParsedFieldTypes == nil {
env.ParsedFieldTypes = inferParsedFieldTypes(env.ParsedFields)
}
return true
}
fields, fieldTypes, ok := ParsedFieldsForEnvelope(*env)
if !ok {
return false
}
env.ParsedFields = fields
env.ParsedFieldTypes = fieldTypes
return true
}
func ParsedFieldsForEnvelope(env envelope.FrameEnvelope) (map[string]any, map[string]string, bool) {
if len(env.ParsedFields) > 0 {
fields := cloneAnyMap(env.ParsedFields)
fieldTypes := cloneStringMap(env.ParsedFieldTypes)
if len(fieldTypes) == 0 {
fieldTypes = inferParsedFieldTypes(fields)
}
return fields, fieldTypes, true
}
rows := realtimeKVFields(env, env.Parsed)
if len(rows) == 0 {
return nil, nil, false
}
fields := make(map[string]any, len(rows))
fieldTypes := make(map[string]string, len(rows))
for _, row := range rows {
key := realtimeKVFieldPath(row.Domain, row.Field)
if key == "" {
continue
}
fields[key] = row.Value
fieldTypes[key] = row.ValueType
}
if len(fields) == 0 {
return nil, nil, false
}
return fields, fieldTypes, true
}
func inferParsedFieldTypes(fields map[string]any) map[string]string {
if len(fields) == 0 {
return nil
}
out := make(map[string]string, len(fields))
for key, value := range fields {
_, valueType, ok := stringifyKVValue(value)
if !ok {
continue
}
out[key] = valueType
}
return out
}
func cloneAnyMap(in map[string]any) map[string]any {
if len(in) == 0 {
return nil
}
out := make(map[string]any, len(in))
for key, value := range in {
out[key] = value
}
return out
}
func cloneStringMap(in map[string]string) map[string]string {
if len(in) == 0 {
return nil
}
out := make(map[string]string, len(in))
for key, value := range in {
out[key] = value
}
return out
}
type protocolFieldMapping struct {
Prefix string
GBDataUnits bool
TopLevelName map[string]string
ArrayName map[string]string
}
var realtimeFieldMappings = map[envelope.Protocol]protocolFieldMapping{
envelope.ProtocolGB32960: {
Prefix: "gb32960",
GBDataUnits: true,
TopLevelName: map[string]string{
"header": "header",
"platform_name": "platform",
"platform_login": "platform_login",
},
ArrayName: map[string]string{
"motors": "motor",
"probes": "probe",
"subsystems": "subsystem",
"summaries": "",
},
},
envelope.ProtocolJT808: {
Prefix: "jt808",
TopLevelName: map[string]string{
"header": "header",
"location": "location",
"registration": "registration",
"authentication": "authentication",
},
ArrayName: map[string]string{
"additional": "additional",
},
},
envelope.ProtocolYutongMQTT: {
Prefix: "yutong_mqtt",
TopLevelName: map[string]string{
"data": "data",
"root": "root",
"endpoint": "metadata",
"topic": "metadata",
},
},
}
const realtimeFieldMappingVersion = "2026-07-03.v1"
func realtimeKVFields(env envelope.FrameEnvelope, parsed map[string]any) []RealtimeKVField {
vin := strings.TrimSpace(env.VIN)
if vin == "" {
return nil
}
eventID := env.StableEventID()
mapping := realtimeMapping(env.Protocol)
add := func(rows []RealtimeKVField, domain string, fields any) []RealtimeKVField {
flat := map[string]any{}
flattenKV("", fields, flat, mapping)
names := make([]string, 0, len(flat))
for name := range flat {
names = append(names, name)
}
sort.Strings(names)
for _, name := range names {
value, valueType, ok := stringifyKVValue(flat[name])
if !ok {
continue
}
rows = append(rows, RealtimeKVField{
Protocol: env.Protocol,
VIN: vin,
Domain: domain,
Field: name,
Value: value,
ValueType: valueType,
EventTimeMS: env.EventTimeMS,
ReceivedAtMS: env.ReceivedAtMS,
EventID: eventID,
})
}
return rows
}
var rows []RealtimeKVField
if mapping.GBDataUnits {
rows = append(rows, gb32960KVFields(env, parsed, mapping, eventID)...)
return rows
}
keys := make([]string, 0, len(parsed))
for key := range parsed {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
domain := mappedTopLevelDomain(mapping, key)
if domain == "" {
continue
}
rows = addMappedValue(rows, add, domain, key, parsed[key])
}
return rows
}
func gb32960KVFields(env envelope.FrameEnvelope, parsed map[string]any, mapping protocolFieldMapping, eventID string) []RealtimeKVField {
addRows := func(rows []RealtimeKVField, domain string, fields any) []RealtimeKVField {
flat := map[string]any{}
flattenKV("", fields, flat, mapping)
names := make([]string, 0, len(flat))
for name := range flat {
names = append(names, name)
}
sort.Strings(names)
for _, name := range names {
value, valueType, ok := stringifyKVValue(flat[name])
if !ok {
continue
}
rows = append(rows, RealtimeKVField{
Protocol: env.Protocol,
VIN: strings.TrimSpace(env.VIN),
Domain: domain,
Field: name,
Value: value,
ValueType: valueType,
EventTimeMS: env.EventTimeMS,
ReceivedAtMS: env.ReceivedAtMS,
EventID: eventID,
})
}
return rows
}
var rows []RealtimeKVField
for key, value := range parsed {
if key == "data_units" {
continue
}
if domain := mappedTopLevelDomain(mapping, key); domain != "" {
rows = addMappedValue(rows, addRows, domain, key, value)
}
}
units, ok := asAnySlice(parsed["data_units"])
if !ok {
return rows
}
for _, unit := range units {
unitMap, ok := unit.(map[string]any)
if !ok {
continue
}
name := sanitizeFieldPart(strconvAny(unitMap["name"]))
if name == "" || name == "<nil>" {
name = strings.ToLower(strings.TrimPrefix(strings.TrimSpace(strconvAny(unitMap["type"])), "0x"))
}
if name == "" {
continue
}
domain := mapping.Prefix + "." + name
if value, ok := unitMap["value"]; ok {
rows = addRows(rows, domain, value)
}
if unitType := strings.TrimSpace(strconvAny(unitMap["type"])); unitType != "" && unitType != "<nil>" {
rows = addRows(rows, domain, map[string]any{"type": unitType})
}
}
return rows
}
func flattenKV(prefix string, value any, out map[string]any, mapping protocolFieldMapping) {
switch typed := value.(type) {
case map[string]any:
keys := make([]string, 0, len(typed))
for key := range typed {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
part := sanitizeFieldPart(key)
if part == "" {
continue
}
if arrayName, ok := mapping.ArrayName[part]; ok && arrayName == "" {
if items, ok := asAnySlice(typed[key]); ok && len(items) == 1 {
flattenKV(prefix, items[0], out, mapping)
continue
}
}
next := part
if prefix != "" {
next = prefix + "." + part
}
flattenKV(next, typed[key], out, mapping)
}
case []any:
for index, item := range typed {
itemMap, ok := item.(map[string]any)
if !ok {
out[prefix] = typed
return
}
arrayName, hasArrayName := mapping.ArrayName[lastPathPart(prefix)]
itemKey := strconv.Itoa(index)
if hasArrayName {
if arrayName == "" && len(typed) == 1 {
flattenKV(prefix, itemMap, out, mapping)
continue
}
if arrayName != "" {
itemKey = arrayName + "_" + strconv.Itoa(index+1)
}
}
next := itemKey
if prefix != "" {
next = prefix + "." + itemKey
}
flattenKV(next, itemMap, out, mapping)
}
case []map[string]any:
items := make([]any, len(typed))
for index, item := range typed {
items[index] = item
}
flattenKV(prefix, items, out, mapping)
default:
if prefix != "" && value != nil {
out[prefix] = value
}
}
}
func addMappedValue(rows []RealtimeKVField, add func([]RealtimeKVField, string, any) []RealtimeKVField, domain string, key string, value any) []RealtimeKVField {
switch typed := value.(type) {
case map[string]any:
return add(rows, domain, typed)
case []any:
return add(rows, domain, typed)
case []map[string]any:
return add(rows, domain, typed)
default:
return add(rows, domain, map[string]any{key: value})
}
}
func realtimeMapping(protocol envelope.Protocol) protocolFieldMapping {
if mapping, ok := realtimeFieldMappings[protocol]; ok {
return mapping
}
return protocolFieldMapping{Prefix: strings.ToLower(string(protocol))}
}
func mappedTopLevelDomain(mapping protocolFieldMapping, key string) string {
key = strings.TrimSpace(key)
if key == "" {
return ""
}
name := mapping.TopLevelName[key]
if name == "" {
name = sanitizeFieldPart(key)
}
if name == "" {
return ""
}
return mapping.Prefix + "." + name
}
func sanitizeFieldPart(value string) string {
value = strings.TrimSpace(value)
if value == "" {
return ""
}
var b strings.Builder
for _, r := range value {
switch {
case r >= 'a' && r <= 'z':
b.WriteRune(r)
case r >= 'A' && r <= 'Z':
b.WriteRune(r + ('a' - 'A'))
case r >= '0' && r <= '9':
b.WriteRune(r)
case r == '_':
b.WriteRune(r)
default:
b.WriteByte('_')
}
}
return strings.Trim(b.String(), "_")
}
func lastPathPart(path string) string {
if index := strings.LastIndex(path, "."); index >= 0 {
return path[index+1:]
}
return path
}
func stringifyKVValue(value any) (string, string, bool) {
switch typed := value.(type) {
case string:
text := strings.TrimSpace(typed)
return text, "string", text != ""
case bool:
return strconv.FormatBool(typed), "bool", true
case float64:
return strconv.FormatFloat(typed, 'f', -1, 64), "number", true
case float32:
return strconv.FormatFloat(float64(typed), 'f', -1, 64), "number", true
case json.Number:
return typed.String(), "number", strings.TrimSpace(typed.String()) != ""
case int:
return strconv.Itoa(typed), "number", true
case int8, int16, int32, int64:
return fmt.Sprintf("%d", typed), "number", true
case uint, uint8, uint16, uint32, uint64:
return fmt.Sprintf("%d", typed), "number", true
default:
data, err := json.Marshal(typed)
if err != nil || string(data) == "null" {
return "", "", false
}
return string(data), "json", true
}
}