feat(go): reuse parsed fields across projections
This commit is contained in:
@@ -33,25 +33,27 @@ const (
|
||||
)
|
||||
|
||||
type FrameEnvelope struct {
|
||||
EventID string `json:"event_id"`
|
||||
TraceID string `json:"trace_id"`
|
||||
Protocol Protocol `json:"protocol"`
|
||||
MessageID string `json:"message_id"`
|
||||
Sequence uint16 `json:"sequence"`
|
||||
VIN string `json:"vin,omitempty"`
|
||||
VehicleKeyHint string `json:"vehicle_key,omitempty"`
|
||||
Phone string `json:"phone,omitempty"`
|
||||
DeviceID string `json:"device_id,omitempty"`
|
||||
Plate string `json:"plate,omitempty"`
|
||||
SourceEndpoint string `json:"source_endpoint,omitempty"`
|
||||
EventTimeMS int64 `json:"event_time_ms"`
|
||||
ReceivedAtMS int64 `json:"received_at_ms"`
|
||||
RawHex string `json:"raw_hex,omitempty"`
|
||||
RawText string `json:"raw_text,omitempty"`
|
||||
Parsed map[string]any `json:"parsed,omitempty"`
|
||||
Fields map[string]any `json:"fields,omitempty"`
|
||||
ParseStatus ParseStatus `json:"parse_status"`
|
||||
ParseError string `json:"parse_error,omitempty"`
|
||||
EventID string `json:"event_id"`
|
||||
TraceID string `json:"trace_id"`
|
||||
Protocol Protocol `json:"protocol"`
|
||||
MessageID string `json:"message_id"`
|
||||
Sequence uint16 `json:"sequence"`
|
||||
VIN string `json:"vin,omitempty"`
|
||||
VehicleKeyHint string `json:"vehicle_key,omitempty"`
|
||||
Phone string `json:"phone,omitempty"`
|
||||
DeviceID string `json:"device_id,omitempty"`
|
||||
Plate string `json:"plate,omitempty"`
|
||||
SourceEndpoint string `json:"source_endpoint,omitempty"`
|
||||
EventTimeMS int64 `json:"event_time_ms"`
|
||||
ReceivedAtMS int64 `json:"received_at_ms"`
|
||||
RawHex string `json:"raw_hex,omitempty"`
|
||||
RawText string `json:"raw_text,omitempty"`
|
||||
Parsed map[string]any `json:"parsed,omitempty"`
|
||||
ParsedFields map[string]any `json:"parsed_fields,omitempty"`
|
||||
ParsedFieldTypes map[string]string `json:"parsed_field_types,omitempty"`
|
||||
Fields map[string]any `json:"fields,omitempty"`
|
||||
ParseStatus ParseStatus `json:"parse_status"`
|
||||
ParseError string `json:"parse_error,omitempty"`
|
||||
}
|
||||
|
||||
func (e FrameEnvelope) VehicleKey() string {
|
||||
|
||||
@@ -222,6 +222,9 @@ func (c *MQTTClient) handleMessage(ctx context.Context, topic string, payload []
|
||||
}
|
||||
frameStatus = env.ParseStatus
|
||||
c.recordFrameMetric(env.ParseStatus)
|
||||
if env.ParseStatus != envelope.ParseBadFrame {
|
||||
realtime.EnsureParsedFields(&env)
|
||||
}
|
||||
if err := c.cfg.Sink.PublishRaw(messageCtx, env); err != nil {
|
||||
c.recordPublishMetric("raw", "error")
|
||||
c.cfg.Logger.Error("publish mqtt raw failed", "topic", topic, "event_id", env.StableEventID(), "error", err)
|
||||
|
||||
@@ -237,6 +237,9 @@ func (s *TCPServer) handleFrame(ctx context.Context, conn net.Conn, raw []byte,
|
||||
}
|
||||
frameStatus = env.ParseStatus
|
||||
s.recordFrameMetric(env.ParseStatus)
|
||||
if env.ParseStatus != envelope.ParseBadFrame {
|
||||
realtime.EnsureParsedFields(&env)
|
||||
}
|
||||
|
||||
if err := s.sink.PublishRaw(frameCtx, env); err != nil {
|
||||
s.recordPublishMetric("raw", "error")
|
||||
|
||||
@@ -318,11 +318,11 @@ func jsonString(value any) string {
|
||||
}
|
||||
|
||||
func parsedFieldsJSONString(env envelope.FrameEnvelope) string {
|
||||
fieldsEnv, ok := realtime.BuildFieldsEnvelope(env)
|
||||
if !ok || len(fieldsEnv.Fields) == 0 {
|
||||
fields, _, ok := realtime.ParsedFieldsForEnvelope(env)
|
||||
if !ok || len(fields) == 0 {
|
||||
return ""
|
||||
}
|
||||
return jsonString(fieldsEnv.Fields)
|
||||
return jsonString(fields)
|
||||
}
|
||||
|
||||
func floatField(env envelope.FrameEnvelope, key string) (float64, bool) {
|
||||
|
||||
@@ -111,8 +111,8 @@ func TestWriterChunksOversizedParsedFields(t *testing.T) {
|
||||
exec := &recordingExec{}
|
||||
writer := NewWriter(exec)
|
||||
env := sampleEnvelope()
|
||||
env.Parsed = map[string]any{
|
||||
"large_payload": strings.Repeat("x", 16_128),
|
||||
env.ParsedFields = map[string]any{
|
||||
"jt808.location.large_payload": strings.Repeat("x", 16_128),
|
||||
}
|
||||
|
||||
if err := writer.AppendRawFrame(context.Background(), env); err != nil {
|
||||
@@ -131,6 +131,30 @@ func TestWriterChunksOversizedParsedFields(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterUsesEnvelopeParsedFieldsWithoutReflattening(t *testing.T) {
|
||||
exec := &recordingExec{}
|
||||
writer := NewWriter(exec)
|
||||
env := sampleEnvelope()
|
||||
env.Parsed = map[string]any{
|
||||
"location": map[string]any{"speed_kmh": 99},
|
||||
}
|
||||
env.ParsedFields = map[string]any{
|
||||
"jt808.location.speed_kmh": "12.3",
|
||||
}
|
||||
|
||||
if err := writer.AppendRawFrame(context.Background(), env); err != nil {
|
||||
t.Fatalf("AppendRawFrame() error = %v", err)
|
||||
}
|
||||
|
||||
rawInsert := findSQL(exec.calls, "INSERT INTO raw_")
|
||||
if !strings.Contains(rawInsert, `"jt808.location.speed_kmh":"12.3"`) {
|
||||
t.Fatalf("raw insert should use precomputed parsed fields: %s", rawInsert)
|
||||
}
|
||||
if strings.Contains(rawInsert, `"jt808.location.speed_kmh":"99"`) {
|
||||
t.Fatalf("raw insert should not re-flatten parsed payload: %s", rawInsert)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterLeavesParsedFieldsEmptyWhenNoParsedPayload(t *testing.T) {
|
||||
exec := &recordingExec{}
|
||||
writer := NewWriter(exec)
|
||||
|
||||
@@ -26,33 +26,27 @@ func BuildFieldsEnvelope(env envelope.FrameEnvelope) (envelope.FrameEnvelope, bo
|
||||
if !envelope.IsRealtimeTelemetryFrame(env) {
|
||||
return envelope.FrameEnvelope{}, false
|
||||
}
|
||||
rows := realtimeKVFields(env, env.Parsed)
|
||||
fields := make(map[string]any, len(rows))
|
||||
for _, row := range rows {
|
||||
key := row.Domain
|
||||
if strings.TrimSpace(row.Field) != "" {
|
||||
key += "." + row.Field
|
||||
}
|
||||
fields[key] = row.Value
|
||||
}
|
||||
if len(fields) == 0 {
|
||||
fields, fieldTypes, ok := ParsedFieldsForEnvelope(env)
|
||||
if !ok {
|
||||
return envelope.FrameEnvelope{}, false
|
||||
}
|
||||
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: fields,
|
||||
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: fields,
|
||||
ParsedFields: fields,
|
||||
ParsedFieldTypes: fieldTypes,
|
||||
Parsed: map[string]any{
|
||||
"source_event_id": env.StableEventID(),
|
||||
"field_mapping": realtimeFieldMappingVersion,
|
||||
@@ -62,6 +56,91 @@ func BuildFieldsEnvelope(env envelope.FrameEnvelope) (envelope.FrameEnvelope, bo
|
||||
return out, true
|
||||
}
|
||||
|
||||
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
|
||||
|
||||
@@ -348,20 +348,10 @@ func (r *Repository) setJSON(ctx context.Context, key string, value any, ttl tim
|
||||
}
|
||||
|
||||
func (r *Repository) setKV(ctx context.Context, vin string, env envelope.FrameEnvelope, parsed map[string]any) error {
|
||||
rows := realtimeKVFields(env, parsed)
|
||||
if len(rows) == 0 {
|
||||
return nil
|
||||
}
|
||||
values := make(map[string]any, len(rows))
|
||||
types := make(map[string]any, len(rows))
|
||||
for _, row := range rows {
|
||||
field := realtimeKVFieldPath(row.Domain, row.Field)
|
||||
if field == "" {
|
||||
continue
|
||||
}
|
||||
values[field] = row.Value
|
||||
types[field] = row.ValueType
|
||||
if len(env.ParsedFields) == 0 && len(parsed) > 0 {
|
||||
env.Parsed = parsed
|
||||
}
|
||||
values, types := realtimeKVMapsForEnvelope(env)
|
||||
if len(values) == 0 {
|
||||
return nil
|
||||
}
|
||||
@@ -383,8 +373,7 @@ func (r *Repository) setKV(ctx context.Context, vin string, env envelope.FrameEn
|
||||
}
|
||||
|
||||
func (r *Repository) setFastProjection(ctx context.Context, vehicleKey string, vin string, env envelope.FrameEnvelope) error {
|
||||
rows := realtimeKVFields(env, env.Parsed)
|
||||
values, types := realtimeKVMaps(rows)
|
||||
values, types := realtimeKVMapsForEnvelope(env)
|
||||
eventTimeMS := eventTimeOrReceivedMS(env)
|
||||
offlineAfterMS := env.ReceivedAtMS + r.cfg.ttl().Milliseconds()
|
||||
online := OnlineStatus{
|
||||
@@ -435,6 +424,30 @@ func (r *Repository) setFastProjection(ctx context.Context, vehicleKey string, v
|
||||
return err
|
||||
}
|
||||
|
||||
func realtimeKVMapsForEnvelope(env envelope.FrameEnvelope) (map[string]any, map[string]any) {
|
||||
if fields, fieldTypes, ok := ParsedFieldsForEnvelope(env); ok {
|
||||
values := make(map[string]any, len(fields))
|
||||
types := make(map[string]any, len(fields))
|
||||
for field, value := range fields {
|
||||
if strings.TrimSpace(field) == "" {
|
||||
continue
|
||||
}
|
||||
stringValue, valueType, ok := stringifyKVValue(value)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
values[field] = stringValue
|
||||
if typed := strings.TrimSpace(fieldTypes[field]); typed != "" {
|
||||
types[field] = typed
|
||||
} else {
|
||||
types[field] = valueType
|
||||
}
|
||||
}
|
||||
return values, types
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func realtimeKVMaps(rows []RealtimeKVField) (map[string]any, map[string]any) {
|
||||
values := make(map[string]any, len(rows))
|
||||
types := make(map[string]any, len(rows))
|
||||
|
||||
@@ -646,6 +646,37 @@ func TestRepositoryOnlineStatus(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRepositoryUsesEnvelopeParsedFieldsForRealtimeKV(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{
|
||||
"location": map[string]any{"speed_kmh": 99},
|
||||
},
|
||||
ParsedFields: map[string]any{
|
||||
"jt808.location.speed_kmh": "12.3",
|
||||
},
|
||||
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
||||
}); err != nil {
|
||||
t.Fatalf("Update() error = %v", err)
|
||||
}
|
||||
|
||||
values, err := repo.client.HGetAll(ctx, realtimeKVValuesKey(envelope.ProtocolJT808, "VIN001")).Result()
|
||||
if err != nil {
|
||||
t.Fatalf("HGetAll() error = %v", err)
|
||||
}
|
||||
if got, want := values["jt808.location.speed_kmh"], "12.3"; got != want {
|
||||
t.Fatalf("speed kv = %#v, want %q; all=%#v", got, want, values)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRepositorySkipsFramesWithoutVIN(t *testing.T) {
|
||||
repo, closeFn := newTestRepository(t)
|
||||
defer closeFn()
|
||||
|
||||
@@ -187,7 +187,10 @@ func (w *SnapshotWriter) snapshotFieldsForEnvelope(ctx context.Context, env enve
|
||||
if len(existing) > 0 {
|
||||
if isStructuredSnapshotParsed(env.Protocol, existing) {
|
||||
parsed = mergeParsedForProtocol(env.Protocol, existing, env.Parsed)
|
||||
return realtimeSnapshotFlatFields(env, parsed), nil
|
||||
mergedEnv := env
|
||||
mergedEnv.ParsedFields = nil
|
||||
mergedEnv.ParsedFieldTypes = nil
|
||||
return realtimeSnapshotFlatFields(mergedEnv, parsed), nil
|
||||
}
|
||||
return mergeRealtimeSnapshotFields(existing, incoming), nil
|
||||
}
|
||||
@@ -196,6 +199,12 @@ func (w *SnapshotWriter) snapshotFieldsForEnvelope(ctx context.Context, env enve
|
||||
}
|
||||
|
||||
func realtimeSnapshotFlatFields(env envelope.FrameEnvelope, parsed map[string]any) map[string]any {
|
||||
if len(env.ParsedFields) > 0 {
|
||||
fields := cloneMap(env.ParsedFields)
|
||||
if len(fields) > 0 {
|
||||
return fields
|
||||
}
|
||||
}
|
||||
rows := realtimeKVFields(env, parsed)
|
||||
if len(rows) == 0 {
|
||||
return nil
|
||||
|
||||
@@ -257,6 +257,57 @@ func TestSnapshotWriterMergesGB32960ParsedJSONAcrossSplitRealtimeFrames(t *testi
|
||||
}
|
||||
}
|
||||
|
||||
func TestSnapshotWriterUsesEnvelopeParsedFieldsWithoutReflattening(t *testing.T) {
|
||||
db, mock, err := sqlmock.New()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
writer := NewSnapshotWriter(db)
|
||||
|
||||
mock.ExpectQuery("SELECT parsed_json FROM vehicle_realtime_snapshot WHERE protocol = \\? AND vin = \\?").
|
||||
WithArgs("GB32960", "VIN001").
|
||||
WillReturnRows(sqlmock.NewRows([]string{"parsed_json"}))
|
||||
mock.ExpectExec("INSERT INTO vehicle_realtime_snapshot").
|
||||
WithArgs(
|
||||
"GB32960",
|
||||
"VIN001",
|
||||
"",
|
||||
"",
|
||||
"",
|
||||
jsonFlatFieldsArg{fields: map[string]string{
|
||||
"gb32960.vehicle.speed_kmh": "12.3",
|
||||
}},
|
||||
sqlmock.AnyArg(),
|
||||
sqlmock.AnyArg(),
|
||||
sqlmock.AnyArg(),
|
||||
).
|
||||
WillReturnResult(sqlmock.NewResult(0, 1))
|
||||
|
||||
err = writer.Update(context.Background(), envelope.FrameEnvelope{
|
||||
Protocol: envelope.ProtocolGB32960,
|
||||
MessageID: "0x02",
|
||||
VIN: "VIN001",
|
||||
EventTimeMS: 1782918600000,
|
||||
ReceivedAtMS: 1782918601000,
|
||||
Parsed: map[string]any{
|
||||
"data_units": []any{
|
||||
map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"speed_kmh": 99}},
|
||||
},
|
||||
},
|
||||
ParsedFields: map[string]any{
|
||||
"gb32960.vehicle.speed_kmh": "12.3",
|
||||
},
|
||||
Fields: map[string]any{envelope.FieldSpeedKMH: 12.3},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Update() error = %v", err)
|
||||
}
|
||||
if err := mock.ExpectationsWereMet(); err != nil {
|
||||
t.Fatalf("sql expectations: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSnapshotWriterCollapsesGB32960StackFragmentsWhenMergingParsedJSON(t *testing.T) {
|
||||
db, mock, err := sqlmock.New()
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user