diff --git a/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-raw-metrics-flow.yaml b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-raw-metrics-flow.yaml new file mode 100644 index 00000000..aac37317 --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-raw-metrics-flow.yaml @@ -0,0 +1,29 @@ +# OTLP Raw Metrics Flow - Local testing +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:53217 + http: + endpoint: 0.0.0.0:53218 + +processors: + batch: + send_batch_size: 5000 + +exporters: + otlphttp: + endpoint: http://localhost:4318 + encoding: proto + compression: gzip +service: + pipelines: + metrics: + receivers: [otlp] + processors: [batch] + exporters: [otlphttp] + telemetry: + logs: + level: debug + resource: + service.name: countminsketchcol-raw-metrics-flow-local diff --git a/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml new file mode 100644 index 00000000..71c059e9 --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-batch.yaml @@ -0,0 +1,35 @@ +# OTLP Sketch Payload Flow - Batch mode - Local testing +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:53217 + http: + endpoint: 0.0.0.0:53218 + +processors: + countmin: + mode: batch + metric_name: "cms_sketch" + rows: 5 + columns: 2000 + transmit_sketch: true + group_by: [] + drop_original: true + +exporters: + otlphttp: + endpoint: http://localhost:4318 + encoding: proto + compression: gzip +service: + pipelines: + metrics: + receivers: [otlp] + processors: [countmin] + exporters: [otlphttp] + telemetry: + logs: + level: debug + resource: + service.name: countminsketchcol-sketch-payload-flow-batch-local diff --git a/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-window.yaml b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-window.yaml new file mode 100644 index 00000000..2dfe6be2 --- /dev/null +++ b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/asap_query_config/config-otlp-sketch-payload-flow-window.yaml @@ -0,0 +1,38 @@ +# OTLP Sketch Payload Flow - Window mode - Local testing +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:53217 + http: + endpoint: 0.0.0.0:53218 + +processors: + countmin: + mode: window + metric_name: "cms_sketch" + rows: 5 + columns: 2000 + transmit_sketch: true + group_by: [] + drop_original: true + window_interval: 60s + +exporters: + otlphttp: + endpoint: http://localhost:4318 + encoding: proto + compression: gzip +service: + pipelines: + metrics: + receivers: [otlp] + processors: [countmin] + exporters: [otlphttp] + telemetry: + metrics: + level: basic + logs: + level: debug + resource: + service.name: countminsketchcol-otlp-sketch-payload-flow-local diff --git a/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/builder-config.yaml b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/builder-config.yaml index c5df90e0..d7b4bcb5 100644 --- a/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/builder-config.yaml +++ b/opentelemetry-collector-contrib-patch/cmd/countminsketchcol/builder-config.yaml @@ -280,4 +280,5 @@ replaces: # - github.com/open-telemetry/opentelemetry-collector-contrib/processor/countminsketchprocessor => /mnt/78D8516BD8512920/GARUDA_ACE/FROOT-LAB/asap-internal/DataCollector/opentelemetry-collector-contrib/processor/countminsketchprocessor - github.com/open-telemetry/opentelemetry-collector-contrib/processor/countminsketchprocessor => ../../../processor/countminsketchprocessor - go.opentelemetry.io/collector/pdata => ../../../../opentelemetry-collector/pdata + # - github.com/ProjectASAP/sketchlib-go => ../../../../../sketchlib-go - go.opentelemetry.io/collector/processor => ../../../../opentelemetry-collector/processor diff --git a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/go.mod b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/go.mod index 92b625fe..cd6f0450 100644 --- a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/go.mod +++ b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/go.mod @@ -3,7 +3,7 @@ module github.com/open-telemetry/opentelemetry-collector-contrib/processor/count go 1.24.0 require ( - github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5 + github.com/ProjectASAP/sketchlib-go v0.0.0-20260310013347-2c5db0c75da8 github.com/stretchr/testify v1.11.1 go.opentelemetry.io/collector/component v1.47.0 go.opentelemetry.io/collector/component/componenttest v0.141.0 @@ -46,13 +46,13 @@ require ( go.yaml.in/yaml/v2 v2.4.3 // indirect golang.org/x/sys v0.38.0 // indirect golang.org/x/text v0.30.0 // indirect - google.golang.org/protobuf v1.36.11 // indirect + google.golang.org/protobuf v1.36.10 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) replace go.opentelemetry.io/collector/pdata => ../../../opentelemetry-collector/pdata -replace go.opentelemetry.io/collector/processor => ../../../opentelemetry-collector/processor +// replace github.com/ProjectASAP/sketchlib-go => ../../../../sketchlib-go replace github.com/open-telemetry/opentelemetry-collector-contrib/pkg/pdatatest => ../../pkg/pdatatest diff --git a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go index f31bbf83..a6797e55 100644 --- a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go +++ b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor.go @@ -20,6 +20,7 @@ import ( "go.opentelemetry.io/collector/pdata/pmetric" "go.opentelemetry.io/collector/processor/selfmonitor" "go.uber.org/zap" + "google.golang.org/protobuf/proto" ) // builderPool recycles strings.Builder instances used in the hot @@ -634,11 +635,15 @@ func (p *windowedCountMinSketchProcessor) inboundDecodeCMS(aggregationKey string } func serializeCMS(s *cms.CountMinSketch) ([]byte, error) { - return s.SerializeProtoBytes() + env, err := s.SerializePortable() + if err != nil { + return nil, err + } + return proto.Marshal(env) } func deserializeCMS(data []byte) (*cms.CountMinSketch, error) { - return cms.DeserializeCountMinSketchFromProtoBytes(data) + return cms.DeserializeCountMinSketchFromBytes(data) } // cloneCMS returns a deep copy of s suitable for use as a delta snapshot. @@ -648,7 +653,7 @@ func cloneCMS(s *cms.CountMinSketch) *cms.CountMinSketch { if err != nil { return nil } - clone, err := cms.DeserializeCountMinSketchFromProtoBytes(data) + clone, err := cms.DeserializeCountMinSketchFromBytes(data) if err != nil { return nil } diff --git a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go index c868650a..cf31e918 100644 --- a/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go +++ b/opentelemetry-collector-contrib-patch/processor/countminsketchprocessor/processor_test.go @@ -180,7 +180,7 @@ func TestProcessor_TumblingWindow_Correctness(t *testing.T) { // ========================================== // 3. Verify Binary Payload (Gob Decode) // ========================================== - payloadVal, ok := dps2[0].Attributes().Get("sketch_payload") + payloadVal, ok := dps2[0].Attributes().Get("cms.sketch_payload") require.True(t, ok, "Sketch payload must exist in attributes") rawBytes := payloadVal.Bytes().AsRaw() @@ -301,7 +301,7 @@ func TestBatchModeQueryMetricsWhenTransmitSketchDisabled(t *testing.T) { dps := getAllDataPoints(out) require.Len(t, dps, 1) assert.Equal(t, 3.0, dps[0].DoubleValue()) - _, hasPayload := dps[0].Attributes().Get("sketch_payload") + _, hasPayload := dps[0].Attributes().Get("cms.sketch_payload") assert.False(t, hasPayload) }