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 9ab42a0d..9a9c781f 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 e075953c..bfd18de5 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/local_resolver/assets/confidence_resolver.wasm b/openfeature-provider/go/confidence/internal/local_resolver/assets/confidence_resolver.wasm index b8f24839..75d31fa5 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/local_resolver/local_resolver.go b/openfeature-provider/go/confidence/internal/local_resolver/local_resolver.go index 1e7675db..f961f4af 100644 --- a/openfeature-provider/go/confidence/internal/local_resolver/local_resolver.go +++ b/openfeature-provider/go/confidence/internal/local_resolver/local_resolver.go @@ -44,7 +44,16 @@ type localResolverImpl struct { } func NewLocalResolverWithPoolSize(ctx context.Context, logSink LogSink, poolSize int) LocalResolver { - factory := NewWasmResolverFactory(logSink) + return NewLocalResolverWithLabels(ctx, logSink, poolSize, nil) +} + +func NewLocalResolverWithLabels(ctx context.Context, logSink LogSink, poolSize int, initLabels map[string]string) LocalResolver { + var factory LocalResolverFactory + if len(initLabels) > 0 { + factory = NewWasmResolverFactoryWithLabels(logSink, initLabels) + } else { + factory = NewWasmResolverFactory(logSink) + } factory = NewRecoveringResolverFactory(factory) if poolSize <= 0 { poolSize = DefaultPoolSize diff --git a/openfeature-provider/go/confidence/internal/local_resolver/wasm.go b/openfeature-provider/go/confidence/internal/local_resolver/wasm.go index 533faa86..e22388e3 100644 --- a/openfeature-provider/go/confidence/internal/local_resolver/wasm.go +++ b/openfeature-provider/go/confidence/internal/local_resolver/wasm.go @@ -42,6 +42,8 @@ type WasmResolver struct { mu *sync.Mutex instanceID string fnCache sync.Map + initLabels map[string]string + firstFlush bool } var _ LocalResolver = (*WasmResolver)(nil) @@ -85,6 +87,16 @@ func (r *WasmResolver) ApplyFlags(request *resolver.ApplyFlagsRequest) error { func (r *WasmResolver) FlushAllLogs() error { resp := &resolverv1.WriteFlagLogsRequest{} err := r.call("wasm_msg_guest_bounded_flush_logs", nil, resp) + if err == nil && r.firstFlush { + r.firstFlush = false + if resp.TelemetryData == nil { + resp.TelemetryData = &resolverv1.TelemetryData{} + } + resp.TelemetryData.ProviderInitRate = append( + resp.TelemetryData.ProviderInitRate, + &resolverv1.TelemetryData_ProviderInitRate{Count: 1, Labels: r.initLabels}, + ) + } if err == nil && proto.Size(resp) > 0 { r.logSink(resp) } @@ -154,9 +166,10 @@ func (r *WasmResolver) call(fnName string, request proto.Message, response proto } type WasmResolverFactory struct { - runtime wazero.Runtime - module wazero.CompiledModule - logSink LogSink + runtime wazero.Runtime + module wazero.CompiledModule + logSink LogSink + initLabels map[string]string } var _ LocalResolverFactory = (*WasmResolverFactory)(nil) @@ -208,6 +221,12 @@ func NewWasmResolverFactory(logSink LogSink) LocalResolverFactory { } } +func NewWasmResolverFactoryWithLabels(logSink LogSink, initLabels map[string]string) LocalResolverFactory { + factory := NewWasmResolverFactory(logSink).(*WasmResolverFactory) + factory.initLabels = initLabels + return factory +} + func (wrf *WasmResolverFactory) New() LocalResolver { ctx := context.Background() config := wazero.NewModuleConfig().WithName("") @@ -221,6 +240,8 @@ func (wrf *WasmResolverFactory) New() LocalResolver { logSink: wrf.logSink, mu: &sync.Mutex{}, instanceID: fmt.Sprintf("%d", id), + initLabels: wrf.initLabels, + firstFlush: true, } } 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 3ca0b257..06eaed3d 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 @@ -190,9 +190,10 @@ func (*WriteFlagLogsResponse) Descriptor() ([]byte, []int) { 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"` + 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() { @@ -232,6 +233,13 @@ func (x *TelemetryData) GetSdk() *resolver.Sdk { return nil } +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"` @@ -1110,6 +1118,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[18] + 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[18] + 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{2, 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"` @@ -1130,7 +1190,7 @@ type FlagAssigned_AppliedFlag struct { func (x *FlagAssigned_AppliedFlag) Reset() { *x = FlagAssigned_AppliedFlag{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[18] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[20] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1142,7 +1202,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[18] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[20] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1258,7 +1318,7 @@ type FlagAssigned_AssignmentInfo struct { func (x *FlagAssigned_AssignmentInfo) Reset() { *x = FlagAssigned_AssignmentInfo{} - 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) } @@ -1270,7 +1330,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[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 { @@ -1309,7 +1369,7 @@ type FlagAssigned_DefaultAssignment struct { func (x *FlagAssigned_DefaultAssignment) Reset() { *x = FlagAssigned_DefaultAssignment{} - 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) } @@ -1321,7 +1381,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[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 { @@ -1355,7 +1415,7 @@ type ClientResolveInfo_EvaluationContextSchemaInstance struct { func (x *ClientResolveInfo_EvaluationContextSchemaInstance) Reset() { *x = ClientResolveInfo_EvaluationContextSchemaInstance{} - 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) } @@ -1367,7 +1427,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[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 { @@ -1403,7 +1463,7 @@ type FlagResolveInfo_VariantResolveInfo struct { func (x *FlagResolveInfo_VariantResolveInfo) Reset() { *x = FlagResolveInfo_VariantResolveInfo{} - mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[23] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[25] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1415,7 +1475,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[23] + mi := &file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[25] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1455,9 +1515,16 @@ const file_confidence_flags_resolver_v1_internal_api_proto_rawDesc = "" + "\x0etelemetry_data\x18\x02 \x01(\v2+.confidence.flags.resolver.v1.TelemetryDataR\rtelemetryData\x12_\n" + "\x13client_resolve_info\x18\x03 \x03(\v2/.confidence.flags.resolver.v1.ClientResolveInfoR\x11clientResolveInfo\x12Y\n" + "\x11flag_resolve_info\x18\x04 \x03(\v2-.confidence.flags.resolver.v1.FlagResolveInfoR\x0fflagResolveInfo\"\x17\n" + - "\x15WriteFlagLogsResponse\"D\n" + + "\x15WriteFlagLogsResponse\"\xfe\x02\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\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" + @@ -1565,7 +1632,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, 24) +var file_confidence_flags_resolver_v1_internal_api_proto_msgTypes = make([]protoimpl.MessageInfo, 26) 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 @@ -1586,50 +1653,54 @@ var file_confidence_flags_resolver_v1_internal_api_proto_goTypes = []any{ (*InclusionData)(nil), // 16: confidence.flags.resolver.v1.InclusionData (*ReadResult)(nil), // 17: confidence.flags.resolver.v1.ReadResult (*ReadOperationsResult)(nil), // 18: confidence.flags.resolver.v1.ReadOperationsResult - (*FlagAssigned_AppliedFlag)(nil), // 19: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag - (*FlagAssigned_AssignmentInfo)(nil), // 20: confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo - (*FlagAssigned_DefaultAssignment)(nil), // 21: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment - (*ClientResolveInfo_EvaluationContextSchemaInstance)(nil), // 22: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance - nil, // 23: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry - (*FlagResolveInfo_VariantResolveInfo)(nil), // 24: confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo - (*resolver.Sdk)(nil), // 25: confidence.flags.resolver.v1.Sdk - (*timestamppb.Timestamp)(nil), // 26: google.protobuf.Timestamp + (*TelemetryData_ProviderInitRate)(nil), // 19: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate + nil, // 20: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntry + (*FlagAssigned_AppliedFlag)(nil), // 21: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag + (*FlagAssigned_AssignmentInfo)(nil), // 22: confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo + (*FlagAssigned_DefaultAssignment)(nil), // 23: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment + (*ClientResolveInfo_EvaluationContextSchemaInstance)(nil), // 24: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance + nil, // 25: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry + (*FlagResolveInfo_VariantResolveInfo)(nil), // 26: confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo + (*resolver.Sdk)(nil), // 27: confidence.flags.resolver.v1.Sdk + (*timestamppb.Timestamp)(nil), // 28: google.protobuf.Timestamp } var file_confidence_flags_resolver_v1_internal_api_proto_depIdxs = []int32{ 5, // 0: confidence.flags.resolver.v1.WriteFlagLogsRequest.flag_assigned:type_name -> confidence.flags.resolver.v1.FlagAssigned 3, // 1: confidence.flags.resolver.v1.WriteFlagLogsRequest.telemetry_data:type_name -> confidence.flags.resolver.v1.TelemetryData 7, // 2: confidence.flags.resolver.v1.WriteFlagLogsRequest.client_resolve_info:type_name -> confidence.flags.resolver.v1.ClientResolveInfo 8, // 3: confidence.flags.resolver.v1.WriteFlagLogsRequest.flag_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo - 25, // 4: confidence.flags.resolver.v1.TelemetryData.sdk:type_name -> confidence.flags.resolver.v1.Sdk - 25, // 5: confidence.flags.resolver.v1.ClientInfo.sdk:type_name -> confidence.flags.resolver.v1.Sdk - 4, // 6: confidence.flags.resolver.v1.FlagAssigned.client_info:type_name -> confidence.flags.resolver.v1.ClientInfo - 19, // 7: confidence.flags.resolver.v1.FlagAssigned.flags:type_name -> confidence.flags.resolver.v1.FlagAssigned.AppliedFlag - 22, // 8: confidence.flags.resolver.v1.ClientResolveInfo.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance - 24, // 9: confidence.flags.resolver.v1.FlagResolveInfo.variant_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo - 15, // 10: confidence.flags.resolver.v1.WriteOperationsRequest.store_variant_op:type_name -> confidence.flags.resolver.v1.VariantData - 11, // 11: confidence.flags.resolver.v1.ReadOp.variant_read_op:type_name -> confidence.flags.resolver.v1.VariantReadOp - 12, // 12: confidence.flags.resolver.v1.ReadOp.inclusion_read_op:type_name -> confidence.flags.resolver.v1.InclusionReadOp - 13, // 13: confidence.flags.resolver.v1.ReadOperationsRequest.ops:type_name -> confidence.flags.resolver.v1.ReadOp - 15, // 14: confidence.flags.resolver.v1.ReadResult.variant_result:type_name -> confidence.flags.resolver.v1.VariantData - 16, // 15: confidence.flags.resolver.v1.ReadResult.inclusion_result:type_name -> confidence.flags.resolver.v1.InclusionData - 17, // 16: confidence.flags.resolver.v1.ReadOperationsResult.results:type_name -> confidence.flags.resolver.v1.ReadResult - 20, // 17: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.assignment_info:type_name -> confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo - 21, // 18: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.default_assignment:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment - 6, // 19: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.fallthrough_assignments:type_name -> confidence.flags.resolver.v1.FallthroughAssignment - 26, // 20: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.apply_time:type_name -> google.protobuf.Timestamp - 0, // 21: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.reason:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason - 23, // 22: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry - 1, // 23: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:input_type -> confidence.flags.resolver.v1.WriteFlagLogsRequest - 9, // 24: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:input_type -> confidence.flags.resolver.v1.WriteOperationsRequest - 14, // 25: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:input_type -> confidence.flags.resolver.v1.ReadOperationsRequest - 2, // 26: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:output_type -> confidence.flags.resolver.v1.WriteFlagLogsResponse - 10, // 27: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:output_type -> confidence.flags.resolver.v1.WriteOperationsResult - 18, // 28: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:output_type -> confidence.flags.resolver.v1.ReadOperationsResult - 26, // [26:29] is the sub-list for method output_type - 23, // [23:26] is the sub-list for method input_type - 23, // [23:23] is the sub-list for extension type_name - 23, // [23:23] is the sub-list for extension extendee - 0, // [0:23] is the sub-list for field type_name + 27, // 4: confidence.flags.resolver.v1.TelemetryData.sdk:type_name -> confidence.flags.resolver.v1.Sdk + 19, // 5: confidence.flags.resolver.v1.TelemetryData.provider_init_rate:type_name -> confidence.flags.resolver.v1.TelemetryData.ProviderInitRate + 27, // 6: confidence.flags.resolver.v1.ClientInfo.sdk:type_name -> confidence.flags.resolver.v1.Sdk + 4, // 7: confidence.flags.resolver.v1.FlagAssigned.client_info:type_name -> confidence.flags.resolver.v1.ClientInfo + 21, // 8: confidence.flags.resolver.v1.FlagAssigned.flags:type_name -> confidence.flags.resolver.v1.FlagAssigned.AppliedFlag + 24, // 9: confidence.flags.resolver.v1.ClientResolveInfo.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance + 26, // 10: confidence.flags.resolver.v1.FlagResolveInfo.variant_resolve_info:type_name -> confidence.flags.resolver.v1.FlagResolveInfo.VariantResolveInfo + 15, // 11: confidence.flags.resolver.v1.WriteOperationsRequest.store_variant_op:type_name -> confidence.flags.resolver.v1.VariantData + 11, // 12: confidence.flags.resolver.v1.ReadOp.variant_read_op:type_name -> confidence.flags.resolver.v1.VariantReadOp + 12, // 13: confidence.flags.resolver.v1.ReadOp.inclusion_read_op:type_name -> confidence.flags.resolver.v1.InclusionReadOp + 13, // 14: confidence.flags.resolver.v1.ReadOperationsRequest.ops:type_name -> confidence.flags.resolver.v1.ReadOp + 15, // 15: confidence.flags.resolver.v1.ReadResult.variant_result:type_name -> confidence.flags.resolver.v1.VariantData + 16, // 16: confidence.flags.resolver.v1.ReadResult.inclusion_result:type_name -> confidence.flags.resolver.v1.InclusionData + 17, // 17: confidence.flags.resolver.v1.ReadOperationsResult.results:type_name -> confidence.flags.resolver.v1.ReadResult + 20, // 18: confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.labels:type_name -> confidence.flags.resolver.v1.TelemetryData.ProviderInitRate.LabelsEntry + 22, // 19: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.assignment_info:type_name -> confidence.flags.resolver.v1.FlagAssigned.AssignmentInfo + 23, // 20: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.default_assignment:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment + 6, // 21: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.fallthrough_assignments:type_name -> confidence.flags.resolver.v1.FallthroughAssignment + 28, // 22: confidence.flags.resolver.v1.FlagAssigned.AppliedFlag.apply_time:type_name -> google.protobuf.Timestamp + 0, // 23: confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.reason:type_name -> confidence.flags.resolver.v1.FlagAssigned.DefaultAssignment.DefaultAssignmentReason + 25, // 24: confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.schema:type_name -> confidence.flags.resolver.v1.ClientResolveInfo.EvaluationContextSchemaInstance.SchemaEntry + 1, // 25: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:input_type -> confidence.flags.resolver.v1.WriteFlagLogsRequest + 9, // 26: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:input_type -> confidence.flags.resolver.v1.WriteOperationsRequest + 14, // 27: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:input_type -> confidence.flags.resolver.v1.ReadOperationsRequest + 2, // 28: confidence.flags.resolver.v1.InternalFlagLoggerService.ClientWriteFlagLogs:output_type -> confidence.flags.resolver.v1.WriteFlagLogsResponse + 10, // 29: confidence.flags.resolver.v1.InternalFlagLoggerService.WriteMaterializedOperations:output_type -> confidence.flags.resolver.v1.WriteOperationsResult + 18, // 30: confidence.flags.resolver.v1.InternalFlagLoggerService.ReadMaterializedOperations:output_type -> confidence.flags.resolver.v1.ReadOperationsResult + 28, // [28:31] is the sub-list for method output_type + 25, // [25:28] is the sub-list for method input_type + 25, // [25:25] is the sub-list for extension type_name + 25, // [25:25] is the sub-list for extension extendee + 0, // [0:25] is the sub-list for field type_name } func init() { file_confidence_flags_resolver_v1_internal_api_proto_init() } @@ -1645,7 +1716,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[18].OneofWrappers = []any{ + file_confidence_flags_resolver_v1_internal_api_proto_msgTypes[20].OneofWrappers = []any{ (*FlagAssigned_AppliedFlag_AssignmentInfo)(nil), (*FlagAssigned_AppliedFlag_DefaultAssignment)(nil), } @@ -1655,7 +1726,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: 24, + NumMessages: 26, NumExtensions: 0, NumServices: 1, }, diff --git a/openfeature-provider/go/confidence/provider_builder.go b/openfeature-provider/go/confidence/provider_builder.go index 8337e42a..16eb2dfa 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" @@ -87,8 +88,11 @@ func NewProvider(ctx context.Context, config ProviderConfig) (*LocalResolverProv materializationStore = newRemoteMaterializationStore(resolverv1.NewInternalFlagLoggerServiceClient(conn), config.ClientSecret) } + initLabels := map[string]string{ + "encryption": strconv.FormatBool(config.EncryptionKey != ""), + } resolverSupplier := func(ctx context.Context, logSink lr.LogSink) lr.LocalResolver { - return lr.NewLocalResolverWithPoolSize(ctx, logSink, config.ResolverPoolSize) + return lr.NewLocalResolverWithLabels(ctx, logSink, config.ResolverPoolSize, initLabels) } resolverSupplierWithMaterialization := wrapResolverSupplierWithMaterializations(resolverSupplier, materializationStore) providerOpts := buildProviderOptions(config.StatePollInterval, config.LogPollInterval) 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 6ae099dd..30205f0c 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; @@ -158,11 +159,15 @@ public OpenFeatureLocalResolveProvider( clientSecret, config.getHttpClientFactory(), config.getEncryptionKey()); final var wasmFlagLogger = new GrpcWasmFlagLogger(clientSecret, config.getChannelFactory()); 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(() -> new WasmLocalResolver(flagLogger::write))); + () -> + new RecoveringResolver( + () -> new WasmLocalResolver(flagLogger::write, initLabels))); this.resolver = new MaterializingResolver(inner, materializationStore); } @@ -189,7 +194,9 @@ public OpenFeatureLocalResolveProvider( final LocalResolver inner = new PooledResolver( numInstances, - () -> new RecoveringResolver(() -> new WasmLocalResolver(wasmFlagLogger::write))); + () -> + new RecoveringResolver( + () -> new WasmLocalResolver(wasmFlagLogger::write, Map.of()))); this.resolver = new MaterializingResolver(inner, materializationStore); } 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 17c8bd1b..036d9451 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 @@ -18,10 +18,12 @@ import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessRequest; import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessResponse; import com.spotify.confidence.sdk.flags.resolver.v1.Sdk; +import com.spotify.confidence.sdk.flags.resolver.v1.TelemetryData; import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest; import com.spotify.confidence.sdk.wasm.Messages; import java.time.Instant; import java.util.List; +import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.atomic.AtomicInteger; @@ -49,6 +51,8 @@ class WasmLocalResolver implements LocalResolver { private final ExportFunction wasmMsgAlloc; private final ExportFunction wasmMsgFree; private final Consumer logSink; + private final Map initLabels; + private boolean firstFlush = true; // api private final ExportFunction wasmMsgGuestSetResolverState; @@ -60,8 +64,9 @@ class WasmLocalResolver implements LocalResolver { private final ExportFunction wasmMsgGuestPrometheusSnapshot; private final ReentrantLock lock = new ReentrantLock(); - public WasmLocalResolver(Consumer logSink) { + public WasmLocalResolver(Consumer logSink, Map initLabels) { this.logSink = logSink; + this.initLabels = initLabels; this.instanceId = String.valueOf(INSTANCE_COUNTER.getAndIncrement()); instance = Instance.builder(ConfidenceResolverModule.load()) @@ -192,7 +197,21 @@ public void flushAllLogs() { final var voidRequest = Messages.Void.getDefaultInstance(); final var reqPtr = transferRequest(voidRequest); final var respPtr = (int) wasmMsgBoundedFlushLogs.apply(reqPtr)[0]; - final var request = consumeResponse(respPtr, WriteFlagLogsRequest::parseFrom); + var request = consumeResponse(respPtr, WriteFlagLogsRequest::parseFrom); + if (firstFlush) { + firstFlush = false; + request = + request.toBuilder() + .setTelemetryData( + request.getTelemetryData().toBuilder() + .addProviderInitRate( + TelemetryData.ProviderInitRate.newBuilder() + .setCount(1) + .putAllLabels(initLabels) + .build()) + .build()) + .build(); + } if (!isEmptyLogRequest(request)) { logSink.accept(request); } @@ -279,9 +298,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); } @@ -350,4 +368,5 @@ private interface ParserFn { T apply(byte[] data) throws InvalidProtocolBufferException; } + } diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ApplyFlagsErrorHandlingTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ApplyFlagsErrorHandlingTest.java index f463d4c6..632fc38f 100644 --- a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ApplyFlagsErrorHandlingTest.java +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ApplyFlagsErrorHandlingTest.java @@ -25,6 +25,7 @@ import java.time.Instant; import java.util.ArrayList; import java.util.List; +import java.util.Map; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -154,7 +155,7 @@ private static Timestamp toTs(Instant t) { @Test void applyWithoutSendTime_isSwallowed_andRecordsNoFlagAssigned() { final List captured = new ArrayList<>(); - final var resolver = new WasmLocalResolver(captured::add); + final var resolver = new WasmLocalResolver(captured::add, Map.of()); resolver.setResolverState(buildState(), ACCOUNT, null); final var resolveResp = resolveWithApplyFalse(resolver); @@ -185,7 +186,7 @@ void applyWithoutSendTime_isSwallowed_andRecordsNoFlagAssigned() { @Test void applyWithSendTime_recordsOneFlagAssigned() { final List captured = new ArrayList<>(); - final var resolver = new WasmLocalResolver(captured::add); + final var resolver = new WasmLocalResolver(captured::add, Map.of()); resolver.setResolverState(buildState(), ACCOUNT, null); final var resolveResp = resolveWithApplyFalse(resolver); @@ -233,7 +234,7 @@ void applyWithSendTime_recordsOneFlagAssigned() { @Test void applyAfterStateRotated_secretRemoved_isSwallowedAndLogged() { final List captured = new ArrayList<>(); - final var resolver = new WasmLocalResolver(captured::add); + final var resolver = new WasmLocalResolver(captured::add, Map.of()); resolver.setResolverState(buildState(SECRET), ACCOUNT, null); final var resolveResp = resolveWithApplyFalse(resolver); diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ResolveTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ResolveTest.java index 70f93095..cbd7f549 100644 --- a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ResolveTest.java +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/ResolveTest.java @@ -36,7 +36,7 @@ class ResolveTest { private final LocalResolver resolver; public ResolveTest() { - resolver = new WasmLocalResolver(request -> {}); + resolver = new WasmLocalResolver(request -> {}, Map.of()); } @BeforeEach diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmMemoryLeakTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmMemoryLeakTest.java index 7982a1e6..5f28c179 100644 --- a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmMemoryLeakTest.java +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmMemoryLeakTest.java @@ -10,6 +10,7 @@ import com.spotify.confidence.sdk.flags.resolver.v1.ResolveProcessRequest; import java.lang.reflect.Field; import java.util.List; +import java.util.Map; import org.junit.jupiter.api.Test; /** @@ -32,7 +33,7 @@ private static int getWasmMemoryPages(WasmLocalResolver resolver) { @Test void wasmMemoryStableOnRepeatedResolveCalls() { - WasmLocalResolver resolver = new WasmLocalResolver(request -> {}); + WasmLocalResolver resolver = new WasmLocalResolver(request -> {}, Map.of()); resolver.setResolverState(ResolveTest.exampleStateBytes, "account", null); ResolveProcessRequest request = diff --git a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmResolveApiFlushCloseRaceTest.java b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmResolveApiFlushCloseRaceTest.java index cdd63258..933a6655 100644 --- a/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmResolveApiFlushCloseRaceTest.java +++ b/openfeature-provider/java/src/test/java/com/spotify/confidence/sdk/WasmResolveApiFlushCloseRaceTest.java @@ -9,6 +9,7 @@ import com.spotify.confidence.sdk.flags.resolver.v1.Sdk; import com.spotify.confidence.sdk.flags.resolver.v1.SdkId; import com.spotify.confidence.sdk.flags.resolver.v1.WriteFlagLogsRequest; +import java.util.Map; import java.util.concurrent.CountDownLatch; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; @@ -44,7 +45,7 @@ void concurrentFlushAndCloseShouldNotLoseAssignments() throws Exception { for (int i = 0; i < iterations; i++) { final var logger = new CapturingWasmFlagLogger(); - final var resolver = new WasmLocalResolver(logger::write); + final var resolver = new WasmLocalResolver(logger::write, Map.of()); resolver.setResolverState(resolverState, accountId, null); // Resolve a flag to create a flag assignment in the WASM buffer diff --git a/openfeature-provider/js/demo.mjs b/openfeature-provider/js/demo.mjs new file mode 100644 index 00000000..56748ec8 --- /dev/null +++ b/openfeature-provider/js/demo.mjs @@ -0,0 +1,246 @@ +#!/usr/bin/env node +/** + * Demo: run the local provider, resolve flags, and inspect telemetry. + * + * Usage: + * node demo.mjs --secret [--encryption-key ] [--flag ] [--duration ] + * + * Examples: + * node demo.mjs --secret abc123 --flag web-sdk-e2e-flag --duration 60 + * node demo.mjs --secret abc123 --encryption-key deadbeef... --flag web-sdk-e2e-flag --duration 300 + */ + +import { OpenFeature } from '@openfeature/server-sdk'; +import { createConfidenceServerProvider } from './dist/index.node.js'; + +function parseArgs() { + const args = process.argv.slice(2); + const opts = { flags: [], duration: 60 }; + for (let i = 0; i < args.length; i++) { + switch (args[i]) { + case '--secret': + opts.secret = args[++i]; + break; + case '--encryption-key': + opts.encryptionKey = args[++i]; + break; + case '--flag': + opts.flags.push(args[++i]); + break; + case '--duration': + opts.duration = parseInt(args[++i], 10); + break; + } + } + if (!opts.secret) { + console.error('Usage: node demo.mjs --secret [--encryption-key ] [--flag ] [--duration ]'); + process.exit(1); + } + if (opts.flags.length === 0) opts.flags.push('web-sdk-e2e-flag'); + return opts; +} + +const opts = parseArgs(); + +console.log(`=== Confidence Provider Demo ===`); +console.log(` encryption: ${opts.encryptionKey ? 'enabled' : 'disabled'}`); +console.log(` flags: ${opts.flags.join(', ')}`); +console.log(` duration: ${opts.duration}s`); +console.log(); + +// Minimal protobuf varint + length-delimited decoder for TelemetryData inspection +function decodeProviderInitRate(buf) { + // Walk WriteFlagLogsRequest looking for field 2 (telemetry_data), + // then inside that look for field 9 (provider_init_rate) + let pos = 0; + while (pos < buf.length) { + const tag = buf[pos++]; + const fieldNum = tag >> 3; + const wireType = tag & 0x7; + if (wireType === 2) { + // length-delimited + let len = 0, + shift = 0; + while (buf[pos] & 0x80) { + len |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + len |= buf[pos++] << shift; + if (fieldNum === 2) { + // telemetry_data — recurse into it + const inner = buf.subarray(pos, pos + len); + const result = findInitRate(inner); + if (result) return result; + } + pos += len; + } else if (wireType === 0) { + // varint + while (buf[pos++] & 0x80) {} + } else { + break; + } + } + return null; +} + +function findInitRate(buf) { + let pos = 0; + while (pos < buf.length) { + const tag = buf[pos++]; + const fieldNum = tag >> 3; + const wireType = tag & 0x7; + if (wireType === 2) { + let len = 0, + shift = 0; + while (buf[pos] & 0x80) { + len |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + len |= buf[pos++] << shift; + if (fieldNum === 9) { + // provider_init_rate + return parseInitRate(buf.subarray(pos, pos + len)); + } + pos += len; + } else if (wireType === 0) { + while (buf[pos++] & 0x80) {} + } else { + break; + } + } + return null; +} + +function parseInitRate(buf) { + let count = 0; + const labels = {}; + let pos = 0; + while (pos < buf.length) { + const tag = buf[pos++]; + const fieldNum = tag >> 3; + const wireType = tag & 0x7; + if (wireType === 0 && fieldNum === 1) { + // count + count = 0; + let shift = 0; + while (buf[pos] & 0x80) { + count |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + count |= buf[pos++] << shift; + } else if (wireType === 2 && fieldNum === 3) { + // labels map entry + let len = 0, + shift = 0; + while (buf[pos] & 0x80) { + len |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + len |= buf[pos++] << shift; + const entry = parseMapEntry(buf.subarray(pos, pos + len)); + if (entry) labels[entry.key] = entry.value; + pos += len; + } else if (wireType === 2) { + let len = 0, + shift = 0; + while (buf[pos] & 0x80) { + len |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + len |= buf[pos++] << shift; + pos += len; + } else if (wireType === 0) { + while (buf[pos++] & 0x80) {} + } else { + break; + } + } + return { count, labels }; +} + +function parseMapEntry(buf) { + let key = '', + value = ''; + let pos = 0; + while (pos < buf.length) { + const tag = buf[pos++]; + const fieldNum = tag >> 3; + const wireType = tag & 0x7; + if (wireType === 2) { + let len = 0, + shift = 0; + while (buf[pos] & 0x80) { + len |= (buf[pos++] & 0x7f) << shift; + shift += 7; + } + len |= buf[pos++] << shift; + const str = new TextDecoder().decode(buf.subarray(pos, pos + len)); + if (fieldNum === 1) key = str; + else if (fieldNum === 2) value = str; + pos += len; + } else { + break; + } + } + return { key, value }; +} + +let flushCount = 0; + +// Intercept outgoing fetches to log telemetry +const originalFetch = globalThis.fetch; +const wrappedFetch = async (url, init) => { + let bodyBytes = null; + if (typeof url === 'string' && url.includes('clientFlagLogs') && init?.body) { + bodyBytes = new Uint8Array(init.body); + } + const resp = await originalFetch(url, init); + if (bodyBytes) { + flushCount++; + const initRate = decodeProviderInitRate(bodyBytes); + if (initRate) { + console.log( + `\n[flush #${flushCount}] → ${resp.status} ★ provider_init_rate: count=${ + initRate.count + }, labels=${JSON.stringify(initRate.labels)}`, + ); + } else { + console.log(`\n[flush #${flushCount}] → ${resp.status} (no provider_init_rate)`); + } + } + return resp; +}; + +const provider = createConfidenceServerProvider({ + flagClientSecret: opts.secret, + encryptionKey: opts.encryptionKey, + flushInterval: 10_000, + fetch: wrappedFetch, +}); + +console.log('Initializing provider...'); +await OpenFeature.setProviderAndWait(provider); +console.log('Provider ready.\n'); + +const client = OpenFeature.getClient(); +const endTime = Date.now() + opts.duration * 1000; +let resolveCount = 0; + +const ctx = { targetingKey: 'demo-user', sticky: false }; + +while (Date.now() < endTime) { + for (const flag of opts.flags) { + const result = client.getBooleanDetails(`${flag}.bool`, false, ctx); + resolveCount++; + } + + const remaining = Math.ceil((endTime - Date.now()) / 1000); + process.stdout.write(`\r resolves: ${resolveCount} | remaining: ${remaining}s `); + + await new Promise(r => setTimeout(r, 2000)); +} + +console.log(`\n\nDone. Total resolves: ${resolveCount}`); +console.log('Shutting down (final flush)...'); +await OpenFeature.close(); +console.log('Bye.'); diff --git a/openfeature-provider/js/proto/test-only.proto b/openfeature-provider/js/proto/test-only.proto index 99f26668..e5d59fef 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.ts b/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts index 7caa249a..59c3e356 100644 --- a/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts +++ b/openfeature-provider/js/src/ConfidenceServerProviderLocal.ts @@ -17,6 +17,7 @@ import { } from './materialization'; import { SetResolverStateRequest } from './proto/confidence/wasm/messages'; import { ClientResolverState } from './proto/confidence/flags/admin/v1/resolver'; +import { WriteFlagLogsRequest } from './proto/confidence/flags/resolver/v1/internal_api'; import FlagBundleType, * as FlagBundle from './flag-bundle'; import { ErrorCode, ResolutionDetails } from './types'; @@ -64,6 +65,8 @@ export class ConfidenceServerProviderLocal implements Provider { private readonly stateUpdateInterval: number; private readonly flushInterval: number; private readonly materializationStore: MaterializationStore | null; + private readonly initLabels: Record; + private firstFlush = true; private stateEtag: string | null = null; private get resolver(): LocalResolver { @@ -142,6 +145,7 @@ export class ConfidenceServerProviderLocal implements Provider { } else { this.materializationStore = null; } + this.initLabels = { encryption: options.encryptionKey ? 'true' : 'false' }; } async initialize(context?: EvaluationContext): Promise { @@ -342,8 +346,17 @@ 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) { + if (this.firstFlush) { + this.firstFlush = false; + const decoded = WriteFlagLogsRequest.decode(writeFlagLogRequest); + if (!decoded.telemetryData) { + decoded.telemetryData = { providerInitRate: [] }; + } + decoded.telemetryData!.providerInitRate = [{ count: 1, labels: this.initLabels }]; + writeFlagLogRequest = WriteFlagLogsRequest.encode(decoded).finish(); + } await this.sendFlagLogs(writeFlagLogRequest, signal); } } 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 ad5e1019..c3105a43 100644 --- a/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto +++ b/openfeature-provider/proto/confidence/flags/resolver/v1/internal_api.proto @@ -44,6 +44,14 @@ message WriteFlagLogsResponse {} message TelemetryData { // Information about the SDK/provider Sdk sdk = 2; + + 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/provider.py b/openfeature-provider/python/src/confidence/provider.py index 990d6f0c..5575672a 100644 --- a/openfeature-provider/python/src/confidence/provider.py +++ b/openfeature-provider/python/src/confidence/provider.py @@ -32,7 +32,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__ @@ -152,6 +156,10 @@ 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._first_flush = True self._state_poll_interval = state_poll_interval self._log_poll_interval = log_poll_interval self._assign_poll_interval = assign_poll_interval @@ -300,7 +308,7 @@ def shutdown(self) -> None: # Flush final logs if self._resolver is not None: try: - log_data = self._resolver.flush_logs() + log_data = self._append_init(self._resolver.flush_logs()) if log_data and self._flag_logger is not None: self._flag_logger.write(log_data) except Exception as e: @@ -744,6 +752,19 @@ def _write_materializations( except Exception as e: logger.error("Failed to write materializations: %s", e) + def _append_init(self, log_data: bytes) -> bytes: + if self._first_flush and log_data: + self._first_flush = False + request = internal_api_pb2.WriteFlagLogsRequest() + request.ParseFromString(log_data) + 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 + return request.SerializeToString() + self._first_flush = False + return log_data + def _flush_assigned(self) -> None: """Flush assigned logs.""" if self._resolver is None or self._flag_logger is None: @@ -804,6 +825,7 @@ def _state_poll_loop(self) -> None: with self._resolver_lock: flushed_logs = self._resolver.flush_logs() self._resolver.set_resolver_state(state, account_id, sdk) + flushed_logs = self._append_init(flushed_logs) if flushed_logs and self._flag_logger is not None: self._flag_logger.write(flushed_logs) logger.debug("Resolver state updated") @@ -831,6 +853,7 @@ def _log_flush_loop(self) -> None: try: with self._resolver_lock: log_data = self._resolver.flush_logs() + log_data = self._append_init(log_data) if log_data and self._flag_logger is not None: self._flag_logger.write(log_data) except Exception as e: diff --git a/openfeature-provider/rust/AGENTS.md b/openfeature-provider/rust/AGENTS.md new file mode 120000 index 00000000..681311eb --- /dev/null +++ b/openfeature-provider/rust/AGENTS.md @@ -0,0 +1 @@ +CLAUDE.md \ No newline at end of file diff --git a/openfeature-provider/rust/src/logger.rs b/openfeature-provider/rust/src/logger.rs index 18615d87..e5a856ed 100644 --- a/openfeature-provider/rust/src/logger.rs +++ b/openfeature-provider/rust/src/logger.rs @@ -1,10 +1,15 @@ //! Log management for sending flag logs to the Confidence API. +use std::collections::BTreeMap; +use std::sync::atomic::{AtomicBool, Ordering}; + use prost::Message; use reqwest_middleware::ClientWithMiddleware; 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; @@ -64,14 +69,23 @@ impl LogSender { pub struct LogManager { sender: LogSender, sdk: Sdk, + init_labels: BTreeMap, + first_flush: AtomicBool, } 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) -> Self { + pub fn new( + client: ClientWithMiddleware, + client_secret: String, + sdk: Sdk, + init_labels: BTreeMap, + ) -> Self { Self { sender: LogSender::new(client, client_secret), sdk, + init_labels, + first_flush: AtomicBool::new(true), } } @@ -86,6 +100,12 @@ impl LogManager { let mut td = TELEMETRY.delta_snapshot(&LAST_FLUSHED); td.sdk = Some(self.sdk.clone()); + if self.first_flush.swap(false, Ordering::Relaxed) { + td.provider_init_rate.push(ProviderInitRate { + count: 1, + labels: self.init_labels.clone(), + }); + } request.telemetry_data = Some(td); let encoded = request.encode_to_vec(); diff --git a/openfeature-provider/rust/src/provider.rs b/openfeature-provider/rust/src/provider.rs index 20e13d77..3ebf252a 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}; @@ -168,16 +168,20 @@ 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(), Some(sdk.clone()), options.encryption_key, )); + 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, + init_labels, )); // Create materialization store if configured