diff --git a/docs/architecture/iot-data-platform-principles.md b/docs/architecture/iot-data-platform-principles.md index e887936a..be9da809 100644 --- a/docs/architecture/iot-data-platform-principles.md +++ b/docs/architecture/iot-data-platform-principles.md @@ -79,3 +79,19 @@ Every received frame should become one `FrameEnvelope`. 2. Add runtime metrics for protocol frame counts, parse failures, Kafka/NATS publish failures, and writer lag. 3. Review TDengine tables and remove unused tables only after export or explicit confirmation. 4. Keep new business tables narrow by default. Add a column only when a query or index proves it is needed. + +## Runtime Metrics Baseline + +Runtime metrics are exposed as Prometheus text from `/metrics`. They are operational signals and must not create new business tables by default. + +Core counters: + +- `vehicle_gateway_frames_total`: received protocol frames by protocol and parse status. +- `vehicle_gateway_publish_total`: raw and unified publish results by protocol. +- `vehicle_realtime_kafka_messages_total`: realtime consumer messages by topic and status. +- `vehicle_realtime_updates_total`: Redis/MySQL realtime projector updates by topic and status. +- `vehicle_history_writes_total`: TDengine history writes by topic and status. +- `vehicle_stat_writes_total`: MySQL metric writes by topic and status. +- `vehicle_bridge_kafka_writes_total`: NATS to Kafka bridge writes by Kafka topic and status. + +If a metric becomes a product requirement, derive a narrow metric table from Kafka replay instead of widening raw or realtime tables. diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index c7f341cc..824ea0aa 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -18,6 +18,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/gateway" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/gb32960" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/jt808" @@ -40,7 +41,8 @@ func main() { os.Exit(1) } defer closeResolver() - health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "vehicle-gateway", nil)) + registry := metrics.NewRegistry() + health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "vehicle-gateway", nil, registry)) protocols := []gateway.TCPProtocol{ { @@ -67,6 +69,7 @@ func main() { Sink: sink, Resolver: resolver, Logger: logger, + Metrics: registry, ReadBufferSize: envInt("TCP_READ_BUFFER_BYTES", 64*1024), IdleTimeout: time.Duration(envInt("TCP_IDLE_TIMEOUT_SECONDS", 180)) * time.Second, MaxConnections: envInt("TCP_MAX_CONNECTIONS", 20_000), @@ -102,6 +105,7 @@ func main() { Sink: sink, Resolver: resolver, Logger: logger, + Metrics: registry, }) if err != nil { logger.Error("build yutong mqtt client failed", "error", err) diff --git a/go/vehicle-gateway/cmd/history-writer/main.go b/go/vehicle-gateway/cmd/history-writer/main.go index efd66766..5594b562 100644 --- a/go/vehicle-gateway/cmd/history-writer/main.go +++ b/go/vehicle-gateway/cmd/history-writer/main.go @@ -16,6 +16,7 @@ import ( "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" ) @@ -35,9 +36,10 @@ func main() { logger.Error("tdengine ping failed", "error", err) os.Exit(1) } + registry := metrics.NewRegistry() health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "vehicle-history-writer", []health.Check{ {Name: "tdengine", Check: db.PingContext}, - })) + }, registry)) writer := history.NewWriter(db) if cfg.EnsureSchema { @@ -70,7 +72,7 @@ func main() { logger.Error("kafka fetch failed", "error", err) continue } - processHistoryMessage(ctx, logger, writer, reader, message) + processHistoryMessage(ctx, logger, registry, writer, reader, message) } } @@ -87,23 +89,37 @@ type kafkaMessageCommitter interface { func processHistoryMessage(ctx context.Context, logger interface { Error(string, ...any) Warn(string, ...any) -}, appender historyAppender, committer kafkaMessageCommitter, message kafka.Message) { +}, registry *metrics.Registry, appender historyAppender, committer kafkaMessageCommitter, message kafka.Message) { messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) defer cancel() + addWriterMetric(registry, "vehicle_history_kafka_messages_total", message, "received") var env envelope.FrameEnvelope if err := json.Unmarshal(message.Value, &env); err != nil { + addWriterMetric(registry, "vehicle_history_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 := appender.AppendAll(messageCtx, env); err != nil { + addWriterMetric(registry, "vehicle_history_writes_total", message, "error") logger.Error("tdengine append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) return } + addWriterMetric(registry, "vehicle_history_writes_total", message, "ok") if err := committer.CommitMessages(messageCtx, message); err != nil { + addWriterMetric(registry, "vehicle_history_kafka_commits_total", message, "error") logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + return } + addWriterMetric(registry, "vehicle_history_kafka_commits_total", message, "ok") +} + +func addWriterMetric(registry *metrics.Registry, name string, message kafka.Message, status string) { + if registry == nil { + return + } + registry.IncCounter(name, metrics.Labels{"topic": message.Topic, "status": status}) } type config struct { diff --git a/go/vehicle-gateway/cmd/history-writer/main_test.go b/go/vehicle-gateway/cmd/history-writer/main_test.go index e6479c4b..1c3c8435 100644 --- a/go/vehicle-gateway/cmd/history-writer/main_test.go +++ b/go/vehicle-gateway/cmd/history-writer/main_test.go @@ -3,11 +3,13 @@ package main import ( "context" "encoding/json" + "strings" "testing" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestProcessHistoryMessageUsesUncancelledContextForAppendAndCommit(t *testing.T) { @@ -29,7 +31,7 @@ func TestProcessHistoryMessageUsesUncancelledContextForAppendAndCommit(t *testin appender := &contextCheckingHistoryAppender{} committer := &contextCheckingHistoryCommitter{} - processHistoryMessage(parent, discardHistoryLogger{}, appender, committer, kafka.Message{Value: payload}) + processHistoryMessage(parent, discardHistoryLogger{}, nil, appender, committer, kafka.Message{Value: payload}) if appender.ctxErr != nil { t.Fatalf("appender saw cancelled context: %v", appender.ctxErr) @@ -42,6 +44,35 @@ func TestProcessHistoryMessageUsesUncancelledContextForAppendAndCommit(t *testin } } +func TestProcessHistoryMessageRecordsMetrics(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", MessageID: "0x01"} + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + registry := metrics.NewRegistry() + + processHistoryMessage( + context.Background(), + discardHistoryLogger{}, + registry, + &contextCheckingHistoryAppender{}, + &contextCheckingHistoryCommitter{}, + kafka.Message{Topic: "vehicle.raw.gb32960.v1", Value: payload}, + ) + + text := registry.Render() + for _, want := range []string{ + `vehicle_history_kafka_messages_total{status="received",topic="vehicle.raw.gb32960.v1"} 1`, + `vehicle_history_writes_total{status="ok",topic="vehicle.raw.gb32960.v1"} 1`, + `vehicle_history_kafka_commits_total{status="ok",topic="vehicle.raw.gb32960.v1"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + type contextCheckingHistoryAppender struct { ctxErr error count int diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go index e4ede85a..98e52583 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go @@ -18,6 +18,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" ) @@ -33,6 +34,7 @@ func main() { os.Exit(1) } defer conn.Close() + registry := metrics.NewRegistry() health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-kafka-bridge", []health.Check{ {Name: "nats", Check: func(context.Context) error { if conn.Status() != nats.CONNECTED { @@ -40,7 +42,7 @@ func main() { } return nil }}, - })) + }, registry)) js, err := conn.JetStream() if err != nil { logger.Error("nats jetstream init failed", "error", err) @@ -78,7 +80,7 @@ func main() { "durable", cfg.NATSDurable, "filter", cfg.NATSFilter, "kafka_brokers", strings.Join(cfg.KafkaBrokers, ",")) - runBridge(ctx, logger, sub, writer, cfg) + runBridge(ctx, logger, registry, sub, writer, cfg) } type config struct { @@ -135,7 +137,7 @@ type bridgeMessage struct { ack func() error } -func runBridge(ctx context.Context, logger *slog.Logger, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { +func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { for { if ctx.Err() != nil { return @@ -161,7 +163,7 @@ func runBridge(ctx context.Context, logger *slog.Logger, sub natsPullSubscriptio }) } operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) - err = bridgeBatch(operationCtx, writer, bridgeMessages, cfg.Route) + err = bridgeBatch(operationCtx, registry, writer, bridgeMessages, cfg.Route) cancel() if err != nil { logger.Error("bridge batch failed", "count", len(bridgeMessages), "error", err) @@ -194,32 +196,56 @@ func ensureStream(js nats.JetStreamContext, cfg config) error { return err } -func bridgeBatch(ctx context.Context, writer kafkaBatchWriter, messages []bridgeMessage, route map[string]string) error { +func bridgeBatch(ctx context.Context, registry *metrics.Registry, writer kafkaBatchWriter, messages []bridgeMessage, route map[string]string) error { if len(messages) == 0 { return nil } kafkaMessages := make([]kafka.Message, 0, len(messages)) for _, message := range messages { + addBridgeSubjectMetric(registry, "vehicle_bridge_messages_total", message.subject, "received") topic, ok := route[message.subject] if !ok || topic == "" { + addBridgeSubjectMetric(registry, "vehicle_bridge_messages_total", message.subject, "route_error") return fmt.Errorf("kafka topic not configured for nats subject %q", message.subject) } kafkaMessages = append(kafkaMessages, kafkaMessage(topic, message.data)) } if err := writer.WriteMessages(ctx, kafkaMessages...); err != nil { + for _, message := range kafkaMessages { + addBridgeTopicMetric(registry, "vehicle_bridge_kafka_writes_total", message.Topic, "error") + } return err } + for _, message := range kafkaMessages { + addBridgeTopicMetric(registry, "vehicle_bridge_kafka_writes_total", message.Topic, "ok") + } for _, message := range messages { if message.ack == nil { continue } if err := message.ack(); err != nil { + addBridgeSubjectMetric(registry, "vehicle_bridge_nats_acks_total", message.subject, "error") return err } + addBridgeSubjectMetric(registry, "vehicle_bridge_nats_acks_total", message.subject, "ok") } return nil } +func addBridgeSubjectMetric(registry *metrics.Registry, name string, subject string, status string) { + if registry == nil { + return + } + registry.IncCounter(name, metrics.Labels{"subject": subject, "status": status}) +} + +func addBridgeTopicMetric(registry *metrics.Registry, name string, topic string, status string) { + if registry == nil { + return + } + registry.IncCounter(name, metrics.Labels{"topic": topic, "status": status}) +} + func kafkaMessage(topic string, data []byte) kafka.Message { var env envelope.FrameEnvelope message := kafka.Message{Topic: topic, Value: data} diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go index 76cd47c2..f4c091a2 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go @@ -4,11 +4,13 @@ import ( "context" "encoding/json" "errors" + "strings" "testing" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestBridgeBatchWritesKafkaThenAcks(t *testing.T) { @@ -20,7 +22,7 @@ func TestBridgeBatchWritesKafkaThenAcks(t *testing.T) { writer := &recordingBridgeWriter{} acked := 0 - err = bridgeBatch(context.Background(), writer, []bridgeMessage{ + err = bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.raw.jt808.v1", data: payload, ack: func() error { acked++ return nil @@ -47,7 +49,7 @@ func TestBridgeBatchDoesNotAckWhenKafkaFails(t *testing.T) { writer := &recordingBridgeWriter{err: errors.New("kafka unavailable")} acked := 0 - err := bridgeBatch(context.Background(), writer, []bridgeMessage{ + err := bridgeBatch(context.Background(), nil, writer, []bridgeMessage{ {subject: "vehicle.event.unified.v1", data: []byte(`{"protocol":"JT808"}`), ack: func() error { acked++ return nil @@ -61,6 +63,33 @@ func TestBridgeBatchDoesNotAckWhenKafkaFails(t *testing.T) { } } +func TestBridgeBatchRecordsMetrics(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + registry := metrics.NewRegistry() + writer := &recordingBridgeWriter{} + + err = bridgeBatch(context.Background(), registry, writer, []bridgeMessage{ + {subject: "vehicle.raw.jt808.v1", data: payload, ack: func() error { return nil }}, + }, map[string]string{"vehicle.raw.jt808.v1": "vehicle.raw.jt808.v1"}) + if err != nil { + t.Fatalf("bridgeBatch() error = %v", err) + } + text := registry.Render() + for _, want := range []string{ + `vehicle_bridge_messages_total{status="received",subject="vehicle.raw.jt808.v1"} 1`, + `vehicle_bridge_kafka_writes_total{status="ok",topic="vehicle.raw.jt808.v1"} 1`, + `vehicle_bridge_nats_acks_total{status="ok",subject="vehicle.raw.jt808.v1"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + type recordingBridgeWriter struct { messages []kafka.Message err error diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 31b78ab0..0fc21375 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -20,6 +20,7 @@ import ( "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" @@ -47,9 +48,11 @@ func main() { }) 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)) @@ -96,7 +99,7 @@ func main() { } defer closeStats() if brokers := splitCSV(os.Getenv("KAFKA_BROKERS")); len(brokers) > 0 { - go consumeKafka(ctx, logger, updater, brokers) + go consumeKafka(ctx, logger, registry, updater, brokers) } else { logger.Warn("KAFKA_BROKERS is empty; realtime api will serve existing redis data only") } @@ -159,7 +162,7 @@ func consumeKafka(ctx context.Context, logger interface { Info(string, ...any) Error(string, ...any) Warn(string, ...any) -}, updater realtimeUpdater, brokers []string) { +}, registry *metrics.Registry, updater realtimeUpdater, brokers []string) { reader := kafka.NewReader(kafka.ReaderConfig{ Brokers: brokers, GroupID: env("KAFKA_GROUP", "go-realtime-api"), @@ -178,7 +181,7 @@ func consumeKafka(ctx context.Context, logger interface { logger.Error("kafka fetch failed", "error", err) continue } - processRealtimeMessage(ctx, logger, updater, reader, message) + processRealtimeMessage(ctx, logger, registry, updater, reader, message) } } @@ -210,23 +213,40 @@ type kafkaMessageCommitter interface { func processRealtimeMessage(ctx context.Context, logger interface { Error(string, ...any) Warn(string, ...any) -}, updater realtimeUpdater, committer kafkaMessageCommitter, message kafka.Message) { +}, 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") 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 env(key string, fallback string) string { diff --git a/go/vehicle-gateway/cmd/realtime-api/main_test.go b/go/vehicle-gateway/cmd/realtime-api/main_test.go index 7a1c2fda..1f5db499 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main_test.go +++ b/go/vehicle-gateway/cmd/realtime-api/main_test.go @@ -3,11 +3,13 @@ package main import ( "context" "encoding/json" + "strings" "testing" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestProcessRealtimeMessageUsesUncancelledContextForUpdateAndCommit(t *testing.T) { @@ -29,7 +31,7 @@ func TestProcessRealtimeMessageUsesUncancelledContextForUpdateAndCommit(t *testi updater := &contextCheckingRealtimeUpdater{} committer := &contextCheckingMessageCommitter{} - processRealtimeMessage(parent, discardRealtimeLogger{}, updater, committer, kafka.Message{Value: payload}) + processRealtimeMessage(parent, discardRealtimeLogger{}, nil, updater, committer, kafka.Message{Value: payload}) if updater.ctxErr != nil { t.Fatalf("updater saw cancelled context: %v", updater.ctxErr) @@ -42,6 +44,42 @@ func TestProcessRealtimeMessageUsesUncancelledContextForUpdateAndCommit(t *testi } } +func TestProcessRealtimeMessageRecordsMetrics(t *testing.T) { + env := envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "13307795425", + EventTimeMS: 1782918600000, + ReceivedAtMS: 1782918600000, + ParseStatus: envelope.ParseOK, + } + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + registry := metrics.NewRegistry() + + processRealtimeMessage( + context.Background(), + discardRealtimeLogger{}, + registry, + &contextCheckingRealtimeUpdater{}, + &contextCheckingMessageCommitter{}, + kafka.Message{Topic: "vehicle.event.go.unified.v1", Value: payload}, + ) + + text := registry.Render() + for _, want := range []string{ + `vehicle_realtime_kafka_messages_total{status="received",topic="vehicle.event.go.unified.v1"} 1`, + `vehicle_realtime_updates_total{status="ok",topic="vehicle.event.go.unified.v1"} 1`, + `vehicle_realtime_kafka_commits_total{status="ok",topic="vehicle.event.go.unified.v1"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + func TestCompositeRealtimeUpdaterUpdatesBothStores(t *testing.T) { env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001"} redisUpdater := &contextCheckingRealtimeUpdater{} diff --git a/go/vehicle-gateway/cmd/stat-writer/main.go b/go/vehicle-gateway/cmd/stat-writer/main.go index 01ff8e5e..1e94f710 100644 --- a/go/vehicle-gateway/cmd/stat-writer/main.go +++ b/go/vehicle-gateway/cmd/stat-writer/main.go @@ -15,6 +15,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats" ) @@ -35,9 +36,10 @@ func main() { logger.Error("mysql ping failed", "error", err) os.Exit(1) } + registry := metrics.NewRegistry() health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "vehicle-stat-writer", []health.Check{ {Name: "mysql", Check: db.PingContext}, - })) + }, registry)) writer := stats.NewWriter(db, cfg.Location) if cfg.EnsureSchema { @@ -66,7 +68,7 @@ func main() { logger.Error("kafka fetch failed", "error", err) continue } - processStatMessage(ctx, logger, writer, reader, message) + processStatMessage(ctx, logger, registry, writer, reader, message) } } @@ -83,23 +85,37 @@ type kafkaMessageCommitter interface { func processStatMessage(ctx context.Context, logger interface { Error(string, ...any) Warn(string, ...any) -}, appender statAppender, committer kafkaMessageCommitter, message kafka.Message) { +}, registry *metrics.Registry, appender statAppender, committer kafkaMessageCommitter, message kafka.Message) { messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) defer cancel() + addStatMetric(registry, "vehicle_stat_kafka_messages_total", message, "received") var env envelope.FrameEnvelope if err := json.Unmarshal(message.Value, &env); err != nil { + addStatMetric(registry, "vehicle_stat_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 := appender.Append(messageCtx, env); err != nil { + addStatMetric(registry, "vehicle_stat_writes_total", message, "error") logger.Error("mysql append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err) return } + addStatMetric(registry, "vehicle_stat_writes_total", message, "ok") if err := committer.CommitMessages(messageCtx, message); err != nil { + addStatMetric(registry, "vehicle_stat_kafka_commits_total", message, "error") logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + return } + addStatMetric(registry, "vehicle_stat_kafka_commits_total", message, "ok") +} + +func addStatMetric(registry *metrics.Registry, name string, message kafka.Message, status string) { + if registry == nil { + return + } + registry.IncCounter(name, metrics.Labels{"topic": message.Topic, "status": status}) } type config struct { diff --git a/go/vehicle-gateway/cmd/stat-writer/main_test.go b/go/vehicle-gateway/cmd/stat-writer/main_test.go index d595f844..cea2d2d4 100644 --- a/go/vehicle-gateway/cmd/stat-writer/main_test.go +++ b/go/vehicle-gateway/cmd/stat-writer/main_test.go @@ -3,11 +3,13 @@ package main import ( "context" "encoding/json" + "strings" "testing" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestProcessStatMessageUsesUncancelledContextForAppendAndCommit(t *testing.T) { @@ -29,7 +31,7 @@ func TestProcessStatMessageUsesUncancelledContextForAppendAndCommit(t *testing.T appender := &contextCheckingStatAppender{} committer := &contextCheckingStatCommitter{} - processStatMessage(parent, discardStatLogger{}, appender, committer, kafka.Message{Value: payload}) + processStatMessage(parent, discardStatLogger{}, nil, appender, committer, kafka.Message{Value: payload}) if appender.ctxErr != nil { t.Fatalf("appender saw cancelled context: %v", appender.ctxErr) @@ -42,6 +44,35 @@ func TestProcessStatMessageUsesUncancelledContextForAppendAndCommit(t *testing.T } } +func TestProcessStatMessageRecordsMetrics(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + registry := metrics.NewRegistry() + + processStatMessage( + context.Background(), + discardStatLogger{}, + registry, + &contextCheckingStatAppender{}, + &contextCheckingStatCommitter{}, + kafka.Message{Topic: "vehicle.raw.jt808.v1", Value: payload}, + ) + + text := registry.Render() + for _, want := range []string{ + `vehicle_stat_kafka_messages_total{status="received",topic="vehicle.raw.jt808.v1"} 1`, + `vehicle_stat_writes_total{status="ok",topic="vehicle.raw.jt808.v1"} 1`, + `vehicle_stat_kafka_commits_total{status="ok",topic="vehicle.raw.jt808.v1"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + type contextCheckingStatAppender struct { ctxErr error count int diff --git a/go/vehicle-gateway/internal/gateway/mqtt_client.go b/go/vehicle-gateway/internal/gateway/mqtt_client.go index 24b4da4f..84dbd37f 100644 --- a/go/vehicle-gateway/internal/gateway/mqtt_client.go +++ b/go/vehicle-gateway/internal/gateway/mqtt_client.go @@ -17,6 +17,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/yutongmqtt" ) @@ -38,6 +39,7 @@ type MQTTClientConfig struct { Sink eventbus.Sink Resolver identity.Resolver Logger *slog.Logger + Metrics *metrics.Registry } type MQTTClient struct { @@ -208,14 +210,41 @@ func (c *MQTTClient) handleMessage(ctx context.Context, topic string, payload [] env = resolved } } + c.recordFrameMetric(env.ParseStatus) if err := c.cfg.Sink.PublishRaw(messageCtx, env); err != nil { + c.recordPublishMetric("raw", "error") c.cfg.Logger.Error("publish mqtt raw failed", "topic", topic, "event_id", env.StableEventID(), "error", err) return } + c.recordPublishMetric("raw", "ok") if env.ParseStatus == envelope.ParseBadFrame { return } if err := c.cfg.Sink.PublishUnified(messageCtx, env); err != nil { + c.recordPublishMetric("unified", "error") c.cfg.Logger.Error("publish mqtt unified failed", "topic", topic, "event_id", env.StableEventID(), "error", err) + return } + c.recordPublishMetric("unified", "ok") +} + +func (c *MQTTClient) recordFrameMetric(status envelope.ParseStatus) { + if c.cfg.Metrics == nil { + return + } + c.cfg.Metrics.IncCounter("vehicle_gateway_frames_total", metrics.Labels{ + "protocol": string(envelope.ProtocolYutongMQTT), + "status": string(status), + }) +} + +func (c *MQTTClient) recordPublishMetric(kind string, status string) { + if c.cfg.Metrics == nil { + return + } + c.cfg.Metrics.IncCounter("vehicle_gateway_publish_total", metrics.Labels{ + "protocol": string(envelope.ProtocolYutongMQTT), + "kind": kind, + "status": status, + }) } diff --git a/go/vehicle-gateway/internal/gateway/mqtt_client_test.go b/go/vehicle-gateway/internal/gateway/mqtt_client_test.go index f026aecb..afd6ef44 100644 --- a/go/vehicle-gateway/internal/gateway/mqtt_client_test.go +++ b/go/vehicle-gateway/internal/gateway/mqtt_client_test.go @@ -11,10 +11,12 @@ import ( "math/big" "os" "path/filepath" + "strings" "testing" "time" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestMQTTClientHandleMessagePublishesRawAndUnified(t *testing.T) { @@ -45,6 +47,39 @@ func TestMQTTClientHandleMessagePublishesRawAndUnified(t *testing.T) { } } +func TestMQTTClientRecordsMessageMetrics(t *testing.T) { + registry := metrics.NewRegistry() + client, err := NewMQTTClient(MQTTClientConfig{ + EndpointName: "endpoint-a", + Broker: "tcp://127.0.0.1:1883", + ClientID: "test-client", + Topics: []string{"/ytforward/shln/+"}, + Sink: &recordingSink{}, + Logger: slog.New(slog.NewTextHandler(testWriter{t: t}, nil)), + Metrics: registry, + }) + if err != nil { + t.Fatalf("NewMQTTClient() error = %v", err) + } + + client.handleMessage(context.Background(), "/ytforward/shln/dev1", []byte(`{ + "device":"LTEST000000000001", + "time":"20260413100000", + "data":{"METER_SPEED":52.3,"TOTAL_MILEAGE":123456.7} + }`)) + + text := registry.Render() + for _, want := range []string{ + `vehicle_gateway_frames_total{protocol="YUTONG_MQTT",status="OK"} 1`, + `vehicle_gateway_publish_total{kind="raw",protocol="YUTONG_MQTT",status="ok"} 1`, + `vehicle_gateway_publish_total{kind="unified",protocol="YUTONG_MQTT",status="ok"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + func TestMQTTClientHandleBadPayloadPublishesOnlyRaw(t *testing.T) { sink := &recordingSink{} client, err := NewMQTTClient(MQTTClientConfig{ diff --git a/go/vehicle-gateway/internal/gateway/tcp_server.go b/go/vehicle-gateway/internal/gateway/tcp_server.go index 844c599f..4861cd77 100644 --- a/go/vehicle-gateway/internal/gateway/tcp_server.go +++ b/go/vehicle-gateway/internal/gateway/tcp_server.go @@ -15,6 +15,7 @@ import ( "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/eventbus" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) type FrameExtractor func([]byte) (frames [][]byte, remainder []byte, err error) @@ -36,6 +37,7 @@ type TCPServer struct { sink eventbus.Sink resolver identity.Resolver logger *slog.Logger + metrics *metrics.Registry readBufferSize int idleTimeout time.Duration maxConnections int @@ -46,6 +48,7 @@ type TCPServerConfig struct { Sink eventbus.Sink Resolver identity.Resolver Logger *slog.Logger + Metrics *metrics.Registry ReadBufferSize int IdleTimeout time.Duration MaxConnections int @@ -89,6 +92,7 @@ func NewTCPServer(cfg TCPServerConfig) (*TCPServer, error) { sink: cfg.Sink, resolver: cfg.Resolver, logger: cfg.Logger, + metrics: cfg.Metrics, readBufferSize: cfg.ReadBufferSize, idleTimeout: cfg.IdleTimeout, maxConnections: cfg.MaxConnections, @@ -209,18 +213,23 @@ func (s *TCPServer) handleFrame(ctx context.Context, conn net.Conn, raw []byte, env = resolved } } + s.recordFrameMetric(env.ParseStatus) if err := s.sink.PublishRaw(frameCtx, env); err != nil { + s.recordPublishMetric("raw", "error") s.logger.Error("publish raw failed", "protocol", s.protocol.Protocol, "event_id", env.StableEventID(), "error", err) return } + s.recordPublishMetric("raw", "ok") if env.ParseStatus == envelope.ParseBadFrame { return } if err := s.sink.PublishUnified(frameCtx, env); err != nil { + s.recordPublishMetric("unified", "error") s.logger.Error("publish unified failed", "protocol", s.protocol.Protocol, "event_id", env.StableEventID(), "error", err) return } + s.recordPublishMetric("unified", "ok") if s.protocol.Respond == nil { return } @@ -238,6 +247,27 @@ func (s *TCPServer) handleFrame(ctx context.Context, conn net.Conn, raw []byte, } } +func (s *TCPServer) recordFrameMetric(status envelope.ParseStatus) { + if s.metrics == nil { + return + } + s.metrics.IncCounter("vehicle_gateway_frames_total", metrics.Labels{ + "protocol": string(s.protocol.Protocol), + "status": string(status), + }) +} + +func (s *TCPServer) recordPublishMetric(kind string, status string) { + if s.metrics == nil { + return + } + s.metrics.IncCounter("vehicle_gateway_publish_total", metrics.Labels{ + "protocol": string(s.protocol.Protocol), + "kind": kind, + "status": status, + }) +} + func (p TCPProtocol) String() string { return fmt.Sprintf("%s@%s", p.Protocol, p.Addr) } diff --git a/go/vehicle-gateway/internal/gateway/tcp_server_test.go b/go/vehicle-gateway/internal/gateway/tcp_server_test.go index cda0b6b3..d6e308ff 100644 --- a/go/vehicle-gateway/internal/gateway/tcp_server_test.go +++ b/go/vehicle-gateway/internal/gateway/tcp_server_test.go @@ -6,10 +6,12 @@ import ( "io" "log/slog" "net" + "strings" "testing" "time" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/gb32960" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/protocol/jt808" ) @@ -45,6 +47,47 @@ func TestTCPServerPublishesGoodFrameToRawAndUnified(t *testing.T) { } } +func TestTCPServerRecordsFrameMetrics(t *testing.T) { + registry := metrics.NewRegistry() + server, err := NewTCPServer(TCPServerConfig{ + Protocol: TCPProtocol{ + Protocol: envelope.ProtocolJT808, + Addr: ":0", + Extract: jt808.ExtractFrames, + Parse: func(_ []byte, receivedAtMS int64, sourceEndpoint string) (envelope.FrameEnvelope, error) { + return envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "13307795425", + SourceEndpoint: sourceEndpoint, + ReceivedAtMS: receivedAtMS, + EventTimeMS: receivedAtMS, + ParseStatus: envelope.ParseOK, + }, nil + }, + }, + Sink: &recordingSink{}, + Logger: slog.New(slog.NewTextHandler(testWriter{t: t}, nil)), + Metrics: registry, + }) + if err != nil { + t.Fatalf("NewTCPServer() error = %v", err) + } + + server.handleFrame(context.Background(), nil, []byte{0x01}, "127.0.0.1:808") + + text := registry.Render() + for _, want := range []string{ + `vehicle_gateway_frames_total{protocol="JT808",status="OK"} 1`, + `vehicle_gateway_publish_total{kind="raw",protocol="JT808",status="ok"} 1`, + `vehicle_gateway_publish_total{kind="unified",protocol="JT808",status="ok"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + func TestTCPServerPublishesBadFrameOnlyToRaw(t *testing.T) { good := buildGBFrame(0x02, 0xfe, "LNBSCB3D4R1234567", nil) good[len(good)-1] ^= 0xff diff --git a/go/vehicle-gateway/internal/health/health.go b/go/vehicle-gateway/internal/health/health.go index d83a5174..0d25476f 100644 --- a/go/vehicle-gateway/internal/health/health.go +++ b/go/vehicle-gateway/internal/health/health.go @@ -7,6 +7,8 @@ import ( "net/http" "strings" "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) type Check struct { @@ -32,22 +34,25 @@ func NewHandler(service string, checks []Check) *Handler { return &Handler{service: service, checks: checks} } -func NewMux(service string, checks []Check) *http.ServeMux { +func NewMux(service string, checks []Check, registry *metrics.Registry) *http.ServeMux { mux := http.NewServeMux() handler := NewHandler(service, checks) mux.Handle("/healthz", handler) mux.Handle("/readyz", handler) + if registry != nil { + mux.Handle("/metrics", metrics.NewHandler(registry)) + } return mux } -func NewServer(addr string, service string, checks []Check) *http.Server { +func NewServer(addr string, service string, checks []Check, registry *metrics.Registry) *http.Server { addr = strings.TrimSpace(addr) if addr == "" { return nil } return &http.Server{ Addr: addr, - Handler: NewMux(service, checks), + Handler: NewMux(service, checks, registry), ReadHeaderTimeout: 5 * time.Second, } } diff --git a/go/vehicle-gateway/internal/health/health_test.go b/go/vehicle-gateway/internal/health/health_test.go index ab74a47b..2764cc24 100644 --- a/go/vehicle-gateway/internal/health/health_test.go +++ b/go/vehicle-gateway/internal/health/health_test.go @@ -7,6 +7,8 @@ import ( "net/http/httptest" "strings" "testing" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestHandlerReturnsHealthyWhenAllChecksPass(t *testing.T) { @@ -63,10 +65,12 @@ func TestHandlerRejectsUnknownPath(t *testing.T) { } func TestNewMuxRegistersHealthAndReadinessRoutes(t *testing.T) { + registry := metrics.NewRegistry() + registry.IncCounter("vehicle_test_total", metrics.Labels{"service": "stat"}) mux := NewMux("vehicle-stat-writer", []Check{ {Name: "mysql", Check: func(context.Context) error { return nil }}, - }) - for _, path := range []string{"/healthz", "/readyz"} { + }, registry) + for _, path := range []string{"/healthz", "/readyz", "/metrics"} { request := httptest.NewRequest(http.MethodGet, path, nil) response := httptest.NewRecorder() @@ -79,13 +83,13 @@ func TestNewMuxRegistersHealthAndReadinessRoutes(t *testing.T) { } func TestNewServerReturnsNilWhenAddressIsEmpty(t *testing.T) { - if server := NewServer("", "vehicle-gateway", nil); server != nil { + if server := NewServer("", "vehicle-gateway", nil, nil); server != nil { t.Fatalf("server = %#v, want nil", server) } } func TestNewServerBuildsConfiguredHTTPServer(t *testing.T) { - server := NewServer(":20290", "vehicle-gateway", nil) + server := NewServer(":20290", "vehicle-gateway", nil, nil) if server == nil { t.Fatal("server is nil") } diff --git a/go/vehicle-gateway/internal/metrics/metrics.go b/go/vehicle-gateway/internal/metrics/metrics.go new file mode 100644 index 00000000..2793756a --- /dev/null +++ b/go/vehicle-gateway/internal/metrics/metrics.go @@ -0,0 +1,141 @@ +package metrics + +import ( + "fmt" + "net/http" + "sort" + "strings" + "sync" +) + +type Labels map[string]string + +type Registry struct { + mu sync.RWMutex + counters map[string]map[string]float64 +} + +func NewRegistry() *Registry { + return &Registry{counters: map[string]map[string]float64{}} +} + +func (r *Registry) IncCounter(name string, labels Labels) { + r.AddCounter(name, labels, 1) +} + +func (r *Registry) AddCounter(name string, labels Labels, value float64) { + name = strings.TrimSpace(name) + if r == nil || name == "" || value == 0 { + return + } + key := labelsKey(labels) + r.mu.Lock() + defer r.mu.Unlock() + if r.counters[name] == nil { + r.counters[name] = map[string]float64{} + } + r.counters[name][key] += value +} + +func (r *Registry) Render() string { + if r == nil { + return "" + } + r.mu.RLock() + defer r.mu.RUnlock() + + var names []string + for name := range r.counters { + names = append(names, name) + } + sort.Strings(names) + + var b strings.Builder + for _, name := range names { + b.WriteString("# TYPE ") + b.WriteString(name) + b.WriteString(" counter\n") + var keys []string + for key := range r.counters[name] { + keys = append(keys, key) + } + sort.Strings(keys) + for _, key := range keys { + b.WriteString(name) + if key != "" { + b.WriteString("{") + b.WriteString(key) + b.WriteString("}") + } + b.WriteString(" ") + b.WriteString(formatNumber(r.counters[name][key])) + b.WriteString("\n") + } + } + return b.String() +} + +func NewHandler(registry *Registry) http.Handler { + if registry == nil { + registry = NewRegistry() + } + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + return + } + if strings.Trim(r.URL.Path, "/") != "metrics" { + http.NotFound(w, r) + return + } + w.Header().Set("Content-Type", "text/plain; version=0.0.4; charset=utf-8") + _, _ = w.Write([]byte(registry.Render())) + }) +} + +func labelsKey(labels Labels) string { + if len(labels) == 0 { + return "" + } + keys := make([]string, 0, len(labels)) + for key := range labels { + keys = append(keys, key) + } + sort.Strings(keys) + parts := make([]string, 0, len(keys)) + for _, key := range keys { + value := labels[key] + parts = append(parts, fmt.Sprintf(`%s="%s"`, sanitizeLabelName(key), escapeLabelValue(value))) + } + return strings.Join(parts, ",") +} + +func sanitizeLabelName(value string) string { + value = strings.TrimSpace(value) + if value == "" { + return "label" + } + var b strings.Builder + for i, r := range value { + ok := r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (i > 0 && r >= '0' && r <= '9') + if ok { + b.WriteRune(r) + } else { + b.WriteByte('_') + } + } + return b.String() +} + +func escapeLabelValue(value string) string { + value = strings.ReplaceAll(value, `\`, `\\`) + value = strings.ReplaceAll(value, "\n", `\n`) + return strings.ReplaceAll(value, `"`, `\"`) +} + +func formatNumber(value float64) string { + if value == float64(int64(value)) { + return fmt.Sprintf("%d", int64(value)) + } + return fmt.Sprintf("%g", value) +} diff --git a/go/vehicle-gateway/internal/metrics/metrics_test.go b/go/vehicle-gateway/internal/metrics/metrics_test.go new file mode 100644 index 00000000..a8cccaa5 --- /dev/null +++ b/go/vehicle-gateway/internal/metrics/metrics_test.go @@ -0,0 +1,60 @@ +package metrics + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +func TestRegistryRendersCountersWithLabels(t *testing.T) { + registry := NewRegistry() + registry.AddCounter("vehicle_frames_total", Labels{"protocol": "JT808", "status": "ok"}, 2) + registry.IncCounter("vehicle_frames_total", Labels{"status": "ok", "protocol": "JT808"}) + registry.IncCounter("vehicle_publish_errors_total", Labels{"target": "kafka"}) + + text := registry.Render() + + for _, want := range []string{ + `# TYPE vehicle_frames_total counter`, + `vehicle_frames_total{protocol="JT808",status="ok"} 3`, + `# TYPE vehicle_publish_errors_total counter`, + `vehicle_publish_errors_total{target="kafka"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + +func TestHandlerServesPrometheusText(t *testing.T) { + registry := NewRegistry() + registry.IncCounter("vehicle_kafka_commits_total", Labels{"service": "history"}) + handler := NewHandler(registry) + request := httptest.NewRequest(http.MethodGet, "/metrics", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", response.Code, response.Body.String()) + } + if got := response.Header().Get("Content-Type"); !strings.Contains(got, "text/plain") { + t.Fatalf("content-type = %q", got) + } + if !strings.Contains(response.Body.String(), `vehicle_kafka_commits_total{service="history"} 1`) { + t.Fatalf("body = %s", response.Body.String()) + } +} + +func TestHandlerRejectsNonMetricsPath(t *testing.T) { + handler := NewHandler(NewRegistry()) + request := httptest.NewRequest(http.MethodGet, "/readyz", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusNotFound { + t.Fatalf("status = %d", response.Code) + } +}