fix(go): store protocol parsed fields in realtime kv

This commit is contained in:
lingniu
2026-07-03 12:08:12 +08:00
parent 7bf78f91ba
commit 112b02e8c5
3 changed files with 83 additions and 94 deletions

View File

@@ -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:

View File

@@ -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 {

View File

@@ -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 {