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

296 lines
10 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", 600)) * 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))
closeStats := func() {}
var updater realtimeUpdater = repository
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")
snapshotWriter := realtime.NewSnapshotWriterWithPlateResolver(db, realtime.NewBindingPlateResolver(db, bindingTable))
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)
}
updater = compositeRealtimeUpdater{primary: repository, secondary: snapshotWriter}
logger.Info("realtime mysql snapshot enabled", "table", "vehicle_realtime_snapshot", "plate_binding_table", bindingTable)
}
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)
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()
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)
}
}
func consumeKafka(ctx context.Context, logger interface {
Info(string, ...any)
Error(string, ...any)
Warn(string, ...any)
}, registry *metrics.Registry, updater realtimeUpdater, brokers []string) {
kafkaTopics := kafkaTopicsFromEnv()
reader := kafka.NewReader(kafka.ReaderConfig{
Brokers: brokers,
GroupID: env("KAFKA_GROUP", "go-realtime-api"),
GroupTopics: kafkaTopics,
MinBytes: 1,
MaxBytes: 10e6,
})
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 kafkaTopicsFromEnv() []string {
return splitCSV(env("KAFKA_TOPICS", topics.Unified))
}
const kafkaMessageOperationTimeout = 30 * time.Second
type realtimeUpdater interface {
Update(context.Context, envelope.FrameEnvelope) error
}
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 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
}