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))) mux.Handle("/api/realtime/kv", realtime.NewKVQueryHandler(realtime.NewKVQueryRepository(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) mux.HandleFunc("/api/realtime/kv", 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 }