diff --git a/confidence-resolver/protos/confidence/flags/resolver/v1/internal_api.proto b/confidence-resolver/protos/confidence/flags/resolver/v1/internal_api.proto index 9ab42a0dd..9a9c781f6 100644 --- a/confidence-resolver/protos/confidence/flags/resolver/v1/internal_api.proto +++ b/confidence-resolver/protos/confidence/flags/resolver/v1/internal_api.proto @@ -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 @@ -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 labels = 3; + } + message ResolveLatency { // Delta sum of observed values since the last flush. uint32 sum = 1; diff --git a/confidence-resolver/src/telemetry.rs b/confidence-resolver/src/telemetry.rs index e075953c9..bfd18de55 100644 --- a/confidence-resolver/src/telemetry.rs +++ b/confidence-resolver/src/telemetry.rs @@ -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(), } } } diff --git a/openfeature-provider/go/confidence/internal/flag_logger/grpc.go b/openfeature-provider/go/confidence/internal/flag_logger/grpc.go index 0edabcdfe..b2949a148 100644 --- a/openfeature-provider/go/confidence/internal/flag_logger/grpc.go +++ b/openfeature-provider/go/confidence/internal/flag_logger/grpc.go @@ -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 } diff --git a/openfeature-provider/go/confidence/internal/flag_logger/multi.go b/openfeature-provider/go/confidence/internal/flag_logger/multi.go index 09152bf75..27afb18ca 100644 --- a/openfeature-provider/go/confidence/internal/flag_logger/multi.go +++ b/openfeature-provider/go/confidence/internal/flag_logger/multi.go @@ -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 } diff --git a/openfeature-provider/go/confidence/internal/flag_logger/multi_test.go b/openfeature-provider/go/confidence/internal/flag_logger/multi_test.go index 789690b0e..6ac3c871f 100644 --- a/openfeature-provider/go/confidence/internal/flag_logger/multi_test.go +++ b/openfeature-provider/go/confidence/internal/flag_logger/multi_test.go @@ -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})) diff --git a/openfeature-provider/go/confidence/internal/local_resolver/assets/confidence_resolver.wasm b/openfeature-provider/go/confidence/internal/local_resolver/assets/confidence_resolver.wasm index 7e8ae7076..ae9ebde18 100755 Binary files a/openfeature-provider/go/confidence/internal/local_resolver/assets/confidence_resolver.wasm and b/openfeature-provider/go/confidence/internal/local_resolver/assets/confidence_resolver.wasm differ diff --git a/openfeature-provider/go/confidence/internal/proto/resolverinternal/internal_api.pb.go b/openfeature-provider/go/confidence/internal/proto/resolverinternal/internal_api.pb.go index e2564c792..0050357cb 100644 --- a/openfeature-provider/go/confidence/internal/proto/resolverinternal/internal_api.pb.go +++ b/openfeature-provider/go/confidence/internal/proto/resolverinternal/internal_api.pb.go @@ -243,9 +243,13 @@ func (x *IngestFlagLogsRequest) GetBatch() *WriteFlagLogsRequest { type TelemetryData struct { state protoimpl.MessageState `protogen:"open.v1"` // Information about the SDK/provider - Sdk *resolver.Sdk `protobuf:"bytes,2,opt,name=sdk,proto3" json:"sdk,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + Sdk *resolver.Sdk `protobuf:"bytes,2,opt,name=sdk,proto3" json:"sdk,omitempty"` + // Version of the embedded resolver (e.g. "0.14.0"). + // This must be preserved when providers decode and re-encode telemetry. + ResolverVersion string `protobuf:"bytes,8,opt,name=resolver_version,json=resolverVersion,proto3" json:"resolver_version,omitempty"` + ProviderInitRate []*TelemetryData_ProviderInitRate `protobuf:"bytes,9,rep,name=provider_init_rate,json=providerInitRate,proto3" json:"provider_init_rate,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *TelemetryData) Reset() { @@ -285,6 +289,20 @@ func (x *TelemetryData) GetSdk() *resolver.Sdk { return nil } +func (x *TelemetryData) GetResolverVersion() string { + if x != nil { + return x.ResolverVersion + } + return "" +} + +func (x *TelemetryData) GetProviderInitRate() []*TelemetryData_ProviderInitRate { + if x != nil { + return x.ProviderInitRate + } + return nil +} + type ClientInfo struct { state protoimpl.MessageState `protogen:"open.v1"` Client string `protobuf:"bytes,1,opt,name=client,proto3" json:"client,omitempty"` @@ -1163,6 +1181,58 @@ func (x *ReadOperationsResult) GetResults() []*ReadResult { return nil } +type TelemetryData_ProviderInitRate struct { + state protoimpl.MessageState `protogen:"open.v1"` + Count uint32 `protobuf:"varint,1,opt,name=count,proto3" json:"count,omitempty"` + Labels map[string]string `protobuf:"bytes,3,rep,name=labels,proto3" json:"labels,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *TelemetryData_ProviderInitRate) Reset() { + *x = TelemetryData_ProviderInitRate{} + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[19] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *TelemetryData_ProviderInitRate) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*TelemetryData_ProviderInitRate) ProtoMessage() {} + +func (x *TelemetryData_ProviderInitRate) ProtoReflect() protoreflect.Message { + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[19] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use TelemetryData_ProviderInitRate.ProtoReflect.Descriptor instead. +func (*TelemetryData_ProviderInitRate) Descriptor() ([]byte, []int) { + return file_confidence_flags_resolver_v1_internal_api_proto_rawDescGZIP(), []int{3, 0} +} + +func (x *TelemetryData_ProviderInitRate) GetCount() uint32 { + if x != nil { + return x.Count + } + return 0 +} + +func (x *TelemetryData_ProviderInitRate) GetLabels() map[string]string { + if x != nil { + return x.Labels + } + return nil +} + type FlagAssigned_AppliedFlag struct { state protoimpl.MessageState `protogen:"open.v1"` Flag string `protobuf:"bytes,1,opt,name=flag,proto3" json:"flag,omitempty"` @@ -1183,7 +1253,7 @@ type FlagAssigned_AppliedFlag struct { func (x *FlagAssigned_AppliedFlag) Reset() { *x = FlagAssigned_AppliedFlag{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[19] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[21] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1195,7 +1265,7 @@ func (x *FlagAssigned_AppliedFlag) String() string { func (*FlagAssigned_AppliedFlag) ProtoMessage() {} func (x *FlagAssigned_AppliedFlag) ProtoReflect() protoreflect.Message { - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[19] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[21] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1311,7 +1381,7 @@ type FlagAssigned_AssignmentInfo struct { func (x *FlagAssigned_AssignmentInfo) Reset() { *x = FlagAssigned_AssignmentInfo{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[20] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[22] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1323,7 +1393,7 @@ func (x *FlagAssigned_AssignmentInfo) String() string { func (*FlagAssigned_AssignmentInfo) ProtoMessage() {} func (x *FlagAssigned_AssignmentInfo) ProtoReflect() protoreflect.Message { - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[20] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[22] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1362,7 +1432,7 @@ type FlagAssigned_DefaultAssignment struct { func (x *FlagAssigned_DefaultAssignment) Reset() { *x = FlagAssigned_DefaultAssignment{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[21] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[23] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1374,7 +1444,7 @@ func (x *FlagAssigned_DefaultAssignment) String() string { func (*FlagAssigned_DefaultAssignment) ProtoMessage() {} func (x *FlagAssigned_DefaultAssignment) ProtoReflect() protoreflect.Message { - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[21] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[23] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1408,7 +1478,7 @@ type ClientResolveInfo_EvaluationContextSchemaInstance struct { func (x *ClientResolveInfo_EvaluationContextSchemaInstance) Reset() { *x = ClientResolveInfo_EvaluationContextSchemaInstance{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[22] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[24] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1420,7 +1490,7 @@ func (x *ClientResolveInfo_EvaluationContextSchemaInstance) String() string { func (*ClientResolveInfo_EvaluationContextSchemaInstance) ProtoMessage() {} func (x *ClientResolveInfo_EvaluationContextSchemaInstance) ProtoReflect() protoreflect.Message { - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[22] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[24] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1456,7 +1526,7 @@ type FlagResolveInfo_VariantResolveInfo struct { func (x *FlagResolveInfo_VariantResolveInfo) Reset() { *x = FlagResolveInfo_VariantResolveInfo{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[24] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[26] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1468,7 +1538,7 @@ func (x *FlagResolveInfo_VariantResolveInfo) String() string { func (*FlagResolveInfo_VariantResolveInfo) ProtoMessage() {} func (x *FlagResolveInfo_VariantResolveInfo) ProtoReflect() protoreflect.Message { - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[24] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[26] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1512,9 +1582,17 @@ const file_confidence_flags_resolver_v1_internal_api_proto_rawDesc = "" + "\x15IngestFlagLogsRequest\x12\x1d\n" + "\n" + "account_id\x18\x01 \x01(\tR\taccountId\x12H\n" + - "\x05batch\x18\x02 \x01(\v22.confidence.flags.resolver.v1.WriteFlagLogsRequestR\x05batch\"D\n" + + "\x05batch\x18\x02 \x01(\v22.confidence.flags.resolver.v1.WriteFlagLogsRequestR\x05batch\"\xa9\x03\n" + "\rTelemetryData\x123\n" + - "\x03sdk\x18\x02 \x01(\v2!.confidence.flags.resolver.v1.SdkR\x03sdk\"\x86\x01\n" + + "\x03sdk\x18\x02 \x01(\v2!.confidence.flags.resolver.v1.SdkR\x03sdk\x12)\n" + + "\x10resolver_version\x18\b \x01(\tR\x0fresolverVersion\x12j\n" + + "\x12provider_init_rate\x18\t \x03(\v2<.confidence.flags.resolver.v1.TelemetryData.ProviderInitRateR\x10providerInitRate\x1a\xcb\x01\n" + + "\x10ProviderInitRate\x12\x14\n" + + "\x05count\x18\x01 \x01(\rR\x05count\x12`\n" + + "\x06labels\x18\x03 \x03(\v2H.confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntryR\x06labels\x1a9\n" + + "\vLabelsEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01J\x04\b\x02\x10\x03\"\x86\x01\n" + "\n" + "ClientInfo\x12\x16\n" + "\x06client\x18\x01 \x01(\tR\x06client\x12+\n" + @@ -1622,7 +1700,7 @@ func file_confidence_flags_resolver_v1_internal_api_proto_rawDescGZIP() []byte { } var file_confidence_flags_resolver_v1_internal_api_proto_enumTypes = make([]protoimpl.EnumInfo, 1) -var file_confidence_flags_resolver_v1_internal_api_proto_msgTypes = make([]protoimpl.MessageInfo, 25) +var file_confidence_flags_resolver_v1_internal_api_proto_msgTypes = make([]protoimpl.MessageInfo, 27) var file_confidence_flags_resolver_v1_internal_api_proto_goTypes = []any{ (FlagAssigned_DefaultAssignment_DefaultAssignmentReason)(0), // 0: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason (*WriteFlagLogsRequest)(nil), // 1: confidence.flags.resolver.v1.WriteFlagLogsRequest @@ -1644,14 +1722,16 @@ var file_confidence_flags_resolver_v1_internal_api_proto_goTypes = []any{ (*InclusionData)(nil), // 17: confidence.flags.resolver.v1.InclusionData (*ReadResult)(nil), // 18: confidence.flags.resolver.v1.ReadResult (*ReadOperationsResult)(nil), // 19: confidence.flags.resolver.v1.ReadOperationsResult - (*FlagAssigned_AppliedFlag)(nil), // 20: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag - (*FlagAssigned_AssignmentInfo)(nil), // 21: confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo - (*FlagAssigned_DefaultAssignment)(nil), // 22: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment - (*ClientResolveInfo_EvaluationContextSchemaInstance)(nil), // 23: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance - nil, // 24: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry - (*FlagResolveInfo_VariantResolveInfo)(nil), // 25: confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo - (*resolver.Sdk)(nil), // 26: confidence.flags.resolver.v1.Sdk - (*timestamppb.Timestamp)(nil), // 27: google.protobuf.Timestamp + (*TelemetryData_ProviderInitRate)(nil), // 20: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate + nil, // 21: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntry + (*FlagAssigned_AppliedFlag)(nil), // 22: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag + (*FlagAssigned_AssignmentInfo)(nil), // 23: confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo + (*FlagAssigned_DefaultAssignment)(nil), // 24: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment + (*ClientResolveInfo_EvaluationContextSchemaInstance)(nil), // 25: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance + nil, // 26: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry + (*FlagResolveInfo_VariantResolveInfo)(nil), // 27: confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo + (*resolver.Sdk)(nil), // 28: confidence.flags.resolver.v1.Sdk + (*timestamppb.Timestamp)(nil), // 29: google.protobuf.Timestamp } var file_confidence_flags_resolver_v1_internal_api_proto_depIdxs = []int32{ 6, // 0: confidence.flags.resolver.v1.WriteFlagLogsRequest.flag_assigned:type_name -> confidence.flags.resolver.v1.FlagAssigned @@ -1659,36 +1739,38 @@ var file_confidence_flags_resolver_v1_internal_api_proto_depIdxs = []int32{ 8, // 2: confidence.flags.resolver.v1.WriteFlagLogsRequest.client_resolve_info:type_name -> confidence.flags.resolver.v1.ClientResolveInfo 9, // 3: confidence.flags.resolver.v1.WriteFlagLogsRequest.flag_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo 1, // 4: confidence.flags.resolver.v1.IngestFlagLogsRequest.batch:type_name -> confidence.flags.resolver.v1.WriteFlagLogsRequest - 26, // 5: confidence.flags.resolver.v1.TelemetryData.sdk:type_name -> confidence.flags.resolver.v1.Sdk - 26, // 6: confidence.flags.resolver.v1.ClientInfo.sdk:type_name -> confidence.flags.resolver.v1.Sdk - 5, // 7: confidence.flags.resolver.v1.FlagAssigned.client_info:type_name -> confidence.flags.resolver.v1.ClientInfo - 20, // 8: confidence.flags.resolver.v1.FlagAssigned.flags:type_name -> confidence.flags.resolver.v1.FlagAssigned.AppliedFlag - 23, // 9: confidence.flags.resolver.v1.ClientResolveInfo.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance - 25, // 10: confidence.flags.resolver.v1.FlagResolveInfo.variant_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo - 16, // 11: confidence.flags.resolver.v1.WriteOperationsRequest.store_variant_op:type_name -> confidence.flags.resolver.v1.VariantData - 12, // 12: confidence.flags.resolver.v1.ReadOp.variant_read_op:type_name -> confidence.flags.resolver.v1.VariantReadOp - 13, // 13: confidence.flags.resolver.v1.ReadOp.inclusion_read_op:type_name -> confidence.flags.resolver.v1.InclusionReadOp - 14, // 14: confidence.flags.resolver.v1.ReadOperationsRequest.ops:type_name -> confidence.flags.resolver.v1.ReadOp - 16, // 15: confidence.flags.resolver.v1.ReadResult.variant_result:type_name -> confidence.flags.resolver.v1.VariantData - 17, // 16: confidence.flags.resolver.v1.ReadResult.inclusion_result:type_name -> confidence.flags.resolver.v1.InclusionData - 18, // 17: confidence.flags.resolver.v1.ReadOperationsResult.results:type_name -> confidence.flags.resolver.v1.ReadResult - 21, // 18: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.assignment_info:type_name -> confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo - 22, // 19: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.default_assignment:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment - 7, // 20: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.fallthrough_assignments:type_name -> confidence.flags.resolver.v1.FallthroughAssignment - 27, // 21: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.apply_time:type_name -> google.protobuf.Timestamp - 0, // 22: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.reason:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason - 24, // 23: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry - 1, // 24: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:input_type -> confidence.flags.resolver.v1.WriteFlagLogsRequest - 10, // 25: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:input_type -> confidence.flags.resolver.v1.WriteOperationsRequest - 15, // 26: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:input_type -> confidence.flags.resolver.v1.ReadOperationsRequest - 2, // 27: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:output_type -> confidence.flags.resolver.v1.WriteFlagLogsResponse - 11, // 28: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:output_type -> confidence.flags.resolver.v1.WriteOperationsResult - 19, // 29: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:output_type -> confidence.flags.resolver.v1.ReadOperationsResult - 27, // [27:30] is the sub-list for method output_type - 24, // [24:27] is the sub-list for method input_type - 24, // [24:24] is the sub-list for extension type_name - 24, // [24:24] is the sub-list for extension extendee - 0, // [0:24] is the sub-list for field type_name + 28, // 5: confidence.flags.resolver.v1.TelemetryData.sdk:type_name -> confidence.flags.resolver.v1.Sdk + 20, // 6: confidence.flags.resolver.v1.TelemetryData.provider_init_rate:type_name -> confidence.flags.resolver.v1.TelemetryData.ProviderInitRate + 28, // 7: confidence.flags.resolver.v1.ClientInfo.sdk:type_name -> confidence.flags.resolver.v1.Sdk + 5, // 8: confidence.flags.resolver.v1.FlagAssigned.client_info:type_name -> confidence.flags.resolver.v1.ClientInfo + 22, // 9: confidence.flags.resolver.v1.FlagAssigned.flags:type_name -> confidence.flags.resolver.v1.FlagAssigned.AppliedFlag + 25, // 10: confidence.flags.resolver.v1.ClientResolveInfo.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance + 27, // 11: confidence.flags.resolver.v1.FlagResolveInfo.variant_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo + 16, // 12: confidence.flags.resolver.v1.WriteOperationsRequest.store_variant_op:type_name -> confidence.flags.resolver.v1.VariantData + 12, // 13: confidence.flags.resolver.v1.ReadOp.variant_read_op:type_name -> confidence.flags.resolver.v1.VariantReadOp + 13, // 14: confidence.flags.resolver.v1.ReadOp.inclusion_read_op:type_name -> confidence.flags.resolver.v1.InclusionReadOp + 14, // 15: confidence.flags.resolver.v1.ReadOperationsRequest.ops:type_name -> confidence.flags.resolver.v1.ReadOp + 16, // 16: confidence.flags.resolver.v1.ReadResult.variant_result:type_name -> confidence.flags.resolver.v1.VariantData + 17, // 17: confidence.flags.resolver.v1.ReadResult.inclusion_result:type_name -> confidence.flags.resolver.v1.InclusionData + 18, // 18: confidence.flags.resolver.v1.ReadOperationsResult.results:type_name -> confidence.flags.resolver.v1.ReadResult + 21, // 19: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.labels:type_name -> confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntry + 23, // 20: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.assignment_info:type_name -> confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo + 24, // 21: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.default_assignment:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment + 7, // 22: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.fallthrough_assignments:type_name -> confidence.flags.resolver.v1.FallthroughAssignment + 29, // 23: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.apply_time:type_name -> google.protobuf.Timestamp + 0, // 24: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.reason:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason + 26, // 25: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry + 1, // 26: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:input_type -> confidence.flags.resolver.v1.WriteFlagLogsRequest + 10, // 27: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:input_type -> confidence.flags.resolver.v1.WriteOperationsRequest + 15, // 28: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:input_type -> confidence.flags.resolver.v1.ReadOperationsRequest + 2, // 29: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:output_type -> confidence.flags.resolver.v1.WriteFlagLogsResponse + 11, // 30: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:output_type -> confidence.flags.resolver.v1.WriteOperationsResult + 19, // 31: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:output_type -> confidence.flags.resolver.v1.ReadOperationsResult + 29, // [29:32] is the sub-list for method output_type + 26, // [26:29] is the sub-list for method input_type + 26, // [26:26] is the sub-list for extension type_name + 26, // [26:26] is the sub-list for extension extendee + 0, // [0:26] is the sub-list for field type_name } func init() { file_confidence_flags_resolver_v1_internal_api_proto_init() } @@ -1704,7 +1786,7 @@ func file_confidence_flags_resolver_v1_internal_api_proto_init() { (*ReadResult_VariantResult)(nil), (*ReadResult_InclusionResult)(nil), } - file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[19].OneofWrappers = []any{ + file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[21].OneofWrappers = []any{ (*FlagAssigned_AppliedFlag_AssignmentInfo)(nil), (*FlagAssigned_AppliedFlag_DefaultAssignment)(nil), } @@ -1714,7 +1796,7 @@ func file_confidence_flags_resolver_v1_internal_api_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_confidence_flags_resolver_v1_internal_api_proto_rawDesc), len(file_confidence_flags_resolver_v1_internal_api_proto_rawDesc)), NumEnums: 1, - NumMessages: 25, + NumMessages: 27, NumExtensions: 0, NumServices: 1, }, diff --git a/openfeature-provider/go/confidence/provider_builder.go b/openfeature-provider/go/confidence/provider_builder.go index 519c0c07b..29d631cc8 100644 --- a/openfeature-provider/go/confidence/provider_builder.go +++ b/openfeature-provider/go/confidence/provider_builder.go @@ -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" @@ -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...) @@ -131,7 +135,7 @@ 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...) @@ -139,13 +143,19 @@ func NewProviderForTest(ctx context.Context, config ProviderTestConfig) (*LocalR 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) + }, + ) } } diff --git a/openfeature-provider/go/confidence/provider_telemetry_resolver.go b/openfeature-provider/go/confidence/provider_telemetry_resolver.go new file mode 100644 index 000000000..f7782653c --- /dev/null +++ b/openfeature-provider/go/confidence/provider_telemetry_resolver.go @@ -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 +} diff --git a/openfeature-provider/go/confidence/provider_telemetry_resolver_test.go b/openfeature-provider/go/confidence/provider_telemetry_resolver_test.go new file mode 100644 index 000000000..82fe5d1f0 --- /dev/null +++ b/openfeature-provider/go/confidence/provider_telemetry_resolver_test.go @@ -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 } diff --git a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/GrpcWasmFlagLogger.java b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/GrpcWasmFlagLogger.java index 8c0618361..ad7f7e342 100644 --- a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/GrpcWasmFlagLogger.java +++ b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/GrpcWasmFlagLogger.java @@ -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; } diff --git a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/OpenFeatureLocalResolveProvider.java b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/OpenFeatureLocalResolveProvider.java index 450852b78..e0c75587b 100644 --- a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/OpenFeatureLocalResolveProvider.java +++ b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/OpenFeatureLocalResolveProvider.java @@ -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; @@ -162,18 +163,25 @@ public OpenFeatureLocalResolveProvider( new GrpcWasmFlagLogger( clientSecret, config.getChannelFactory(), config.getHttpClientFactory()); this.flagLogger = wasmFlagLogger; + final Map initLabels = + Map.of("encryption", String.valueOf(config.getEncryptionKey() != null)); final int numInstances = PooledResolver.getNumInstances(config.getResolverPoolSize()); - final LocalResolver inner = - new PooledResolver( - numInstances, - () -> - new RecoveringResolver( + final LocalResolver telemetryResolver = + new ProviderTelemetryResolver( + flagLogger::write, + SDK, + initLabels, + providerLogSink -> + new PooledResolver( + numInstances, () -> - new WasmLocalResolver( - flagLogger::write, - config.isEnableApplyDedup(), - config.isDisableExposureCollection()))); - this.resolver = new MaterializingResolver(inner, materializationStore); + new RecoveringResolver( + () -> + new WasmLocalResolver( + providerLogSink, + config.isEnableApplyDedup(), + config.isDisableExposureCollection())))); + this.resolver = new MaterializingResolver(telemetryResolver, materializationStore); } /** @@ -224,15 +232,22 @@ public OpenFeatureLocalResolveProvider( this.flagLogger = wasmFlagLogger; final int numInstances = PooledResolver.getNumInstances(LocalProviderConfig.DEFAULT_RESOLVER_POOL_SIZE); - final LocalResolver inner = - new PooledResolver( - numInstances, - () -> - new RecoveringResolver( + final LocalResolver telemetryResolver = + new ProviderTelemetryResolver( + wasmFlagLogger::write, + SDK, + Map.of(), + providerLogSink -> + new PooledResolver( + numInstances, () -> - new WasmLocalResolver( - wasmFlagLogger::write, enableApplyDedup, disableExposureCollection))); - this.resolver = new MaterializingResolver(inner, materializationStore); + new RecoveringResolver( + () -> + new WasmLocalResolver( + providerLogSink, + enableApplyDedup, + disableExposureCollection)))); + this.resolver = new MaterializingResolver(telemetryResolver, materializationStore); } @Override diff --git a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/ProviderTelemetryResolver.java b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/ProviderTelemetryResolver.java new file mode 100644 index 000000000..cf9dbc80d --- /dev/null +++ b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/ProviderTelemetryResolver.java @@ -0,0 +1,105 @@ +package com.spotify.confidence.sdk; + +import com.spotify.confidence.sdk.flags.resolver.v1.ApplyFlagsRequest; +import com.spotify.confidence.sdk.flags.resolver.v1.RegisterResolveRequest; +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 java.util.Map; +import java.util.concurrent.CompletionStage; +import java.util.function.Consumer; +import java.util.function.Function; + +/** Owns provider-scoped telemetry above the resolver pool and recovery layers. */ +final class ProviderTelemetryResolver implements LocalResolver { + private final Consumer logSink; + private final Sdk sdk; + private final Map labels; + private final LocalResolver delegate; + private boolean initSent; + + ProviderTelemetryResolver( + Consumer logSink, + Sdk sdk, + Map labels, + Function, LocalResolver> innerFactory) { + this.logSink = logSink; + this.sdk = sdk; + this.labels = Map.copyOf(labels); + this.delegate = innerFactory.apply(this::writeLogs); + } + + private synchronized void writeLogs(WriteFlagLogsRequest request) { + final WriteFlagLogsRequest outgoing = initSent ? request : addInitTelemetry(request); + logSink.accept(outgoing); + initSent = true; + } + + private synchronized void emitInitIfPending() { + if (initSent) { + return; + } + logSink.accept(addInitTelemetry(WriteFlagLogsRequest.getDefaultInstance())); + initSent = true; + } + + private WriteFlagLogsRequest addInitTelemetry(WriteFlagLogsRequest request) { + return request.toBuilder() + .setTelemetryData( + request.getTelemetryData().toBuilder() + .setSdk(sdk) + .addProviderInitRate( + TelemetryData.ProviderInitRate.newBuilder() + .setCount(1) + .putAllLabels(labels) + .build()) + .build()) + .build(); + } + + @Override + public void setResolverState(byte[] state, String accountId, Sdk sdk) { + delegate.setResolverState(state, accountId, sdk); + } + + @Override + public CompletionStage resolveProcess(ResolveProcessRequest request) { + return delegate.resolveProcess(request); + } + + @Override + public void applyFlags(ApplyFlagsRequest request) { + delegate.applyFlags(request); + } + + @Override + public void registerResolve(RegisterResolveRequest request) { + delegate.registerResolve(request); + } + + @Override + public void flushAllLogs() { + delegate.flushAllLogs(); + } + + @Override + public void flushAssignLogs() { + delegate.flushAssignLogs(); + } + + @Override + public void close() { + try { + delegate.close(); + } finally { + emitInitIfPending(); + } + } + + @Override + public String prometheusSnapshot() { + return delegate.prometheusSnapshot(); + } +} diff --git a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/WasmLocalResolver.java b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/WasmLocalResolver.java index f599ec039..cb85f23aa 100644 --- a/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/WasmLocalResolver.java +++ b/openfeature-provider/java/src/main/java/com/spotify/confidence/sdk/WasmLocalResolver.java @@ -63,7 +63,7 @@ class WasmLocalResolver implements LocalResolver { private final boolean disableExposureCollection; public WasmLocalResolver(Consumer logSink) { - this(logSink, false); + this(logSink, false, false); } public WasmLocalResolver(Consumer logSink, boolean enableApplyDedup) { @@ -288,7 +288,8 @@ public String prometheusSnapshot() { private static boolean isEmptyLogRequest(WriteFlagLogsRequest request) { return request.getFlagAssignedCount() == 0 && request.getClientResolveInfoCount() == 0 - && request.getFlagResolveInfoCount() == 0; + && request.getFlagResolveInfoCount() == 0 + && !request.hasTelemetryData(); } private T consumeResponse(int addr, ParserFn codec) { @@ -296,9 +297,8 @@ private T consumeResponse(int addr, ParserFn 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); } diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/GrpcWasmFlagLoggerTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/GrpcWasmFlagLoggerTest.java index 4cd294302..bbfef2b72 100644 --- a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/GrpcWasmFlagLoggerTest.java +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/GrpcWasmFlagLoggerTest.java @@ -18,6 +18,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; @@ -39,6 +40,21 @@ void testEmptyRequest_shouldSkip() { logger.shutdown(); } + @Test + void testTelemetryOnlyRequest_shouldSend() { + final var capturedRequest = new AtomicReference(); + final var logger = new GrpcWasmFlagLogger("test-client-secret", capturedRequest::set); + final var request = + WriteFlagLogsRequest.newBuilder() + .setTelemetryData(TelemetryData.getDefaultInstance()) + .build(); + + logger.write(request); + logger.shutdown(); + + assertEquals(request, capturedRequest.get()); + } + @Test void testSmallRequest_shouldSendAsIs() { // Given diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ProviderTelemetryResolverTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ProviderTelemetryResolverTest.java new file mode 100644 index 000000000..7513a4d87 --- /dev/null +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ProviderTelemetryResolverTest.java @@ -0,0 +1,124 @@ +package com.spotify.confidence.sdk; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import com.spotify.confidence.sdk.flags.resolver.v1.ApplyFlagsRequest; +import com.spotify.confidence.sdk.flags.resolver.v1.RegisterResolveRequest; +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.SdkId; +import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest; +import java.util.ArrayList; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Consumer; +import org.junit.jupiter.api.Test; + +class ProviderTelemetryResolverTest { + private static final Sdk SDK = + Sdk.newBuilder().setId(SdkId.SDK_ID_JAVA_LOCAL_PROVIDER).setVersion("test-version").build(); + + @Test + void emitsOnceAcrossPooledFlushes() { + final var captured = new ArrayList(); + final var resolver = + new ProviderTelemetryResolver( + captured::add, + SDK, + Map.of("encryption", "true"), + sink -> new TelemetryTestResolver(sink, 3)); + + resolver.flushAllLogs(); + + assertThat(providerInitEventCount(captured)).isEqualTo(1); + final var telemetry = captured.get(0).getTelemetryData(); + assertThat(telemetry.getSdk()).isEqualTo(SDK); + assertThat(telemetry.getProviderInitRate(0).getLabelsMap()).containsEntry("encryption", "true"); + } + + @Test + void closeEmitsWithoutResolve() { + final var captured = new ArrayList(); + final var resolver = + new ProviderTelemetryResolver( + captured::add, SDK, Map.of(), sink -> new TelemetryTestResolver(sink, 0)); + + resolver.close(); + + assertThat(providerInitEventCount(captured)).isEqualTo(1); + } + + @Test + void retriesAfterSinkFailure() { + final var attempts = new AtomicInteger(); + final var captured = new ArrayList(); + final var resolver = + new ProviderTelemetryResolver( + request -> { + if (attempts.getAndIncrement() == 0) { + throw new IllegalStateException("send failed"); + } + captured.add(request); + }, + SDK, + Map.of(), + sink -> new TelemetryTestResolver(sink, 1)); + + assertThatThrownBy(resolver::flushAllLogs).isInstanceOf(IllegalStateException.class); + resolver.flushAllLogs(); + + assertThat(providerInitEventCount(captured)).isEqualTo(1); + } + + private static int providerInitEventCount(ArrayList requests) { + return requests.stream() + .mapToInt(request -> request.getTelemetryData().getProviderInitRateCount()) + .sum(); + } + + private static final class TelemetryTestResolver implements LocalResolver { + private final Consumer sink; + private final int flushCount; + + private TelemetryTestResolver(Consumer sink, int flushCount) { + this.sink = sink; + this.flushCount = flushCount; + } + + @Override + public void setResolverState(byte[] state, String accountId, Sdk sdk) {} + + @Override + public CompletionStage resolveProcess(ResolveProcessRequest request) { + return CompletableFuture.completedFuture(ResolveProcessResponse.getDefaultInstance()); + } + + @Override + public void applyFlags(ApplyFlagsRequest request) {} + + @Override + public void registerResolve(RegisterResolveRequest request) {} + + @Override + public void flushAllLogs() { + for (int i = 0; i < flushCount; i++) { + sink.accept(WriteFlagLogsRequest.getDefaultInstance()); + } + } + + @Override + public void flushAssignLogs() {} + + @Override + public void close() {} + + @Override + public String prometheusSnapshot() { + return ""; + } + } +} diff --git a/openfeature-provider/js/proto/test-only.proto b/openfeature-provider/js/proto/test-only.proto index 99f266687..e5d59fef8 100644 --- a/openfeature-provider/js/proto/test-only.proto +++ b/openfeature-provider/js/proto/test-only.proto @@ -35,6 +35,14 @@ message TelemetryData { string resolver_version = 8; + repeated ProviderInitRate provider_init_rate = 9; + + message ProviderInitRate { + uint32 count = 1; + reserved 2; // status — tbd + map labels = 3; + } + message ResolveLatency { uint32 sum = 1; uint32 count = 2; diff --git a/openfeature-provider/js/src/ConfidenceServerProviderLocal.test.ts b/openfeature-provider/js/src/ConfidenceServerProviderLocal.test.ts index c26298aff..1244ef9db 100644 --- a/openfeature-provider/js/src/ConfidenceServerProviderLocal.test.ts +++ b/openfeature-provider/js/src/ConfidenceServerProviderLocal.test.ts @@ -9,6 +9,8 @@ import { abortableSleep, TimeUnit, timeoutSignal } from './util'; import { advanceTimersUntil, NetworkMock } from './test-helpers'; import { sha256Hex } from './hash'; import { ResolveReason } from './proto/confidence/flags/resolver/v1/types'; +import { WriteFlagLogsRequest } from './proto/test-only'; +import { VERSION } from './version'; vi.mock(import('./hash'), async () => { const { sha256Hex } = await import('./test-helpers'); @@ -156,6 +158,59 @@ describe('state update scheduling', () => { }); describe('flush behavior', () => { + it('preserves existing telemetry and sets provider SDK metadata when adding provider init telemetry', async () => { + let sentBody: Uint8Array | undefined; + net.resolver.flagLogs.handler = async (req: Request) => { + sentBody = new Uint8Array(await req.arrayBuffer()); + return new Response(null, { status: 200 }); + }; + mockedWasmResolver.flushLogs.mockReturnValueOnce( + WriteFlagLogsRequest.encode( + WriteFlagLogsRequest.create({ + telemetryData: { + resolverVersion: '0.20.0', + sdk: { id: 25, version: 'resolver-version' }, + providerInitRate: [{ count: 2, labels: { existing: 'true' } }], + }, + }), + ).finish(), + ); + + await advanceTimersUntil(provider.flush()); + + expect(sentBody).toBeDefined(); + 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' } }, + ]); + }); + + it('does not retry provider init telemetry after an HTTP failure', async () => { + const sentBodies: Uint8Array[] = []; + let attempts = 0; + net.resolver.flagLogs.handler = async (req: Request) => { + sentBodies.push(new Uint8Array(await req.arrayBuffer())); + attempts++; + return new Response(null, { status: attempts <= 3 ? 503 : 200 }); + }; + mockedWasmResolver.flushLogs.mockReturnValue( + WriteFlagLogsRequest.encode( + WriteFlagLogsRequest.create({ telemetryData: { resolverVersion: '0.20.0' } }), + ).finish(), + ); + + await advanceTimersUntil(provider.flush()); + await advanceTimersUntil(provider.flush()); + + const firstAttempt = WriteFlagLogsRequest.decode(sentBodies[0]); + const retryAttempt = WriteFlagLogsRequest.decode(sentBodies[3]); + expect(firstAttempt.telemetryData?.providerInitRate).toHaveLength(1); + expect(retryAttempt.telemetryData?.providerInitRate).toHaveLength(0); + }); + it('flushes periodically at the configured interval', async () => { await advanceTimersUntil(expect(provider.initialize()).resolves.toBeUndefined()); @@ -190,6 +245,27 @@ describe('flush behavior', () => { expect(net.resolver.flagLogs.calls).toBe(start + 1); }); + it('emits provider init telemetry on close when there are no resolver logs', async () => { + let sentBody: Uint8Array | undefined; + net.resolver.flagLogs.handler = async (req: Request) => { + sentBody = new Uint8Array(await req.arrayBuffer()); + return new Response(null, { status: 200 }); + }; + mockedWasmResolver.flushLogs.mockReturnValueOnce(new Uint8Array(0)); + + await advanceTimersUntil(expect(provider.onClose()).resolves.toBeUndefined()); + + expect(sentBody).toBeDefined(); + const decoded = WriteFlagLogsRequest.decode(sentBody!); + expect(decoded.telemetryData?.sdk).toEqual({ id: 22, customId: undefined, version: VERSION }); + expect(decoded.telemetryData?.providerInitRate).toEqual([{ count: 1, labels: { encryption: 'false' } }]); + }); + it('keeps close best-effort when provider init telemetry cannot be sent', async () => { + mockedWasmResolver.flushLogs.mockReturnValueOnce(new Uint8Array(0)); + net.resolver.flagLogs.status = 'No network'; + + await advanceTimersUntil(expect(provider.onClose()).resolves.toBeUndefined()); + }); it('skips flush if there are no logs to send', async () => { await advanceTimersUntil(expect(provider.initialize()).resolves.toBeUndefined()); diff --git a/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts b/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts index 22cf30df7..6daf14e1c 100644 --- a/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts +++ b/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts @@ -77,6 +77,8 @@ export class ConfidenceServerProviderLocal implements Provider { private readonly stateUpdateInterval: number; private readonly flushInterval: number; private readonly materializationStore: MaterializationStore | null; + private readonly initLabels: Record; + private initTelemetryState: 'pending' | 'sending' | 'sent' = 'pending'; private stateEtag: string | null = null; private logDestinations: LogDestination[] = []; private accountId = ''; @@ -164,6 +166,7 @@ export class ConfidenceServerProviderLocal implements Provider { } else { this.materializationStore = null; } + this.initLabels = { encryption: options.encryptionKey ? 'true' : 'false' }; } async initialize(context?: EvaluationContext): Promise { @@ -195,8 +198,25 @@ export class ConfidenceServerProviderLocal implements Provider { } async onClose(): Promise { - await this.flush(timeoutSignal(3000)); - this.main.abort(); + const signal = timeoutSignal(3000); + try { + try { + await this.flush(signal); + } catch { + // best-effort: try an init-only request below + } + if (this.initTelemetryState !== 'sent') { + try { + const request = this.addProviderInitTelemetry(new Uint8Array()); + await this.sendFlagLogs(request, signal); + this.initTelemetryState = 'sent'; + } catch { + // best-effort: provider is shutting down + } + } + } finally { + this.main.abort(); + } } async resolve(context: EvaluationContext, flagNames: string[], apply = false): Promise { @@ -375,9 +395,24 @@ export class ConfidenceServerProviderLocal implements Provider { // TODO should this return success/failure, or even throw? async flush(signal?: AbortSignal): Promise { - const writeFlagLogRequest = this.resolver.flushLogs(); + let writeFlagLogRequest = this.resolver.flushLogs(); if (writeFlagLogRequest.length > 0) { - await this.sendFlagLogs(writeFlagLogRequest, signal); + const includeInit = this.initTelemetryState === 'pending'; + if (includeInit) { + this.initTelemetryState = 'sending'; + writeFlagLogRequest = this.addProviderInitTelemetry(writeFlagLogRequest); + } + try { + await this.sendFlagLogs(writeFlagLogRequest, signal); + if (includeInit) { + this.initTelemetryState = 'sent'; + } + } catch (error) { + if (includeInit) { + this.initTelemetryState = 'pending'; + } + throw error; + } } } @@ -414,6 +449,19 @@ export class ConfidenceServerProviderLocal implements Provider { } } + private addProviderInitTelemetry(encodedWriteFlagLogRequest: Uint8Array): Uint8Array { + const request = WriteFlagLogsRequest.decode(encodedWriteFlagLogRequest); + if (!request.telemetryData) { + request.telemetryData = { resolverVersion: '', providerInitRate: [] }; + } + request.telemetryData.sdk = { + id: SdkId.SDK_ID_JS_LOCAL_SERVER_PROVIDER, + version: VERSION, + }; + request.telemetryData.providerInitRate.push({ count: 1, labels: this.initLabels }); + return WriteFlagLogsRequest.encode(request).finish(); + } + /** * Send flag logs to a specific destination. Returns true on success (HTTP 2xx), * false on a non-OK HTTP response. Throws on network errors. diff --git a/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto b/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto index b82063440..8ba7bc835 100644 --- a/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto +++ b/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto @@ -50,6 +50,18 @@ message IngestFlagLogsRequest { message TelemetryData { // Information about the SDK/provider Sdk sdk = 2; + + // Version of the embedded resolver (e.g. "0.14.0"). + // This must be preserved when providers decode and re-encode telemetry. + string resolver_version = 8; + + repeated ProviderInitRate provider_init_rate = 9; + + message ProviderInitRate { + uint32 count = 1; + reserved 2; // status — tbd + map labels = 3; + } } message ClientInfo { diff --git a/openfeature-provider/python/src/confidence/proto/confidence/flags/resolver/v1/internal_api_pb2.py b/openfeature-provider/python/src/confidence/proto/confidence/flags/resolver/v1/internal_api_pb2.py index b1c029b86..a835c476f 100644 --- a/openfeature-provider/python/src/confidence/proto/confidence/flags/resolver/v1/internal_api_pb2.py +++ b/openfeature-provider/python/src/confidence/proto/confidence/flags/resolver/v1/internal_api_pb2.py @@ -26,7 +26,7 @@ from . import types_pb2 as confidence_dot_flags_dot_resolver_dot_v1_dot_types__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n/confidence/flags/resolver/v1/internal_api.proto\x12\x1c\x63onfidence.flags.resolver.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a(confidence/flags/resolver/v1/types.proto\"\xb6\x02\n\x14WriteFlagLogsRequest\x12\x41\n\rflag_assigned\x18\x01 \x03(\x0b\x32*.confidence.flags.resolver.v1.FlagAssigned\x12\x43\n\x0etelemetry_data\x18\x02 \x01(\x0b\x32+.confidence.flags.resolver.v1.TelemetryData\x12L\n\x13\x63lient_resolve_info\x18\x03 \x03(\x0b\x32/.confidence.flags.resolver.v1.ClientResolveInfo\x12H\n\x11\x66lag_resolve_info\x18\x04 \x03(\x0b\x32-.confidence.flags.resolver.v1.FlagResolveInfo\"\x17\n\x15WriteFlagLogsResponse\"n\n\x15IngestFlagLogsRequest\x12\x12\n\naccount_id\x18\x01 \x01(\t\x12\x41\n\x05\x62\x61tch\x18\x02 \x01(\x0b\x32\x32.confidence.flags.resolver.v1.WriteFlagLogsRequest\"?\n\rTelemetryData\x12.\n\x03sdk\x18\x02 \x01(\x0b\x32!.confidence.flags.resolver.v1.Sdk\"g\n\nClientInfo\x12\x0e\n\x06\x63lient\x18\x01 \x01(\t\x12\x19\n\x11\x63lient_credential\x18\x02 \x01(\t\x12.\n\x03sdk\x18\x03 \x01(\x0b\x32!.confidence.flags.resolver.v1.Sdk\"\xa4\x07\n\x0c\x46lagAssigned\x12\x12\n\nresolve_id\x18\n \x01(\t\x12=\n\x0b\x63lient_info\x18\x03 \x01(\x0b\x32(.confidence.flags.resolver.v1.ClientInfo\x12\x45\n\x05\x66lags\x18\x0f \x03(\x0b\x32\x36.confidence.flags.resolver.v1.FlagAssigned.AppliedFlag\x1a\xbd\x03\n\x0b\x41ppliedFlag\x12\x0c\n\x04\x66lag\x18\x01 \x01(\t\x12\x15\n\rtargeting_key\x18\x02 \x01(\t\x12\x1e\n\x16targeting_key_selector\x18\x03 \x01(\t\x12T\n\x0f\x61ssignment_info\x18\x04 \x01(\x0b\x32\x39.confidence.flags.resolver.v1.FlagAssigned.AssignmentInfoH\x00\x12Z\n\x12\x64\x65\x66\x61ult_assignment\x18\x05 \x01(\x0b\x32<.confidence.flags.resolver.v1.FlagAssigned.DefaultAssignmentH\x00\x12\x15\n\rassignment_id\x18\x06 \x01(\t\x12\x0c\n\x04rule\x18\x07 \x01(\t\x12T\n\x17\x66\x61llthrough_assignments\x18\x08 \x03(\x0b\x32\x33.confidence.flags.resolver.v1.FallthroughAssignment\x12.\n\napply_time\x18\t \x01(\x0b\x32\x1a.google.protobuf.TimestampB\x0c\n\nassignment\x1a\x32\n\x0e\x41ssignmentInfo\x12\x0f\n\x07segment\x18\x01 \x01(\t\x12\x0f\n\x07variant\x18\x02 \x01(\t\x1a\x85\x02\n\x11\x44\x65\x66\x61ultAssignment\x12\x64\n\x06reason\x18\x01 \x01(\x0e\x32T.confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason\"\x89\x01\n\x17\x44\x65\x66\x61ultAssignmentReason\x12)\n%DEFAULT_ASSIGNMENT_REASON_UNSPECIFIED\x10\x00\x12\x14\n\x10NO_SEGMENT_MATCH\x10\x01\x12\x1a\n\x12NO_TREATMENT_MATCH\x10\x02\x1a\x02\x08\x01\x12\x11\n\rFLAG_ARCHIVED\x10\x03\"s\n\x15\x46\x61llthroughAssignment\x12\x0c\n\x04rule\x18\x01 \x01(\t\x12\x15\n\rassignment_id\x18\x02 \x01(\t\x12\x15\n\rtargeting_key\x18\x03 \x01(\t\x12\x1e\n\x16targeting_key_selector\x18\x04 \x01(\t\"\xdf\x02\n\x11\x43lientResolveInfo\x12\x0e\n\x06\x63lient\x18\x01 \x01(\t\x12\x19\n\x11\x63lient_credential\x18\x02 \x01(\t\x12_\n\x06schema\x18\x03 \x03(\x0b\x32O.confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance\x1a\xbd\x01\n\x1f\x45valuationContextSchemaInstance\x12k\n\x06schema\x18\x01 \x03(\x0b\x32[.confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry\x1a-\n\x0bSchemaEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\x05:\x02\x38\x01\"\xb5\x01\n\x0f\x46lagResolveInfo\x12\x0c\n\x04\x66lag\x18\x01 \x01(\t\x12^\n\x14variant_resolve_info\x18\x02 \x03(\x0b\x32@.confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo\x1a\x34\n\x12VariantResolveInfo\x12\x0f\n\x07variant\x18\x01 \x01(\t\x12\r\n\x05\x63ount\x18\x03 \x01(\x03\"]\n\x16WriteOperationsRequest\x12\x43\n\x10store_variant_op\x18\x01 \x03(\x0b\x32).confidence.flags.resolver.v1.VariantData\"\x17\n\x15WriteOperationsResult\"D\n\rVariantReadOp\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x0c\n\x04rule\x18\x03 \x01(\t\"8\n\x0fInclusionReadOp\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\"\xa2\x01\n\x06ReadOp\x12\x46\n\x0fvariant_read_op\x18\x01 \x01(\x0b\x32+.confidence.flags.resolver.v1.VariantReadOpH\x00\x12J\n\x11inclusion_read_op\x18\x02 \x01(\x0b\x32-.confidence.flags.resolver.v1.InclusionReadOpH\x00\x42\x04\n\x02op\"J\n\x15ReadOperationsRequest\x12\x31\n\x03ops\x18\x03 \x03(\x0b\x32$.confidence.flags.resolver.v1.ReadOp\"S\n\x0bVariantData\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x0c\n\x04rule\x18\x03 \x01(\t\x12\x0f\n\x07variant\x18\x04 \x01(\t\"K\n\rInclusionData\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x13\n\x0bis_included\x18\x03 \x01(\x08\"\xa4\x01\n\nReadResult\x12\x43\n\x0evariant_result\x18\x01 \x01(\x0b\x32).confidence.flags.resolver.v1.VariantDataH\x00\x12G\n\x10inclusion_result\x18\x02 \x01(\x0b\x32+.confidence.flags.resolver.v1.InclusionDataH\x00\x42\x08\n\x06result\"Q\n\x14ReadOperationsResult\x12\x39\n\x07results\x18\x01 \x03(\x0b\x32(.confidence.flags.resolver.v1.ReadResult2\xb2\x03\n\x19InternalFlagLoggerService\x12~\n\x13\x43lientWriteFlagLogs\x12\x32.confidence.flags.resolver.v1.WriteFlagLogsRequest\x1a\x33.confidence.flags.resolver.v1.WriteFlagLogsResponse\x12\x8a\x01\n\x1bWriteMaterializedOperations\x12\x34.confidence.flags.resolver.v1.WriteOperationsRequest\x1a\x33.confidence.flags.resolver.v1.WriteOperationsResult\"\x00\x12\x87\x01\n\x1aReadMaterializedOperations\x12\x33.confidence.flags.resolver.v1.ReadOperationsRequest\x1a\x32.confidence.flags.resolver.v1.ReadOperationsResult\"\x00\x42\xad\x01\n,com.spotify.confidence.sdk.flags.resolver.v1B\x10InternalApiProtoP\x01Zigithub.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternalb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n/confidence/flags/resolver/v1/internal_api.proto\x12\x1c\x63onfidence.flags.resolver.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a(confidence/flags/resolver/v1/types.proto\"\xb6\x02\n\x14WriteFlagLogsRequest\x12\x41\n\rflag_assigned\x18\x01 \x03(\x0b\x32*.confidence.flags.resolver.v1.FlagAssigned\x12\x43\n\x0etelemetry_data\x18\x02 \x01(\x0b\x32+.confidence.flags.resolver.v1.TelemetryData\x12L\n\x13\x63lient_resolve_info\x18\x03 \x03(\x0b\x32/.confidence.flags.resolver.v1.ClientResolveInfo\x12H\n\x11\x66lag_resolve_info\x18\x04 \x03(\x0b\x32-.confidence.flags.resolver.v1.FlagResolveInfo\"\x17\n\x15WriteFlagLogsResponse\"n\n\x15IngestFlagLogsRequest\x12\x12\n\naccount_id\x18\x01 \x01(\t\x12\x41\n\x05\x62\x61tch\x18\x02 \x01(\x0b\x32\x32.confidence.flags.resolver.v1.WriteFlagLogsRequest\"\xe6\x02\n\rTelemetryData\x12.\n\x03sdk\x18\x02 \x01(\x0b\x32!.confidence.flags.resolver.v1.Sdk\x12\x18\n\x10resolver_version\x18\x08 \x01(\t\x12X\n\x12provider_init_rate\x18\t \x03(\x0b\x32<.confidence.flags.resolver.v1.TelemetryData.ProviderInitRate\x1a\xb0\x01\n\x10ProviderInitRate\x12\r\n\x05\x63ount\x18\x01 \x01(\r\x12X\n\x06labels\x18\x03 \x03(\x0b\x32H.confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntry\x1a-\n\x0bLabelsEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01J\x04\x08\x02\x10\x03\"g\n\nClientInfo\x12\x0e\n\x06\x63lient\x18\x01 \x01(\t\x12\x19\n\x11\x63lient_credential\x18\x02 \x01(\t\x12.\n\x03sdk\x18\x03 \x01(\x0b\x32!.confidence.flags.resolver.v1.Sdk\"\xa4\x07\n\x0c\x46lagAssigned\x12\x12\n\nresolve_id\x18\n \x01(\t\x12=\n\x0b\x63lient_info\x18\x03 \x01(\x0b\x32(.confidence.flags.resolver.v1.ClientInfo\x12\x45\n\x05\x66lags\x18\x0f \x03(\x0b\x32\x36.confidence.flags.resolver.v1.FlagAssigned.AppliedFlag\x1a\xbd\x03\n\x0b\x41ppliedFlag\x12\x0c\n\x04\x66lag\x18\x01 \x01(\t\x12\x15\n\rtargeting_key\x18\x02 \x01(\t\x12\x1e\n\x16targeting_key_selector\x18\x03 \x01(\t\x12T\n\x0f\x61ssignment_info\x18\x04 \x01(\x0b\x32\x39.confidence.flags.resolver.v1.FlagAssigned.AssignmentInfoH\x00\x12Z\n\x12\x64\x65\x66\x61ult_assignment\x18\x05 \x01(\x0b\x32<.confidence.flags.resolver.v1.FlagAssigned.DefaultAssignmentH\x00\x12\x15\n\rassignment_id\x18\x06 \x01(\t\x12\x0c\n\x04rule\x18\x07 \x01(\t\x12T\n\x17\x66\x61llthrough_assignments\x18\x08 \x03(\x0b\x32\x33.confidence.flags.resolver.v1.FallthroughAssignment\x12.\n\napply_time\x18\t \x01(\x0b\x32\x1a.google.protobuf.TimestampB\x0c\n\nassignment\x1a\x32\n\x0e\x41ssignmentInfo\x12\x0f\n\x07segment\x18\x01 \x01(\t\x12\x0f\n\x07variant\x18\x02 \x01(\t\x1a\x85\x02\n\x11\x44\x65\x66\x61ultAssignment\x12\x64\n\x06reason\x18\x01 \x01(\x0e\x32T.confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason\"\x89\x01\n\x17\x44\x65\x66\x61ultAssignmentReason\x12)\n%DEFAULT_ASSIGNMENT_REASON_UNSPECIFIED\x10\x00\x12\x14\n\x10NO_SEGMENT_MATCH\x10\x01\x12\x1a\n\x12NO_TREATMENT_MATCH\x10\x02\x1a\x02\x08\x01\x12\x11\n\rFLAG_ARCHIVED\x10\x03\"s\n\x15\x46\x61llthroughAssignment\x12\x0c\n\x04rule\x18\x01 \x01(\t\x12\x15\n\rassignment_id\x18\x02 \x01(\t\x12\x15\n\rtargeting_key\x18\x03 \x01(\t\x12\x1e\n\x16targeting_key_selector\x18\x04 \x01(\t\"\xdf\x02\n\x11\x43lientResolveInfo\x12\x0e\n\x06\x63lient\x18\x01 \x01(\t\x12\x19\n\x11\x63lient_credential\x18\x02 \x01(\t\x12_\n\x06schema\x18\x03 \x03(\x0b\x32O.confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance\x1a\xbd\x01\n\x1f\x45valuationContextSchemaInstance\x12k\n\x06schema\x18\x01 \x03(\x0b\x32[.confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry\x1a-\n\x0bSchemaEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\x05:\x02\x38\x01\"\xb5\x01\n\x0f\x46lagResolveInfo\x12\x0c\n\x04\x66lag\x18\x01 \x01(\t\x12^\n\x14variant_resolve_info\x18\x02 \x03(\x0b\x32@.confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo\x1a\x34\n\x12VariantResolveInfo\x12\x0f\n\x07variant\x18\x01 \x01(\t\x12\r\n\x05\x63ount\x18\x03 \x01(\x03\"]\n\x16WriteOperationsRequest\x12\x43\n\x10store_variant_op\x18\x01 \x03(\x0b\x32).confidence.flags.resolver.v1.VariantData\"\x17\n\x15WriteOperationsResult\"D\n\rVariantReadOp\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x0c\n\x04rule\x18\x03 \x01(\t\"8\n\x0fInclusionReadOp\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\"\xa2\x01\n\x06ReadOp\x12\x46\n\x0fvariant_read_op\x18\x01 \x01(\x0b\x32+.confidence.flags.resolver.v1.VariantReadOpH\x00\x12J\n\x11inclusion_read_op\x18\x02 \x01(\x0b\x32-.confidence.flags.resolver.v1.InclusionReadOpH\x00\x42\x04\n\x02op\"J\n\x15ReadOperationsRequest\x12\x31\n\x03ops\x18\x03 \x03(\x0b\x32$.confidence.flags.resolver.v1.ReadOp\"S\n\x0bVariantData\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x0c\n\x04rule\x18\x03 \x01(\t\x12\x0f\n\x07variant\x18\x04 \x01(\t\"K\n\rInclusionData\x12\x0c\n\x04unit\x18\x01 \x01(\t\x12\x17\n\x0fmaterialization\x18\x02 \x01(\t\x12\x13\n\x0bis_included\x18\x03 \x01(\x08\"\xa4\x01\n\nReadResult\x12\x43\n\x0evariant_result\x18\x01 \x01(\x0b\x32).confidence.flags.resolver.v1.VariantDataH\x00\x12G\n\x10inclusion_result\x18\x02 \x01(\x0b\x32+.confidence.flags.resolver.v1.InclusionDataH\x00\x42\x08\n\x06result\"Q\n\x14ReadOperationsResult\x12\x39\n\x07results\x18\x01 \x03(\x0b\x32(.confidence.flags.resolver.v1.ReadResult2\xb2\x03\n\x19InternalFlagLoggerService\x12~\n\x13\x43lientWriteFlagLogs\x12\x32.confidence.flags.resolver.v1.WriteFlagLogsRequest\x1a\x33.confidence.flags.resolver.v1.WriteFlagLogsResponse\x12\x8a\x01\n\x1bWriteMaterializedOperations\x12\x34.confidence.flags.resolver.v1.WriteOperationsRequest\x1a\x33.confidence.flags.resolver.v1.WriteOperationsResult\"\x00\x12\x87\x01\n\x1aReadMaterializedOperations\x12\x33.confidence.flags.resolver.v1.ReadOperationsRequest\x1a\x32.confidence.flags.resolver.v1.ReadOperationsResult\"\x00\x42\xad\x01\n,com.spotify.confidence.sdk.flags.resolver.v1B\x10InternalApiProtoP\x01Zigithub.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternalb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) @@ -34,6 +34,8 @@ if not _descriptor._USE_C_DESCRIPTORS: _globals['DESCRIPTOR']._loaded_options = None _globals['DESCRIPTOR']._serialized_options = b'\n,com.spotify.confidence.sdk.flags.resolver.v1B\020InternalApiProtoP\001Zigithub.com/spotify/confidence-resolver/openfeature-provider/go/confidence/internal/proto/resolverinternal' + _globals['_TELEMETRYDATA_PROVIDERINITRATE_LABELSENTRY']._loaded_options = None + _globals['_TELEMETRYDATA_PROVIDERINITRATE_LABELSENTRY']._serialized_options = b'8\001' _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON'].values_by_name["NO_TREATMENT_MATCH"]._loaded_options = None _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON'].values_by_name["NO_TREATMENT_MATCH"]._serialized_options = b'\010\001' _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE_SCHEMAENTRY']._loaded_options = None @@ -44,52 +46,56 @@ _globals['_WRITEFLAGLOGSRESPONSE']._serialized_end=492 _globals['_INGESTFLAGLOGSREQUEST']._serialized_start=494 _globals['_INGESTFLAGLOGSREQUEST']._serialized_end=604 - _globals['_TELEMETRYDATA']._serialized_start=606 - _globals['_TELEMETRYDATA']._serialized_end=669 - _globals['_CLIENTINFO']._serialized_start=671 - _globals['_CLIENTINFO']._serialized_end=774 - _globals['_FLAGASSIGNED']._serialized_start=777 - _globals['_FLAGASSIGNED']._serialized_end=1709 - _globals['_FLAGASSIGNED_APPLIEDFLAG']._serialized_start=948 - _globals['_FLAGASSIGNED_APPLIEDFLAG']._serialized_end=1393 - _globals['_FLAGASSIGNED_ASSIGNMENTINFO']._serialized_start=1395 - _globals['_FLAGASSIGNED_ASSIGNMENTINFO']._serialized_end=1445 - _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT']._serialized_start=1448 - _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT']._serialized_end=1709 - _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON']._serialized_start=1572 - _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON']._serialized_end=1709 - _globals['_FALLTHROUGHASSIGNMENT']._serialized_start=1711 - _globals['_FALLTHROUGHASSIGNMENT']._serialized_end=1826 - _globals['_CLIENTRESOLVEINFO']._serialized_start=1829 - _globals['_CLIENTRESOLVEINFO']._serialized_end=2180 - _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE']._serialized_start=1991 - _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE']._serialized_end=2180 - _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE_SCHEMAENTRY']._serialized_start=2135 - _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE_SCHEMAENTRY']._serialized_end=2180 - _globals['_FLAGRESOLVEINFO']._serialized_start=2183 - _globals['_FLAGRESOLVEINFO']._serialized_end=2364 - _globals['_FLAGRESOLVEINFO_VARIANTRESOLVEINFO']._serialized_start=2312 - _globals['_FLAGRESOLVEINFO_VARIANTRESOLVEINFO']._serialized_end=2364 - _globals['_WRITEOPERATIONSREQUEST']._serialized_start=2366 - _globals['_WRITEOPERATIONSREQUEST']._serialized_end=2459 - _globals['_WRITEOPERATIONSRESULT']._serialized_start=2461 - _globals['_WRITEOPERATIONSRESULT']._serialized_end=2484 - _globals['_VARIANTREADOP']._serialized_start=2486 - _globals['_VARIANTREADOP']._serialized_end=2554 - _globals['_INCLUSIONREADOP']._serialized_start=2556 - _globals['_INCLUSIONREADOP']._serialized_end=2612 - _globals['_READOP']._serialized_start=2615 - _globals['_READOP']._serialized_end=2777 - _globals['_READOPERATIONSREQUEST']._serialized_start=2779 - _globals['_READOPERATIONSREQUEST']._serialized_end=2853 - _globals['_VARIANTDATA']._serialized_start=2855 - _globals['_VARIANTDATA']._serialized_end=2938 - _globals['_INCLUSIONDATA']._serialized_start=2940 - _globals['_INCLUSIONDATA']._serialized_end=3015 - _globals['_READRESULT']._serialized_start=3018 - _globals['_READRESULT']._serialized_end=3182 - _globals['_READOPERATIONSRESULT']._serialized_start=3184 - _globals['_READOPERATIONSRESULT']._serialized_end=3265 - _globals['_INTERNALFLAGLOGGERSERVICE']._serialized_start=3268 - _globals['_INTERNALFLAGLOGGERSERVICE']._serialized_end=3702 + _globals['_TELEMETRYDATA']._serialized_start=607 + _globals['_TELEMETRYDATA']._serialized_end=965 + _globals['_TELEMETRYDATA_PROVIDERINITRATE']._serialized_start=789 + _globals['_TELEMETRYDATA_PROVIDERINITRATE']._serialized_end=965 + _globals['_TELEMETRYDATA_PROVIDERINITRATE_LABELSENTRY']._serialized_start=914 + _globals['_TELEMETRYDATA_PROVIDERINITRATE_LABELSENTRY']._serialized_end=959 + _globals['_CLIENTINFO']._serialized_start=967 + _globals['_CLIENTINFO']._serialized_end=1070 + _globals['_FLAGASSIGNED']._serialized_start=1073 + _globals['_FLAGASSIGNED']._serialized_end=2005 + _globals['_FLAGASSIGNED_APPLIEDFLAG']._serialized_start=1244 + _globals['_FLAGASSIGNED_APPLIEDFLAG']._serialized_end=1689 + _globals['_FLAGASSIGNED_ASSIGNMENTINFO']._serialized_start=1691 + _globals['_FLAGASSIGNED_ASSIGNMENTINFO']._serialized_end=1741 + _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT']._serialized_start=1744 + _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT']._serialized_end=2005 + _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON']._serialized_start=1868 + _globals['_FLAGASSIGNED_DEFAULTASSIGNMENT_DEFAULTASSIGNMENTREASON']._serialized_end=2005 + _globals['_FALLTHROUGHASSIGNMENT']._serialized_start=2007 + _globals['_FALLTHROUGHASSIGNMENT']._serialized_end=2122 + _globals['_CLIENTRESOLVEINFO']._serialized_start=2125 + _globals['_CLIENTRESOLVEINFO']._serialized_end=2476 + _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE']._serialized_start=2287 + _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE']._serialized_end=2476 + _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE_SCHEMAENTRY']._serialized_start=2431 + _globals['_CLIENTRESOLVEINFO_EVALUATIONCONTEXTSCHEMAINSTANCE_SCHEMAENTRY']._serialized_end=2476 + _globals['_FLAGRESOLVEINFO']._serialized_start=2479 + _globals['_FLAGRESOLVEINFO']._serialized_end=2660 + _globals['_FLAGRESOLVEINFO_VARIANTRESOLVEINFO']._serialized_start=2608 + _globals['_FLAGRESOLVEINFO_VARIANTRESOLVEINFO']._serialized_end=2660 + _globals['_WRITEOPERATIONSREQUEST']._serialized_start=2662 + _globals['_WRITEOPERATIONSREQUEST']._serialized_end=2755 + _globals['_WRITEOPERATIONSRESULT']._serialized_start=2757 + _globals['_WRITEOPERATIONSRESULT']._serialized_end=2780 + _globals['_VARIANTREADOP']._serialized_start=2782 + _globals['_VARIANTREADOP']._serialized_end=2850 + _globals['_INCLUSIONREADOP']._serialized_start=2852 + _globals['_INCLUSIONREADOP']._serialized_end=2908 + _globals['_READOP']._serialized_start=2911 + _globals['_READOP']._serialized_end=3073 + _globals['_READOPERATIONSREQUEST']._serialized_start=3075 + _globals['_READOPERATIONSREQUEST']._serialized_end=3149 + _globals['_VARIANTDATA']._serialized_start=3151 + _globals['_VARIANTDATA']._serialized_end=3234 + _globals['_INCLUSIONDATA']._serialized_start=3236 + _globals['_INCLUSIONDATA']._serialized_end=3311 + _globals['_READRESULT']._serialized_start=3314 + _globals['_READRESULT']._serialized_end=3478 + _globals['_READOPERATIONSRESULT']._serialized_start=3480 + _globals['_READOPERATIONSRESULT']._serialized_end=3561 + _globals['_INTERNALFLAGLOGGERSERVICE']._serialized_start=3564 + _globals['_INTERNALFLAGLOGGERSERVICE']._serialized_end=3998 # @@protoc_insertion_point(module_scope) diff --git a/openfeature-provider/python/src/confidence/provider.py b/openfeature-provider/python/src/confidence/provider.py index 79e5bfd80..eeac2275d 100644 --- a/openfeature-provider/python/src/confidence/provider.py +++ b/openfeature-provider/python/src/confidence/provider.py @@ -35,7 +35,11 @@ VariantReadResult, VariantWriteOp, ) -from confidence.proto.confidence.flags.resolver.v1 import api_pb2, types_pb2 +from confidence.proto.confidence.flags.resolver.v1 import ( + api_pb2, + internal_api_pb2, + types_pb2, +) from confidence.proto.confidence.wasm import wasm_api_pb2 from confidence.state_fetcher import StateFetcher from confidence.version import __version__ @@ -165,6 +169,11 @@ def __init__( """ self._client_secret = client_secret self._encryption_key = encryption_key + self._init_labels: Dict[str, str] = { + "encryption": str(bool(encryption_key)).lower() + } + self._init_telemetry_state = "pending" + self._init_telemetry_lock = threading.Lock() self._state_poll_interval = state_poll_interval self._log_poll_interval = log_poll_interval self._assign_poll_interval = assign_poll_interval @@ -322,9 +331,7 @@ def shutdown(self) -> None: # Flush final logs if self._resolver is not None: try: - log_data = self._resolver.flush_logs() - if log_data and self._flag_logger is not None: - self._flag_logger.write(log_data) + self._write_logs(self._resolver.flush_logs()) except Exception as e: logger.error("Failed to flush final logs: %s", e) @@ -770,6 +777,42 @@ def _write_materializations( except Exception as e: logger.error("Failed to write materializations: %s", e) + def _write_logs(self, log_data: bytes) -> None: + if not log_data or self._flag_logger is None: + return + + include_init = False + with self._init_telemetry_lock: + if self._init_telemetry_state == "pending": + self._init_telemetry_state = "sending" + include_init = True + + if include_init: + request = internal_api_pb2.WriteFlagLogsRequest.FromString(log_data) + request.telemetry_data.sdk.CopyFrom( + types_pb2.Sdk( + id=types_pb2.SdkId.SDK_ID_PYTHON_PROVIDER, + version=__version__, + ) + ) + init_rate = request.telemetry_data.provider_init_rate.add() + init_rate.count = 1 + for k, v in self._init_labels.items(): + init_rate.labels[k] = v + log_data = request.SerializeToString() + + try: + self._flag_logger.write(log_data) + except Exception: + if include_init: + with self._init_telemetry_lock: + self._init_telemetry_state = "pending" + raise + else: + if include_init: + with self._init_telemetry_lock: + self._init_telemetry_state = "sent" + def _create_flag_logger( self, account_id: str, log_destinations: List[int] ) -> FlagLogger: @@ -872,8 +915,7 @@ def _state_poll_loop(self) -> None: self._enable_apply_dedup, self._disable_exposure_collection, ) - if flushed_logs and self._flag_logger is not None: - self._flag_logger.write(flushed_logs) + self._write_logs(flushed_logs) # Update account ID and destinations on the logger if self._flag_logger is not None: @@ -911,8 +953,7 @@ def _log_flush_loop(self) -> None: try: with self._resolver_lock: log_data = self._resolver.flush_logs() - if log_data and self._flag_logger is not None: - self._flag_logger.write(log_data) + self._write_logs(log_data) except Exception as e: logger.error("Failed to flush logs: %s", e) last_full_flush = now diff --git a/openfeature-provider/python/tests/test_provider.py b/openfeature-provider/python/tests/test_provider.py index e07aaa03d..683dc45f8 100644 --- a/openfeature-provider/python/tests/test_provider.py +++ b/openfeature-provider/python/tests/test_provider.py @@ -5,6 +5,8 @@ from openfeature.flag_evaluation import FlagResolutionDetails, Reason from confidence.provider import ConfidenceProvider +from confidence.proto.confidence.flags.resolver.v1 import internal_api_pb2, types_pb2 +from confidence.version import __version__ from tests.conftest import MockFlagLogger, MockStateFetcher @@ -36,6 +38,65 @@ def test_get_metadata( class TestInitialize: """Tests for provider initialization.""" + def test_init_telemetry_includes_sdk( + self, + wasm_bytes: bytes, + test_client_secret: str, + ) -> None: + provider = ConfidenceProvider( + client_secret=test_client_secret, + flag_logger=MockFlagLogger(), + wasm_bytes=wasm_bytes, + ) + request = internal_api_pb2.WriteFlagLogsRequest() + request.telemetry_data.SetInParent() + + provider._write_logs(request.SerializeToString()) + decoded = internal_api_pb2.WriteFlagLogsRequest.FromString( + provider._flag_logger.writes[0] + ) + + 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 + + def test_init_telemetry_retries_after_failed_write( + self, + wasm_bytes: bytes, + test_client_secret: str, + ) -> None: + class FailOnceLogger(MockFlagLogger): + def __init__(self) -> None: + super().__init__() + self.attempts = 0 + + def write(self, request_bytes: bytes) -> None: + self.attempts += 1 + if self.attempts == 1: + raise RuntimeError("send failed") + super().write(request_bytes) + + mock_logger = FailOnceLogger() + provider = ConfidenceProvider( + client_secret=test_client_secret, + flag_logger=mock_logger, + wasm_bytes=wasm_bytes, + ) + request = internal_api_pb2.WriteFlagLogsRequest() + request.telemetry_data.SetInParent() + encoded = request.SerializeToString() + + try: + provider._write_logs(encoded) + except RuntimeError: + pass + provider._write_logs(encoded) + + decoded = internal_api_pb2.WriteFlagLogsRequest.FromString( + mock_logger.writes[0] + ) + assert len(decoded.telemetry_data.provider_init_rate) == 1 + def test_initialize_fetches_state( self, wasm_bytes: bytes, diff --git a/openfeature-provider/rust/src/logger.rs b/openfeature-provider/rust/src/logger.rs index 57721e6ea..b492e0b40 100644 --- a/openfeature-provider/rust/src/logger.rs +++ b/openfeature-provider/rust/src/logger.rs @@ -1,5 +1,7 @@ //! Log management for sending flag logs to the Confidence API. +use std::collections::BTreeMap; +use std::sync::atomic::{AtomicU8, Ordering}; use std::sync::Arc; use prost::Message; @@ -7,7 +9,9 @@ use reqwest_middleware::ClientWithMiddleware; use tokio::sync::RwLock; use confidence_resolver::assign_logger::AssignLogger; -use confidence_resolver::proto::confidence::flags::resolver::v1::{Sdk, WriteFlagLogsRequest}; +use confidence_resolver::proto::confidence::flags::resolver::v1::{ + telemetry_data::ProviderInitRate, Sdk, WriteFlagLogsRequest, +}; use confidence_resolver::resolve_logger::ResolveLogger; use crate::error::Result; @@ -20,6 +24,35 @@ const CLOUDFLARE_URL: &str = /// Target size for log batches (4 MB). const LOG_TARGET_BYTES: usize = 4 * 1024 * 1024; +const INIT_PENDING: u8 = 0; +const INIT_SENDING: u8 = 1; +const INIT_SENT: u8 = 2; + +struct InitTelemetryState(AtomicU8); + +impl InitTelemetryState { + fn new() -> Self { + Self(AtomicU8::new(INIT_PENDING)) + } + + fn claim(&self) -> bool { + self.0 + .compare_exchange( + INIT_PENDING, + INIT_SENDING, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_ok() + } + + fn complete(&self, success: bool) { + self.0.store( + if success { INIT_SENT } else { INIT_PENDING }, + Ordering::Release, + ); + } +} /// Log sender that sends flag logs to the Confidence API. pub struct LogSender { @@ -156,19 +189,25 @@ fn encode_ingest_request(account_id: &str, batch: &[u8]) -> Vec { pub struct LogManager { sender: LogSender, sdk: Sdk, + init_labels: BTreeMap, + init_state: InitTelemetryState, } impl LogManager { + /// Create a new log manager with the given client, client secret, and SDK identity. pub fn new( client: ClientWithMiddleware, client_secret: String, sdk: Sdk, account_id: Arc>>, destinations: Arc>>, + init_labels: BTreeMap, ) -> Self { Self { sender: LogSender::new(client, client_secret, account_id, destinations), sdk, + init_labels, + init_state: InitTelemetryState::new(), } } @@ -183,11 +222,26 @@ impl LogManager { let mut td = TELEMETRY.delta_snapshot(&LAST_FLUSHED); 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(), + }); + } request.telemetry_data = Some(td); let encoded = request.encode_to_vec(); if !encoded.is_empty() && has_logs(&request) { - self.sender.send(&encoded).await?; + if let Err(error) = self.sender.send(&encoded).await { + if include_init { + self.init_state.complete(false); + } + return Err(error); + } + if include_init { + self.init_state.complete(true); + } } Ok(()) @@ -247,4 +301,16 @@ mod tests { assert_eq!(decoded.account_id, "acct"); assert!(decoded.batch.is_empty()); } + + #[test] + fn init_telemetry_state_retries_failure_and_stops_after_success() { + let state = InitTelemetryState::new(); + + assert!(state.claim()); + assert!(!state.claim()); + state.complete(false); + assert!(state.claim()); + state.complete(true); + assert!(!state.claim()); + } } diff --git a/openfeature-provider/rust/src/provider.rs b/openfeature-provider/rust/src/provider.rs index e2d1d3a9d..32edb0be4 100644 --- a/openfeature-provider/rust/src/provider.rs +++ b/openfeature-provider/rust/src/provider.rs @@ -1,6 +1,6 @@ //! OpenFeature provider implementation for Confidence. -use std::collections::HashMap; +use std::collections::{BTreeMap, HashMap}; use std::sync::Arc; use std::time::{Duration, Instant}; @@ -182,6 +182,7 @@ impl ConfidenceProvider { let client = client_builder.build(); let sdk = provider_sdk(); + let encryption_enabled = options.encryption_key.is_some(); let state_fetcher = Arc::new(StateFetcher::new( client.clone(), options.client_secret.clone(), @@ -189,12 +190,15 @@ impl ConfidenceProvider { options.encryption_key, )); let shared_state = Arc::new(SharedState::new()); + let init_labels = + BTreeMap::from([("encryption".to_string(), encryption_enabled.to_string())]); let log_manager = Arc::new(LogManager::new( client.clone(), options.client_secret.clone(), sdk, Arc::clone(&shared_state.account_id), Arc::clone(&shared_state.log_destinations), + init_labels, )); // Create materialization store if configured