Files
lingniu-vehicle-ingest/go/vehicle-gateway/cmd/realtime-api/main.go
2026-07-02 00:29:49 +08:00

199 lines
6.5 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/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)))
mux.Handle("/api/history/mileage-points", history.NewMileagePointHandler(history.NewMileagePointRepository(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)
mux.HandleFunc("/api/history/mileage-points", 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
}