Files
lingniu-vehicle-ingest/go/vehicle-gateway/cmd/realtime-api/main.go
2026-07-03 15:52:53 +08:00

513 lines
16 KiB
Go

package main
import (
"context"
"database/sql"
"encoding/json"
"net/http"
"os"
"os/signal"
"strconv"
"strings"
"syscall"
"time"
_ "github.com/go-sql-driver/mysql"
"github.com/redis/go-redis/v9"
"github.com/segmentio/kafka-go"
_ "github.com/taosdata/driver-go/v3/taosWS"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/health"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/history"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics"
)
func main() {
logger := observability.NewLogger("vehicle-realtime-api")
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
client := redis.NewClient(&redis.Options{
Addr: env("REDIS_ADDR", "127.0.0.1:6379"),
Username: env("REDIS_USERNAME", ""),
Password: env("REDIS_PASSWORD", ""),
DB: envInt("REDIS_DB", 0),
})
defer client.Close()
if err := client.Ping(ctx).Err(); err != nil {
logger.Error("redis ping failed", "error", err)
os.Exit(1)
}
repository := realtime.NewRepository(client, realtime.Config{
OnlineTTL: time.Duration(envInt("ONLINE_TTL_SECONDS", envInt("REALTIME_ONLINE_TTL_SECONDS", 60))) * time.Second,
})
mux := http.NewServeMux()
registry := metrics.NewRegistry()
healthChecks := []health.Check{
{Name: "redis", Check: func(ctx context.Context) error { return client.Ping(ctx).Err() }},
}
mux.Handle("/metrics", metrics.NewHandler(registry))
mux.Handle("/openapi.json", realtime.NewOpenAPIHandler())
mux.Handle("/swagger-ui/", realtime.NewSwaggerUIHandler())
mux.Handle("/api/realtime/vehicles/", realtime.NewHandler(repository))
mux.Handle("/api/realtime/online", realtime.NewOnlineIndexHandler(repository))
mux.Handle("/api/debug/pipeline", realtime.NewPipelineDebugHandler(repository))
closeStats := func() {}
var updater realtimeUpdater = storeMetricUpdater{
store: "redis",
delegate: fastRealtimeUpdater{repository: repository},
registry: registry,
}
if dsn := strings.TrimSpace(os.Getenv("MYSQL_DSN")); dsn != "" {
db, err := sql.Open("mysql", dsn)
if err != nil {
logger.Error("stats mysql open failed", "error", err)
os.Exit(1)
}
if err := db.PingContext(ctx); err != nil {
_ = db.Close()
logger.Error("stats mysql ping failed", "error", err)
os.Exit(1)
}
healthChecks = append(healthChecks, health.Check{Name: "mysql", Check: db.PingContext})
closeStats = func() { _ = db.Close() }
bindingTable := env("VEHICLE_IDENTITY_TABLE", "vehicle_identity_binding")
plateResolver := realtime.NewCachedPlateResolver(
realtime.NewBindingPlateResolver(db, bindingTable),
time.Duration(envInt("PLATE_CACHE_TTL_SECONDS", 600))*time.Second,
)
snapshotWriter := realtime.NewSnapshotWriterWithPlateResolver(db, plateResolver)
if env("MYSQL_REALTIME_SNAPSHOT_ENABLED", "true") != "false" {
if err := snapshotWriter.EnsureSchema(ctx); err != nil {
_ = db.Close()
logger.Error("realtime snapshot mysql schema bootstrap failed", "error", err)
os.Exit(1)
}
mysqlDelegate := realtimeUpdater(retryRealtimeUpdater{
delegate: snapshotWriter,
attempts: envInt("MYSQL_REALTIME_RETRY_ATTEMPTS", 3),
delay: time.Duration(envInt("MYSQL_REALTIME_RETRY_DELAY_MS", 20)) * time.Millisecond,
})
mysqlUpdater := realtimeUpdater(storeMetricUpdater{store: "mysql", delegate: mysqlDelegate, registry: registry})
if env("MYSQL_REALTIME_ASYNC_ENABLED", "true") != "false" {
mysqlUpdater = newAsyncSecondaryRealtimeUpdater(
nil,
storeMetricUpdater{store: "mysql", delegate: mysqlDelegate, registry: registry},
envInt("MYSQL_REALTIME_ASYNC_QUEUE_SIZE", 20000),
envInt("MYSQL_REALTIME_ASYNC_WORKERS", 4),
registry,
logger,
)
}
updater = compositeRealtimeUpdater{
primary: storeMetricUpdater{store: "redis", delegate: fastRealtimeUpdater{repository: repository}, registry: registry},
secondary: mysqlUpdater,
}
logger.Info("realtime mysql snapshot enabled", "table", "vehicle_realtime_snapshot", "plate_binding_table", bindingTable, "plate_cache_ttl_seconds", envInt("PLATE_CACHE_TTL_SECONDS", 600))
}
mux.Handle("/api/stats/daily-metrics", stats.NewMetricHandler(stats.NewMetricRepository(db)))
mux.Handle("/api/realtime/snapshots", realtime.NewSnapshotQueryHandler(realtime.NewSnapshotQueryRepository(db)))
mux.Handle("/api/realtime/locations", realtime.NewLocationQueryHandler(realtime.NewLocationQueryRepository(db)))
logger.Info("stats mysql query enabled")
} else {
mysqlUnavailable := func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusServiceUnavailable)
_ = json.NewEncoder(w).Encode(map[string]any{"error": "MYSQL_DSN is not configured"})
}
mux.HandleFunc("/api/stats/daily-metrics", mysqlUnavailable)
mux.HandleFunc("/api/realtime/snapshots", mysqlUnavailable)
mux.HandleFunc("/api/realtime/locations", mysqlUnavailable)
logger.Warn("MYSQL_DSN is empty; stats query api disabled")
}
defer closeStats()
if brokers := splitCSV(os.Getenv("KAFKA_BROKERS")); len(brokers) > 0 {
go consumeKafka(ctx, logger, registry, updater, brokers)
} else {
logger.Warn("KAFKA_BROKERS is empty; realtime api will serve existing redis data only")
}
closeHistory := func() {}
if dsn := strings.TrimSpace(os.Getenv("TDENGINE_DSN")); dsn != "" {
driver := env("TDENGINE_DRIVER", "taosWS")
db, err := sql.Open(driver, dsn)
if err != nil {
logger.Error("history tdengine open failed", "driver", driver, "error", err)
os.Exit(1)
}
if err := db.PingContext(ctx); err != nil {
_ = db.Close()
logger.Error("history tdengine ping failed", "driver", driver, "error", err)
os.Exit(1)
}
healthChecks = append(healthChecks, health.Check{Name: "tdengine", Check: db.PingContext})
closeHistory = func() { _ = db.Close() }
database := env("TDENGINE_DATABASE", history.DefaultDatabase)
rawFrameHandler := history.NewRawFrameHandler(history.NewRawFrameRepository(db, database))
mux.Handle("/api/history/raw-frames", rawFrameHandler)
mux.Handle("/api/history/raw-frames/query", rawFrameHandler)
mux.Handle("/api/history/locations", history.NewLocationHandler(history.NewLocationRepository(db, database)))
logger.Info("history query enabled", "driver", driver, "database", database)
} else {
historyUnavailable := func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusServiceUnavailable)
_ = json.NewEncoder(w).Encode(map[string]any{"error": "TDENGINE_DSN is not configured"})
}
mux.HandleFunc("/api/history/raw-frames", historyUnavailable)
mux.HandleFunc("/api/history/raw-frames/query", historyUnavailable)
mux.HandleFunc("/api/history/locations", historyUnavailable)
logger.Warn("TDENGINE_DSN is empty; history query api disabled")
}
defer closeHistory()
healthHandler := health.NewHandler("vehicle-realtime-api", healthChecks)
mux.Handle("/healthz", healthHandler)
mux.Handle("/readyz", healthHandler)
server := &http.Server{
Addr: env("HTTP_ADDR", ":20210"),
Handler: mux,
ReadHeaderTimeout: 5 * time.Second,
}
go func() {
<-ctx.Done()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
_ = server.Shutdown(shutdownCtx)
}()
logger.Info("realtime api started", "addr", server.Addr)
if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
logger.Error("http server failed", "error", err)
os.Exit(1)
}
}
type realtimeLogger interface {
Info(string, ...any)
Error(string, ...any)
Warn(string, ...any)
}
func consumeKafka(ctx context.Context, logger realtimeLogger, registry *metrics.Registry, updater realtimeUpdater, brokers []string) {
kafkaTopics := kafkaTopicsFromEnv()
reader := kafka.NewReader(realtimeReaderConfig(brokers, kafkaTopics))
defer reader.Close()
logger.Info("realtime kafka consumer started", "topics", strings.Join(kafkaTopics, ","))
for {
message, err := reader.FetchMessage(ctx)
if err != nil {
if ctx.Err() != nil {
return
}
logger.Error("kafka fetch failed", "error", err)
continue
}
processRealtimeMessage(ctx, logger, registry, updater, reader, message)
}
}
func realtimeReaderConfig(brokers []string, kafkaTopics []string) kafka.ReaderConfig {
return kafka.ReaderConfig{
Brokers: brokers,
GroupID: env("KAFKA_GROUP", "go-realtime-api"),
GroupTopics: kafkaTopics,
MinBytes: 1,
MaxBytes: 10e6,
StartOffset: kafka.LastOffset,
}
}
func kafkaTopicsFromEnv() []string {
defaultTopics := strings.Join([]string{topics.RawGB32960, topics.RawJT808, topics.RawYutongMQTT}, ",")
return splitCSV(env("KAFKA_TOPICS", defaultTopics))
}
const kafkaMessageOperationTimeout = 30 * time.Second
type realtimeUpdater interface {
Update(context.Context, envelope.FrameEnvelope) error
}
type fastRealtimeUpdater struct {
repository *realtime.Repository
}
func (u fastRealtimeUpdater) Update(ctx context.Context, env envelope.FrameEnvelope) error {
return u.repository.FastUpdate(ctx, env)
}
type storeMetricUpdater struct {
store string
delegate realtimeUpdater
registry *metrics.Registry
}
func (u storeMetricUpdater) Update(ctx context.Context, env envelope.FrameEnvelope) error {
if u.delegate == nil {
return nil
}
started := time.Now()
err := u.delegate.Update(ctx, env)
elapsed := time.Since(started)
status := "ok"
if err != nil {
status = "error"
}
if u.registry != nil {
u.registry.IncCounter("vehicle_realtime_store_updates_total", metrics.Labels{
"store": u.store,
"protocol": string(env.Protocol),
"status": status,
})
u.registry.SetGauge("vehicle_realtime_store_update_duration_ms", metrics.Labels{
"store": u.store,
"protocol": string(env.Protocol),
"status": status,
}, float64(elapsed.Milliseconds()))
}
return err
}
type retryRealtimeUpdater struct {
delegate realtimeUpdater
attempts int
delay time.Duration
}
func (u retryRealtimeUpdater) Update(ctx context.Context, env envelope.FrameEnvelope) error {
if u.delegate == nil {
return nil
}
attempts := u.attempts
if attempts <= 0 {
attempts = 1
}
var err error
for attempt := 1; attempt <= attempts; attempt++ {
err = u.delegate.Update(ctx, env)
if err == nil || !isTransientMySQLWriteError(err) || attempt == attempts {
return err
}
if u.delay <= 0 {
continue
}
timer := time.NewTimer(u.delay)
select {
case <-ctx.Done():
timer.Stop()
return ctx.Err()
case <-timer.C:
}
}
return err
}
func isTransientMySQLWriteError(err error) bool {
text := strings.ToLower(strings.TrimSpace(err.Error()))
return strings.Contains(text, "deadlock") ||
strings.Contains(text, "error 1213") ||
strings.Contains(text, "40001") ||
strings.Contains(text, "lock wait timeout") ||
strings.Contains(text, "error 1205")
}
type compositeRealtimeUpdater struct {
primary realtimeUpdater
secondary realtimeUpdater
}
func (u compositeRealtimeUpdater) Update(ctx context.Context, env envelope.FrameEnvelope) error {
if err := u.primary.Update(ctx, env); err != nil {
return err
}
if u.secondary == nil {
return nil
}
return u.secondary.Update(ctx, env)
}
type asyncSecondaryRealtimeUpdater struct {
primary realtimeUpdater
secondary realtimeUpdater
queue chan envelope.FrameEnvelope
cancel context.CancelFunc
registry *metrics.Registry
logger realtimeLogger
}
func newAsyncSecondaryRealtimeUpdater(primary realtimeUpdater, secondary realtimeUpdater, queueSize int, workers int, registry *metrics.Registry, logger realtimeLogger) *asyncSecondaryRealtimeUpdater {
if queueSize <= 0 {
queueSize = 20000
}
if workers <= 0 {
workers = 4
}
ctx, cancel := context.WithCancel(context.Background())
updater := &asyncSecondaryRealtimeUpdater{
primary: primary,
secondary: secondary,
queue: make(chan envelope.FrameEnvelope, queueSize),
cancel: cancel,
registry: registry,
logger: logger,
}
for i := 0; i < workers; i++ {
go updater.run(ctx)
}
return updater
}
func (u *asyncSecondaryRealtimeUpdater) Update(ctx context.Context, env envelope.FrameEnvelope) error {
if u.primary != nil {
if err := u.primary.Update(ctx, env); err != nil {
return err
}
}
if u.secondary == nil {
return nil
}
select {
case u.queue <- env:
u.recordQueueMetric(env, "queued")
u.recordQueueDepth(env)
return nil
default:
u.recordQueueMetric(env, "dropped")
u.recordQueueDepth(env)
return nil
}
}
func (u *asyncSecondaryRealtimeUpdater) run(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case env := <-u.queue:
u.recordQueueDepth(env)
secondaryCtx, cancel := context.WithTimeout(context.Background(), kafkaMessageOperationTimeout)
if err := u.secondary.Update(secondaryCtx, env); err != nil && u.logger != nil {
u.logger.Error("async realtime secondary update failed", "protocol", env.Protocol, "vin", env.VIN, "event_id", env.StableEventID(), "error", err)
}
cancel()
u.recordQueueDepth(env)
}
}
}
func (u *asyncSecondaryRealtimeUpdater) recordQueueMetric(env envelope.FrameEnvelope, status string) {
if u.registry == nil {
return
}
u.registry.IncCounter("vehicle_realtime_async_queue_total", metrics.Labels{
"store": "mysql",
"protocol": string(env.Protocol),
"status": status,
})
}
func (u *asyncSecondaryRealtimeUpdater) recordQueueDepth(env envelope.FrameEnvelope) {
if u.registry == nil {
return
}
u.registry.SetGauge("vehicle_realtime_async_queue_depth", metrics.Labels{
"store": "mysql",
"protocol": string(env.Protocol),
}, float64(len(u.queue)))
}
func (u *asyncSecondaryRealtimeUpdater) Close() {
if u.cancel != nil {
u.cancel()
}
}
type kafkaMessageCommitter interface {
CommitMessages(context.Context, ...kafka.Message) error
}
func processRealtimeMessage(ctx context.Context, logger interface {
Error(string, ...any)
Warn(string, ...any)
}, registry *metrics.Registry, updater realtimeUpdater, committer kafkaMessageCommitter, message kafka.Message) {
messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout)
defer cancel()
addRealtimeMetric(registry, "vehicle_realtime_kafka_messages_total", message, "received")
addRealtimeLagMetric(registry, message)
var env envelope.FrameEnvelope
if err := json.Unmarshal(message.Value, &env); err != nil {
addRealtimeMetric(registry, "vehicle_realtime_kafka_messages_total", message, "invalid_json")
logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
_ = committer.CommitMessages(messageCtx, message)
return
}
if err := updater.Update(messageCtx, env); err != nil {
addRealtimeMetric(registry, "vehicle_realtime_updates_total", message, "error")
logger.Error("redis realtime update failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err)
return
}
addRealtimeMetric(registry, "vehicle_realtime_updates_total", message, "ok")
if err := committer.CommitMessages(messageCtx, message); err != nil {
addRealtimeMetric(registry, "vehicle_realtime_kafka_commits_total", message, "error")
logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
return
}
addRealtimeMetric(registry, "vehicle_realtime_kafka_commits_total", message, "ok")
}
func addRealtimeMetric(registry *metrics.Registry, name string, message kafka.Message, status string) {
if registry == nil {
return
}
registry.IncCounter(name, metrics.Labels{
"topic": message.Topic,
"status": status,
})
}
func addRealtimeLagMetric(registry *metrics.Registry, message kafka.Message) {
if registry == nil {
return
}
registry.SetKafkaLag("vehicle_realtime_kafka_lag", message.Topic, message.Partition, message.Offset, message.HighWaterMark)
}
func env(key string, fallback string) string {
value := strings.TrimSpace(os.Getenv(key))
if value == "" {
return fallback
}
return value
}
func envInt(key string, fallback int) int {
value := strings.TrimSpace(os.Getenv(key))
if value == "" {
return fallback
}
parsed, err := strconv.Atoi(value)
if err != nil {
return fallback
}
return parsed
}
func splitCSV(value string) []string {
var out []string
for _, item := range strings.Split(value, ",") {
item = strings.TrimSpace(item)
if item != "" {
out = append(out, item)
}
}
return out
}