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) } var historyWriter fastAppender var tdCheck health.Check if cfg.TDengineEnabled { 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(cfg.TDengineMaxOpenConns) tdDB.SetMaxIdleConns(cfg.TDengineMaxIdleConns) if err := tdDB.PingContext(ctx); err != nil { logger.Error("tdengine ping failed", "error", err) os.Exit(1) } writer := history.NewWriterWithDatabase(tdDB, cfg.TDengineDatabase) if cfg.TDengineEnsureSchema { if err := writer.EnsureSchema(ctx, cfg.TDengineDatabase); err != nil { logger.Error("tdengine schema bootstrap failed", "error", err) os.Exit(1) } } historyWriter = writer tdCheck = health.Check{Name: "tdengine", Check: tdDB.PingContext} } else { logger.Info("nats fast writer tdengine stage disabled") } 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() recordFastWriterConfigMetrics(registry, cfg) registry.SetGauge("vehicle_fast_writer_tdengine_enabled", nil, boolGauge(cfg.TDengineEnabled)) healthChecks := []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: "redis", Check: func(ctx context.Context) error { return redisClient.Ping(ctx).Err() }}, } if tdCheck.Name != "" { healthChecks = append(healthChecks, tdCheck) } health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-fast-writer", healthChecks, registry)) logger.Info("nats fast writer started", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "filter", cfg.NATSFilter, "batch_size", cfg.BatchSize, "fetch_wait_ms", cfg.FetchWait.Milliseconds(), "workers", cfg.Workers, "tdengine_enabled", cfg.TDengineEnabled, "tdengine_max_open_conns", cfg.TDengineMaxOpenConns, "tdengine_max_idle_conns", cfg.TDengineMaxIdleConns, "operation_timeout_ms", cfg.OperationWait.Milliseconds()) for i := 0; i < cfg.Workers; i++ { go runFastWorker(ctx, logger, registry, js, 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 StreamMaxBytes int64 StreamEnsureWait time.Duration Workers int TDengineEnabled bool TDengineDriver string TDengineDSN string TDengineDatabase string TDengineEnsureSchema bool TDengineMaxOpenConns int TDengineMaxIdleConns int 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)), env("NATS_SUBJECT_GB32960_FIELDS", env("KAFKA_TOPIC_GB32960_FIELDS", topics.FieldsGB32960)), env("NATS_SUBJECT_JT808_FIELDS", env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808)), env("NATS_SUBJECT_YUTONG_MQTT_FIELDS", env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT)), }, ","))) workers := envInt("FAST_WRITER_WORKERS", 8) maxOpenConns := envInt("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", 1) maxIdleConns := envInt("FAST_WRITER_TDENGINE_MAX_IDLE_CONNS", maxOpenConns) 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", 20)) * time.Millisecond, OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 1000)) * 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, StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024), StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second, Workers: workers, TDengineEnabled: envBool("FAST_WRITER_TDENGINE_ENABLED", false), TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: env("TDENGINE_DSN", ""), TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase), TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false", TDengineMaxOpenConns: maxOpenConns, TDengineMaxIdleConns: maxIdleConns, 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 AppendAllBatch(context.Context, []envelope.FrameEnvelope) error } type fastUpdater interface { FastUpdate(context.Context, envelope.FrameEnvelope) error } type fastResultUpdater interface { FastUpdateWithResult(context.Context, envelope.FrameEnvelope) (realtime.FastUpdateResult, error) } type fastBatchUpdater interface { FastUpdateBatch(context.Context, []envelope.FrameEnvelope) error } type fastBatchResultUpdater interface { FastUpdateBatchWithResult(context.Context, []envelope.FrameEnvelope) (realtime.FastUpdateResult, error) } type fastMessage struct { subject string data []byte ack func() error } func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, infoReader natsConsumerInfoReader, sub natsPullSubscription, appender fastAppender, updater fastUpdater, cfg config) { var lastConsumerInfoAt time.Time for { if ctx.Err() != nil { return } if time.Since(lastConsumerInfoAt) >= time.Duration(envInt("NATS_CONSUMER_METRICS_INTERVAL_SECONDS", 10))*time.Second { lastConsumerInfoAt = time.Now() if infoReader != nil { info, err := infoReader.ConsumerInfo(cfg.NATSStream, cfg.NATSDurable, nats.Context(ctx)) if err != nil { logger.Warn("nats consumer info failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) } else { recordFastNATSConsumerInfoMetrics(registry, cfg, info) } } } msgs, err := sub.Fetch(cfg.BatchSize, nats.MaxWait(cfg.FetchWait)) if err != nil { if isFastWorkerShutdownFetchError(ctx, err) { return } if errors.Is(err, nats.ErrTimeout) { continue } if isTransientFastFetchError(err) { logger.Warn("nats fetch interrupted", "error", err) time.Sleep(time.Second) continue } logger.Error("nats fetch failed", "error", err) time.Sleep(time.Second) continue } fastMessages := make([]*fastMessage, 0, len(msgs)) for _, msg := range msgs { natsMsg := msg fastMessages = append(fastMessages, &fastMessage{ subject: natsMsg.Subject, data: natsMsg.Data, ack: func() error { return natsMsg.Ack() }, }) } operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) err = processFastBatch(operationCtx, registry, appender, updater, fastMessages) cancel() if err != nil { addFastMetric(registry, fastBatchSubject(fastMessages), "error") logger.Error("fast write batch failed", "messages", len(fastMessages), "error", err) continue } } } func isFastWorkerShutdownFetchError(ctx context.Context, err error) bool { if err == nil { return false } if ctx.Err() != nil { return true } return errors.Is(err, nats.ErrConnectionClosed) } func isTransientFastFetchError(err error) bool { if err == nil { return false } text := strings.ToLower(strings.TrimSpace(err.Error())) return strings.Contains(text, "disconnected during fetch") || strings.Contains(text, "connection closed") || strings.Contains(text, "connection reset") || strings.Contains(text, "broken pipe") || strings.Contains(text, "temporary") || strings.Contains(text, "temporarily") || strings.Contains(text, "timeout") } type natsPullSubscription interface { Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error) } type natsConsumerInfoReader interface { ConsumerInfo(stream string, name string, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) } var fastWriterStageDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000} var fastWriterRedisE2EDurationBucketsMS = []float64{10, 25, 50, 100, 250, 500, 1000, 2500, 5000, 10000} var fastWriterRedisE2ERecent = metrics.NewRecentLatencyByKey(512) var fastBatchPending = metrics.PendingPairGauge{} func processFastBatch(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, messages []*fastMessage) error { if len(messages) == 0 { return nil } for _, msg := range messages { addFastMetric(registry, msg.subject, "received") } envelopes := make([]envelope.FrameEnvelope, 0, len(messages)) validMessages := make([]*fastMessage, 0, len(messages)) for _, msg := range messages { var env envelope.FrameEnvelope if err := json.Unmarshal(msg.data, &env); err != nil { addFastMetric(registry, msg.subject, "invalid_json") if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack invalid json: %w", err) } } continue } if status, err := topics.ValidateRawEnvelope(msg.subject, env); err != nil { addFastMetric(registry, msg.subject, status) if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack mismatched raw envelope: %w", err) } } continue } envelopes = append(envelopes, env) validMessages = append(validMessages, msg) } if len(envelopes) == 0 { return nil } addFastBatchPending(registry, len(messages), len(envelopes)) defer addFastBatchPending(registry, -len(messages), -len(envelopes)) subject := fastBatchSubject(validMessages) if appender != nil { started := time.Now() err := appender.AppendAllBatch(ctx, envelopes) recordFastWriterStageDuration(registry, subject, "tdengine", statusFromError(err), time.Since(started)) if err != nil { if shouldFallbackFastBatchError(err) { if fallbackErr := processFastMessagesIndividually(ctx, registry, appender, updater, validMessages); fallbackErr != nil { addFastBatchFallbackMetric(registry, "tdengine", "error") return fmt.Errorf("tdengine batch fallback: %w", fallbackErr) } addFastBatchFallbackMetric(registry, "tdengine", "ok") return nil } addFastBatchFallbackMetric(registry, "tdengine", "skipped_transient") if catchupErr := updateFastRedisBatchWithoutAck(ctx, registry, updater, validMessages, envelopes); catchupErr != nil { addFastDecoupledUpdateMetric(registry, "tdengine_transient", "error") return fmt.Errorf("tdengine batch append: %w; redis catchup: %v", err, catchupErr) } addFastDecoupledUpdateMetric(registry, "tdengine_transient", "ok") return fmt.Errorf("tdengine batch append: %w", err) } } if batchUpdater, ok := updater.(fastBatchResultUpdater); ok { for _, group := range fastSubjectGroups(validMessages, envelopes) { started := time.Now() result, err := batchUpdater.FastUpdateBatchWithResult(ctx, group.envelopes) recordFastWriterStageDuration(registry, group.subject, "redis", statusFromError(err), time.Since(started)) if err != nil { if shouldFallbackFastBatchError(err) { if fallbackErr := processFastMessagesIndividually(ctx, registry, nil, updater, validMessages); fallbackErr != nil { addFastBatchFallbackMetric(registry, "redis", "error") return fmt.Errorf("redis fast batch fallback: %w", fallbackErr) } addFastBatchFallbackMetric(registry, "redis", "ok") return nil } addFastBatchFallbackMetric(registry, "redis", "skipped_transient") return fmt.Errorf("redis fast batch update: %w", err) } recordFastWriterRedisEnvelopeMetrics(registry, group.subject, result) recordFastWriterRedisFieldMetrics(registry, group.subject, result) recordFastWriterRedisE2EDurationMessages(registry, group.messages, group.envelopes, group.subject) } for _, msg := range validMessages { if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack: %w", err) } } addFastMetric(registry, msg.subject, "ok") } return nil } if batchUpdater, ok := updater.(fastBatchUpdater); ok { for _, group := range fastSubjectGroups(validMessages, envelopes) { started := time.Now() err := batchUpdater.FastUpdateBatch(ctx, group.envelopes) recordFastWriterStageDuration(registry, group.subject, "redis", statusFromError(err), time.Since(started)) if err != nil { if shouldFallbackFastBatchError(err) { if fallbackErr := processFastMessagesIndividually(ctx, registry, nil, updater, validMessages); fallbackErr != nil { addFastBatchFallbackMetric(registry, "redis", "error") return fmt.Errorf("redis fast batch fallback: %w", fallbackErr) } addFastBatchFallbackMetric(registry, "redis", "ok") return nil } addFastBatchFallbackMetric(registry, "redis", "skipped_transient") return fmt.Errorf("redis fast batch update: %w", err) } recordFastWriterRedisE2EDurationMessages(registry, group.messages, group.envelopes, group.subject) } for _, msg := range validMessages { if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack: %w", err) } } addFastMetric(registry, msg.subject, "ok") } return nil } for i, env := range envelopes { msg := validMessages[i] started := time.Now() var result realtime.FastUpdateResult var err error if resultUpdater, ok := updater.(fastResultUpdater); ok { result, err = resultUpdater.FastUpdateWithResult(ctx, env) } else { err = updater.FastUpdate(ctx, env) } recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started)) if err != nil { return fmt.Errorf("redis fast update: %w", err) } recordFastWriterRedisEnvelopeMetrics(registry, msg.subject, result) recordFastWriterRedisFieldMetrics(registry, msg.subject, result) recordFastWriterRedisE2EDuration(registry, msg.subject, env) if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack: %w", err) } } addFastMetric(registry, msg.subject, "ok") } return nil } func processFastMessagesIndividually(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, messages []*fastMessage) error { for _, msg := range messages { if err := processFastMessageWithReceived(ctx, registry, appender, updater, msg, false); err != nil { return err } } return nil } func processFastMessage(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage) error { return processFastMessageWithReceived(ctx, registry, appender, updater, msg, true) } func processFastMessageWithReceived(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage, recordReceived bool) error { if recordReceived { addFastMetric(registry, msg.subject, "received") } var env envelope.FrameEnvelope if err := json.Unmarshal(msg.data, &env); err != nil { addFastMetric(registry, msg.subject, "invalid_json") if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack invalid json: %w", err) } } return nil } if status, err := topics.ValidateRawEnvelope(msg.subject, env); err != nil { addFastMetric(registry, msg.subject, status) if msg.ack != nil { started := time.Now() err := msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack mismatched raw envelope: %w", err) } } return nil } if appender != nil { started := time.Now() err := appender.AppendAll(ctx, env) recordFastWriterStageDuration(registry, msg.subject, "tdengine", statusFromError(err), time.Since(started)) if err != nil { if isTransientFastBatchError(err) { if catchupErr := updateFastRedisSingleWithoutAck(ctx, registry, updater, msg.subject, env); catchupErr != nil { addFastDecoupledUpdateMetric(registry, "tdengine_transient", "error") return fmt.Errorf("tdengine append: %w; redis catchup: %v", err, catchupErr) } addFastDecoupledUpdateMetric(registry, "tdengine_transient", "ok") } return fmt.Errorf("tdengine append: %w", err) } } started := time.Now() var result realtime.FastUpdateResult var err error if resultUpdater, ok := updater.(fastResultUpdater); ok { result, err = resultUpdater.FastUpdateWithResult(ctx, env) } else { err = updater.FastUpdate(ctx, env) } recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started)) if err != nil { return fmt.Errorf("redis fast update: %w", err) } recordFastWriterRedisEnvelopeMetrics(registry, msg.subject, result) recordFastWriterRedisFieldMetrics(registry, msg.subject, result) recordFastWriterRedisE2EDuration(registry, msg.subject, env) if msg.ack != nil { started = time.Now() err = msg.ack() recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) if err != nil { addFastMetric(registry, msg.subject, "ack_error") return fmt.Errorf("nats ack: %w", err) } } addFastMetric(registry, msg.subject, "ok") return nil } func updateFastRedisBatchWithoutAck(ctx context.Context, registry *metrics.Registry, updater fastUpdater, messages []*fastMessage, envelopes []envelope.FrameEnvelope) error { if len(envelopes) == 0 { return nil } subject := fastBatchSubject(messages) if batchUpdater, ok := updater.(fastBatchResultUpdater); ok { started := time.Now() result, err := batchUpdater.FastUpdateBatchWithResult(ctx, envelopes) recordFastWriterStageDuration(registry, subject, "redis", statusFromError(err), time.Since(started)) if err == nil { recordFastWriterRedisEnvelopeMetrics(registry, subject, result) recordFastWriterRedisFieldMetrics(registry, subject, result) recordFastWriterRedisE2EDurationMessages(registry, messages, envelopes, subject) return nil } if shouldFallbackFastBatchError(err) { addFastBatchFallbackMetric(registry, "redis_catchup", "attempted") return updateFastRedisSinglesWithoutAck(ctx, registry, updater, messages, envelopes) } return err } if batchUpdater, ok := updater.(fastBatchUpdater); ok { started := time.Now() err := batchUpdater.FastUpdateBatch(ctx, envelopes) recordFastWriterStageDuration(registry, subject, "redis", statusFromError(err), time.Since(started)) if err == nil { recordFastWriterRedisE2EDurationMessages(registry, messages, envelopes, subject) return nil } if shouldFallbackFastBatchError(err) { addFastBatchFallbackMetric(registry, "redis_catchup", "attempted") return updateFastRedisSinglesWithoutAck(ctx, registry, updater, messages, envelopes) } return err } return updateFastRedisSinglesWithoutAck(ctx, registry, updater, messages, envelopes) } func updateFastRedisSinglesWithoutAck(ctx context.Context, registry *metrics.Registry, updater fastUpdater, messages []*fastMessage, envelopes []envelope.FrameEnvelope) error { for index, env := range envelopes { subject := fastMessageSubject(messages, index, fastBatchSubject(messages)) if err := updateFastRedisSingleWithoutAck(ctx, registry, updater, subject, env); err != nil { return err } } return nil } func updateFastRedisSingleWithoutAck(ctx context.Context, registry *metrics.Registry, updater fastUpdater, subject string, env envelope.FrameEnvelope) error { started := time.Now() var result realtime.FastUpdateResult var err error if resultUpdater, ok := updater.(fastResultUpdater); ok { result, err = resultUpdater.FastUpdateWithResult(ctx, env) } else { err = updater.FastUpdate(ctx, env) } recordFastWriterStageDuration(registry, subject, "redis", statusFromError(err), time.Since(started)) if err != nil { return err } recordFastWriterRedisEnvelopeMetrics(registry, subject, result) recordFastWriterRedisFieldMetrics(registry, subject, result) recordFastWriterRedisE2EDuration(registry, subject, env) return nil } func shouldFallbackFastBatchError(err error) bool { if err == nil || errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { return false } return !isTransientFastBatchError(err) } func isTransientFastBatchError(err error) bool { if err == nil { return false } text := strings.ToLower(strings.TrimSpace(err.Error())) return strings.Contains(text, "timeout") || strings.Contains(text, "temporary") || strings.Contains(text, "temporarily") || strings.Contains(text, "connection refused") || strings.Contains(text, "connection reset") || strings.Contains(text, "connection closed") || strings.Contains(text, "broken pipe") || strings.Contains(text, "bad connection") || strings.Contains(text, "i/o timeout") || text == "eof" || strings.Contains(text, "unexpected eof") || strings.Contains(text, "server is down") || strings.Contains(text, "network is unreachable") || strings.Contains(text, "no route to host") } func recordFastWriterConfigMetrics(registry *metrics.Registry, cfg config) { if registry == nil { return } registry.SetGauge("vehicle_fast_writer_config", metrics.Labels{"setting": "workers"}, float64(cfg.Workers)) registry.SetGauge("vehicle_fast_writer_config", metrics.Labels{"setting": "batch_size"}, float64(cfg.BatchSize)) registry.SetGauge("vehicle_fast_writer_config", metrics.Labels{"setting": "fetch_wait_ms"}, float64(cfg.FetchWait.Milliseconds())) registry.SetGauge("vehicle_fast_writer_config", metrics.Labels{"setting": "operation_timeout_ms"}, float64(cfg.OperationWait.Milliseconds())) } func boolGauge(value bool) float64 { if value { return 1 } return 0 } type fastSubjectGroup struct { subject string messages []*fastMessage envelopes []envelope.FrameEnvelope } func fastSubjectGroups(messages []*fastMessage, envelopes []envelope.FrameEnvelope) []fastSubjectGroup { if len(envelopes) == 0 { return nil } fallback := fastBatchSubject(messages) groups := make([]fastSubjectGroup, 0, len(envelopes)) indexBySubject := make(map[string]int, len(envelopes)) for index, env := range envelopes { subject := fastMessageSubject(messages, index, fallback) groupIndex, exists := indexBySubject[subject] if !exists { groupIndex = len(groups) indexBySubject[subject] = groupIndex groups = append(groups, fastSubjectGroup{subject: subject}) } groups[groupIndex].envelopes = append(groups[groupIndex].envelopes, env) if index >= 0 && index < len(messages) { groups[groupIndex].messages = append(groups[groupIndex].messages, messages[index]) } } return groups } func fastBatchSubject(messages []*fastMessage) string { if len(messages) == 0 { return "unknown" } subject := messages[0].subject for _, msg := range messages[1:] { if msg.subject != subject { return "mixed" } } if strings.TrimSpace(subject) == "" { return "unknown" } return subject } func fastMessageSubject(messages []*fastMessage, index int, fallback string) string { if index >= 0 && index < len(messages) && strings.TrimSpace(messages[index].subject) != "" { return messages[index].subject } if strings.TrimSpace(fallback) != "" { return fallback } return "unknown" } func ensureStream(js nats.JetStreamContext, cfg config) error { ctx, cancel := context.WithTimeout(context.Background(), cfg.StreamEnsureWait) defer cancel() opts := []nats.JSOpt{nats.Context(ctx)} stream := &nats.StreamConfig{ Name: cfg.NATSStream, Subjects: cfg.NATSSubjects, Storage: nats.FileStorage, Retention: nats.LimitsPolicy, MaxAge: cfg.StreamMaxAge, MaxBytes: cfg.StreamMaxBytes, Duplicates: 2 * time.Minute, } if _, err := js.StreamInfo(cfg.NATSStream, opts...); err == nil { _, err = js.UpdateStream(stream, opts...) return err } _, err := js.AddStream(stream, opts...) if err == nil { return nil } _, updateErr := js.UpdateStream(stream, opts...) if updateErr == nil { return nil } return err } func addFastMetric(registry *metrics.Registry, subject string, status string) { if registry == nil { return } labels := metrics.Labels{"subject": subject, "status": status} registry.IncCounter("vehicle_fast_writer_messages_total", labels) metrics.RecordLastActivity(registry, "vehicle_fast_writer_last_message_unix_seconds", labels) } func addFastBatchFallbackMetric(registry *metrics.Registry, stage string, status string) { if registry == nil { return } registry.IncCounter("vehicle_fast_writer_batch_fallback_total", metrics.Labels{ "stage": stage, "status": status, }) } func addFastDecoupledUpdateMetric(registry *metrics.Registry, reason string, status string) { if registry == nil { return } registry.IncCounter("vehicle_fast_writer_decoupled_updates_total", metrics.Labels{ "reason": reason, "status": status, }) } func addFastBatchPending(registry *metrics.Registry, messages int, envelopes int) { if registry == nil { return } fastBatchPending.Add(registry, "vehicle_fast_writer_batch_pending_messages", "vehicle_fast_writer_batch_pending_envelopes", messages, envelopes) } func recordFastNATSConsumerInfoMetrics(registry *metrics.Registry, cfg config, info *nats.ConsumerInfo) { if registry == nil || info == nil { return } labels := metrics.Labels{"stream": cfg.NATSStream, "consumer": cfg.NATSDurable} registry.SetGauge("vehicle_fast_writer_nats_consumer_pending", labels, float64(info.NumPending)) registry.SetGauge("vehicle_fast_writer_nats_consumer_ack_pending", labels, float64(info.NumAckPending)) registry.SetGauge("vehicle_fast_writer_nats_consumer_waiting", labels, float64(info.NumWaiting)) } func recordFastWriterStageDuration(registry *metrics.Registry, subject string, stage string, status string, elapsed time.Duration) { if registry == nil { return } labels := metrics.Labels{ "subject": subject, "stage": stage, "status": status, } registry.ObserveHistogram("vehicle_fast_writer_stage_duration_ms_histogram", labels, fastWriterStageDurationBucketsMS, float64(elapsed.Milliseconds())) metrics.RecordLastActivity(registry, "vehicle_fast_writer_last_stage_unix_seconds", labels) } func recordFastWriterRedisE2EDurationMessages(registry *metrics.Registry, messages []*fastMessage, envelopes []envelope.FrameEnvelope, fallback string) { for index, env := range envelopes { recordFastWriterRedisE2EDuration(registry, fastMessageSubject(messages, index, fallback), env) } } func recordFastWriterRedisE2EDuration(registry *metrics.Registry, subject string, env envelope.FrameEnvelope) { if registry == nil || env.ReceivedAtMS <= 0 { return } elapsedMS := time.Since(time.UnixMilli(env.ReceivedAtMS)).Milliseconds() if elapsedMS < 0 { elapsedMS = 0 } labels := metrics.Labels{ "subject": subject, } registry.ObserveHistogram("vehicle_fast_writer_redis_e2e_duration_ms_histogram", labels, fastWriterRedisE2EDurationBucketsMS, float64(elapsedMS)) p99, samples := fastWriterRedisE2ERecent.Observe(subject, float64(elapsedMS)) registry.SetGauge("vehicle_fast_writer_redis_e2e_recent_p99_ms", labels, p99) registry.SetGauge("vehicle_fast_writer_redis_e2e_recent_samples", labels, float64(samples)) metrics.RecordLastActivity(registry, "vehicle_fast_writer_last_redis_e2e_unix_seconds", labels) } func recordFastWriterRedisFieldMetrics(registry *metrics.Registry, subject string, result realtime.FastUpdateResult) { if registry == nil || result.FieldsSeen == 0 { return } addFastWriterRedisFieldMetric(registry, subject, "seen", result.FieldsSeen) addFastWriterRedisFieldMetric(registry, subject, "written", result.FieldsWritten) addFastWriterRedisFieldMetric(registry, subject, "skipped_stale", result.FieldsSkippedStale) } func recordFastWriterRedisEnvelopeMetrics(registry *metrics.Registry, subject string, result realtime.FastUpdateResult) { if registry == nil || result.EnvelopesSeen == 0 { return } addFastWriterRedisEnvelopeMetric(registry, subject, "seen", result.EnvelopesSeen) addFastWriterRedisEnvelopeMetric(registry, subject, "updated", result.EnvelopesUpdated) addFastWriterRedisEnvelopeMetric(registry, subject, "skipped_non_realtime", result.EnvelopesSkippedNonRealtime) addFastWriterRedisEnvelopeMetric(registry, subject, "skipped_missing_vin", result.EnvelopesSkippedMissingVIN) addFastWriterRedisEnvelopeMetric(registry, subject, "skipped_missing_vehicle_key", result.EnvelopesSkippedMissingVehicleKey) addFastWriterRedisEnvelopeMetric(registry, subject, "skipped_missing_fields", result.EnvelopesSkippedMissingFields) } func addFastWriterRedisEnvelopeMetric(registry *metrics.Registry, subject string, status string, value int) { if value <= 0 { return } registry.AddCounter("vehicle_fast_writer_redis_envelopes_total", metrics.Labels{ "subject": subject, "status": status, }, float64(value)) } func addFastWriterRedisFieldMetric(registry *metrics.Registry, subject string, status string, value int) { if value <= 0 { return } registry.AddCounter("vehicle_fast_writer_redis_fields_total", metrics.Labels{ "subject": subject, "status": status, }, float64(value)) } func statusFromError(err error) string { if err != nil { return "error" } return "ok" } 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 envInt64(key string, fallback int64) int64 { value := strings.TrimSpace(os.Getenv(key)) if value == "" { return fallback } var parsed int64 if _, err := fmt.Sscanf(value, "%d", &parsed); err != nil || parsed <= 0 { return fallback } return parsed } func envBool(key string, fallback bool) bool { value := strings.ToLower(strings.TrimSpace(os.Getenv(key))) switch value { case "": return fallback case "1", "true", "yes", "y", "on": return true case "0", "false", "no", "n", "off": return false default: return fallback } } 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 }