Просмотр исходного кода

fix: add transport label to AlertsReceived metric

AlertsReceived counter now uses two labels (transport, result)
instead of just result, so Prometheus queries can filter by transport
(grpc vs http). All call sites updated:
- pipeline.go: all branches (accepted, rate_limited, circuit_open,
  broker_unavailable, deduped, invalid, quarantined) now pass d.Transport
- cmd/ingestd/http.go: adds "http" transport label
- internal/grpcserver/handler.go: adds "grpc" transport label for
  fast-reject rate_limited paths (before pipeline.Process call)

This fixes the m11_smoke.py query:
  sum(rate(ba_ingestd_alerts_received_total{transport="grpc",result="ok"}[30s]))
Luis Rosales 1 месяц назад
Родитель
Сommit
1528d511b7
4 измененных файлов с 20 добавлено и 17 удалено
  1. 2 2
      cmd/ingestd/http.go
  2. 2 0
      internal/grpcserver/handler.go
  3. 3 3
      internal/observability/metrics.go
  4. 13 12
      internal/pipeline/pipeline.go

+ 2 - 2
cmd/ingestd/http.go

@@ -67,11 +67,11 @@ func (d *httpDeps) handleIngest(w http.ResponseWriter, r *http.Request) {
 	if err != nil {
 	if err != nil {
 		var maxErr *http.MaxBytesError
 		var maxErr *http.MaxBytesError
 		if errors.As(err, &maxErr) {
 		if errors.As(err, &maxErr) {
-			d.Metrics.AlertsReceived.WithLabelValues("payload_too_large").Inc()
+			d.Metrics.AlertsReceived.WithLabelValues("http", "payload_too_large").Inc()
 			writeErr(w, http.StatusRequestEntityTooLarge, "payload_too_large", "")
 			writeErr(w, http.StatusRequestEntityTooLarge, "payload_too_large", "")
 			return
 			return
 		}
 		}
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues("http", "invalid").Inc()
 		writeErr(w, http.StatusBadRequest, "bad_request", err.Error())
 		writeErr(w, http.StatusBadRequest, "bad_request", err.Error())
 		return
 		return
 	}
 	}

+ 2 - 0
internal/grpcserver/handler.go

@@ -106,6 +106,7 @@ func (s *Server) StreamAlerts(stream pbv1.Ingest_StreamAlertsServer) error {
 			allowed, retryAfter, _ := s.Dep.Limiter.Allow(ctx, rateLimitKey, src.RateLimitPerSec)
 			allowed, retryAfter, _ := s.Dep.Limiter.Allow(ctx, rateLimitKey, src.RateLimitPerSec)
 			if !allowed {
 			if !allowed {
 				s.Dep.Metrics.GRPCRateLimited.WithLabelValues(src.SourceID).Inc()
 				s.Dep.Metrics.GRPCRateLimited.WithLabelValues(src.SourceID).Inc()
+				s.Dep.Metrics.AlertsReceived.WithLabelValues("grpc", "rate_limited").Inc()
 				ackCh <- &pbv1.Ack{
 				ackCh <- &pbv1.Ack{
 					AlertId:   alert.DedupeKey,
 					AlertId:   alert.DedupeKey,
 					DedupeKey: alert.DedupeKey,
 					DedupeKey: alert.DedupeKey,
@@ -126,6 +127,7 @@ func (s *Server) StreamAlerts(stream pbv1.Ingest_StreamAlertsServer) error {
 				// Proceed.
 				// Proceed.
 			default:
 			default:
 				s.Dep.Metrics.GRPCRateLimited.WithLabelValues(src.SourceID).Inc()
 				s.Dep.Metrics.GRPCRateLimited.WithLabelValues(src.SourceID).Inc()
+				s.Dep.Metrics.AlertsReceived.WithLabelValues("grpc", "rate_limited").Inc()
 				ackCh <- &pbv1.Ack{
 				ackCh <- &pbv1.Ack{
 					AlertId:   alert.DedupeKey,
 					AlertId:   alert.DedupeKey,
 					DedupeKey: alert.DedupeKey,
 					DedupeKey: alert.DedupeKey,

+ 3 - 3
internal/observability/metrics.go

@@ -75,9 +75,9 @@ func NewIngestdMetrics(reg prometheus.Registerer, serviceName string) *IngestdMe
 			Namespace: "ba",
 			Namespace: "ba",
 			Subsystem: "ingestd",
 			Subsystem: "ingestd",
 			Name:      "alerts_received_total",
 			Name:      "alerts_received_total",
-			Help:      "Number of inbound alerts by result.",
+			Help:      "Number of inbound alerts by result and transport.",
 			ConstLabels: prometheus.Labels{"service": serviceName},
 			ConstLabels: prometheus.Labels{"service": serviceName},
-		}, []string{"result"}),
+		}, []string{"transport", "result"}),
 		PayloadBytes: prometheus.NewHistogram(prometheus.HistogramOpts{
 		PayloadBytes: prometheus.NewHistogram(prometheus.HistogramOpts{
 			Namespace: "ba",
 			Namespace: "ba",
 			Subsystem: "ingestd",
 			Subsystem: "ingestd",
@@ -224,7 +224,7 @@ func NewIngestdMetrics(reg prometheus.Registerer, serviceName string) *IngestdMe
 		m.GRPCRateLimited,
 		m.GRPCRateLimited,
 		m.GRPCAckLatency,
 		m.GRPCAckLatency,
 	)
 	)
-	m.AlertsReceived.WithLabelValues("accepted")
+	m.AlertsReceived.WithLabelValues("internal", "accepted")
 	return m
 	return m
 }
 }
 
 

+ 13 - 12
internal/pipeline/pipeline.go

@@ -47,6 +47,7 @@ import (
 // populate the Sources map with entries keyed by "company_id:source_id".
 // populate the Sources map with entries keyed by "company_id:source_id".
 type SourceConfig struct {
 type SourceConfig struct {
 	CompanyID       string
 	CompanyID       string
+	SourceID        string
 	HMACSecret      []byte   // may be empty for transports that don't use HMAC
 	HMACSecret      []byte   // may be empty for transports that don't use HMAC
 	RateLimitPerSec int
 	RateLimitPerSec int
 	AllowedTargets  []string // M2: allowed routing targets (pipeline ignores; caller enforces)
 	AllowedTargets  []string // M2: allowed routing targets (pipeline ignores; caller enforces)
@@ -143,18 +144,18 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	// 5. Parse + validate.
 	// 5. Parse + validate.
 	var a alert.Alert
 	var a alert.Alert
 	if err := json.Unmarshal(body, &a); err != nil {
 	if err := json.Unmarshal(body, &a); err != nil {
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "invalid").Inc()
 		return Reject("invalid_json", 400, err.Error())
 		return Reject("invalid_json", 400, err.Error())
 	}
 	}
 	if err := a.Validate(); err != nil {
 	if err := a.Validate(); err != nil {
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "invalid").Inc()
 		return Reject("invalid", 400, err.Error())
 		return Reject("invalid", 400, err.Error())
 	}
 	}
 
 
 	// Source lookup.
 	// Source lookup.
 	src, ok := d.Sources[a.CompanyID+":"+a.SourceID]
 	src, ok := d.Sources[a.CompanyID+":"+a.SourceID]
 	if !ok {
 	if !ok {
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "invalid").Inc()
 		return Reject("unknown_source", 401,
 		return Reject("unknown_source", 401,
 			fmt.Sprintf("no such source %s/%s", a.CompanyID, a.SourceID))
 			fmt.Sprintf("no such source %s/%s", a.CompanyID, a.SourceID))
 	}
 	}
@@ -162,7 +163,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	// 2. Quarantine check (M9 layer 7). Before we spend any CPU.
 	// 2. Quarantine check (M9 layer 7). Before we spend any CPU.
 	if d.Quarantine != nil {
 	if d.Quarantine != nil {
 		if banned, remaining, err := d.Quarantine.IsBanned(ctx, a.SourceID); err == nil && banned {
 		if banned, remaining, err := d.Quarantine.IsBanned(ctx, a.SourceID); err == nil && banned {
-			d.Metrics.AlertsReceived.WithLabelValues("quarantined").Inc()
+			d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "quarantined").Inc()
 			d.Logger.Warn("source quarantined",
 			d.Logger.Warn("source quarantined",
 				"source_id", a.SourceID,
 				"source_id", a.SourceID,
 				"company_id", a.CompanyID,
 				"company_id", a.CompanyID,
@@ -175,7 +176,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 
 
 	// 6. Auth (transport-specific; gRPC skips by passing "").
 	// 6. Auth (transport-specific; gRPC skips by passing "").
 	if sig != "" && !verifyHMAC(sig, src.HMACSecret, body, now) {
 	if sig != "" && !verifyHMAC(sig, src.HMACSecret, body, now) {
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "invalid").Inc()
 		return Reject("bad_signature", 401, "")
 		return Reject("bad_signature", 401, "")
 	}
 	}
 
 
@@ -213,7 +214,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 		if ok, ttl, err := d.Limiter.Allow(ctx, "source:"+a.CompanyID+":"+a.SourceID, src.RateLimitPerSec); err != nil {
 		if ok, ttl, err := d.Limiter.Allow(ctx, "source:"+a.CompanyID+":"+a.SourceID, src.RateLimitPerSec); err != nil {
 			d.Logger.Warn("ratelimit redis error (failing open)", "err", err, "scope", "source")
 			d.Logger.Warn("ratelimit redis error (failing open)", "err", err, "scope", "source")
 		} else if !ok {
 		} else if !ok {
-			d.Metrics.AlertsReceived.WithLabelValues("rate_limited").Inc()
+			d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "rate_limited").Inc()
 			d.Metrics.RateLimitHits.WithLabelValues("source").Inc()
 			d.Metrics.RateLimitHits.WithLabelValues("source").Inc()
 			recordHit()
 			recordHit()
 			return Reject("rate_limited_source", 429, strconv.Itoa(int(ttl.Seconds())))
 			return Reject("rate_limited_source", 429, strconv.Itoa(int(ttl.Seconds())))
@@ -223,7 +224,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	// 4. Per-company rate limit (new alerts only).
 	// 4. Per-company rate limit (new alerts only).
 	if isNew {
 	if isNew {
 		if ok, ttl, _ := d.Limiter.Allow(ctx, "company:"+a.CompanyID, d.CompanyRatePerSec); !ok {
 		if ok, ttl, _ := d.Limiter.Allow(ctx, "company:"+a.CompanyID, d.CompanyRatePerSec); !ok {
-			d.Metrics.AlertsReceived.WithLabelValues("rate_limited").Inc()
+			d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "rate_limited").Inc()
 			d.Metrics.RateLimitHits.WithLabelValues("company").Inc()
 			d.Metrics.RateLimitHits.WithLabelValues("company").Inc()
 			recordHit()
 			recordHit()
 			return Reject("rate_limited_company", 429, strconv.Itoa(int(ttl.Seconds())))
 			return Reject("rate_limited_company", 429, strconv.Itoa(int(ttl.Seconds())))
@@ -239,7 +240,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	subject := broker.AlertsSubject(a.CompanyID)
 	subject := broker.AlertsSubject(a.CompanyID)
 	payload, err := json.Marshal(a)
 	payload, err := json.Marshal(a)
 	if err != nil {
 	if err != nil {
-		d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "invalid").Inc()
 		return Reject("marshal_failed", 500, err.Error())
 		return Reject("marshal_failed", 500, err.Error())
 	}
 	}
 	start := time.Now()
 	start := time.Now()
@@ -253,7 +254,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	}
 	}
 	if publishErr != nil {
 	if publishErr != nil {
 		if errors.Is(publishErr, circuitbreaker.ErrCircuitOpen) {
 		if errors.Is(publishErr, circuitbreaker.ErrCircuitOpen) {
-			d.Metrics.AlertsReceived.WithLabelValues("circuit_open").Inc()
+			d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "circuit_open").Inc()
 			d.Metrics.CBState.WithLabelValues("nats").Set(circuitbreaker.StateOpen)
 			d.Metrics.CBState.WithLabelValues("nats").Set(circuitbreaker.StateOpen)
 			d.Logger.Warn("circuit breaker open",
 			d.Logger.Warn("circuit breaker open",
 				"subject", subject,
 				"subject", subject,
@@ -263,7 +264,7 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 			recordHit()
 			recordHit()
 			return Reject("circuit_open", 503, "broker circuit breaker open")
 			return Reject("circuit_open", 503, "broker circuit breaker open")
 		}
 		}
-		d.Metrics.AlertsReceived.WithLabelValues("broker_unavailable").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "broker_unavailable").Inc()
 		d.Logger.Error("nats publish", "err", publishErr, "subject", subject)
 		d.Logger.Error("nats publish", "err", publishErr, "subject", subject)
 		recordHit()
 		recordHit()
 		return Reject("broker_unavailable", 503, publishErr.Error())
 		return Reject("broker_unavailable", 503, publishErr.Error())
@@ -272,9 +273,9 @@ func (d *Deps) Process(ctx context.Context, body []byte, sig string) Result {
 	d.Metrics.PayloadBytes.Observe(float64(len(payload)))
 	d.Metrics.PayloadBytes.Observe(float64(len(payload)))
 
 
 	if isNew {
 	if isNew {
-		d.Metrics.AlertsReceived.WithLabelValues("accepted").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "accepted").Inc()
 	} else {
 	} else {
-		d.Metrics.AlertsReceived.WithLabelValues("deduped").Inc()
+		d.Metrics.AlertsReceived.WithLabelValues(d.Transport, "deduped").Inc()
 	}
 	}
 
 
 	d.Logger.Info("alert accepted",
 	d.Logger.Info("alert accepted",