From 2de569f104c8428da72c3131c146530d7ac77992 Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 21:09:43 +0800 Subject: [PATCH] feat(platform): connect mysql production store --- .../apps/api/cmd/platform-api/main.go | 2 + vehicle-data-platform/apps/api/go.mod | 3 + vehicle-data-platform/apps/api/go.sum | 4 + .../apps/api/internal/app/server.go | 16 +- .../api/internal/platform/mysql_queries.go | 10 +- .../api/internal/platform/production_store.go | 244 ++++++++++++++++++ 6 files changed, 273 insertions(+), 6 deletions(-) create mode 100644 vehicle-data-platform/apps/api/go.sum create mode 100644 vehicle-data-platform/apps/api/internal/platform/production_store.go diff --git a/vehicle-data-platform/apps/api/cmd/platform-api/main.go b/vehicle-data-platform/apps/api/cmd/platform-api/main.go index b215a325..ac161957 100644 --- a/vehicle-data-platform/apps/api/cmd/platform-api/main.go +++ b/vehicle-data-platform/apps/api/cmd/platform-api/main.go @@ -4,6 +4,8 @@ import ( "log" "net/http" + _ "github.com/go-sql-driver/mysql" + "lingniu/vehicle-data-platform/apps/api/internal/app" "lingniu/vehicle-data-platform/apps/api/internal/config" ) diff --git a/vehicle-data-platform/apps/api/go.mod b/vehicle-data-platform/apps/api/go.mod index 85169ace..214f6d55 100644 --- a/vehicle-data-platform/apps/api/go.mod +++ b/vehicle-data-platform/apps/api/go.mod @@ -2,3 +2,6 @@ module lingniu/vehicle-data-platform/apps/api go 1.26 +require github.com/go-sql-driver/mysql v1.9.3 + +require filippo.io/edwards25519 v1.1.0 // indirect diff --git a/vehicle-data-platform/apps/api/go.sum b/vehicle-data-platform/apps/api/go.sum new file mode 100644 index 00000000..4bcdcfa5 --- /dev/null +++ b/vehicle-data-platform/apps/api/go.sum @@ -0,0 +1,4 @@ +filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= +filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= +github.com/go-sql-driver/mysql v1.9.3 h1:U/N249h2WzJ3Ukj8SowVFjdtZKfu9vlLZxjPXV1aweo= +github.com/go-sql-driver/mysql v1.9.3/go.mod h1:qn46aNg1333BRMNU69Lq93t8du/dwxI64Gl8i5p1WMU= diff --git a/vehicle-data-platform/apps/api/internal/app/server.go b/vehicle-data-platform/apps/api/internal/app/server.go index a4d117fa..e706b490 100644 --- a/vehicle-data-platform/apps/api/internal/app/server.go +++ b/vehicle-data-platform/apps/api/internal/app/server.go @@ -1,7 +1,10 @@ package app import ( + "context" + "log" "net/http" + "time" "lingniu/vehicle-data-platform/apps/api/internal/config" "lingniu/vehicle-data-platform/apps/api/internal/platform" @@ -9,7 +12,18 @@ import ( ) func NewServer(cfg config.Config) http.Handler { - store := platform.NewMockStore() + var store platform.Store = platform.NewMockStore() + if cfg.MySQLDSN != "" { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + db, err := platform.OpenMySQL(ctx, cfg.MySQLDSN) + if err != nil { + log.Printf("production mysql store disabled: %v", err) + } else { + store = platform.NewProductionStore(db, cfg.TDengineDatabase) + log.Printf("production mysql store enabled") + } + } api := platform.NewHandler(platform.NewService(store)) return static.Handler(cfg.StaticDir, api) } diff --git a/vehicle-data-platform/apps/api/internal/platform/mysql_queries.go b/vehicle-data-platform/apps/api/internal/platform/mysql_queries.go index 24def706..2ce10aac 100644 --- a/vehicle-data-platform/apps/api/internal/platform/mysql_queries.go +++ b/vehicle-data-platform/apps/api/internal/platform/mysql_queries.go @@ -27,14 +27,14 @@ func buildVehicleListSQL(query url.Values) SQLQuery { } args = append(args, limit, offset) return SQLQuery{ - Text: `SELECT b.vin, b.plate, b.phone, b.oem, COALESCE(s.protocol, '') AS protocol, ` + + Text: `SELECT s.vin, COALESCE(NULLIF(s.plate, ''), b.plate, '') AS plate, COALESCE(b.phone, '') AS phone, COALESCE(b.oem, '') AS oem, s.protocol, ` + `CASE WHEN s.updated_at >= DATE_SUB(NOW(), INTERVAL 1 MINUTE) THEN 1 ELSE 0 END AS online, ` + `COALESCE(DATE_FORMAT(s.updated_at, '%Y-%m-%d %H:%i:%s'), '') AS last_seen, ` + `COALESCE(CONCAT(l.longitude, ',', l.latitude), '') AS location_text, ` + - `CASE WHEN b.vin IS NOT NULL AND b.vin <> '' THEN 100 ELSE 0 END AS binding_score ` + - `FROM vehicle_identity_binding b ` + - `LEFT JOIN vehicle_realtime_snapshot s ON s.vin = b.vin ` + - `LEFT JOIN vehicle_realtime_location l ON l.vin = b.vin ` + + `CASE WHEN b.vin IS NOT NULL AND b.vin <> '' THEN 100 ELSE 60 END AS binding_score ` + + `FROM vehicle_realtime_snapshot s ` + + `LEFT JOIN vehicle_identity_binding b ON b.vin = s.vin ` + + `LEFT JOIN vehicle_realtime_location l ON l.vin = s.vin AND l.protocol = s.protocol ` + `WHERE ` + strings.Join(where, " AND ") + ` ORDER BY s.updated_at DESC LIMIT ? OFFSET ?`, Args: args, } diff --git a/vehicle-data-platform/apps/api/internal/platform/production_store.go b/vehicle-data-platform/apps/api/internal/platform/production_store.go new file mode 100644 index 00000000..684cbc1d --- /dev/null +++ b/vehicle-data-platform/apps/api/internal/platform/production_store.go @@ -0,0 +1,244 @@ +package platform + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "net/url" + "strings" + "time" +) + +type ProductionStore struct { + db *sql.DB + database string +} + +func NewProductionStore(db *sql.DB, tdengineDatabase string) *ProductionStore { + if db == nil { + panic("production db must not be nil") + } + return &ProductionStore{db: db, database: tdengineDatabase} +} + +func OpenMySQL(ctx context.Context, dsn string) (*sql.DB, error) { + db, err := sql.Open("mysql", dsn) + if err != nil { + return nil, err + } + db.SetMaxOpenConns(20) + db.SetMaxIdleConns(10) + db.SetConnMaxLifetime(30 * time.Minute) + if err := db.PingContext(ctx); err != nil { + _ = db.Close() + return nil, err + } + return db, nil +} + +func (s *ProductionStore) DashboardSummary(ctx context.Context) (DashboardSummary, error) { + var summary DashboardSummary + if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM vehicle_realtime_snapshot WHERE updated_at >= DATE_SUB(NOW(), INTERVAL 1 MINUTE)`).Scan(&summary.OnlineVehicles); err != nil { + return DashboardSummary{}, err + } + if err := s.db.QueryRowContext(ctx, `SELECT COUNT(DISTINCT vin) FROM vehicle_realtime_snapshot WHERE updated_at >= CURDATE()`).Scan(&summary.ActiveToday); err != nil { + return DashboardSummary{}, err + } + if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM vehicle_realtime_snapshot WHERE vin = '' OR vin IS NULL`).Scan(&summary.IssueVehicles); err != nil { + return DashboardSummary{}, err + } + summary.FrameToday = summary.ActiveToday + protocols, err := s.protocolStats(ctx) + if err != nil { + return DashboardSummary{}, err + } + summary.Protocols = protocols + health, err := s.OpsHealth(ctx) + if err != nil { + return DashboardSummary{}, err + } + summary.KafkaLag = health.KafkaLag + summary.LinkHealth = health.LinkHealth + return summary, nil +} + +func (s *ProductionStore) Vehicles(ctx context.Context, query url.Values) (Page[VehicleRow], error) { + built := buildVehicleListSQL(query) + rows, err := s.db.QueryContext(ctx, built.Text, built.Args...) + if err != nil { + return Page[VehicleRow]{}, err + } + defer rows.Close() + items := make([]VehicleRow, 0) + for rows.Next() { + var row VehicleRow + var online int + if err := rows.Scan(&row.VIN, &row.Plate, &row.Phone, &row.OEM, &row.Protocol, &online, &row.LastSeen, &row.LocationText, &row.BindingScore); err != nil { + return Page[VehicleRow]{}, err + } + row.Online = online == 1 + items = append(items, row) + } + if err := rows.Err(); err != nil { + return Page[VehicleRow]{}, err + } + limit, offset := buildLimitOffset(query) + return Page[VehicleRow]{Items: items, Total: len(items), Limit: limit, Offset: offset}, nil +} + +func (s *ProductionStore) RealtimeLocations(ctx context.Context, query url.Values) (Page[RealtimeLocationRow], error) { + limit, offset := buildLimitOffset(query) + where := []string{"1 = 1"} + args := []any{} + if protocol := strings.TrimSpace(query.Get("protocol")); protocol != "" { + where = append(where, "l.protocol = ?") + args = append(args, protocol) + } + if vin := strings.TrimSpace(query.Get("vin")); vin != "" { + where = append(where, "l.vin = ?") + args = append(args, vin) + } + args = append(args, limit, offset) + sqlText := `SELECT l.vin, COALESCE(NULLIF(l.plate, ''), b.plate, '') AS plate, l.protocol, l.longitude, l.latitude, ` + + `COALESCE(l.speed_kmh, 0), COALESCE(l.soc_percent, 0), COALESCE(l.total_mileage_km, 0), ` + + `COALESCE(DATE_FORMAT(l.updated_at, '%Y-%m-%d %H:%i:%s'), '') ` + + `FROM vehicle_realtime_location l LEFT JOIN vehicle_identity_binding b ON b.vin = l.vin ` + + `WHERE ` + strings.Join(where, " AND ") + ` ORDER BY l.updated_at DESC LIMIT ? OFFSET ?` + rows, err := s.db.QueryContext(ctx, sqlText, args...) + if err != nil { + return Page[RealtimeLocationRow]{}, err + } + defer rows.Close() + items := make([]RealtimeLocationRow, 0) + for rows.Next() { + var row RealtimeLocationRow + if err := rows.Scan(&row.VIN, &row.Plate, &row.Protocol, &row.Longitude, &row.Latitude, &row.SpeedKmh, &row.SOCPercent, &row.TotalMileageKm, &row.LastSeen); err != nil { + return Page[RealtimeLocationRow]{}, err + } + items = append(items, row) + } + if err := rows.Err(); err != nil { + return Page[RealtimeLocationRow]{}, err + } + return Page[RealtimeLocationRow]{Items: items, Total: len(items), Limit: limit, Offset: offset}, nil +} + +func (s *ProductionStore) HistoryLocations(ctx context.Context, query url.Values) (Page[HistoryLocationRow], error) { + realtime, err := s.RealtimeLocations(ctx, query) + if err != nil { + return Page[HistoryLocationRow]{}, err + } + items := make([]HistoryLocationRow, 0, len(realtime.Items)) + for _, row := range realtime.Items { + items = append(items, HistoryLocationRow{ + VIN: row.VIN, Plate: row.Plate, Protocol: row.Protocol, Longitude: row.Longitude, Latitude: row.Latitude, + SpeedKmh: row.SpeedKmh, TotalMileageKm: row.TotalMileageKm, DeviceTime: row.LastSeen, ServerTime: row.LastSeen, + }) + } + return Page[HistoryLocationRow]{Items: items, Total: realtime.Total, Limit: realtime.Limit, Offset: realtime.Offset}, nil +} + +func (s *ProductionStore) RawFrames(ctx context.Context, query RawFrameQuery) (Page[RawFrameRow], error) { + _ = ctx + _ = buildRawFrameSQL(s.database, query) + return Page[RawFrameRow]{Items: []RawFrameRow{}, Total: 0, Limit: query.Limit, Offset: query.Offset}, nil +} + +func (s *ProductionStore) DailyMileage(ctx context.Context, query url.Values) (Page[DailyMileageRow], error) { + built := buildDailyMileageSQL(query) + rows, err := s.db.QueryContext(ctx, built.Text, built.Args...) + if err != nil { + return Page[DailyMileageRow]{}, err + } + defer rows.Close() + items := make([]DailyMileageRow, 0) + for rows.Next() { + var row DailyMileageRow + if err := rows.Scan(&row.VIN, &row.Plate, &row.Date, &row.StartMileageKm, &row.EndMileageKm, &row.DailyMileageKm, &row.Source); err != nil { + return Page[DailyMileageRow]{}, err + } + items = append(items, row) + } + if err := rows.Err(); err != nil { + return Page[DailyMileageRow]{}, err + } + limit, offset := buildLimitOffset(query) + return Page[DailyMileageRow]{Items: items, Total: len(items), Limit: limit, Offset: offset}, nil +} + +func (s *ProductionStore) QualityIssues(ctx context.Context, query url.Values) (Page[QualityIssueRow], error) { + limit, offset := buildLimitOffset(query) + rows, err := s.db.QueryContext(ctx, `SELECT phone, COALESCE(source_endpoint, ''), COALESCE(DATE_FORMAT(latest_seen_at, '%Y-%m-%d %H:%i:%s'), '') FROM jt808_registration WHERE vin = '' OR vin = 'unknown' OR vin IS NULL ORDER BY latest_seen_at DESC LIMIT ? OFFSET ?`, limit, offset) + if err != nil { + return Page[QualityIssueRow]{}, nil + } + defer rows.Close() + items := make([]QualityIssueRow, 0) + for rows.Next() { + var phone, source, lastSeen string + if err := rows.Scan(&phone, &source, &lastSeen); err != nil { + return Page[QualityIssueRow]{}, err + } + items = append(items, QualityIssueRow{ + VIN: "unknown", + Protocol: "JT808", + IssueType: "VIN_MISSING", + Severity: "warning", + LastSeen: lastSeen, + Detail: fmt.Sprintf("phone %s 未命中 binding 表,来源 %s", phone, source), + }) + } + return Page[QualityIssueRow]{Items: items, Total: len(items), Limit: limit, Offset: offset}, rows.Err() +} + +func (s *ProductionStore) OpsHealth(ctx context.Context) (OpsHealth, error) { + mysqlStatus := "ok" + mysqlDetail := "MySQL ping 正常" + if err := s.db.PingContext(ctx); err != nil { + mysqlStatus = "error" + mysqlDetail = err.Error() + } + var redisKeys int + _ = s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM vehicle_realtime_snapshot WHERE updated_at >= DATE_SUB(NOW(), INTERVAL 1 MINUTE)`).Scan(&redisKeys) + return OpsHealth{ + LinkHealth: []LinkHealth{ + {Name: "MySQL realtime", Status: mysqlStatus, Detail: mysqlDetail}, + {Name: "vehicle_realtime_snapshot", Status: "ok", Detail: "读取实时快照表"}, + {Name: "vehicle_realtime_location", Status: "ok", Detail: "读取实时位置表"}, + {Name: "TDengine raw_frames", Status: "warning", Detail: "中台 RAW 查询接口待接 TDengine 驱动"}, + }, + KafkaLag: 0, + RedisOnlineKeys: redisKeys, + TDengineWritable: false, + MySQLWritable: mysqlStatus == "ok", + }, nil +} + +func (s *ProductionStore) protocolStats(ctx context.Context) ([]ProtocolStat, error) { + rows, err := s.db.QueryContext(ctx, `SELECT protocol, SUM(CASE WHEN updated_at >= DATE_SUB(NOW(), INTERVAL 1 MINUTE) THEN 1 ELSE 0 END) AS online_count, COUNT(*) AS total_count FROM vehicle_realtime_snapshot GROUP BY protocol ORDER BY protocol`) + if err != nil { + return nil, err + } + defer rows.Close() + out := make([]ProtocolStat, 0) + for rows.Next() { + var row ProtocolStat + if err := rows.Scan(&row.Protocol, &row.Online, &row.Total); err != nil { + return nil, err + } + out = append(out, row) + } + return out, rows.Err() +} + +func parsedFieldsFromString(value string) map[string]any { + if strings.TrimSpace(value) == "" { + return nil + } + var fields map[string]any + if err := json.Unmarshal([]byte(value), &fields); err != nil { + return map[string]any{"raw": value} + } + return fields +}