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 @@ -69,10 +69,9 @@ func (r *providerTelemetryResolver) addInitTelemetry(logs *resolverv1.WriteFlagL
logs.TelemetryData = &resolverv1.TelemetryData{}
}
logs.TelemetryData.Sdk = r.sdk
logs.TelemetryData.ProviderInitRate = append(
logs.TelemetryData.ProviderInitRate,
&resolverv1.TelemetryData_ProviderInitRate{Count: 1, Labels: r.labels},
)
logs.TelemetryData.ProviderInitRate = []*resolverv1.TelemetryData_ProviderInitRate{
{Count: 1, Labels: r.labels},
}
}

func (r *providerTelemetryResolver) SetResolverState(request *wasm.SetResolverStateRequest) error {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,30 @@ func TestProviderTelemetryResolverCloseEmitsWithoutResolve(t *testing.T) {
}
}

func TestProviderTelemetryResolverOwnsSingleInitSample(t *testing.T) {
resolver := &providerTelemetryResolver{
labels: map[string]string{"encryption": "true"},
sdk: &resolvertypes.Sdk{Version: "test-version"},
}
logs := &resolverv1.WriteFlagLogsRequest{
TelemetryData: &resolverv1.TelemetryData{
ProviderInitRate: []*resolverv1.TelemetryData_ProviderInitRate{
{Count: 1, Labels: map[string]string{"existing": "true"}},
},
},
}

resolver.addInitTelemetry(logs)

got := logs.GetTelemetryData().GetProviderInitRate()
if len(got) != 1 {
t.Fatalf("expected exactly one provider init sample, got %d", len(got))
}
if got[0].GetLabels()["encryption"] != "true" {
t.Fatalf("expected provider-owned labels, got %v", got[0].GetLabels())
}
}

func TestProviderTelemetryResolverRetriesAfterSinkFailure(t *testing.T) {
attempts := 0
var captured []*resolverv1.WriteFlagLogsRequest
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ private WriteFlagLogsRequest addInitTelemetry(WriteFlagLogsRequest request) {
.setTelemetryData(
request.getTelemetryData().toBuilder()
.setSdk(sdk)
.clearProviderInitRate()
.addProviderInitRate(
TelemetryData.ProviderInitRate.newBuilder()
.setCount(1)
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.ResolveProcessResponse;
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.TelemetryData;
import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest;
import java.util.ArrayList;
import java.util.Map;
Expand Down Expand Up @@ -52,6 +53,32 @@ void closeEmitsWithoutResolve() {
assertThat(providerInitEventCount(captured)).isEqualTo(1);
}

@Test
void ownsSingleInitSample() {
final var captured = new ArrayList<WriteFlagLogsRequest>();
final var existingRequest =
WriteFlagLogsRequest.newBuilder()
.setTelemetryData(
TelemetryData.newBuilder()
.addProviderInitRate(
TelemetryData.ProviderInitRate.newBuilder()
.setCount(1)
.putLabels("existing", "true")))
.build();
final var resolver =
new ProviderTelemetryResolver(
captured::add,
SDK,
Map.of("encryption", "true"),
sink -> new TelemetryTestResolver(sink, 1, existingRequest));

resolver.flushAllLogs();

assertThat(captured.get(0).getTelemetryData().getProviderInitRateList())
.singleElement()
.satisfies(sample -> assertThat(sample.getLabelsMap()).containsOnlyKeys("encryption"));
}

@Test
void retriesAfterSinkFailure() {
final var attempts = new AtomicInteger();
Expand Down Expand Up @@ -83,10 +110,17 @@ private static int providerInitEventCount(ArrayList<WriteFlagLogsRequest> reques
private static final class TelemetryTestResolver implements LocalResolver {
private final Consumer<WriteFlagLogsRequest> sink;
private final int flushCount;
private final WriteFlagLogsRequest request;

private TelemetryTestResolver(Consumer<WriteFlagLogsRequest> sink, int flushCount) {
this(sink, flushCount, WriteFlagLogsRequest.getDefaultInstance());
}

private TelemetryTestResolver(
Consumer<WriteFlagLogsRequest> sink, int flushCount, WriteFlagLogsRequest request) {
this.sink = sink;
this.flushCount = flushCount;
this.request = request;
}

@Override
Expand All @@ -106,7 +140,7 @@ public void registerResolve(RegisterResolveRequest request) {}
@Override
public void flushAllLogs() {
for (int i = 0; i < flushCount; i++) {
sink.accept(WriteFlagLogsRequest.getDefaultInstance());
sink.accept(request);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -170,7 +170,7 @@ describe('flush behavior', () => {
telemetryData: {
resolverVersion: '0.20.0',
sdk: { id: 25, version: 'resolver-version' },
providerInitRate: [{ count: 2, labels: { existing: 'true' } }],
providerInitRate: [{ count: 1, labels: { existing: 'true' } }],
},
}),
).finish(),
Expand All @@ -182,10 +182,7 @@ describe('flush behavior', () => {
const decoded = WriteFlagLogsRequest.decode(sentBody!);
expect(decoded.telemetryData?.resolverVersion).toBe('0.20.0');
expect(decoded.telemetryData?.sdk).toEqual({ id: 22, customId: undefined, version: VERSION });
expect(decoded.telemetryData?.providerInitRate).toEqual([
{ count: 2, labels: { existing: 'true' } },
{ count: 1, labels: { encryption: 'false' } },
]);
expect(decoded.telemetryData?.providerInitRate).toEqual([{ count: 1, labels: { encryption: 'false' } }]);
});

it('does not retry provider init telemetry after an HTTP failure', async () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -458,7 +458,7 @@ export class ConfidenceServerProviderLocal implements Provider {
id: SdkId.SDK_ID_JS_LOCAL_SERVER_PROVIDER,
version: VERSION,
};
request.telemetryData.providerInitRate.push({ count: 1, labels: this.initLabels });
request.telemetryData.providerInitRate = [{ count: 1, labels: this.initLabels }];
return WriteFlagLogsRequest.encode(request).finish();
}

Expand Down
1 change: 1 addition & 0 deletions openfeature-provider/python/src/confidence/provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -795,6 +795,7 @@ def _write_logs(self, log_data: bytes) -> None:
version=__version__,
)
)
request.telemetry_data.ClearField("provider_init_rate")
init_rate = request.telemetry_data.provider_init_rate.add()
init_rate.count = 1
for k, v in self._init_labels.items():
Expand Down
6 changes: 6 additions & 0 deletions openfeature-provider/python/tests/test_provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,9 @@ def test_init_telemetry_includes_sdk(
)
request = internal_api_pb2.WriteFlagLogsRequest()
request.telemetry_data.SetInParent()
existing = request.telemetry_data.provider_init_rate.add()
existing.count = 1
existing.labels["existing"] = "true"

provider._write_logs(request.SerializeToString())
decoded = internal_api_pb2.WriteFlagLogsRequest.FromString(
Expand All @@ -59,6 +62,9 @@ def test_init_telemetry_includes_sdk(
assert decoded.telemetry_data.sdk.id == types_pb2.SdkId.SDK_ID_PYTHON_PROVIDER
assert decoded.telemetry_data.sdk.version == __version__
assert len(decoded.telemetry_data.provider_init_rate) == 1
assert dict(decoded.telemetry_data.provider_init_rate[0].labels) == {
"encryption": "false"
}

def test_init_telemetry_retries_after_failed_write(
self,
Expand Down
31 changes: 26 additions & 5 deletions openfeature-provider/rust/src/logger.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use tokio::sync::RwLock;

use confidence_resolver::assign_logger::AssignLogger;
use confidence_resolver::proto::confidence::flags::resolver::v1::{
telemetry_data::ProviderInitRate, Sdk, WriteFlagLogsRequest,
telemetry_data::ProviderInitRate, Sdk, TelemetryData, WriteFlagLogsRequest,
};
use confidence_resolver::resolve_logger::ResolveLogger;

Expand Down Expand Up @@ -224,10 +224,7 @@ impl LogManager {
td.sdk = Some(self.sdk.clone());
let include_init = self.init_state.claim();
if include_init {
td.provider_init_rate.push(ProviderInitRate {
count: 1,
labels: self.init_labels.clone(),
});
set_provider_init_telemetry(&mut td, &self.init_labels);
}
request.telemetry_data = Some(td);

Expand Down Expand Up @@ -260,6 +257,13 @@ impl LogManager {
}
}

fn set_provider_init_telemetry(telemetry: &mut TelemetryData, labels: &BTreeMap<String, String>) {
telemetry.provider_init_rate = vec![ProviderInitRate {
count: 1,
labels: labels.clone(),
}];
}

/// Check if a WriteFlagLogsRequest has any logs to send.
fn has_logs(request: &WriteFlagLogsRequest) -> bool {
!request.flag_assigned.is_empty()
Expand All @@ -272,6 +276,23 @@ fn has_logs(request: &WriteFlagLogsRequest) -> bool {
mod tests {
use super::*;

#[test]
fn provider_owns_single_init_sample() {
let mut telemetry = TelemetryData {
provider_init_rate: vec![ProviderInitRate {
count: 1,
labels: BTreeMap::from([("existing".to_string(), "true".to_string())]),
}],
..Default::default()
};
let labels = BTreeMap::from([("encryption".to_string(), "true".to_string())]);

set_provider_init_telemetry(&mut telemetry, &labels);

assert_eq!(telemetry.provider_init_rate.len(), 1);
assert_eq!(telemetry.provider_init_rate[0].labels, labels);
}

#[test]
fn test_encode_ingest_request_roundtrip() {
let account_id = "test-account-123";
Expand Down