From 80474eeccbe4b57fbb7b767a768898770d1e41e9 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 21:57:18 +0800 Subject: [PATCH] feat: add go vehicle identity resolver --- deploy/portainer/docker-compose-go.yml | 2 + .../go-vehicle-gateway-verification.md | 2 + go/vehicle-gateway/cmd/gateway/main.go | 31 +++++ go/vehicle-gateway/go.mod | 1 + go/vehicle-gateway/go.sum | 3 + .../internal/gateway/mqtt_client.go | 17 +++ .../internal/gateway/tcp_server.go | 19 +++ .../internal/identity/resolver.go | 112 ++++++++++++++++++ .../internal/identity/resolver_test.go | 97 +++++++++++++++ 9 files changed, 284 insertions(+) create mode 100644 go/vehicle-gateway/internal/identity/resolver.go create mode 100644 go/vehicle-gateway/internal/identity/resolver_test.go diff --git a/deploy/portainer/docker-compose-go.yml b/deploy/portainer/docker-compose-go.yml index 78fb0d26..88dc3aa5 100644 --- a/deploy/portainer/docker-compose-go.yml +++ b/deploy/portainer/docker-compose-go.yml @@ -44,6 +44,8 @@ services: YUTONG_MQTT_CLIENT_ID: ${YUTONG_MQTT_CLIENT_ID:-lingniu-go-yutong-mqtt} YUTONG_MQTT_USERNAME: ${YUTONG_MQTT_USERNAME:-} YUTONG_MQTT_PASSWORD: ${YUTONG_MQTT_PASSWORD:-} + IDENTITY_MYSQL_DSN: ${IDENTITY_MYSQL_DSN:-} + VEHICLE_IDENTITY_TABLE: ${VEHICLE_IDENTITY_TABLE:-vehicle_identity_binding} ports: - "${GO_GB32960_TCP_PORT:-32960}:32960" - "${GO_JT808_TCP_PORT:-808}:808" diff --git a/docs/operations/go-vehicle-gateway-verification.md b/docs/operations/go-vehicle-gateway-verification.md index 28370dc0..80ec5e09 100644 --- a/docs/operations/go-vehicle-gateway-verification.md +++ b/docs/operations/go-vehicle-gateway-verification.md @@ -47,6 +47,8 @@ LINGNIU_GO_IMAGE_VERSION= KAFKA_BROKERS=172.17.111.56:9092 TDENGINE_DSN=root:@ws(172.17.111.57:6041)/lingniu_vehicle_ts MYSQL_DSN=lingniu_vehicle:@tcp(rm-bp179zbv481rnw3e2.mysql.rds.aliyuncs.com:3306)/lingniu_vehicle_data?parseTime=true&charset=utf8mb4,utf8&loc=Asia%2FShanghai +IDENTITY_MYSQL_DSN=${MYSQL_DSN} +VEHICLE_IDENTITY_TABLE=vehicle_identity_binding REDIS_ADDR=r-bp1u741kij7e51i481.redis.rds.aliyuncs.com:6379 REDIS_PASSWORD= REDIS_DB=50 diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 1069fb1b..abfdc625 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -2,6 +2,7 @@ package main import ( "context" + "database/sql" "log/slog" "os" "os/signal" @@ -10,9 +11,12 @@ import ( "syscall" "time" + _ "github.com/go-sql-driver/mysql" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/gateway" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/gb32960" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/jt808" @@ -29,6 +33,12 @@ func main() { os.Exit(1) } defer sink.Close() + resolver, closeResolver, err := buildIdentityResolver(ctx, logger) + if err != nil { + logger.Error("build identity resolver failed", "error", err) + os.Exit(1) + } + defer closeResolver() protocols := []gateway.TCPProtocol{ { @@ -51,6 +61,7 @@ func main() { server, err := gateway.NewTCPServer(gateway.TCPServerConfig{ Protocol: protocol, Sink: sink, + Resolver: resolver, Logger: logger, ReadBufferSize: envInt("TCP_READ_BUFFER_BYTES", 64*1024), IdleTimeout: time.Duration(envInt("TCP_IDLE_TIMEOUT_SECONDS", 180)) * time.Second, @@ -78,6 +89,7 @@ func main() { Topics: splitCSV(env("YUTONG_MQTT_TOPICS", env("YUTONG_MQTT_TOPIC", "/ytforward/shln/+"))), QoS: byte(envInt("YUTONG_MQTT_QOS", 2)), Sink: sink, + Resolver: resolver, Logger: logger, }) if err != nil { @@ -101,6 +113,25 @@ func main() { wg.Wait() } +func buildIdentityResolver(ctx context.Context, logger *slog.Logger) (identity.Resolver, func(), error) { + dsn := env("IDENTITY_MYSQL_DSN", strings.TrimSpace(os.Getenv("MYSQL_DSN"))) + if strings.TrimSpace(dsn) == "" { + logger.Warn("identity mysql dsn is empty; using noop identity resolver") + return identity.NoopResolver{}, func() {}, nil + } + db, err := sql.Open("mysql", dsn) + if err != nil { + return nil, nil, err + } + if err := db.PingContext(ctx); err != nil { + _ = db.Close() + return nil, nil, err + } + table := env("VEHICLE_IDENTITY_TABLE", "vehicle_identity_binding") + logger.Info("identity mysql resolver enabled", "table", table) + return identity.NewMySQLResolver(db, table), func() { _ = db.Close() }, nil +} + func buildSink(logger *slog.Logger) (eventbus.Sink, error) { brokers := splitCSV(os.Getenv("KAFKA_BROKERS")) if len(brokers) == 0 { diff --git a/go/vehicle-gateway/go.mod b/go/vehicle-gateway/go.mod index 022aa215..34455b63 100644 --- a/go/vehicle-gateway/go.mod +++ b/go/vehicle-gateway/go.mod @@ -3,6 +3,7 @@ module lingniu-vehicle-ingest/go/vehicle-gateway go 1.26 require ( + github.com/DATA-DOG/go-sqlmock v1.5.2 github.com/alicebob/miniredis/v2 v2.35.0 github.com/eclipse/paho.mqtt.golang v1.5.1 github.com/go-sql-driver/mysql v1.9.3 diff --git a/go/vehicle-gateway/go.sum b/go/vehicle-gateway/go.sum index 00fee118..1dafdbf9 100644 --- a/go/vehicle-gateway/go.sum +++ b/go/vehicle-gateway/go.sum @@ -1,5 +1,7 @@ filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= +github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7OputlJIzU= +github.com/DATA-DOG/go-sqlmock v1.5.2/go.mod h1:88MAG/4G7SMwSE3CeA0ZKzrT5CiOU3OJ+JlNzwDqpNU= github.com/alicebob/miniredis/v2 v2.35.0 h1:QwLphYqCEAo1eu1TqPRN2jgVMPBweeQcR21jeqDCONI= github.com/alicebob/miniredis/v2 v2.35.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM= github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= @@ -25,6 +27,7 @@ github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aN github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/kisielk/sqlstruct v0.0.0-20201105191214-5f3e10d3ab46/go.mod h1:yyMNCyc/Ib3bDTKd379tNMpB/7/H5TjM2Y9QJ5THLbE= github.com/klauspost/compress v1.15.9 h1:wKRjX6JRtDdrE9qwa4b/Cip7ACOshUI4smpCQanqjSY= github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 h1:ZqeYNhU3OHLH3mGKHDcjJRFFRrJa6eAM5H+CtDdOsPc= diff --git a/go/vehicle-gateway/internal/gateway/mqtt_client.go b/go/vehicle-gateway/internal/gateway/mqtt_client.go index 9403ffa3..4c798d72 100644 --- a/go/vehicle-gateway/internal/gateway/mqtt_client.go +++ b/go/vehicle-gateway/internal/gateway/mqtt_client.go @@ -12,6 +12,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/yutongmqtt" ) @@ -24,6 +25,7 @@ type MQTTClientConfig struct { Topics []string QoS byte Sink eventbus.Sink + Resolver identity.Resolver Logger *slog.Logger } @@ -48,6 +50,9 @@ func NewMQTTClient(cfg MQTTClientConfig) (*MQTTClient, error) { if cfg.Logger == nil { cfg.Logger = slog.Default() } + if cfg.Resolver == nil { + cfg.Resolver = identity.NoopResolver{} + } if cfg.EndpointName == "" { cfg.EndpointName = "yutong" } @@ -115,6 +120,18 @@ func (c *MQTTClient) handleMessage(ctx context.Context, topic string, payload [] ParseError: err.Error(), } env.EventID = env.StableEventID() + } else { + resolved, resolveErr := c.cfg.Resolver.Resolve(ctx, env) + if resolveErr != nil { + c.cfg.Logger.Warn("mqtt identity resolve failed", "topic", topic, "event_id", env.StableEventID(), "error", resolveErr) + if env.Parsed == nil { + env.Parsed = map[string]any{} + } + env.Parsed["identity"] = map[string]any{"resolved": false, "error": resolveErr.Error()} + env.ParseStatus = envelope.ParsePartial + } else { + env = resolved + } } if err := c.cfg.Sink.PublishRaw(ctx, env); err != nil { c.cfg.Logger.Error("publish mqtt raw failed", "topic", topic, "event_id", env.StableEventID(), "error", err) diff --git a/go/vehicle-gateway/internal/gateway/tcp_server.go b/go/vehicle-gateway/internal/gateway/tcp_server.go index 05a3a426..18f52484 100644 --- a/go/vehicle-gateway/internal/gateway/tcp_server.go +++ b/go/vehicle-gateway/internal/gateway/tcp_server.go @@ -14,6 +14,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" ) type FrameExtractor func([]byte) (frames [][]byte, remainder []byte, err error) @@ -30,6 +31,7 @@ type TCPProtocol struct { type TCPServer struct { protocol TCPProtocol sink eventbus.Sink + resolver identity.Resolver logger *slog.Logger readBufferSize int idleTimeout time.Duration @@ -39,6 +41,7 @@ type TCPServer struct { type TCPServerConfig struct { Protocol TCPProtocol Sink eventbus.Sink + Resolver identity.Resolver Logger *slog.Logger ReadBufferSize int IdleTimeout time.Duration @@ -64,6 +67,9 @@ func NewTCPServer(cfg TCPServerConfig) (*TCPServer, error) { if cfg.Logger == nil { cfg.Logger = slog.Default() } + if cfg.Resolver == nil { + cfg.Resolver = identity.NoopResolver{} + } if cfg.ReadBufferSize <= 0 { cfg.ReadBufferSize = 32 * 1024 } @@ -76,6 +82,7 @@ func NewTCPServer(cfg TCPServerConfig) (*TCPServer, error) { return &TCPServer{ protocol: cfg.Protocol, sink: cfg.Sink, + resolver: cfg.Resolver, logger: cfg.Logger, readBufferSize: cfg.ReadBufferSize, idleTimeout: cfg.IdleTimeout, @@ -181,6 +188,18 @@ func (s *TCPServer) handleFrame(ctx context.Context, raw []byte, source string) ParseError: err.Error(), } env.EventID = env.StableEventID() + } else { + resolved, resolveErr := s.resolver.Resolve(ctx, env) + if resolveErr != nil { + s.logger.Warn("identity resolve failed", "protocol", s.protocol.Protocol, "event_id", env.StableEventID(), "error", resolveErr) + if env.Parsed == nil { + env.Parsed = map[string]any{} + } + env.Parsed["identity"] = map[string]any{"resolved": false, "error": resolveErr.Error()} + env.ParseStatus = envelope.ParsePartial + } else { + env = resolved + } } if err := s.sink.PublishRaw(ctx, env); err != nil { diff --git a/go/vehicle-gateway/internal/identity/resolver.go b/go/vehicle-gateway/internal/identity/resolver.go new file mode 100644 index 00000000..1f569dab --- /dev/null +++ b/go/vehicle-gateway/internal/identity/resolver.go @@ -0,0 +1,112 @@ +package identity + +import ( + "context" + "database/sql" + "errors" + "strings" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +type Resolver interface { + Resolve(context.Context, envelope.FrameEnvelope) (envelope.FrameEnvelope, error) +} + +type NoopResolver struct{} + +func (NoopResolver) Resolve(_ context.Context, env envelope.FrameEnvelope) (envelope.FrameEnvelope, error) { + return env, nil +} + +type MySQLResolver struct { + db *sql.DB + table string +} + +func NewMySQLResolver(db *sql.DB, table string) *MySQLResolver { + if db == nil { + panic("identity db must not be nil") + } + table = strings.TrimSpace(table) + if table == "" || !safeIdentifier(table) { + table = "vehicle_identity_binding" + } + return &MySQLResolver{db: db, table: table} +} + +func safeIdentifier(value string) bool { + if value == "" { + return false + } + for _, r := range value { + if (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') || r == '_' { + continue + } + return false + } + return true +} + +func (r *MySQLResolver) Resolve(ctx context.Context, env envelope.FrameEnvelope) (envelope.FrameEnvelope, error) { + if strings.TrimSpace(env.VIN) != "" { + return env, nil + } + candidates := CandidateKeys(env) + for _, candidate := range candidates { + vin, err := r.lookup(ctx, candidate.Column, candidate.Value) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + continue + } + return env, err + } + if strings.TrimSpace(vin) != "" { + env.VIN = strings.TrimSpace(vin) + if env.Parsed == nil { + env.Parsed = map[string]any{} + } + env.Parsed["identity"] = map[string]any{ + "resolved": true, + "source": candidate.Column, + "value": candidate.Value, + } + env.EventID = env.StableEventID() + return env, nil + } + } + return env, nil +} + +func (r *MySQLResolver) lookup(ctx context.Context, column string, value string) (string, error) { + query := "SELECT vin FROM " + r.table + " WHERE " + column + " = ? AND vin IS NOT NULL AND vin <> '' ORDER BY updated_at DESC LIMIT 1" + var vin string + err := r.db.QueryRowContext(ctx, query, value).Scan(&vin) + return vin, err +} + +type CandidateKey struct { + Column string + Value string +} + +func CandidateKeys(env envelope.FrameEnvelope) []CandidateKey { + var out []CandidateKey + add := func(column string, value string) { + value = strings.TrimSpace(value) + if value == "" || strings.EqualFold(value, "unknown") { + return + } + for _, existing := range out { + if existing.Column == column && existing.Value == value { + return + } + } + out = append(out, CandidateKey{Column: column, Value: value}) + } + add("phone", env.Phone) + add("device_id", env.DeviceID) + add("plate", env.Plate) + add("vin", env.VehicleKeyHint) + return out +} diff --git a/go/vehicle-gateway/internal/identity/resolver_test.go b/go/vehicle-gateway/internal/identity/resolver_test.go new file mode 100644 index 00000000..2bbf1960 --- /dev/null +++ b/go/vehicle-gateway/internal/identity/resolver_test.go @@ -0,0 +1,97 @@ +package identity + +import ( + "context" + "database/sql" + "testing" + + "github.com/DATA-DOG/go-sqlmock" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestCandidateKeysPrioritizesPhoneDevicePlate(t *testing.T) { + keys := CandidateKeys(envelope.FrameEnvelope{ + Phone: "013307795425", + DeviceID: "D1", + Plate: "豫A12345", + VehicleKeyHint: "D1", + }) + got := []string{} + for _, key := range keys { + got = append(got, key.Column+"="+key.Value) + } + want := []string{"phone=013307795425", "device_id=D1", "plate=豫A12345", "vin=D1"} + if len(got) != len(want) { + t.Fatalf("candidate count = %d got=%#v", len(got), got) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("candidate[%d] = %q, want %q", i, got[i], want[i]) + } + } +} + +func TestMySQLResolverFillsVINFromPhone(t *testing.T) { + db, mock := newMockDB(t) + defer db.Close() + mock.ExpectQuery("SELECT vin FROM vehicle_identity_binding WHERE phone = \\?"). + WithArgs("013307795425"). + WillReturnRows(sqlmock.NewRows([]string{"vin"}).AddRow("LNBVIN00000000001")) + + resolver := NewMySQLResolver(db, "vehicle_identity_binding") + env, err := resolver.Resolve(context.Background(), envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + Phone: "013307795425", + Parsed: map[string]any{}, + }) + if err != nil { + t.Fatalf("Resolve() error = %v", err) + } + if env.VIN != "LNBVIN00000000001" { + t.Fatalf("vin = %q", env.VIN) + } + identity, ok := env.Parsed["identity"].(map[string]any) + if !ok || identity["source"] != "phone" { + t.Fatalf("identity metadata = %#v", env.Parsed["identity"]) + } + if err := mock.ExpectationsWereMet(); err != nil { + t.Fatalf("sql expectations: %v", err) + } +} + +func TestMySQLResolverFallsBackToDeviceID(t *testing.T) { + db, mock := newMockDB(t) + defer db.Close() + mock.ExpectQuery("SELECT vin FROM vehicle_identity_binding WHERE phone = \\?"). + WithArgs("013307795425"). + WillReturnError(sql.ErrNoRows) + mock.ExpectQuery("SELECT vin FROM vehicle_identity_binding WHERE device_id = \\?"). + WithArgs("D1"). + WillReturnRows(sqlmock.NewRows([]string{"vin"}).AddRow("LNBVIN00000000002")) + + resolver := NewMySQLResolver(db, "vehicle_identity_binding") + env, err := resolver.Resolve(context.Background(), envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + Phone: "013307795425", + DeviceID: "D1", + }) + if err != nil { + t.Fatalf("Resolve() error = %v", err) + } + if env.VIN != "LNBVIN00000000002" { + t.Fatalf("vin = %q", env.VIN) + } + if err := mock.ExpectationsWereMet(); err != nil { + t.Fatalf("sql expectations: %v", err) + } +} + +func newMockDB(t *testing.T) (*sql.DB, sqlmock.Sqlmock) { + t.Helper() + db, mock, err := sqlmock.New() + if err != nil { + t.Fatalf("sqlmock.New() error = %v", err) + } + return db, mock +}