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/history" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats" ) 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", 600)) * time.Second, }) if brokers := splitCSV(os.Getenv("KAFKA_BROKERS")); len(brokers) > 0 { go consumeKafka(ctx, logger, repository, brokers) } else { logger.Warn("KAFKA_BROKERS is empty; realtime api will serve existing redis data only") } mux := http.NewServeMux() mux.Handle("/api/realtime/vehicles/", realtime.NewHandler(repository)) closeStats := func() {} 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) } closeStats = func() { _ = db.Close() } mux.Handle("/api/stats/daily-metrics", stats.NewMetricHandler(stats.NewMetricRepository(db))) logger.Info("stats mysql query enabled") } else { mux.HandleFunc("/api/stats/daily-metrics", 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"}) }) logger.Warn("MYSQL_DSN is empty; stats query api disabled") } defer closeStats() 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) } closeHistory = func() { _ = db.Close() } database := env("TDENGINE_DATABASE", history.DefaultDatabase) mux.Handle("/api/history/raw-frames", history.NewRawFrameHandler(history.NewRawFrameRepository(db, database))) 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/locations", historyUnavailable) logger.Warn("TDENGINE_DSN is empty; history query api disabled") } defer closeHistory() 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) } } func consumeKafka(ctx context.Context, logger interface { Info(string, ...any) Error(string, ...any) Warn(string, ...any) }, repository *realtime.Repository, brokers []string) { reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: brokers, GroupID: env("KAFKA_GROUP", "go-realtime-api"), GroupTopics: splitCSV(env("KAFKA_TOPICS", "vehicle.event.unified.v1")), MinBytes: 1, MaxBytes: 10e6, }) defer reader.Close() logger.Info("realtime kafka consumer started", "topics", env("KAFKA_TOPICS", "vehicle.event.unified.v1")) for { message, err := reader.FetchMessage(ctx) if err != nil { if ctx.Err() != nil { return } logger.Error("kafka fetch failed", "error", err) continue } var env envelope.FrameEnvelope if err := json.Unmarshal(message.Value, &env); err != nil { logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) _ = reader.CommitMessages(ctx, message) continue } if err := repository.Update(ctx, env); err != nil { logger.Error("redis realtime update failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) continue } if err := reader.CommitMessages(ctx, message); err != nil { logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) } } } 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 }