Skip to content

Commit c0e7dbe

Browse files
nicklaslclaudecodex
authored
feat: add provider_init_rate to telemetry (#485)
Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Co-authored-by: Codex <noreply@openai.com>
1 parent 1ac8d80 commit c0e7dbe

25 files changed

Lines changed: 1071 additions & 149 deletions

File tree

confidence-resolver/protos/confidence/flags/resolver/v1/internal_api.proto

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,8 @@ message WriteFlagAssignedResponse {
112112
// consumers everything needed to reconstruct bucket boundaries and resample
113113
// between different histogram configurations if needed.
114114
message TelemetryData {
115+
reserved 1; // was: int64 dropped_events
116+
115117
// Information about the SDK/provider
116118
Sdk sdk = 2 [
117119
(google.api.field_behavior) = OPTIONAL
@@ -128,6 +130,14 @@ message TelemetryData {
128130
// Set from the confidence-resolver crate version at build time.
129131
string resolver_version = 8;
130132

133+
repeated ProviderInitRate provider_init_rate = 9;
134+
135+
message ProviderInitRate {
136+
uint32 count = 1;
137+
reserved 2; // status — tbd
138+
map<string, string> labels = 3;
139+
}
140+
131141
message ResolveLatency {
132142
// Delta sum of observed values since the last flush.
133143
uint32 sum = 1;

confidence-resolver/src/telemetry.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -498,6 +498,7 @@ impl Telemetry {
498498
state_age,
499499
memory_bytes: (self.memory_provider)(),
500500
resolver_version: crate::version::VERSION.to_string(),
501+
provider_init_rate: Vec::new(),
501502
}
502503
}
503504
}

openfeature-provider/go/confidence/internal/flag_logger/grpc.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,7 @@ func (g *GrpcFlagLogger) Write(request *resolverv1.WriteFlagLogsRequest) {
3535
clientResolveCount := len(request.ClientResolveInfo)
3636
flagResolveCount := len(request.FlagResolveInfo)
3737

38-
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 {
38+
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 && request.TelemetryData == nil {
3939
g.logger.Debug("Skipping empty flag log request")
4040
return
4141
}

openfeature-provider/go/confidence/internal/flag_logger/multi.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ func (m *MultiDestinationFlagLogger) Write(request *resolverv1.WriteFlagLogsRequ
7373
clientResolveCount := len(request.ClientResolveInfo)
7474
flagResolveCount := len(request.FlagResolveInfo)
7575

76-
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 {
76+
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 && request.TelemetryData == nil {
7777
m.logger.Debug("Skipping empty flag log request")
7878
return
7979
}

openfeature-provider/go/confidence/internal/flag_logger/multi_test.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -173,6 +173,29 @@ func TestMultiDestinationFlagLogger_SkipsEmptyRequest(t *testing.T) {
173173
}
174174
}
175175

176+
func TestMultiDestinationFlagLogger_SendsTelemetryOnlyRequest(t *testing.T) {
177+
var called atomic.Int32
178+
logger := &MultiDestinationFlagLogger{
179+
senders: map[admin.LogDestination]logSender{
180+
admin.LogDestination_LOG_DESTINATION_SPOTIFY_EDGE: func(context.Context, *resolverv1.WriteFlagLogsRequest) error {
181+
called.Add(1)
182+
return nil
183+
},
184+
},
185+
destinations: func() []admin.LogDestination {
186+
return []admin.LogDestination{admin.LogDestination_LOG_DESTINATION_SPOTIFY_EDGE}
187+
},
188+
logger: slog.New(slog.NewTextHandler(&bytes.Buffer{}, nil)),
189+
}
190+
191+
logger.Write(&resolverv1.WriteFlagLogsRequest{TelemetryData: &resolverv1.TelemetryData{}})
192+
logger.Shutdown()
193+
194+
if called.Load() != 1 {
195+
t.Fatalf("expected telemetry-only request to be sent once, got %d", called.Load())
196+
}
197+
}
198+
176199
func TestMultiDestinationFlagLogger_AllDestinationsFail(t *testing.T) {
177200
var buf bytes.Buffer
178201
testLogger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{Level: slog.LevelWarn}))
Binary file not shown.

openfeature-provider/go/confidence/internal/proto/resolverinternal/internal_api.pb.go

Lines changed: 138 additions & 56 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

openfeature-provider/go/confidence/provider_builder.go

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import (
66
"log/slog"
77
"net/http"
88
"os"
9+
"strconv"
910
"time"
1011

1112
fl "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/flag_logger"
@@ -104,7 +105,10 @@ func NewProvider(ctx context.Context, config ProviderConfig) (*LocalResolverProv
104105
materializationStore = newRemoteMaterializationStore(resolverv1.NewInternalFlagLoggerServiceClient(conn), config.ClientSecret)
105106
}
106107

107-
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter)
108+
initLabels := map[string]string{
109+
"encryption": strconv.FormatBool(config.EncryptionKey != ""),
110+
}
111+
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter, initLabels)
108112
resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore)
109113
providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval, config.EnableApplyDedup, config.DisableExposureCollection)
110114
provider := NewLocalResolverProvider(resolverSupplierWithMaterialization, stateProvider, flagLogger, config.ClientSecret, logger, providerOpts...)
@@ -131,21 +135,27 @@ func NewProviderForTest(ctx context.Context, config ProviderTestConfig) (*LocalR
131135
if materializationStore == nil {
132136
materializationStore = newUnsupportedMaterializationStore()
133137
}
134-
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter)
138+
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter, nil)
135139
resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore)
136140
providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval, false, config.DisableExposureCollection)
137141
provider := NewLocalResolverProvider(resolverSupplierWithMaterialization, config.StateProvider, config.FlagLogger, config.ClientSecret, logger, providerOpts...)
138142

139143
return provider, nil
140144
}
141145

142-
func newLocalResolverSupplier(poolSize int, useWasmInterpreter bool) func(context.Context, lr.LogSink) lr.LocalResolver {
146+
func newLocalResolverSupplier(poolSize int, useWasmInterpreter bool, initLabels map[string]string) func(context.Context, lr.LogSink) lr.LocalResolver {
143147
cfg := lr.LocalResolverConfig{
144148
PoolSize: poolSize,
145149
UseWasmInterpreter: useWasmInterpreter,
146150
}
147151
return func(ctx context.Context, logSink lr.LogSink) lr.LocalResolver {
148-
return lr.NewLocalResolver(ctx, logSink, cfg)
152+
return newProviderTelemetryResolver(
153+
logSink,
154+
initLabels,
155+
func(providerLogSink lr.LogSink) lr.LocalResolver {
156+
return lr.NewLocalResolver(ctx, providerLogSink, cfg)
157+
},
158+
)
149159
}
150160
}
151161

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,104 @@
1+
package confidence
2+
3+
import (
4+
"context"
5+
"sync"
6+
7+
lr "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/local_resolver"
8+
resolvertypes "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolver"
9+
resolverv1 "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternal"
10+
"github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/wasm"
11+
)
12+
13+
// providerTelemetryResolver owns provider-scoped telemetry above pooling and recovery. The inner
14+
// resolver stack is constructed with writeLogs, so every pooled WASM instance shares one init
15+
// state. Close forces an init-only request when no earlier full flush produced one.
16+
type providerTelemetryResolver struct {
17+
delegate lr.LocalResolver
18+
logSink lr.LogSink
19+
labels map[string]string
20+
sdk *resolvertypes.Sdk
21+
22+
mu sync.Mutex
23+
initSent bool
24+
}
25+
26+
func newProviderTelemetryResolver(
27+
logSink lr.LogSink,
28+
labels map[string]string,
29+
innerFactory func(lr.LogSink) lr.LocalResolver,
30+
) lr.LocalResolver {
31+
r := &providerTelemetryResolver{
32+
logSink: logSink,
33+
labels: labels,
34+
sdk: &resolvertypes.Sdk{
35+
Sdk: &resolvertypes.Sdk_Id{Id: resolvertypes.SdkId_SDK_ID_GO_LOCAL_PROVIDER},
36+
Version: Version,
37+
},
38+
}
39+
r.delegate = innerFactory(r.writeLogs)
40+
return r
41+
}
42+
43+
func (r *providerTelemetryResolver) writeLogs(logs *resolverv1.WriteFlagLogsRequest) {
44+
r.mu.Lock()
45+
defer r.mu.Unlock()
46+
47+
if !r.initSent {
48+
r.addInitTelemetry(logs)
49+
}
50+
r.logSink(logs)
51+
r.initSent = true
52+
}
53+
54+
func (r *providerTelemetryResolver) emitInitIfPending() {
55+
r.mu.Lock()
56+
defer r.mu.Unlock()
57+
58+
if r.initSent {
59+
return
60+
}
61+
logs := &resolverv1.WriteFlagLogsRequest{}
62+
r.addInitTelemetry(logs)
63+
r.logSink(logs)
64+
r.initSent = true
65+
}
66+
67+
func (r *providerTelemetryResolver) addInitTelemetry(logs *resolverv1.WriteFlagLogsRequest) {
68+
if logs.TelemetryData == nil {
69+
logs.TelemetryData = &resolverv1.TelemetryData{}
70+
}
71+
logs.TelemetryData.Sdk = r.sdk
72+
logs.TelemetryData.ProviderInitRate = append(
73+
logs.TelemetryData.ProviderInitRate,
74+
&resolverv1.TelemetryData_ProviderInitRate{Count: 1, Labels: r.labels},
75+
)
76+
}
77+
78+
func (r *providerTelemetryResolver) SetResolverState(request *wasm.SetResolverStateRequest) error {
79+
return r.delegate.SetResolverState(request)
80+
}
81+
82+
func (r *providerTelemetryResolver) ResolveProcess(request *wasm.ResolveProcessRequest) (*wasm.ResolveProcessResponse, error) {
83+
return r.delegate.ResolveProcess(request)
84+
}
85+
86+
func (r *providerTelemetryResolver) RegisterResolve(request *wasm.RegisterResolveRequest) {
87+
r.delegate.RegisterResolve(request)
88+
}
89+
90+
func (r *providerTelemetryResolver) ApplyFlags(request *resolvertypes.ApplyFlagsRequest) error {
91+
return r.delegate.ApplyFlags(request)
92+
}
93+
94+
func (r *providerTelemetryResolver) FlushAllLogs() error { return r.delegate.FlushAllLogs() }
95+
func (r *providerTelemetryResolver) FlushAssignLogs() error { return r.delegate.FlushAssignLogs() }
96+
func (r *providerTelemetryResolver) PrometheusSnapshot(bucketsPerDecade uint32, openmetrics bool) string {
97+
return r.delegate.PrometheusSnapshot(bucketsPerDecade, openmetrics)
98+
}
99+
100+
func (r *providerTelemetryResolver) Close(ctx context.Context) error {
101+
err := r.delegate.Close(ctx)
102+
r.emitInitIfPending()
103+
return err
104+
}
Lines changed: 109 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,109 @@
1+
package confidence
2+
3+
import (
4+
"context"
5+
"testing"
6+
7+
lr "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/local_resolver"
8+
resolvertypes "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolver"
9+
resolverv1 "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternal"
10+
"github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/wasm"
11+
)
12+
13+
func TestProviderTelemetryResolverEmitsOnceAcrossPooledFlushes(t *testing.T) {
14+
var captured []*resolverv1.WriteFlagLogsRequest
15+
resolver := newProviderTelemetryResolver(
16+
func(logs *resolverv1.WriteFlagLogsRequest) { captured = append(captured, logs) },
17+
map[string]string{"encryption": "true"},
18+
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink, flushCount: 3} },
19+
)
20+
21+
if err := resolver.FlushAllLogs(); err != nil {
22+
t.Fatal(err)
23+
}
24+
25+
if got := providerInitEventCount(captured); got != 1 {
26+
t.Fatalf("expected one provider init event across pooled flushes, got %d", got)
27+
}
28+
telemetry := captured[0].GetTelemetryData()
29+
if telemetry.GetSdk().GetId() != resolvertypes.SdkId_SDK_ID_GO_LOCAL_PROVIDER {
30+
t.Fatalf("unexpected SDK ID: %v", telemetry.GetSdk().GetId())
31+
}
32+
if telemetry.GetSdk().GetVersion() != Version {
33+
t.Fatalf("unexpected SDK version: %q", telemetry.GetSdk().GetVersion())
34+
}
35+
}
36+
37+
func TestProviderTelemetryResolverCloseEmitsWithoutResolve(t *testing.T) {
38+
var captured []*resolverv1.WriteFlagLogsRequest
39+
resolver := newProviderTelemetryResolver(
40+
func(logs *resolverv1.WriteFlagLogsRequest) { captured = append(captured, logs) },
41+
map[string]string{"encryption": "true"},
42+
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink} },
43+
)
44+
45+
if err := resolver.Close(context.Background()); err != nil {
46+
t.Fatal(err)
47+
}
48+
49+
if got := providerInitEventCount(captured); got != 1 {
50+
t.Fatalf("expected shutdown to emit one provider init event, got %d", got)
51+
}
52+
}
53+
54+
func TestProviderTelemetryResolverRetriesAfterSinkFailure(t *testing.T) {
55+
attempts := 0
56+
var captured []*resolverv1.WriteFlagLogsRequest
57+
resolver := newProviderTelemetryResolver(
58+
func(logs *resolverv1.WriteFlagLogsRequest) {
59+
attempts++
60+
if attempts == 1 {
61+
panic("send failed")
62+
}
63+
captured = append(captured, logs)
64+
},
65+
nil,
66+
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink, flushCount: 1} },
67+
)
68+
69+
func() {
70+
defer func() { _ = recover() }()
71+
_ = resolver.FlushAllLogs()
72+
}()
73+
if err := resolver.FlushAllLogs(); err != nil {
74+
t.Fatal(err)
75+
}
76+
77+
if got := providerInitEventCount(captured); got != 1 {
78+
t.Fatalf("expected init telemetry to be retried, got %d events", got)
79+
}
80+
}
81+
82+
func providerInitEventCount(requests []*resolverv1.WriteFlagLogsRequest) int {
83+
count := 0
84+
for _, request := range requests {
85+
count += len(request.GetTelemetryData().GetProviderInitRate())
86+
}
87+
return count
88+
}
89+
90+
type telemetryTestResolver struct {
91+
sink lr.LogSink
92+
flushCount int
93+
}
94+
95+
func (r *telemetryTestResolver) SetResolverState(*wasm.SetResolverStateRequest) error { return nil }
96+
func (r *telemetryTestResolver) ResolveProcess(*wasm.ResolveProcessRequest) (*wasm.ResolveProcessResponse, error) {
97+
return &wasm.ResolveProcessResponse{}, nil
98+
}
99+
func (r *telemetryTestResolver) RegisterResolve(*wasm.RegisterResolveRequest) {}
100+
func (r *telemetryTestResolver) ApplyFlags(*resolvertypes.ApplyFlagsRequest) error { return nil }
101+
func (r *telemetryTestResolver) FlushAllLogs() error {
102+
for range r.flushCount {
103+
r.sink(&resolverv1.WriteFlagLogsRequest{})
104+
}
105+
return nil
106+
}
107+
func (r *telemetryTestResolver) FlushAssignLogs() error { return nil }
108+
func (r *telemetryTestResolver) PrometheusSnapshot(uint32, bool) string { return "" }
109+
func (r *telemetryTestResolver) Close(context.Context) error { return nil }

0 commit comments

Comments
 (0)