fix(go): flatten jt808 additional fields

This commit is contained in:
lingniu
2026-07-02 16:23:03 +08:00
parent 5b2e1abcdd
commit d0289b1b48
4 changed files with 170 additions and 93 deletions

View File

@@ -30,14 +30,7 @@ type Location struct {
DirectionDeg uint16 DirectionDeg uint16
Time time.Time Time time.Time
TotalMileageKM *float64 TotalMileageKM *float64
Additional []AdditionalItem Additional map[string]any
}
type AdditionalItem struct {
ID byte
Length int
ValueHex string
Parsed any
} }
// ExtractFrames splits JT/T 808 TCP streams by 0x7e delimiters and unescapes // ExtractFrames splits JT/T 808 TCP streams by 0x7e delimiters and unescapes
@@ -311,88 +304,70 @@ func parseLocation(body []byte) (Location, error) {
return location, nil return location, nil
} }
func parseAdditional(data []byte, location *Location) []AdditionalItem { func parseAdditional(data []byte, location *Location) map[string]any {
var out []AdditionalItem out := map[string]any{}
for len(data) >= 2 { for len(data) >= 2 {
id := data[0] id := data[0]
size := int(data[1]) size := int(data[1])
data = data[2:] data = data[2:]
if len(data) < size { if len(data) < size {
out = append(out, AdditionalItem{
ID: id,
Length: size,
ValueHex: strings.ToUpper(hex.EncodeToString(data)),
Parsed: "truncated",
})
return out return out
} }
value := data[:size] value := data[:size]
item := AdditionalItem{ var parsed map[string]any
ID: id,
Length: size,
ValueHex: strings.ToUpper(hex.EncodeToString(value)),
}
switch { switch {
case id == 0x01 && size == 4: case id == 0x01 && size == 4:
rawMileage := binary.BigEndian.Uint32(value) rawMileage := binary.BigEndian.Uint32(value)
mileage := float64(rawMileage) / 10 mileage := float64(rawMileage) / 10
location.TotalMileageKM = &mileage location.TotalMileageKM = &mileage
item.Parsed = map[string]any{ parsed = map[string]any{
"name": envelope.FieldTotalMileageKM, "name": envelope.FieldTotalMileageKM,
"raw_value_tenth_km": rawMileage, "value": mileage,
"value": mileage,
"unit": "km",
} }
case id == 0x02 && size == 2: case id == 0x02 && size == 2:
rawFuel := binary.BigEndian.Uint16(value) rawFuel := binary.BigEndian.Uint16(value)
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "fuel_l", "name": "fuel_l",
"raw_value_tenth_l": rawFuel, "value": float64(rawFuel) / 10,
"value": float64(rawFuel) / 10,
"unit": "L",
} }
case id == 0x03 && size == 2: case id == 0x03 && size == 2:
rawSpeed := binary.BigEndian.Uint16(value) rawSpeed := binary.BigEndian.Uint16(value)
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "recorder_speed_kmh", "name": "recorder_speed_kmh",
"raw_value_tenth_kmh": rawSpeed, "value": float64(rawSpeed) / 10,
"value": float64(rawSpeed) / 10,
"unit": "km/h",
} }
case id == 0x04 && size == 2: case id == 0x04 && size == 2:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "manual_alarm_event_id", "name": "manual_alarm_event_id",
"value": binary.BigEndian.Uint16(value), "value": binary.BigEndian.Uint16(value),
} }
case id == 0x05 && size == 30: case id == 0x05 && size == 30:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "tire_pressure", "name": "tire_pressure",
"values": parseTirePressure(value), "values": parseTirePressure(value),
} }
case id == 0x06 && size == 2: case id == 0x06 && size == 2:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "carriage_temperature_c", "name": "carriage_temperature_c",
"value": int16(binary.BigEndian.Uint16(value)), "value": int16(binary.BigEndian.Uint16(value)),
"unit": "C",
} }
case id == 0x11 && (size == 1 || size == 5): case id == 0x11 && (size == 1 || size == 5):
parsed := map[string]any{ parsed = map[string]any{
"name": "overspeed_alarm", "name": "overspeed_alarm",
"location_type": value[0], "location_type": value[0],
} }
if size == 5 { if size == 5 {
parsed["area_id"] = binary.BigEndian.Uint32(value[1:5]) parsed["area_id"] = binary.BigEndian.Uint32(value[1:5])
} }
item.Parsed = parsed
case id == 0x12 && size == 6: case id == 0x12 && size == 6:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "area_route_alarm", "name": "area_route_alarm",
"location_type": value[0], "location_type": value[0],
"area_id": binary.BigEndian.Uint32(value[1:5]), "area_id": binary.BigEndian.Uint32(value[1:5]),
"direction": value[5], "direction": value[5],
} }
case id == 0x13 && size == 7: case id == 0x13 && size == 7:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "road_section_driving_time_alarm", "name": "road_section_driving_time_alarm",
"road_section_id": binary.BigEndian.Uint32(value[0:4]), "road_section_id": binary.BigEndian.Uint32(value[0:4]),
"driving_time_seconds": binary.BigEndian.Uint16(value[4:6]), "driving_time_seconds": binary.BigEndian.Uint16(value[4:6]),
@@ -402,38 +377,38 @@ func parseAdditional(data []byte, location *Location) []AdditionalItem {
} }
case id == 0x25 && size == 4: case id == 0x25 && size == 4:
rawStatus := binary.BigEndian.Uint32(value) rawStatus := binary.BigEndian.Uint32(value)
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "extended_vehicle_signal_status", "name": "extended_vehicle_signal_status",
"value": rawStatus, "value": rawStatus,
"bits": parseExtendedVehicleSignal(rawStatus), "bits": parseExtendedVehicleSignal(rawStatus),
} }
case id == 0x2a && size == 2: case id == 0x2a && size == 2:
rawStatus := binary.BigEndian.Uint16(value) rawStatus := binary.BigEndian.Uint16(value)
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "io_status", "name": "io_status",
"value": rawStatus, "value": rawStatus,
"bits": parseIOStatus(rawStatus), "bits": parseIOStatus(rawStatus),
} }
case id == 0x2b && size == 4: case id == 0x2b && size == 4:
analog := binary.BigEndian.Uint32(value) analog := binary.BigEndian.Uint32(value)
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "analog", "name": "analog",
"ad0": analog & 0xffff, "ad0": analog & 0xffff,
"ad1": analog >> 16, "ad1": analog >> 16,
"value": analog, "value": analog,
} }
case id == 0x30 && size == 1: case id == 0x30 && size == 1:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "network_signal_strength", "name": "network_signal_strength",
"value": value[0], "value": value[0],
} }
case id == 0x31 && size == 1: case id == 0x31 && size == 1:
item.Parsed = map[string]any{ parsed = map[string]any{
"name": "gnss_satellite_count", "name": "gnss_satellite_count",
"value": value[0], "value": value[0],
} }
} }
out = append(out, item) addParsedAdditional(out, parsed)
data = data[size:] data = data[size:]
} }
return out return out
@@ -448,7 +423,6 @@ func locationToMap(location Location) map[string]any {
"altitude_m": location.AltitudeM, "altitude_m": location.AltitudeM,
"speed_kmh": location.SpeedKMH, "speed_kmh": location.SpeedKMH,
"direction_deg": location.DirectionDeg, "direction_deg": location.DirectionDeg,
"additional": additionalToMaps(location.Additional),
} }
if !location.Time.IsZero() { if !location.Time.IsZero() {
out["device_time"] = location.Time.Format(time.RFC3339) out["device_time"] = location.Time.Format(time.RFC3339)
@@ -456,23 +430,38 @@ func locationToMap(location Location) map[string]any {
if location.TotalMileageKM != nil { if location.TotalMileageKM != nil {
out[envelope.FieldTotalMileageKM] = *location.TotalMileageKM out[envelope.FieldTotalMileageKM] = *location.TotalMileageKM
} }
for key, value := range location.Additional {
out[key] = value
}
return out return out
} }
func additionalToMaps(items []AdditionalItem) []map[string]any { func addParsedAdditional(out map[string]any, parsed map[string]any) {
out := make([]map[string]any, 0, len(items)) if parsed == nil {
for _, item := range items { return
value := map[string]any{
"id": fmt.Sprintf("0x%02X", item.ID),
"length": item.Length,
"value_hex": item.ValueHex,
}
if item.Parsed != nil {
value["parsed"] = item.Parsed
}
out = append(out, value)
} }
return out name, ok := parsed["name"].(string)
if !ok || name == "" {
return
}
if value, ok := parsed["value"]; ok && len(parsed) <= 2 {
out[name] = value
return
}
value := make(map[string]any, len(parsed)-1)
for key, item := range parsed {
if key == "name" {
continue
}
value[key] = item
}
if len(value) == 1 {
if singleValue, ok := value["value"]; ok {
out[name] = singleValue
return
}
}
out[name] = value
} }
func parseTirePressure(data []byte) []map[string]any { func parseTirePressure(data []byte) []map[string]any {

View File

@@ -57,19 +57,14 @@ func TestExtractFramesUnescapesAndParsesLocationWithTotalMileage(t *testing.T) {
if !ok { if !ok {
t.Fatalf("parsed location missing: %#v", env.Parsed) t.Fatalf("parsed location missing: %#v", env.Parsed)
} }
additional, ok := location["additional"].([]map[string]any) if _, ok := location["additional"]; ok {
if !ok || len(additional) == 0 { t.Fatalf("additional wrapper should not be exposed: %#v", location["additional"])
t.Fatalf("additional fields missing: %#v", location["additional"])
} }
if additional[0]["id"] != "0x01" || additional[0]["value_hex"] != "0001900C" { if location[envelope.FieldTotalMileageKM] != 10241.2 {
t.Fatalf("unexpected first additional item: %#v", additional[0]) t.Fatalf("mileage should be flattened on location: %#v", location)
} }
parsed, ok := additional[0]["parsed"].(map[string]any) if _, ok := location["raw_value_tenth_km"]; ok {
if !ok { t.Fatalf("raw mileage helper should not be exposed: %#v", location)
t.Fatalf("mileage parsed field missing: %#v", additional[0])
}
if parsed["raw_value_tenth_km"] != uint32(102412) || parsed["value"] != 10241.2 {
t.Fatalf("unexpected mileage parsed value: %#v", parsed)
} }
} }
@@ -148,32 +143,38 @@ func TestParseFrameParsesCommonLocationAdditionalItems(t *testing.T) {
} }
assertFloatField(t, env, envelope.FieldTotalMileageKM, 1.1) assertFloatField(t, env, envelope.FieldTotalMileageKM, 1.1)
location := env.Parsed["location"].(map[string]any) location, ok := env.Parsed["location"].(map[string]any)
additional := location["additional"].([]map[string]any) if !ok {
byID := map[string]map[string]any{} t.Fatalf("location missing: %#v", env.Parsed)
for _, item := range additional {
byID[item["id"].(string)] = item["parsed"].(map[string]any)
} }
if byID["0x01"]["raw_value_tenth_km"] != uint32(11) || byID["0x01"]["value"] != 1.1 { if _, ok := location["additional"]; ok {
t.Fatalf("mileage parsed = %#v", byID["0x01"]) t.Fatalf("additional wrapper should not be exposed: %#v", location["additional"])
} }
if byID["0x05"]["name"] != "tire_pressure" { if location[envelope.FieldTotalMileageKM] != 1.1 {
t.Fatalf("tire pressure missing: %#v", byID["0x05"]) t.Fatalf("mileage parsed = %#v", location[envelope.FieldTotalMileageKM])
} }
tireValues := byID["0x05"]["values"].([]map[string]any) tirePressure, ok := location["tire_pressure"].(map[string]any)
if !ok {
t.Fatalf("tire pressure missing: %#v", location["tire_pressure"])
}
tireValues := tirePressure["values"].([]map[string]any)
if tireValues[0]["value"] != byte(0x37) || tireValues[0]["valid"] != true { if tireValues[0]["value"] != byte(0x37) || tireValues[0]["valid"] != true {
t.Fatalf("tire pressure first value = %#v", tireValues[0]) t.Fatalf("tire pressure first value = %#v", tireValues[0])
} }
if byID["0x11"]["location_type"] != byte(0x42) || byID["0x11"]["area_id"] != uint32(0x42) { overspeed := location["overspeed_alarm"].(map[string]any)
t.Fatalf("overspeed parsed = %#v", byID["0x11"]) if overspeed["location_type"] != byte(0x42) || overspeed["area_id"] != uint32(0x42) {
t.Fatalf("overspeed parsed = %#v", overspeed)
} }
if byID["0x12"]["location_type"] != byte(0x4d) || byID["0x12"]["area_id"] != uint32(0x4d) || byID["0x12"]["direction"] != byte(0x4d) { areaRoute := location["area_route_alarm"].(map[string]any)
t.Fatalf("area route parsed = %#v", byID["0x12"]) if areaRoute["location_type"] != byte(0x4d) || areaRoute["area_id"] != uint32(0x4d) || areaRoute["direction"] != byte(0x4d) {
t.Fatalf("area route parsed = %#v", areaRoute)
} }
if byID["0x13"]["road_section_id"] != uint32(0x58) || byID["0x13"]["driving_time_seconds"] != uint16(0x58) { roadSection := location["road_section_driving_time_alarm"].(map[string]any)
t.Fatalf("road section parsed = %#v", byID["0x13"]) if roadSection["road_section_id"] != uint32(0x58) || roadSection["driving_time_seconds"] != uint16(0x58) {
t.Fatalf("road section parsed = %#v", roadSection)
} }
ioBits := byID["0x2A"]["bits"].(map[string]bool) ioStatus := location["io_status"].(map[string]any)
ioBits := ioStatus["bits"].(map[string]bool)
if ioBits["deep_sleep"] || !ioBits["sleep"] { if ioBits["deep_sleep"] || !ioBits["sleep"] {
t.Fatalf("io bits = %#v", ioBits) t.Fatalf("io bits = %#v", ioBits)
} }

View File

@@ -61,7 +61,7 @@ func (r *Repository) Update(ctx context.Context, env envelope.FrameEnvelope) err
protocolSnapshot.VIN = vin protocolSnapshot.VIN = vin
} }
mergeFields(&protocolSnapshot, env.Fields, eventMS) mergeFields(&protocolSnapshot, env.Fields, eventMS)
protocolSnapshot.Parsed = mergeParsed(protocolSnapshot.Parsed, env.Parsed) protocolSnapshot.Parsed = mergeParsedForProtocol(env.Protocol, protocolSnapshot.Parsed, env.Parsed)
if eventMS >= protocolSnapshot.EventTimeMS { if eventMS >= protocolSnapshot.EventTimeMS {
protocolSnapshot.EventTimeMS = eventMS protocolSnapshot.EventTimeMS = eventMS
protocolSnapshot.EventID = env.StableEventID() protocolSnapshot.EventID = env.StableEventID()
@@ -290,6 +290,27 @@ func mergeParsed(existing map[string]any, incoming map[string]any) map[string]an
return out return out
} }
func mergeParsedForProtocol(protocol envelope.Protocol, existing map[string]any, incoming map[string]any) map[string]any {
out := mergeParsed(existing, incoming)
for _, key := range replaceParsedKeys(protocol, incoming) {
out[key] = cloneAny(incoming[key])
}
return out
}
func replaceParsedKeys(protocol envelope.Protocol, incoming map[string]any) []string {
keys := make([]string, 0, 2)
if _, ok := incoming["header"]; ok {
keys = append(keys, "header")
}
if protocol == envelope.ProtocolJT808 {
if _, ok := incoming["location"]; ok {
keys = append(keys, "location")
}
}
return keys
}
func mergeAny(existing any, incoming any) any { func mergeAny(existing any, incoming any) any {
switch incomingTyped := incoming.(type) { switch incomingTyped := incoming.(type) {
case map[string]any: case map[string]any:

View File

@@ -196,6 +196,72 @@ func TestRepositoryMergesNestedGB32960MotorSlicesBySerialNo(t *testing.T) {
} }
} }
func TestRepositoryReplacesJT808LocationWhenUpdatingRealtimeRaw(t *testing.T) {
repo, closeFn := newTestRepository(t)
defer closeFn()
ctx := context.Background()
if err := repo.Update(ctx, envelope.FrameEnvelope{
Protocol: envelope.ProtocolJT808,
VIN: "VIN001",
EventTimeMS: 1000,
ReceivedAtMS: 1100,
MessageID: "0x0200",
Parsed: map[string]any{
"header": map[string]any{
"message_id": "0x0200",
"phone_raw": "013307795425",
"phone_bcd_hex": "013307795425",
},
"location": map[string]any{
"additional": []any{
map[string]any{"id": "0x01", "parsed": map[string]any{"raw_value_tenth_km": 102412}},
},
"total_mileage_km": 10241.2,
},
},
Fields: map[string]any{envelope.FieldTotalMileageKM: 10241.2},
}); err != nil {
t.Fatalf("Update() initial error = %v", err)
}
if err := repo.Update(ctx, envelope.FrameEnvelope{
Protocol: envelope.ProtocolJT808,
VIN: "VIN001",
EventTimeMS: 2000,
ReceivedAtMS: 2100,
MessageID: "0x0200",
Parsed: map[string]any{
"header": map[string]any{
"message_id": "0x0200",
"phone": "13307795425",
},
"location": map[string]any{
"total_mileage_km": 10241.3,
"network_signal_strength": byte(31),
},
},
Fields: map[string]any{envelope.FieldTotalMileageKM: 10241.3},
}); err != nil {
t.Fatalf("Update() replacement error = %v", err)
}
realtimeRaw, err := repo.GetRealtimeRaw(ctx, "VIN001", envelope.ProtocolJT808)
if err != nil {
t.Fatalf("GetRealtimeRaw() error = %v", err)
}
header := realtimeRaw["header"].(map[string]any)
if _, ok := header["phone_raw"]; ok {
t.Fatalf("header should be replaced, got %#v", header)
}
location := realtimeRaw["location"].(map[string]any)
if _, ok := location["additional"]; ok {
t.Fatalf("jt808 location should be replaced, got %#v", location)
}
if location["total_mileage_km"] != float64(10241.3) {
t.Fatalf("mileage = %#v", location["total_mileage_km"])
}
}
func TestRepositoryOnlineStatus(t *testing.T) { func TestRepositoryOnlineStatus(t *testing.T) {
repo, closeFn := newTestRepository(t) repo, closeFn := newTestRepository(t)
defer closeFn() defer closeFn()