Skip to content
Open
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
Binary file not shown.
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,16 @@ type localResolverImpl struct {
}

func NewLocalResolverWithPoolSize(ctx context.Context, logSink LogSink, poolSize int) LocalResolver {
factory := NewWasmResolverFactory(logSink)
return NewLocalResolverWithLabels(ctx, logSink, poolSize, nil)
}

func NewLocalResolverWithLabels(ctx context.Context, logSink LogSink, poolSize int, initLabels map[string]string) LocalResolver {
var factory LocalResolverFactory
if len(initLabels) > 0 {
factory = NewWasmResolverFactoryWithLabels(logSink, initLabels)
} else {
factory = NewWasmResolverFactory(logSink)
}
factory = NewRecoveringResolverFactory(factory)
if poolSize <= 0 {
poolSize = DefaultPoolSize
Expand Down
27 changes: 24 additions & 3 deletions openfeature-provider/go/confidence/internal/local_resolver/wasm.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,8 @@ type WasmResolver struct {
mu *sync.Mutex
instanceID string
fnCache sync.Map
initLabels map[string]string
firstFlush bool
}

var _ LocalResolver = (*WasmResolver)(nil)
Expand Down Expand Up @@ -85,6 +87,16 @@ func (r *WasmResolver) ApplyFlags(request *resolver.ApplyFlagsRequest) error {
func (r *WasmResolver) FlushAllLogs() error {
resp := &resolverv1.WriteFlagLogsRequest{}
err := r.call("wasm_msg_guest_bounded_flush_logs", nil, resp)
if err == nil && r.firstFlush {
r.firstFlush = false
if resp.TelemetryData == nil {
resp.TelemetryData = &resolverv1.TelemetryData{}
}
resp.TelemetryData.ProviderInitRate = append(
resp.TelemetryData.ProviderInitRate,
&resolverv1.TelemetryData_ProviderInitRate{Count: 1, Labels: r.initLabels},
)
}
if err == nil && proto.Size(resp) > 0 {
r.logSink(resp)
}
Expand Down Expand Up @@ -154,9 +166,10 @@ func (r *WasmResolver) call(fnName string, request proto.Message, response proto
}

type WasmResolverFactory struct {
runtime wazero.Runtime
module wazero.CompiledModule
logSink LogSink
runtime wazero.Runtime
module wazero.CompiledModule
logSink LogSink
initLabels map[string]string
}

var _ LocalResolverFactory = (*WasmResolverFactory)(nil)
Expand Down Expand Up @@ -208,6 +221,12 @@ func NewWasmResolverFactory(logSink LogSink) LocalResolverFactory {
}
}

func NewWasmResolverFactoryWithLabels(logSink LogSink, initLabels map[string]string) LocalResolverFactory {
factory := NewWasmResolverFactory(logSink).(*WasmResolverFactory)
factory.initLabels = initLabels
return factory
}

func (wrf *WasmResolverFactory) New() LocalResolver {
ctx := context.Background()
config := wazero.NewModuleConfig().WithName("")
Expand All @@ -221,6 +240,8 @@ func (wrf *WasmResolverFactory) New() LocalResolver {
logSink: wrf.logSink,
mu: &sync.Mutex{},
instanceID: fmt.Sprintf("%d", id),
initLabels: wrf.initLabels,
firstFlush: true,
}
}

Expand Down

Large diffs are not rendered by default.

6 changes: 5 additions & 1 deletion 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 @@ -87,8 +88,11 @@ func NewProvider(ctx context.Context, config ProviderConfig) (*LocalResolverProv
materializationStore = newRemoteMaterializationStore(resolverv1.NewInternalFlagLoggerServiceClient(conn), config.ClientSecret)
}

initLabels := map[string]string{
"encryption": strconv.FormatBool(config.EncryptionKey != ""),
}
resolverSupplier := func(ctx context.Context, logSink lr.LogSink) lr.LocalResolver {
return lr.NewLocalResolverWithPoolSize(ctx, logSink, config.ResolverPoolSize)
return lr.NewLocalResolverWithLabels(ctx, logSink, config.ResolverPoolSize, initLabels)
}
resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore)
providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import io.grpc.StatusRuntimeException;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicReference;
Expand Down Expand Up @@ -158,11 +159,15 @@ public OpenFeatureLocalResolveProvider(
clientSecret, config.getHttpClientFactory(), config.getEncryptionKey());
final var wasmFlagLogger = new GrpcWasmFlagLogger(clientSecret, config.getChannelFactory());
this.flagLogger = wasmFlagLogger;
final Map<String, String> initLabels =
Map.of("encryption", String.valueOf(config.getEncryptionKey() != null));
final int numInstances = PooledResolver.getNumInstances(config.getResolverPoolSize());
final LocalResolver inner =
new PooledResolver(
numInstances,
() -> new RecoveringResolver(() -> new WasmLocalResolver(flagLogger::write)));
() ->
new RecoveringResolver(
() -> new WasmLocalResolver(flagLogger::write, initLabels)));
this.resolver = new MaterializingResolver(inner, materializationStore);
}

Expand All @@ -189,7 +194,9 @@ public OpenFeatureLocalResolveProvider(
final LocalResolver inner =
new PooledResolver(
numInstances,
() -> new RecoveringResolver(() -> new WasmLocalResolver(wasmFlagLogger::write)));
() ->
new RecoveringResolver(
() -> new WasmLocalResolver(wasmFlagLogger::write, Map.of())));
this.resolver = new MaterializingResolver(inner, materializationStore);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,12 @@
import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessRequest;
import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessResponse;
import com.spotify.confidence.sdk.flags.resolver.v1.Sdk;
import com.spotify.confidence.sdk.flags.resolver.v1.TelemetryData;
import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest;
import com.spotify.confidence.sdk.wasm.Messages;
import java.time.Instant;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.atomic.AtomicInteger;
Expand Down Expand Up @@ -49,6 +51,8 @@ class WasmLocalResolver implements LocalResolver {
private final ExportFunction wasmMsgAlloc;
private final ExportFunction wasmMsgFree;
private final Consumer<WriteFlagLogsRequest> logSink;
private final Map<String, String> initLabels;
private boolean firstFlush = true;

// api
private final ExportFunction wasmMsgGuestSetResolverState;
Expand All @@ -60,8 +64,9 @@ class WasmLocalResolver implements LocalResolver {
private final ExportFunction wasmMsgGuestPrometheusSnapshot;
private final ReentrantLock lock = new ReentrantLock();

public WasmLocalResolver(Consumer<WriteFlagLogsRequest> logSink) {
public WasmLocalResolver(Consumer<WriteFlagLogsRequest> logSink, Map<String, String> initLabels) {
this.logSink = logSink;
this.initLabels = initLabels;
this.instanceId = String.valueOf(INSTANCE_COUNTER.getAndIncrement());
instance =
Instance.builder(ConfidenceResolverModule.load())
Expand Down Expand Up @@ -192,7 +197,21 @@ public void flushAllLogs() {
final var voidRequest = Messages.Void.getDefaultInstance();
final var reqPtr = transferRequest(voidRequest);
final var respPtr = (int) wasmMsgBoundedFlushLogs.apply(reqPtr)[0];
final var request = consumeResponse(respPtr, WriteFlagLogsRequest::parseFrom);
var request = consumeResponse(respPtr, WriteFlagLogsRequest::parseFrom);
if (firstFlush) {
firstFlush = false;
request =
request.toBuilder()
.setTelemetryData(
request.getTelemetryData().toBuilder()
.addProviderInitRate(
TelemetryData.ProviderInitRate.newBuilder()
.setCount(1)
.putAllLabels(initLabels)
.build())
.build())
.build();
}
if (!isEmptyLogRequest(request)) {
logSink.accept(request);
}
Expand Down Expand Up @@ -279,9 +298,8 @@ private <T extends Message> T consumeResponse(int addr, ParserFn<T> codec) {
final Messages.Response response = Messages.Response.parseFrom(consume(addr));
if (response.hasError()) {
throw new RuntimeException(response.getError());
} else {
return codec.apply(response.getData().toByteArray());
}
return codec.apply(response.getData().toByteArray());
} catch (InvalidProtocolBufferException e) {
throw new RuntimeException(e);
}
Expand Down Expand Up @@ -350,4 +368,5 @@ private interface ParserFn<T> {

T apply(byte[] data) throws InvalidProtocolBufferException;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -154,7 +155,7 @@ private static Timestamp toTs(Instant t) {
@Test
void applyWithoutSendTime_isSwallowed_andRecordsNoFlagAssigned() {
final List<WriteFlagLogsRequest> captured = new ArrayList<>();
final var resolver = new WasmLocalResolver(captured::add);
final var resolver = new WasmLocalResolver(captured::add, Map.of());
resolver.setResolverState(buildState(), ACCOUNT, null);

final var resolveResp = resolveWithApplyFalse(resolver);
Expand Down Expand Up @@ -185,7 +186,7 @@ void applyWithoutSendTime_isSwallowed_andRecordsNoFlagAssigned() {
@Test
void applyWithSendTime_recordsOneFlagAssigned() {
final List<WriteFlagLogsRequest> captured = new ArrayList<>();
final var resolver = new WasmLocalResolver(captured::add);
final var resolver = new WasmLocalResolver(captured::add, Map.of());
resolver.setResolverState(buildState(), ACCOUNT, null);

final var resolveResp = resolveWithApplyFalse(resolver);
Expand Down Expand Up @@ -233,7 +234,7 @@ void applyWithSendTime_recordsOneFlagAssigned() {
@Test
void applyAfterStateRotated_secretRemoved_isSwallowedAndLogged() {
final List<WriteFlagLogsRequest> captured = new ArrayList<>();
final var resolver = new WasmLocalResolver(captured::add);
final var resolver = new WasmLocalResolver(captured::add, Map.of());
resolver.setResolverState(buildState(SECRET), ACCOUNT, null);

final var resolveResp = resolveWithApplyFalse(resolver);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ class ResolveTest {
private final LocalResolver resolver;

public ResolveTest() {
resolver = new WasmLocalResolver(request -> {});
resolver = new WasmLocalResolver(request -> {}, Map.of());
}

@BeforeEach
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessRequest;
import java.lang.reflect.Field;
import java.util.List;
import java.util.Map;
import org.junit.jupiter.api.Test;

/**
Expand All @@ -32,7 +33,7 @@ private static int getWasmMemoryPages(WasmLocalResolver resolver) {

@Test
void wasmMemoryStableOnRepeatedResolveCalls() {
WasmLocalResolver resolver = new WasmLocalResolver(request -> {});
WasmLocalResolver resolver = new WasmLocalResolver(request -> {}, Map.of());
resolver.setResolverState(ResolveTest.exampleStateBytes, "account", null);

ResolveProcessRequest request =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import com.spotify.confidence.sdk.flags.resolver.v1.Sdk;
import com.spotify.confidence.sdk.flags.resolver.v1.SdkId;
import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -44,7 +45,7 @@ void concurrentFlushAndCloseShouldNotLoseAssignments() throws Exception {

for (int i = 0; i < iterations; i++) {
final var logger = new CapturingWasmFlagLogger();
final var resolver = new WasmLocalResolver(logger::write);
final var resolver = new WasmLocalResolver(logger::write, Map.of());
resolver.setResolverState(resolverState, accountId, null);

// Resolve a flag to create a flag assignment in the WASM buffer
Expand Down
Loading
Loading