From 3c9678491a818cc262d3c1881123eb2a6bfc00f6 Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 20:06:38 +0800 Subject: [PATCH] feat(go): bootstrap minimal identity schema --- docs/architecture/storage-minimal-contract.md | 9 +++- go/vehicle-gateway/cmd/gateway/main.go | 14 +++-- .../internal/identity/resolver.go | 53 +++++++++++++++++++ .../internal/identity/resolver_test.go | 17 ++++++ 4 files changed, 87 insertions(+), 6 deletions(-) diff --git a/docs/architecture/storage-minimal-contract.md b/docs/architecture/storage-minimal-contract.md index 0bde0b77..809bffe6 100644 --- a/docs/architecture/storage-minimal-contract.md +++ b/docs/architecture/storage-minimal-contract.md @@ -24,8 +24,8 @@ | MySQL | `vehicle_realtime_snapshot` | 每个协议+VIN 的最新事件时间、接收时间、车牌、事件 ID | phone、device、source endpoint、完整 JSON | | MySQL | `vehicle_realtime_location` | 每个协议+VIN 的最新位置核心字段 | raw、parsed JSON、消息头内部字段 | | MySQL | `vehicle_daily_mileage` | 日期、VIN/临时 vehicle key、协议、日里程、首末总里程、样本数 | 泛化 metric key/value、每帧细节、位置点列表 | -| MySQL | `vehicle_identity_binding` | 人工维护或导入的 VIN、车牌、phone、device_id 映射 | 注册历史 | -| MySQL | `jt808_registration` | JT808 phone 主键下的注册、鉴权、VIN 匹配状态 | GB32960/MQTT 注册信息 | +| MySQL | `vehicle_identity_binding` | 人工维护或导入的 VIN、车牌、phone、device_id 映射 | 注册历史、协议状态 | +| MySQL | `jt808_registration` | JT808 phone 主键下的注册、鉴权、VIN 匹配状态、来源端点 | GB32960/MQTT 注册信息、位置历史 | | Redis | `vehicle:realtime-raw:{protocol}:{vin}` | 每协议最新完整 parsed 状态 | 历史数据、统计结果 | ## 当前应收敛的重复点 @@ -45,6 +45,11 @@ - 收敛方向:业务查询以 VIN 为主;没有 VIN 的数据只用于排查和待绑定,不作为正式车辆指标。 - 生产:RDS 已在上线前删除泛化 `vehicle_daily_metric`,改为专用 `vehicle_daily_mileage`。 +4. identity 层只保留两张表。 + - `vehicle_identity_binding`:VIN 与 plate/phone/device_id 的映射,供 808 等协议反查 VIN。 + - `jt808_registration`:808 注册、鉴权、最新活跃和 VIN 匹配状态。 + - 生产:RDS 已在上线前删除旧 `vehicle_identity_binding_registration`、`vehicle_identity_bindings`,gateway 启动会自动创建最小 schema。 + ## 字段提升规则 协议字段进入系统后分三层: diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 824ea0aa..07e9dfcf 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -143,11 +143,17 @@ func buildIdentityResolver(ctx context.Context, logger *slog.Logger) (identity.R return nil, nil, err } table := env("VEHICLE_IDENTITY_TABLE", "vehicle_identity_binding") + resolver := identity.NewMySQLResolverWithOptions(db, table, identity.MySQLResolverOptions{ + LocationTouchInterval: time.Duration(envInt("JT808_REGISTRATION_LOCATION_TOUCH_INTERVAL_SECONDS", 600)) * time.Second, + }) + if envBool("IDENTITY_MYSQL_ENSURE_SCHEMA", true) { + if err := resolver.EnsureSchema(ctx); err != nil { + _ = db.Close() + return nil, nil, err + } + } logger.Info("identity mysql resolver enabled", "table", table) - locationTouchInterval := time.Duration(envInt("JT808_REGISTRATION_LOCATION_TOUCH_INTERVAL_SECONDS", 600)) * time.Second - return identity.NewMySQLResolverWithOptions(db, table, identity.MySQLResolverOptions{ - LocationTouchInterval: locationTouchInterval, - }), func() { _ = db.Close() }, nil + return resolver, func() { _ = db.Close() }, nil } func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) { diff --git a/go/vehicle-gateway/internal/identity/resolver.go b/go/vehicle-gateway/internal/identity/resolver.go index 31135d1d..4ed21dd4 100644 --- a/go/vehicle-gateway/internal/identity/resolver.go +++ b/go/vehicle-gateway/internal/identity/resolver.go @@ -58,6 +58,59 @@ func NewMySQLResolverWithOptions(db *sql.DB, table string, opts MySQLResolverOpt } } +func (r *MySQLResolver) EnsureSchema(ctx context.Context) error { + if _, err := r.db.ExecContext(ctx, identityBindingTableSQL(r.table)); err != nil { + return err + } + _, err := r.db.ExecContext(ctx, jt808RegistrationTableSQL) + return err +} + +func identityBindingTableSQL(table string) string { + if table == "" || !safeIdentifier(table) { + table = "vehicle_identity_binding" + } + return `CREATE TABLE IF NOT EXISTS ` + table + ` ( + id BIGINT PRIMARY KEY AUTO_INCREMENT, + vin VARCHAR(32) NOT NULL, + plate VARCHAR(32) NULL, + phone VARCHAR(32) NULL, + device_id VARCHAR(64) NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + UNIQUE KEY uk_identity_plate (plate), + UNIQUE KEY uk_identity_phone (phone), + UNIQUE KEY uk_identity_device (device_id), + KEY idx_identity_vin (vin) +)` +} + +const jt808RegistrationTableSQL = `CREATE TABLE IF NOT EXISTS jt808_registration ( + phone VARCHAR(32) PRIMARY KEY, + device_id VARCHAR(64) NULL, + plate VARCHAR(32) NULL, + vin VARCHAR(32) NULL DEFAULT 'unknown', + province VARCHAR(16) NULL, + city VARCHAR(16) NULL, + manufacturer VARCHAR(64) NULL, + device_type VARCHAR(64) NULL, + plate_color VARCHAR(16) NULL, + auth_token VARCHAR(128) NULL, + auth_imei VARCHAR(64) NULL, + auth_software_version VARCHAR(64) NULL, + source_endpoint VARCHAR(64) NULL, + first_registered_at DATETIME NULL, + latest_registered_at DATETIME NULL, + latest_authenticated_at DATETIME NULL, + latest_seen_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, + KEY idx_jt808_registration_device (device_id), + KEY idx_jt808_registration_plate (plate), + KEY idx_jt808_registration_vin (vin), + KEY idx_jt808_registration_seen (latest_seen_at) +)` + func safeIdentifier(value string) bool { if value == "" { return false diff --git a/go/vehicle-gateway/internal/identity/resolver_test.go b/go/vehicle-gateway/internal/identity/resolver_test.go index 085d4d39..66033371 100644 --- a/go/vehicle-gateway/internal/identity/resolver_test.go +++ b/go/vehicle-gateway/internal/identity/resolver_test.go @@ -185,6 +185,23 @@ func TestMySQLResolverTracksFirstJT808LocationThenThrottles(t *testing.T) { } } +func TestMySQLResolverEnsuresMinimalIdentitySchema(t *testing.T) { + db, mock := newMockDB(t) + defer db.Close() + mock.ExpectExec("CREATE TABLE IF NOT EXISTS vehicle_identity_binding"). + WillReturnResult(sqlmock.NewResult(0, 0)) + mock.ExpectExec("CREATE TABLE IF NOT EXISTS jt808_registration"). + WillReturnResult(sqlmock.NewResult(0, 0)) + + resolver := NewMySQLResolver(db, "vehicle_identity_binding") + if err := resolver.EnsureSchema(context.Background()); err != nil { + t.Fatalf("EnsureSchema() error = %v", err) + } + 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()