From 57d721354a0127b938eb9081dfe4d75fca49e6ac Mon Sep 17 00:00:00 2001 From: Julien Harbulot Date: Wed, 6 May 2026 21:41:58 +1200 Subject: [PATCH 1/3] Track pubsub topic validator receives --- network/topics/controller.go | 8 ++++- network/topics/observability.go | 18 +++++++++++ network/topics/observability_test.go | 48 ++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+), 1 deletion(-) create mode 100644 network/topics/observability_test.go diff --git a/network/topics/controller.go b/network/topics/controller.go index be69588a57..2181eddb45 100644 --- a/network/topics/controller.go +++ b/network/topics/controller.go @@ -308,7 +308,13 @@ func (ctrl *topicsCtrl) setupTopicValidator(name string) error { opts := []pubsub.ValidatorOpt{pubsub.WithValidatorTimeout(topicValidatorTimeout)} - err = ctrl.ps.RegisterTopicValidator(name, ctrl.msgValidator.ValidatorForTopic(name), opts...) + validator := ctrl.msgValidator.ValidatorForTopic(name) + wrappedValidator := func(ctx context.Context, p peer.ID, pmsg *pubsub.Message) pubsub.ValidationResult { + recordPubsubMessageReceived(ctx, name) + return validator(ctx, p, pmsg) + } + + err = ctrl.ps.RegisterTopicValidator(name, wrappedValidator, opts...) if err != nil { return fmt.Errorf("could not register topic validator: %w", err) } diff --git a/network/topics/observability.go b/network/topics/observability.go index 08ac2a5808..1c69f6fcd6 100644 --- a/network/topics/observability.go +++ b/network/topics/observability.go @@ -1,6 +1,8 @@ package topics import ( + "context" + "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/metric" @@ -12,6 +14,8 @@ import ( const ( observabilityName = "github.com/ssvlabs/ssv/network/topics" observabilityNamespace = "ssv.p2p.messages" + + pubsubObservabilityNamespace = "ssv.p2p.pubsub.messages" ) var ( @@ -29,6 +33,12 @@ var ( metric.WithUnit("{message}"), metric.WithDescription("total number of outbound(broadcasted) messages"))) + pubsubMessagesReceivedCounter = metrics.New( + meter.Int64Counter( + observability.InstrumentName(pubsubObservabilityNamespace, "received"), + metric.WithUnit("{message}"), + metric.WithDescription("total number of messages received by the pubsub topic validator"))) + msgIDHandlerBufferFallbackCounter = metrics.New( meter.Int64Counter( observability.InstrumentName(observabilityNamespace, "msg_id_buffer_fallback"), @@ -36,6 +46,10 @@ var ( metric.WithDescription("total number of msg_id add operations processed synchronously because the async buffer was full"))) ) +func pubsubTopicAttribute(value string) attribute.KeyValue { + return attribute.String("topic", value) +} + func messageTopicAttribute(value string) attribute.KeyValue { return attribute.String("ssv.p2p.message.topic", value) } @@ -46,3 +60,7 @@ func messageTypeAttribute(value uint64) attribute.KeyValue { Value: observability.Uint64AttributeValue(value), } } + +func recordPubsubMessageReceived(ctx context.Context, topic string) { + pubsubMessagesReceivedCounter.Add(ctx, 1, metric.WithAttributes(pubsubTopicAttribute(topic))) +} diff --git a/network/topics/observability_test.go b/network/topics/observability_test.go new file mode 100644 index 0000000000..f6ce4ed419 --- /dev/null +++ b/network/topics/observability_test.go @@ -0,0 +1,48 @@ +package topics + +import ( + "testing" + + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" +) + +func TestRecordPubsubMessageReceived(t *testing.T) { + reader := metric.NewManualReader() + provider := metric.NewMeterProvider(metric.WithReader(reader)) + previousProvider := otel.GetMeterProvider() + otel.SetMeterProvider(provider) + t.Cleanup(func() { + otel.SetMeterProvider(previousProvider) + require.NoError(t, provider.Shutdown(t.Context())) + }) + + const topic = "ssv.v2.42" + recordPubsubMessageReceived(t.Context(), topic) + recordPubsubMessageReceived(t.Context(), topic) + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(t.Context(), &rm)) + + for _, scopeMetrics := range rm.ScopeMetrics { + for _, metric := range scopeMetrics.Metrics { + if metric.Name != "ssv.p2p.pubsub.messages.received" { + continue + } + + sum, ok := metric.Data.(metricdata.Sum[int64]) + require.True(t, ok) + require.Len(t, sum.DataPoints, 1) + require.EqualValues(t, 2, sum.DataPoints[0].Value) + + topicAttr, ok := sum.DataPoints[0].Attributes.Value("topic") + require.True(t, ok) + require.Equal(t, topic, topicAttr.AsString()) + return + } + } + + t.Fatal("pubsub received metric was not collected") +} From 716103ec622e780e447d5ec88fa96cedec906d70 Mon Sep 17 00:00:00 2001 From: Julien Harbulot Date: Wed, 6 May 2026 23:50:23 +1200 Subject: [PATCH 2/3] Namespace pubsub topic metric attribute --- network/topics/observability.go | 3 ++- network/topics/observability_test.go | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/network/topics/observability.go b/network/topics/observability.go index 1c69f6fcd6..5b357054cc 100644 --- a/network/topics/observability.go +++ b/network/topics/observability.go @@ -16,6 +16,7 @@ const ( observabilityNamespace = "ssv.p2p.messages" pubsubObservabilityNamespace = "ssv.p2p.pubsub.messages" + pubsubTopicAttributeKey = "ssv.p2p.pubsub.topic" ) var ( @@ -47,7 +48,7 @@ var ( ) func pubsubTopicAttribute(value string) attribute.KeyValue { - return attribute.String("topic", value) + return attribute.String(pubsubTopicAttributeKey, value) } func messageTopicAttribute(value string) attribute.KeyValue { diff --git a/network/topics/observability_test.go b/network/topics/observability_test.go index f6ce4ed419..41f1dc650a 100644 --- a/network/topics/observability_test.go +++ b/network/topics/observability_test.go @@ -37,7 +37,7 @@ func TestRecordPubsubMessageReceived(t *testing.T) { require.Len(t, sum.DataPoints, 1) require.EqualValues(t, 2, sum.DataPoints[0].Value) - topicAttr, ok := sum.DataPoints[0].Attributes.Value("topic") + topicAttr, ok := sum.DataPoints[0].Attributes.Value(pubsubTopicAttributeKey) require.True(t, ok) require.Equal(t, topic, topicAttr.AsString()) return From a925ed505a03b4c69da3d4c5b442dd0e84206a39 Mon Sep 17 00:00:00 2001 From: iurii-ssv <183610124+iurii-ssv@users.noreply.github.com> Date: Thu, 25 Jun 2026 11:20:17 +0300 Subject: [PATCH 3/3] network/topics: clarify pubsub-received metric semantics, rename shadowed test var (#2862) - Tighten the counter description to spell out that it counts deliveries to the topic validator *before* SSV validation runs, and reference inboundMessageCounter so operators can see the post-validation counterpart at a glance. - Add a brief comment on recordPubsubMessageReceived noting it is invoked from the validator wrapper before the inner validator, so all outcomes (including reject/timeout) are counted. - Rename the test loop variable that shadowed the imported metric package. --- network/topics/observability.go | 5 ++++- network/topics/observability_test.go | 6 +++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/network/topics/observability.go b/network/topics/observability.go index 5b357054cc..2c05a9327a 100644 --- a/network/topics/observability.go +++ b/network/topics/observability.go @@ -38,7 +38,7 @@ var ( meter.Int64Counter( observability.InstrumentName(pubsubObservabilityNamespace, "received"), metric.WithUnit("{message}"), - metric.WithDescription("total number of messages received by the pubsub topic validator"))) + metric.WithDescription("total number of messages delivered to the pubsub topic validator, before SSV validation runs (compare with ssv_p2p_messages_in_total for the post-validation rate)"))) msgIDHandlerBufferFallbackCounter = metrics.New( meter.Int64Counter( @@ -62,6 +62,9 @@ func messageTypeAttribute(value uint64) attribute.KeyValue { } } +// recordPubsubMessageReceived is called from the topic validator wrapper before the inner SSV +// validator runs, so the counter increments for every message libp2p hands to the validator +// regardless of validation outcome (accept/ignore/reject/timeout). func recordPubsubMessageReceived(ctx context.Context, topic string) { pubsubMessagesReceivedCounter.Add(ctx, 1, metric.WithAttributes(pubsubTopicAttribute(topic))) } diff --git a/network/topics/observability_test.go b/network/topics/observability_test.go index 41f1dc650a..8e9b30ce60 100644 --- a/network/topics/observability_test.go +++ b/network/topics/observability_test.go @@ -27,12 +27,12 @@ func TestRecordPubsubMessageReceived(t *testing.T) { require.NoError(t, reader.Collect(t.Context(), &rm)) for _, scopeMetrics := range rm.ScopeMetrics { - for _, metric := range scopeMetrics.Metrics { - if metric.Name != "ssv.p2p.pubsub.messages.received" { + for _, m := range scopeMetrics.Metrics { + if m.Name != "ssv.p2p.pubsub.messages.received" { continue } - sum, ok := metric.Data.(metricdata.Sum[int64]) + sum, ok := m.Data.(metricdata.Sum[int64]) require.True(t, ok) require.Len(t, sum.DataPoints, 1) require.EqualValues(t, 2, sum.DataPoints[0].Value)