feat(go): bootstrap minimal identity schema
This commit is contained in:
@@ -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。
|
||||
|
||||
## 字段提升规则
|
||||
|
||||
协议字段进入系统后分三层:
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user