Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,8 @@ message WriteFlagAssignedResponse {
// consumers everything needed to reconstruct bucket boundaries and resample
// between different histogram configurations if needed.
message TelemetryData {
reserved 1; // was: int64 dropped_events

// Information about the SDK/provider
Sdk sdk = 2 [
(google.api.field_behavior) = OPTIONAL
Expand All @@ -128,6 +130,14 @@ message TelemetryData {
// Set from the confidence-resolver crate version at build time.
string resolver_version = 8;

repeated ProviderInitRate provider_init_rate = 9;

message ProviderInitRate {
uint32 count = 1;
reserved 2; // status — tbd
map<string, string> labels = 3;
}

message ResolveLatency {
// Delta sum of observed values since the last flush.
uint32 sum = 1;
Expand Down
1 change: 1 addition & 0 deletions confidence-resolver/src/telemetry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -498,6 +498,7 @@ impl Telemetry {
state_age,
memory_bytes: (self.memory_provider)(),
resolver_version: crate::version::VERSION.to_string(),
provider_init_rate: Vec::new(),
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ func (g *GrpcFlagLogger) Write(request *resolverv1.WriteFlagLogsRequest) {
clientResolveCount := len(request.ClientResolveInfo)
flagResolveCount := len(request.FlagResolveInfo)

if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 {
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 && request.TelemetryData == nil {
g.logger.Debug("Skipping empty flag log request")
return
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ func (m *MultiDestinationFlagLogger) Write(request *resolverv1.WriteFlagLogsRequ
clientResolveCount := len(request.ClientResolveInfo)
flagResolveCount := len(request.FlagResolveInfo)

if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 {
if clientResolveCount == 0 && flagAssignedCount == 0 && flagResolveCount == 0 && request.TelemetryData == nil {
m.logger.Debug("Skipping empty flag log request")
return
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,29 @@ func TestMultiDestinationFlagLogger_SkipsEmptyRequest(t *testing.T) {
}
}

func TestMultiDestinationFlagLogger_SendsTelemetryOnlyRequest(t *testing.T) {
var called atomic.Int32
logger := &MultiDestinationFlagLogger{
senders: map[admin.LogDestination]logSender{
admin.LogDestination_LOG_DESTINATION_SPOTIFY_EDGE: func(context.Context, *resolverv1.WriteFlagLogsRequest) error {
called.Add(1)
return nil
},
},
destinations: func() []admin.LogDestination {
return []admin.LogDestination{admin.LogDestination_LOG_DESTINATION_SPOTIFY_EDGE}
},
logger: slog.New(slog.NewTextHandler(&bytes.Buffer{}, nil)),
}

logger.Write(&resolverv1.WriteFlagLogsRequest{TelemetryData: &resolverv1.TelemetryData{}})
logger.Shutdown()

if called.Load() != 1 {
t.Fatalf("expected telemetry-only request to be sent once, got %d", called.Load())
}
}

func TestMultiDestinationFlagLogger_AllDestinationsFail(t *testing.T) {
var buf bytes.Buffer
testLogger := slog.New(slog.NewTextHandler(&buf, &slog.HandlerOptions{Level: slog.LevelWarn}))
Expand Down
Binary file not shown.

Large diffs are not rendered by default.

18 changes: 14 additions & 4 deletions openfeature-provider/go/confidence/provider_builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"log/slog"
"net/http"
"os"
"strconv"
"time"

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

resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter)
initLabels := map[string]string{
"encryption": strconv.FormatBool(config.EncryptionKey != ""),
}
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter, initLabels)
resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore)
providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval, config.EnableApplyDedup, config.DisableExposureCollection)
provider := NewLocalResolverProvider(resolverSupplierWithMaterialization, stateProvider, flagLogger, config.ClientSecret, logger, providerOpts...)
Expand All @@ -131,21 +135,27 @@ func NewProviderForTest(ctx context.Context, config ProviderTestConfig) (*LocalR
if materializationStore == nil {
materializationStore = newUnsupportedMaterializationStore()
}
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter)
resolverSupplier := newLocalResolverSupplier(config.ResolverPoolSize, config.UseWasmInterpreter, nil)
resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore)
providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval, false, config.DisableExposureCollection)
provider := NewLocalResolverProvider(resolverSupplierWithMaterialization, config.StateProvider, config.FlagLogger, config.ClientSecret, logger, providerOpts...)

return provider, nil
}

func newLocalResolverSupplier(poolSize int, useWasmInterpreter bool) func(context.Context, lr.LogSink) lr.LocalResolver {
func newLocalResolverSupplier(poolSize int, useWasmInterpreter bool, initLabels map[string]string) func(context.Context, lr.LogSink) lr.LocalResolver {
cfg := lr.LocalResolverConfig{
PoolSize: poolSize,
UseWasmInterpreter: useWasmInterpreter,
}
return func(ctx context.Context, logSink lr.LogSink) lr.LocalResolver {
return lr.NewLocalResolver(ctx, logSink, cfg)
return newProviderTelemetryResolver(
logSink,
initLabels,
func(providerLogSink lr.LogSink) lr.LocalResolver {
return lr.NewLocalResolver(ctx, providerLogSink, cfg)
},
)
}
}

Expand Down
104 changes: 104 additions & 0 deletions openfeature-provider/go/confidence/provider_telemetry_resolver.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
package confidence

import (
"context"
"sync"

lr "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/local_resolver"
resolvertypes "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolver"
resolverv1 "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternal"
"github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/wasm"
)

// providerTelemetryResolver owns provider-scoped telemetry above pooling and recovery. The inner
// resolver stack is constructed with writeLogs, so every pooled WASM instance shares one init
// state. Close forces an init-only request when no earlier full flush produced one.
type providerTelemetryResolver struct {
delegate lr.LocalResolver
logSink lr.LogSink
labels map[string]string
sdk *resolvertypes.Sdk

mu sync.Mutex
initSent bool
}

func newProviderTelemetryResolver(
logSink lr.LogSink,
labels map[string]string,
innerFactory func(lr.LogSink) lr.LocalResolver,
) lr.LocalResolver {
r := &providerTelemetryResolver{
logSink: logSink,
labels: labels,
sdk: &resolvertypes.Sdk{
Sdk: &resolvertypes.Sdk_Id{Id: resolvertypes.SdkId_SDK_ID_GO_LOCAL_PROVIDER},
Version: Version,
},
}
r.delegate = innerFactory(r.writeLogs)
return r
}

func (r *providerTelemetryResolver) writeLogs(logs *resolverv1.WriteFlagLogsRequest) {
r.mu.Lock()
defer r.mu.Unlock()

if !r.initSent {
r.addInitTelemetry(logs)
}
r.logSink(logs)
r.initSent = true
}

func (r *providerTelemetryResolver) emitInitIfPending() {
r.mu.Lock()
defer r.mu.Unlock()

if r.initSent {
return
}
logs := &resolverv1.WriteFlagLogsRequest{}
r.addInitTelemetry(logs)
r.logSink(logs)
r.initSent = true
}

func (r *providerTelemetryResolver) addInitTelemetry(logs *resolverv1.WriteFlagLogsRequest) {
if logs.TelemetryData == nil {
logs.TelemetryData = &resolverv1.TelemetryData{}
}
logs.TelemetryData.Sdk = r.sdk
logs.TelemetryData.ProviderInitRate = append(
logs.TelemetryData.ProviderInitRate,
&resolverv1.TelemetryData_ProviderInitRate{Count: 1, Labels: r.labels},
)
}

func (r *providerTelemetryResolver) SetResolverState(request *wasm.SetResolverStateRequest) error {
return r.delegate.SetResolverState(request)
}

func (r *providerTelemetryResolver) ResolveProcess(request *wasm.ResolveProcessRequest) (*wasm.ResolveProcessResponse, error) {
return r.delegate.ResolveProcess(request)
}

func (r *providerTelemetryResolver) RegisterResolve(request *wasm.RegisterResolveRequest) {
r.delegate.RegisterResolve(request)
}

func (r *providerTelemetryResolver) ApplyFlags(request *resolvertypes.ApplyFlagsRequest) error {
return r.delegate.ApplyFlags(request)
}

func (r *providerTelemetryResolver) FlushAllLogs() error { return r.delegate.FlushAllLogs() }
func (r *providerTelemetryResolver) FlushAssignLogs() error { return r.delegate.FlushAssignLogs() }
func (r *providerTelemetryResolver) PrometheusSnapshot(bucketsPerDecade uint32, openmetrics bool) string {
return r.delegate.PrometheusSnapshot(bucketsPerDecade, openmetrics)
}

func (r *providerTelemetryResolver) Close(ctx context.Context) error {
err := r.delegate.Close(ctx)
r.emitInitIfPending()
return err
}
109 changes: 109 additions & 0 deletions openfeature-provider/go/confidence/provider_telemetry_resolver_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
package confidence

import (
"context"
"testing"

lr "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/local_resolver"
resolvertypes "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolver"
resolverv1 "github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternal"
"github.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/wasm"
)

func TestProviderTelemetryResolverEmitsOnceAcrossPooledFlushes(t *testing.T) {
var captured []*resolverv1.WriteFlagLogsRequest
resolver := newProviderTelemetryResolver(
func(logs *resolverv1.WriteFlagLogsRequest) { captured = append(captured, logs) },
map[string]string{"encryption": "true"},
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink, flushCount: 3} },
)

if err := resolver.FlushAllLogs(); err != nil {
t.Fatal(err)
}

if got := providerInitEventCount(captured); got != 1 {
t.Fatalf("expected one provider init event across pooled flushes, got %d", got)
}
telemetry := captured[0].GetTelemetryData()
if telemetry.GetSdk().GetId() != resolvertypes.SdkId_SDK_ID_GO_LOCAL_PROVIDER {
t.Fatalf("unexpected SDK ID: %v", telemetry.GetSdk().GetId())
}
if telemetry.GetSdk().GetVersion() != Version {
t.Fatalf("unexpected SDK version: %q", telemetry.GetSdk().GetVersion())
}
}

func TestProviderTelemetryResolverCloseEmitsWithoutResolve(t *testing.T) {
var captured []*resolverv1.WriteFlagLogsRequest
resolver := newProviderTelemetryResolver(
func(logs *resolverv1.WriteFlagLogsRequest) { captured = append(captured, logs) },
map[string]string{"encryption": "true"},
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink} },
)

if err := resolver.Close(context.Background()); err != nil {
t.Fatal(err)
}

if got := providerInitEventCount(captured); got != 1 {
t.Fatalf("expected shutdown to emit one provider init event, got %d", got)
}
}

func TestProviderTelemetryResolverRetriesAfterSinkFailure(t *testing.T) {
attempts := 0
var captured []*resolverv1.WriteFlagLogsRequest
resolver := newProviderTelemetryResolver(
func(logs *resolverv1.WriteFlagLogsRequest) {
attempts++
if attempts == 1 {
panic("send failed")
}
captured = append(captured, logs)
},
nil,
func(sink lr.LogSink) lr.LocalResolver { return &telemetryTestResolver{sink: sink, flushCount: 1} },
)

func() {
defer func() { _ = recover() }()
_ = resolver.FlushAllLogs()
}()
if err := resolver.FlushAllLogs(); err != nil {
t.Fatal(err)
}

if got := providerInitEventCount(captured); got != 1 {
t.Fatalf("expected init telemetry to be retried, got %d events", got)
}
}

func providerInitEventCount(requests []*resolverv1.WriteFlagLogsRequest) int {
count := 0
for _, request := range requests {
count += len(request.GetTelemetryData().GetProviderInitRate())
}
return count
}

type telemetryTestResolver struct {
sink lr.LogSink
flushCount int
}

func (r *telemetryTestResolver) SetResolverState(*wasm.SetResolverStateRequest) error { return nil }
func (r *telemetryTestResolver) ResolveProcess(*wasm.ResolveProcessRequest) (*wasm.ResolveProcessResponse, error) {
return &wasm.ResolveProcessResponse{}, nil
}
func (r *telemetryTestResolver) RegisterResolve(*wasm.RegisterResolveRequest) {}
func (r *telemetryTestResolver) ApplyFlags(*resolvertypes.ApplyFlagsRequest) error { return nil }
func (r *telemetryTestResolver) FlushAllLogs() error {
for range r.flushCount {
r.sink(&resolverv1.WriteFlagLogsRequest{})
}
return nil
}
func (r *telemetryTestResolver) FlushAssignLogs() error { return nil }
func (r *telemetryTestResolver) PrometheusSnapshot(uint32, bool) string { return "" }
func (r *telemetryTestResolver) Close(context.Context) error { return nil }
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,8 @@ private InternalFlagLoggerServiceGrpc.InternalFlagLoggerServiceBlockingStub crea
public void write(WriteFlagLogsRequest request) {
if (request.getClientResolveInfoList().isEmpty()
&& request.getFlagAssignedList().isEmpty()
&& request.getFlagResolveInfoList().isEmpty()) {
&& request.getFlagResolveInfoList().isEmpty()
&& !request.hasTelemetryData()) {
logger.debug("Skipping empty flag log request");
return;
}
Expand Down
Loading