package main import ( "context" "database/sql" "encoding/json" "errors" "fmt" "log/slog" "os" "os/signal" "strings" "syscall" "time" "github.com/nats-io/nats.go" "github.com/redis/go-redis/v9" _ "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/topics" ) func main() { logger := observability.NewLogger("nats-fast-writer") ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() cfg := loadConfig() conn, err := nats.Connect(cfg.NATSURL, nats.Name(cfg.NATSClientName), nats.Timeout(5*time.Second)) if err != nil { logger.Error("nats connect failed", "error", err) os.Exit(1) } defer conn.Close() js, err := conn.JetStream() if err != nil { logger.Error("nats jetstream init failed", "error", err) os.Exit(1) } if err := ensureStream(js, cfg); err != nil { logger.Error("nats stream ensure failed", "stream", cfg.NATSStream, "error", err) os.Exit(1) } sub, err := js.PullSubscribe( cfg.NATSFilter, cfg.NATSDurable, nats.BindStream(cfg.NATSStream), nats.ManualAck(), nats.AckWait(cfg.AckWait), nats.MaxDeliver(-1), ) if err != nil { logger.Error("nats pull consumer init failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) os.Exit(1) } tdDB, err := sql.Open(cfg.TDengineDriver, cfg.TDengineDSN) if err != nil { logger.Error("tdengine open failed", "error", err) os.Exit(1) } defer tdDB.Close() tdDB.SetMaxOpenConns(1) tdDB.SetMaxIdleConns(1) if err := tdDB.PingContext(ctx); err != nil { logger.Error("tdengine ping failed", "error", err) os.Exit(1) } historyWriter := history.NewWriter(tdDB) if cfg.TDengineEnsureSchema { if err := historyWriter.EnsureSchema(ctx, cfg.TDengineDatabase); err != nil { logger.Error("tdengine schema bootstrap failed", "error", err) os.Exit(1) } } if strings.TrimSpace(cfg.TDengineDatabase) != "" { if _, err := tdDB.ExecContext(ctx, "USE "+cfg.TDengineDatabase); err != nil { logger.Error("tdengine use database failed", "database", cfg.TDengineDatabase, "error", err) os.Exit(1) } } redisClient := redis.NewClient(&redis.Options{ Addr: cfg.RedisAddr, Username: cfg.RedisUsername, Password: cfg.RedisPassword, DB: cfg.RedisDB, }) defer redisClient.Close() if err := redisClient.Ping(ctx).Err(); err != nil { logger.Error("redis ping failed", "error", err) os.Exit(1) } realtimeRepo := realtime.NewRepository(redisClient, realtime.Config{OnlineTTL: cfg.OnlineTTL}) registry := metrics.NewRegistry() health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-fast-writer", []health.Check{ {Name: "nats", Check: func(context.Context) error { if conn.Status() != nats.CONNECTED { return fmt.Errorf("nats status is %s", conn.Status().String()) } return nil }}, {Name: "tdengine", Check: tdDB.PingContext}, {Name: "redis", Check: func(ctx context.Context) error { return redisClient.Ping(ctx).Err() }}, }, registry)) logger.Info("nats fast writer started", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "filter", cfg.NATSFilter, "batch_size", cfg.BatchSize, "workers", cfg.Workers, "operation_timeout_ms", cfg.OperationWait.Milliseconds()) for i := 0; i < cfg.Workers; i++ { go runFastWorker(ctx, logger, registry, sub, historyWriter, realtimeRepo, cfg) } <-ctx.Done() } type config struct { NATSURL string NATSClientName string NATSStream string NATSDurable string NATSFilter string NATSSubjects []string BatchSize int FetchWait time.Duration OperationWait time.Duration AckWait time.Duration StreamMaxAge time.Duration Workers int TDengineDriver string TDengineDSN string TDengineDatabase string TDengineEnsureSchema bool RedisAddr string RedisUsername string RedisPassword string RedisDB int OnlineTTL time.Duration } func loadConfig() config { subjects := splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join([]string{ env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)), env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)), env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)), }, ","))) return config{ NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-fast-writer"), NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"), NATSDurable: env("NATS_DURABLE", "vehicle-fast-writer"), NATSFilter: env("NATS_FILTER", "vehicle.raw.go.>"), NATSSubjects: subjects, BatchSize: envInt("FAST_WRITER_BATCH_SIZE", 100), FetchWait: time.Duration(envInt("FAST_WRITER_FETCH_WAIT_MS", 100)) * time.Millisecond, OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 100)) * time.Millisecond, AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second, StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour, Workers: envInt("FAST_WRITER_WORKERS", 8), TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: env("TDENGINE_DSN", ""), TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase), TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false", RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"), RedisUsername: env("REDIS_USERNAME", ""), RedisPassword: env("REDIS_PASSWORD", ""), RedisDB: envInt("REDIS_DB", 0), OnlineTTL: time.Duration(envInt("REALTIME_ONLINE_TTL_SECONDS", 60)) * time.Second, } } type fastAppender interface { AppendAll(context.Context, envelope.FrameEnvelope) error } type fastUpdater interface { FastUpdate(context.Context, envelope.FrameEnvelope) error } type fastMessage struct { subject string data []byte ack func() error } func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, sub natsPullSubscription, appender fastAppender, updater fastUpdater, cfg config) { for { if ctx.Err() != nil { return } msgs, err := sub.Fetch(cfg.BatchSize, nats.MaxWait(cfg.FetchWait)) if err != nil { if errors.Is(err, nats.ErrTimeout) { continue } logger.Error("nats fetch failed", "error", err) time.Sleep(time.Second) continue } for _, msg := range msgs { msg := msg operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) err := processFastMessage(operationCtx, appender, updater, &fastMessage{ subject: msg.Subject, data: msg.Data, ack: func() error { return msg.Ack() }, }) cancel() if err != nil { addFastMetric(registry, msg.Subject, "error") logger.Error("fast write failed", "subject", msg.Subject, "error", err) continue } addFastMetric(registry, msg.Subject, "ok") } } } type natsPullSubscription interface { Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error) } func processFastMessage(ctx context.Context, appender fastAppender, updater fastUpdater, msg *fastMessage) error { var env envelope.FrameEnvelope if err := json.Unmarshal(msg.data, &env); err != nil { if msg.ack != nil { _ = msg.ack() } return nil } if err := appender.AppendAll(ctx, env); err != nil { return fmt.Errorf("tdengine append: %w", err) } if err := updater.FastUpdate(ctx, env); err != nil { return fmt.Errorf("redis fast update: %w", err) } if msg.ack != nil { if err := msg.ack(); err != nil { return fmt.Errorf("nats ack: %w", err) } } return nil } func ensureStream(js nats.JetStreamContext, cfg config) error { stream := &nats.StreamConfig{ Name: cfg.NATSStream, Subjects: cfg.NATSSubjects, Storage: nats.FileStorage, Retention: nats.LimitsPolicy, MaxAge: cfg.StreamMaxAge, Duplicates: 2 * time.Minute, } if _, err := js.StreamInfo(cfg.NATSStream); err == nil { _, err = js.UpdateStream(stream) return err } _, err := js.AddStream(stream) if err == nil { return nil } _, updateErr := js.UpdateStream(stream) if updateErr == nil { return nil } return err } func addFastMetric(registry *metrics.Registry, subject string, status string) { if registry == nil { return } registry.IncCounter("vehicle_fast_writer_messages_total", metrics.Labels{"subject": subject, "status": status}) } 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 } var parsed int if _, err := fmt.Sscanf(value, "%d", &parsed); 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 }