Skip to content

Commit b179025

Browse files
authored
Merge pull request #22 from InjectiveLabs/f/otel-counters
- feat: allow otel counters - feat: configure metrics
2 parents a84c7a4 + a67f6cd commit b179025

8 files changed

Lines changed: 767 additions & 19 deletions

File tree

client.go

Lines changed: 87 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,9 @@ package metrics
22

33
import (
44
"context"
5+
"os"
56
"runtime"
7+
"strings"
68
"sync"
79
"time"
810

@@ -38,23 +40,24 @@ var (
3840
)
3941

4042
type StatterConfig struct {
41-
Addr string // localhost:8125
42-
Prefix string // metrics prefix
43-
Agent string // telegraf/datadog
44-
EnvName string // dev/test/staging/prod
45-
HostName string // hostname
46-
Version string // version
47-
DefaultTags []interface{} // default tags for all metrics
48-
StuckFunctionTimeout time.Duration // stuck time
49-
MockingThreshold time.Duration // mocking threshold
50-
MockingEnabled bool // whether to enable mock statter, which only produce logs
51-
Disabled bool // whether to disable metrics completely
52-
TracingEnabled bool // whether tracing should be enabled
53-
ProfilingEnabled bool // whether Datadog profiling should be enabled
54-
MixPanelEnabled bool // whether MixPanel should be enabled
55-
MixPanelProjectToken string // MixPanel project token
56-
OTELInsecure bool // disable TLS (use for self-hosted SigNoz without TLS)
57-
OTELHeaders map[string]string // extra headers, e.g. {"signoz-access-token": "<token>"} for SigNoz Cloud
43+
Addr string // localhost:8125
44+
Prefix string // metrics prefix
45+
Agent string // telegraf/datadog
46+
EnvName string // dev/test/staging/prod
47+
HostName string // hostname
48+
Version string // version
49+
DefaultTags []interface{} // default tags for all metrics
50+
StuckFunctionTimeout time.Duration // stuck time
51+
MockingThreshold time.Duration // mocking threshold
52+
MockingEnabled bool // whether to enable mock statter, which only produce logs
53+
Disabled bool // whether to disable metrics completely
54+
TracingEnabled bool // whether tracing should be enabled
55+
ProfilingEnabled bool // whether Datadog profiling should be enabled
56+
MixPanelEnabled bool // whether MixPanel should be enabled
57+
MixPanelProjectToken string // MixPanel project token
58+
OTELInsecure bool // disable TLS (use for self-hosted SigNoz without TLS)
59+
OTELHeaders map[string]string // extra headers, e.g. {"signoz-access-token": "<token>"} for SigNoz Cloud
60+
OTELUseCounterForCount bool // use monotonic OTel counters for Count/Incr instead of UpDownCounters
5861
}
5962

6063
func (m *StatterConfig) BaseTags() []string {
@@ -177,6 +180,7 @@ func Init(addr string, prefix string, cfg *StatterConfig) error {
177180
cfg.OTELInsecure,
178181
cfg.OTELHeaders,
179182
config.BaseTags(),
183+
cfg.OTELUseCounterForCount,
180184
)
181185

182186
default:
@@ -229,6 +233,72 @@ func Init(addr string, prefix string, cfg *StatterConfig) error {
229233
return nil
230234
}
231235

236+
type ServiceConfig struct {
237+
Disabled bool
238+
ServiceName string
239+
AgentID string
240+
AgentAddress string
241+
MetricsPrefix string
242+
EnvName string
243+
MockingEnabled bool
244+
MockingThreshold time.Duration
245+
MixPanelEnabled bool
246+
MixPanelProjectToken string
247+
OTelInsecure bool
248+
OTelUseCounters bool
249+
TracingEnabled bool
250+
RetryInitialInterval time.Duration
251+
}
252+
253+
func (c ServiceConfig) normalizePrefix() string {
254+
return strings.TrimRight(c.MetricsPrefix, ".") + "."
255+
}
256+
257+
func InitService(ctx context.Context, cfg ServiceConfig) (func(timeout time.Duration), error) {
258+
closeFn := func(time.Duration) {}
259+
if err := cfg.Validate(); err != nil {
260+
return closeFn, err
261+
}
262+
if cfg.Disabled {
263+
return closeFn, nil
264+
}
265+
266+
if cfg.RetryInitialInterval <= 0 {
267+
cfg.RetryInitialInterval = 10 * time.Second
268+
}
269+
270+
for {
271+
hostname, _ := os.Hostname()
272+
err := Init(cfg.AgentAddress, cfg.normalizePrefix(), &StatterConfig{
273+
Agent: cfg.AgentID,
274+
EnvName: cfg.EnvName,
275+
HostName: hostname,
276+
MockingEnabled: cfg.MockingEnabled,
277+
MockingThreshold: cfg.MockingThreshold,
278+
MixPanelEnabled: cfg.MixPanelEnabled,
279+
MixPanelProjectToken: cfg.MixPanelProjectToken,
280+
OTELInsecure: cfg.OTelInsecure,
281+
OTELUseCounterForCount: cfg.OTelUseCounters,
282+
TracingEnabled: cfg.TracingEnabled,
283+
DefaultTags: []interface{}{"service.name", cfg.ServiceName},
284+
})
285+
if err != nil {
286+
log.WithError(err).Warningf("metrics init failed, will retry in %s seconds", cfg.RetryInitialInterval)
287+
select {
288+
case <-ctx.Done():
289+
return closeFn, nil
290+
case <-time.After(cfg.RetryInitialInterval):
291+
}
292+
continue
293+
}
294+
log.Debugf("metrics %s client initialized at %s with prefix %s (mocking %v)",
295+
cfg.AgentID, cfg.AgentAddress, cfg.normalizePrefix(), cfg.MockingEnabled)
296+
break
297+
}
298+
299+
return CloseWithTimeout, nil
300+
}
301+
232302
func StartMixPanel(projectToken string) {
233303
clientMux.Lock()
234304
defer clientMux.Unlock()

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ require (
99
github.com/cosmos/cosmos-sdk v0.50.6
1010
github.com/mixpanel/mixpanel-go v1.2.1
1111
github.com/pkg/errors v0.9.1
12+
github.com/spf13/pflag v1.0.5
1213
github.com/stretchr/testify v1.11.1
1314
go.opentelemetry.io/otel v1.43.0
1415
go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.43.0
@@ -120,7 +121,6 @@ require (
120121
github.com/spaolacci/murmur3 v1.1.0 // indirect
121122
github.com/spf13/cast v1.6.0 // indirect
122123
github.com/spf13/cobra v1.8.0 // indirect
123-
github.com/spf13/pflag v1.0.5 // indirect
124124
github.com/syndtr/goleveldb v1.0.1-0.20220721030215-126854af5e6d // indirect
125125
github.com/tendermint/go-amino v0.16.0 // indirect
126126
github.com/tinylib/msgp v1.1.8 // indirect

otel.go

Lines changed: 82 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package metrics
22

33
import (
44
"context"
5+
"fmt"
56
"strings"
67
"sync"
78
"time"
@@ -20,8 +21,10 @@ type otelStatter struct {
2021
meter otelmetric.Meter
2122
meterProvider *sdkmetric.MeterProvider
2223
prefix string
24+
useCounters bool
2325

2426
mu sync.RWMutex
27+
counters map[string]otelmetric.Int64Counter
2528
updownCounters map[string]otelmetric.Int64UpDownCounter
2629
gauges map[string]otelmetric.Float64Gauge
2730
histograms map[string]otelmetric.Float64Histogram
@@ -41,7 +44,7 @@ func newOTELResource(baseTags []string) *resource.Resource {
4144
return res
4245
}
4346

44-
func newOTELStatter(endpoint, prefix string, insecure bool, headers map[string]string, baseTags []string) (Statter, error) {
47+
func newOTELStatter(endpoint, prefix string, insecure bool, headers map[string]string, baseTags []string, useCounters bool) (Statter, error) {
4548
ctx := context.Background()
4649

4750
metricOpts := []otlpmetricgrpc.Option{
@@ -70,6 +73,8 @@ func newOTELStatter(endpoint, prefix string, insecure bool, headers map[string]s
7073
meter: mp.Meter(prefix),
7174
meterProvider: mp,
7275
prefix: prefix,
76+
useCounters: useCounters,
77+
counters: make(map[string]otelmetric.Int64Counter),
7378
updownCounters: make(map[string]otelmetric.Int64UpDownCounter),
7479
gauges: make(map[string]otelmetric.Float64Gauge),
7580
histograms: make(map[string]otelmetric.Float64Histogram),
@@ -114,6 +119,27 @@ func (s *otelStatter) tagsToAttrs(tags []string) []attribute.KeyValue {
114119
return attrs
115120
}
116121

122+
func (s *otelStatter) getCounter(name string) (otelmetric.Int64Counter, error) {
123+
fullName := s.prefix + name
124+
s.mu.RLock()
125+
c, ok := s.counters[fullName]
126+
s.mu.RUnlock()
127+
if ok {
128+
return c, nil
129+
}
130+
s.mu.Lock()
131+
defer s.mu.Unlock()
132+
if c, ok = s.counters[fullName]; ok {
133+
return c, nil
134+
}
135+
c, err := s.meter.Int64Counter(fullName)
136+
if err != nil {
137+
return nil, err
138+
}
139+
s.counters[fullName] = c
140+
return c, nil
141+
}
142+
117143
func (s *otelStatter) getUpDownCounter(name string) (otelmetric.Int64UpDownCounter, error) {
118144
fullName := s.prefix + name
119145
s.mu.RLock()
@@ -178,6 +204,17 @@ func (s *otelStatter) getHistogram(name string) (otelmetric.Float64Histogram, er
178204
}
179205

180206
func (s *otelStatter) Count(name string, value int64, tags []string, rate float64) error {
207+
if s.useCounters {
208+
if value < 0 {
209+
return fmt.Errorf("counter value must be non-negative")
210+
}
211+
c, err := s.getCounter(name)
212+
if err != nil {
213+
return err
214+
}
215+
c.Add(context.Background(), value, otelmetric.WithAttributes(s.tagsToAttrs(tags)...))
216+
return nil
217+
}
181218
c, err := s.getUpDownCounter(name)
182219
if err != nil {
183220
return err
@@ -187,6 +224,14 @@ func (s *otelStatter) Count(name string, value int64, tags []string, rate float6
187224
}
188225

189226
func (s *otelStatter) Incr(name string, tags []string, rate float64) error {
227+
if s.useCounters {
228+
c, err := s.getCounter(name)
229+
if err != nil {
230+
return err
231+
}
232+
c.Add(context.Background(), 1, otelmetric.WithAttributes(s.tagsToAttrs(tags)...))
233+
return nil
234+
}
190235
c, err := s.getUpDownCounter(name)
191236
if err != nil {
192237
return err
@@ -239,3 +284,39 @@ func (s *otelStatter) Close() error {
239284
func (s *otelStatter) CloseCtx(ctx context.Context) error {
240285
return s.meterProvider.Shutdown(ctx)
241286
}
287+
288+
func OTelSetupHistogram(name string, opts ...otelmetric.Float64HistogramOption) error {
289+
clientMux.RLock()
290+
c := client
291+
clientMux.RUnlock()
292+
293+
if c == nil {
294+
return nil
295+
}
296+
297+
s, ok := c.(*otelStatter)
298+
if !ok {
299+
return nil
300+
}
301+
302+
fullName := s.prefix + name
303+
304+
s.mu.Lock()
305+
defer s.mu.Unlock()
306+
307+
if _, exists := s.histograms[fullName]; exists {
308+
return fmt.Errorf("histogram %s already initialized", fullName)
309+
}
310+
311+
finalOpts := append([]otelmetric.Float64HistogramOption{
312+
otelmetric.WithUnit("ms"),
313+
}, opts...)
314+
315+
h, err := s.meter.Float64Histogram(fullName, finalOpts...)
316+
if err != nil {
317+
return err
318+
}
319+
320+
s.histograms[fullName] = h
321+
return nil
322+
}

otel_test.go

Lines changed: 102 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,102 @@
1+
package metrics
2+
3+
import (
4+
"testing"
5+
6+
"github.com/stretchr/testify/require"
7+
otelmetric "go.opentelemetry.io/otel/metric"
8+
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
9+
"go.opentelemetry.io/otel/sdk/metric/metricdata"
10+
)
11+
12+
func newTestOTELStatter(useCounters bool) (*otelStatter, *sdkmetric.ManualReader) {
13+
reader := sdkmetric.NewManualReader()
14+
mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
15+
return &otelStatter{
16+
meter: mp.Meter("test"),
17+
meterProvider: mp,
18+
useCounters: useCounters,
19+
counters: make(map[string]otelmetric.Int64Counter),
20+
updownCounters: make(map[string]otelmetric.Int64UpDownCounter),
21+
gauges: make(map[string]otelmetric.Float64Gauge),
22+
histograms: make(map[string]otelmetric.Float64Histogram),
23+
}, reader
24+
}
25+
26+
func collectOTELMetric(t *testing.T, reader *sdkmetric.ManualReader, name string) metricdata.Metrics {
27+
t.Helper()
28+
29+
var rm metricdata.ResourceMetrics
30+
require.NoError(t, reader.Collect(t.Context(), &rm))
31+
for _, sm := range rm.ScopeMetrics {
32+
for _, m := range sm.Metrics {
33+
if m.Name == name {
34+
return m
35+
}
36+
}
37+
}
38+
require.Failf(t, "metric not found", "metric %q not found", name)
39+
return metricdata.Metrics{}
40+
}
41+
42+
func requireOTELInt64Sum(t *testing.T, m metricdata.Metrics) metricdata.Sum[int64] {
43+
t.Helper()
44+
45+
sum, ok := m.Data.(metricdata.Sum[int64])
46+
require.Truef(t, ok, "expected int64 sum, got %T", m.Data)
47+
return sum
48+
}
49+
50+
func TestOTELCountDefaultsToUpDownCounter(t *testing.T) {
51+
statter, reader := newTestOTELStatter(false)
52+
t.Cleanup(func() {
53+
require.NoError(t, statter.Close())
54+
})
55+
56+
require.NoError(t, statter.Count("events.total", 5, []string{"status=ok"}, 1))
57+
58+
sum := requireOTELInt64Sum(t, collectOTELMetric(t, reader, "events.total"))
59+
require.False(t, sum.IsMonotonic)
60+
require.Len(t, sum.DataPoints, 1)
61+
require.EqualValues(t, 5, sum.DataPoints[0].Value)
62+
}
63+
64+
func TestOTELCountCanUseCounter(t *testing.T) {
65+
statter, reader := newTestOTELStatter(true)
66+
t.Cleanup(func() {
67+
require.NoError(t, statter.Close())
68+
})
69+
70+
require.NoError(t, statter.Count("events.total", 5, []string{"status=ok"}, 1))
71+
require.NoError(t, statter.Incr("events.total", []string{"status=ok"}, 1))
72+
73+
sum := requireOTELInt64Sum(t, collectOTELMetric(t, reader, "events.total"))
74+
require.True(t, sum.IsMonotonic)
75+
require.Len(t, sum.DataPoints, 1)
76+
require.EqualValues(t, 6, sum.DataPoints[0].Value)
77+
}
78+
79+
func TestOTELCountRejectsNegativeValueWithCounterOption(t *testing.T) {
80+
statter, _ := newTestOTELStatter(true)
81+
t.Cleanup(func() {
82+
require.NoError(t, statter.Close())
83+
})
84+
85+
err := statter.Count("events.total", -1, []string{"status=ok"}, 1)
86+
require.EqualError(t, err, "counter value must be non-negative")
87+
require.Empty(t, statter.counters)
88+
}
89+
90+
func TestOTELDecrUsesUpDownCounterWithCounterOption(t *testing.T) {
91+
statter, reader := newTestOTELStatter(true)
92+
t.Cleanup(func() {
93+
require.NoError(t, statter.Close())
94+
})
95+
96+
require.NoError(t, statter.Decr("active.requests", []string{"status=ok"}, 1))
97+
98+
sum := requireOTELInt64Sum(t, collectOTELMetric(t, reader, "active.requests"))
99+
require.False(t, sum.IsMonotonic)
100+
require.Len(t, sum.DataPoints, 1)
101+
require.EqualValues(t, -1, sum.DataPoints[0].Value)
102+
}

0 commit comments

Comments
 (0)