diff --git a/go/vehicle-gateway/internal/realtime/kv.go b/go/vehicle-gateway/internal/realtime/kv.go index 3f527ca7..ddc6c1cd 100644 --- a/go/vehicle-gateway/internal/realtime/kv.go +++ b/go/vehicle-gateway/internal/realtime/kv.go @@ -28,7 +28,7 @@ func realtimeKVFields(env envelope.FrameEnvelope, parsed map[string]any) []Realt return nil } eventID := env.StableEventID() - add := func(rows []RealtimeKVField, domain string, fields map[string]any) []RealtimeKVField { + add := func(rows []RealtimeKVField, domain string, fields any) []RealtimeKVField { flat := map[string]any{} flattenKV("", fields, flat) names := make([]string, 0, len(flat)) @@ -57,70 +57,35 @@ func realtimeKVFields(env envelope.FrameEnvelope, parsed map[string]any) []Realt } var rows []RealtimeKVField - switch env.Protocol { - case envelope.ProtocolGB32960: - units, ok := asAnySlice(parsed["data_units"]) - if !ok { - return rows + keys := make([]string, 0, len(parsed)) + for key := range parsed { + keys = append(keys, key) + } + sort.Strings(keys) + metadata := map[string]any{} + for _, key := range keys { + domain := strings.TrimSpace(key) + if domain == "" { + continue } - for _, unit := range units { - unitMap, ok := unit.(map[string]any) - if !ok { - continue - } - domain := strings.TrimSpace(strconvAny(unitMap["name"])) - if domain == "" { - domain = strings.TrimSpace(strconvAny(unitMap["type"])) - } - if domain == "" { - continue - } - value, ok := unitMap["value"].(map[string]any) - if !ok { - continue - } - rows = add(rows, domain, domainKVFields(domain, value)) + value := parsed[key] + switch typed := value.(type) { + case map[string]any: + rows = add(rows, domain, typed) + case []any: + rows = add(rows, domain, typed) + case []map[string]any: + rows = add(rows, domain, typed) + default: + metadata[domain] = value } - case envelope.ProtocolJT808: - fields := cloneFields(env.Fields) - if location, ok := parsed["location"].(map[string]any); ok { - for key, value := range location { - fields[key] = value - } - } - rows = add(rows, "location", fields) - case envelope.ProtocolYutongMQTT: - fields := cloneFields(env.Fields) - if data, ok := parsed["data"].(map[string]any); ok { - for key, value := range data { - fields[key] = value - } - } - rows = add(rows, "vehicle", fields) + } + if len(metadata) > 0 { + rows = add(rows, "metadata", metadata) } return rows } -func domainKVFields(domain string, value map[string]any) map[string]any { - out := cloneMap(value) - if domain != "gd_fc_stack" { - return out - } - summaries, ok := asAnySlice(out["summaries"]) - if !ok || len(summaries) != 1 { - return out - } - summary, ok := summaries[0].(map[string]any) - if !ok { - return out - } - delete(out, "summaries") - for key, value := range summary { - out[key] = value - } - return out -} - func flattenKV(prefix string, value any, out map[string]any) { switch typed := value.(type) { case map[string]any: diff --git a/go/vehicle-gateway/internal/realtime/kv_test.go b/go/vehicle-gateway/internal/realtime/kv_test.go index 48e80e42..fe6793d5 100644 --- a/go/vehicle-gateway/internal/realtime/kv_test.go +++ b/go/vehicle-gateway/internal/realtime/kv_test.go @@ -14,6 +14,7 @@ func TestRealtimeKVFieldsFromGB32960ParsedDomains(t *testing.T) { ReceivedAtMS: 1100, EventID: "event-1", }, map[string]any{ + "header": map[string]any{"vin": "VIN001", "command": "0x02"}, "data_units": []any{ map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}}, map[string]any{"type": "0x30", "name": "gd_fc_stack", "value": map[string]any{ @@ -26,21 +27,19 @@ func TestRealtimeKVFieldsFromGB32960ParsedDomains(t *testing.T) { }) values := kvMap(rows) - if values["vehicle/soc_percent"] != "88" { + if values["header/vin"] != "VIN001" || values["header/command"] != "0x02" { + t.Fatalf("gb32960 header kv missing: %#v", values) + } + if values["data_units/0.name"] != "vehicle" || values["data_units/0.type"] != "0x01" || values["data_units/0.value.soc_percent"] != "88" { t.Fatalf("vehicle soc kv missing: %#v", values) } - if values["gd_fc_stack/stack_water_outlet_temp_c"] != "63" { - t.Fatalf("stack temp kv missing: %#v", values) - } - if values["gd_fc_stack/hydrogen_inlet_pressure_kpa"] != "130" { + if values["data_units/1.value.summaries.0.stack_water_outlet_temp_c"] != "63" || + values["data_units/1.value.summaries.0.hydrogen_inlet_pressure_kpa"] != "130" { t.Fatalf("stack pressure kv missing: %#v", values) } - if values["gd_fc_stack/summaries.0.stack_water_outlet_temp_c"] != "" { - t.Fatalf("stack single summary should be flattened onto domain, got %#v", values) - } } -func TestRealtimeKVFieldsFromJT808LocationAndMQTTFields(t *testing.T) { +func TestRealtimeKVFieldsFromJT808ParsedFields(t *testing.T) { jtRows := realtimeKVFields(envelope.FrameEnvelope{ Protocol: envelope.ProtocolJT808, VIN: "VIN001", @@ -49,12 +48,30 @@ func TestRealtimeKVFieldsFromJT808LocationAndMQTTFields(t *testing.T) { envelope.FieldLatitude: 30.2, envelope.FieldTotalMileageKM: 10241.2, }, - }, map[string]any{}) + }, map[string]any{ + "header": map[string]any{"message_id": "0x0200", "phone": "13307795425"}, + "location": map[string]any{ + "longitude": 121.1, + "latitude": 30.2, + "total_mileage_km": 10241.2, + "additional": []any{map[string]any{"id": "0x01", "value_hex": "00077235"}}, + }, + }) jtValues := kvMap(jtRows) - if jtValues["location/total_mileage_km"] != "10241.2" || jtValues["location/longitude"] != "121.1" { + if jtValues["header/message_id"] != "0x0200" || jtValues["header/phone"] != "13307795425" { + t.Fatalf("jt808 header kv missing: %#v", jtValues) + } + if jtValues["location/total_mileage_km"] != "10241.2" || + jtValues["location/longitude"] != "121.1" || + jtValues["location/additional.0.id"] != "0x01" { t.Fatalf("jt808 location kv missing: %#v", jtValues) } + if jtValues["location/soc_percent"] != "" { + t.Fatalf("jt808 kv should not include standardized env.Fields-only values: %#v", jtValues) + } +} +func TestRealtimeKVFieldsFromYutongMQTTParsedFields(t *testing.T) { mqttRows := realtimeKVFields(envelope.FrameEnvelope{ Protocol: envelope.ProtocolYutongMQTT, VIN: "VIN002", @@ -63,6 +80,8 @@ func TestRealtimeKVFieldsFromJT808LocationAndMQTTFields(t *testing.T) { "gear": 3, }, }, map[string]any{ + "endpoint": "yutong", + "topic": "/ytforward/shln/3", "data": map[string]any{ "ACC_PEDAL_APT": 14, "BATTERY_CAPACITY_SOC": 78.4, @@ -71,18 +90,27 @@ func TestRealtimeKVFieldsFromJT808LocationAndMQTTFields(t *testing.T) { "TOTAL_MILEAGE": 119925000, "fuelCellCoolInTempt": 59, }, + "root": map[string]any{ + "device": "LMRKH9AC2R1004087", + "version": "1.0", + }, }) mqttValues := kvMap(mqttRows) - if mqttValues["vehicle/soc_percent"] != "76" || mqttValues["vehicle/gear"] != "3" { - t.Fatalf("mqtt vehicle kv missing: %#v", mqttValues) - } - if mqttValues["vehicle/ACC_PEDAL_APT"] != "14" || - mqttValues["vehicle/CURRENT_OF_FC"] != "54.9" || - mqttValues["vehicle/HYDROGEN_LOW_PRESSURE"] != "1.41" || - mqttValues["vehicle/TOTAL_MILEAGE"] != "119925000" || - mqttValues["vehicle/fuelCellCoolInTempt"] != "59" { + if mqttValues["data/ACC_PEDAL_APT"] != "14" || + mqttValues["data/CURRENT_OF_FC"] != "54.9" || + mqttValues["data/HYDROGEN_LOW_PRESSURE"] != "1.41" || + mqttValues["data/TOTAL_MILEAGE"] != "119925000" || + mqttValues["data/fuelCellCoolInTempt"] != "59" { t.Fatalf("mqtt raw data kv missing: %#v", mqttValues) } + if mqttValues["root/device"] != "LMRKH9AC2R1004087" || + mqttValues["metadata/endpoint"] != "yutong" || + mqttValues["metadata/topic"] != "/ytforward/shln/3" { + t.Fatalf("mqtt metadata kv missing: %#v", mqttValues) + } + if mqttValues["data/soc_percent"] != "" || mqttValues["data/gear"] != "" { + t.Fatalf("mqtt kv should not include standardized env.Fields-only values: %#v", mqttValues) + } } func kvMap(rows []RealtimeKVField) map[string]string { diff --git a/go/vehicle-gateway/internal/realtime/repository_test.go b/go/vehicle-gateway/internal/realtime/repository_test.go index d359f1d5..c53eed8e 100644 --- a/go/vehicle-gateway/internal/realtime/repository_test.go +++ b/go/vehicle-gateway/internal/realtime/repository_test.go @@ -294,21 +294,17 @@ func TestRepositoryWritesRealtimeKVHashes(t *testing.T) { t.Fatalf("Update() error = %v", err) } - vehicleKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Result() + dataUnitsKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:data_units").Result() if err != nil { - t.Fatalf("vehicle kv HGetAll error = %v", err) + t.Fatalf("data_units kv HGetAll error = %v", err) } - if vehicleKV["soc_percent"] != "88" || vehicleKV["_event_id"] == "" || vehicleKV["_event_time_ms"] != "1000" { - t.Fatalf("vehicle kv = %#v", vehicleKV) + if dataUnitsKV["0.value.soc_percent"] != "88" || dataUnitsKV["_event_id"] == "" || dataUnitsKV["_event_time_ms"] != "1000" { + t.Fatalf("data_units kv = %#v", dataUnitsKV) } - stackKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:gd_fc_stack").Result() - if err != nil { - t.Fatalf("stack kv HGetAll error = %v", err) + if dataUnitsKV["1.value.summaries.0.stack_water_outlet_temp_c"] != "63" { + t.Fatalf("data_units stack kv = %#v", dataUnitsKV) } - if stackKV["stack_water_outlet_temp_c"] != "63" { - t.Fatalf("stack kv = %#v", stackKV) - } - if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Val(); ttl != -1 { + if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:data_units").Val(); ttl != -1 { t.Fatalf("kv ttl should not expire, got %v", ttl) } } @@ -332,14 +328,14 @@ func TestRepositoryFastUpdateOnlyWritesPermanentKVAndMinuteOnline(t *testing.T) t.Fatalf("FastUpdate() error = %v", err) } - vehicleKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Result() + dataUnitsKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:data_units").Result() if err != nil { - t.Fatalf("vehicle kv HGetAll error = %v", err) + t.Fatalf("data_units kv HGetAll error = %v", err) } - if vehicleKV["soc_percent"] != "88" || vehicleKV["_event_time_ms"] != "1000" { - t.Fatalf("vehicle kv = %#v", vehicleKV) + if dataUnitsKV["0.value.soc_percent"] != "88" || dataUnitsKV["_event_time_ms"] != "1000" { + t.Fatalf("data_units kv = %#v", dataUnitsKV) } - if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Val(); ttl != -1 { + if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:data_units").Val(); ttl != -1 { t.Fatalf("kv ttl should not expire, got %v", ttl) } if ttl := repo.client.TTL(ctx, "vehicle:online:VIN001").Val(); ttl <= 0 || ttl > time.Minute {