@@ -12,6 +12,7 @@ import (
1212 "encoding/json"
1313 "fmt"
1414 "io"
15+ "maps"
1516 "net/http"
1617 "strings"
1718 "sync"
@@ -442,11 +443,6 @@ func sendTracesComputeTopLevelBySpanKind(t *testing.T, endpoint string) {
442443}
443444
444445func TestIntegrationLogs (t * testing.T ) {
445- require .NoError (t , featuregate .GlobalRegistry ().Set ("exporter.datadogexporter.metricexportserializerclient" , false ))
446- defer func () {
447- require .NoError (t , featuregate .GlobalRegistry ().Set ("exporter.datadogexporter.metricexportserializerclient" , true ))
448- }()
449-
450446 // 1. Set up mock Datadog server
451447 // See also https://github.com/DataDog/datadog-agent/blob/49c16e0d4deab396626238fa1d572b684475a53f/cmd/trace-agent/test/backend.go
452448 seriesRec := & testutil.HTTPRequestRecorderWithChan {Pattern : testutil .MetricV2Endpoint , ReqChan : make (chan []byte )}
@@ -499,11 +495,9 @@ func TestIntegrationLogs(t *testing.T) {
499495 case <- doneChannel :
500496 assert .Len (t , logsData , 5 )
501497 case metricsBytes := <- seriesRec .ReqChan :
502- var smap seriesSlice
503- gz := getGzipReader (t , metricsBytes )
504- dec := json .NewDecoder (gz )
505- assert .NoError (t , dec .Decode (& smap ))
506- for _ , s := range smap .Series {
498+ newSeries , err := seriesSliceFromAny (t , metricsBytes )
499+ require .NoError (t , err )
500+ for _ , s := range newSeries {
507501 if s .Metric == "otelcol_receiver_accepted_log_records" || s .Metric == "otelcol_exporter_sent_log_records" {
508502 metricMap .Series = append (metricMap .Series , s )
509503 }
@@ -732,11 +726,6 @@ func seriesFromAPIClient(t *testing.T, metricsBytes []byte, expectedMetrics map[
732726}
733727
734728func TestIntegrationInternalMetrics (t * testing.T ) {
735- t .Skip ("flaky test http://github.com/open-telemetry/opentelemetry-collector-contrib/issues/40056" )
736- require .NoError (t , featuregate .GlobalRegistry ().Set ("exporter.datadogexporter.metricexportserializerclient" , false ))
737- defer func () {
738- require .NoError (t , featuregate .GlobalRegistry ().Set ("exporter.datadogexporter.metricexportserializerclient" , true ))
739- }()
740729 expectedMetrics := map [string ]struct {}{
741730 // Datadog internal metrics on trace and stats writers
742731 "datadog.otlp_translator.resources.missing_source" : {},
@@ -807,19 +796,82 @@ func testIntegrationInternalMetrics(t *testing.T, expectedMetrics map[string]str
807796 case <- tracesRec .ReqChan :
808797 // Drain the channel, no need to look into the traces
809798 case metricsBytes := <- seriesRec .ReqChan :
810- var metrics seriesSlice
811- gz := getGzipReader (t , metricsBytes )
812- dec := json .NewDecoder (gz )
813- assert .NoError (t , dec .Decode (& metrics ))
814- for _ , s := range metrics .Series {
815- if _ , ok := expectedMetrics [s .Metric ]; ok {
816- metricMap [s .Metric ] = s
817- }
818- }
799+ // The mock server receives two different compression formats:
800+ // - gzip JSON (from the DD trace agent's internal metrics reporter, magic bytes 0x1f 0x8b)
801+ // - zlib-compressed protobuf (from the serializer exporter, magic bytes 0x78 0x??)
802+ newMetrics , err := seriesFromAny (t , metricsBytes , expectedMetrics )
803+ require .NoError (t , err )
804+ maps .Copy (metricMap , newMetrics )
819805 case <- time .After (60 * time .Second ):
820- t .Fatalf ("did not receive expected metrics after 1m" )
806+ t .Fatalf ("did not receive expected metrics after 1m, missing: %v" , missingMetrics (metricMap , expectedMetrics ))
807+ }
808+ }
809+ }
810+
811+ // seriesFromAny decodes a metric payload from the mock Datadog server, handling both
812+ // gzip-compressed JSON (DD trace agent format) and zlib-compressed protobuf (serializer exporter format).
813+ func seriesFromAny (t * testing.T , metricsBytes []byte , expectedMetrics map [string ]struct {}) (map [string ]series , error ) {
814+ if len (metricsBytes ) >= 2 && metricsBytes [0 ] == 0x1f && metricsBytes [1 ] == 0x8b {
815+ // gzip magic bytes: trace agent sends gzip-compressed JSON
816+ return seriesFromAPIClient (t , metricsBytes , expectedMetrics )
817+ }
818+ // zlib magic bytes (0x78 0x??): serializer exporter sends zlib-compressed protobuf
819+ return seriesFromSerializer (metricsBytes , expectedMetrics )
820+ }
821+
822+ // seriesSliceFromAny decodes all series from a metric payload, handling both
823+ // gzip-compressed JSON (DD trace agent format) and zlib-compressed protobuf (serializer exporter format).
824+ // Unlike seriesFromAny, it returns a slice so duplicate metric names across different resources are preserved.
825+ func seriesSliceFromAny (t * testing.T , metricsBytes []byte ) ([]series , error ) {
826+ if len (metricsBytes ) >= 2 && metricsBytes [0 ] == 0x1f && metricsBytes [1 ] == 0x8b {
827+ // gzip magic bytes: trace agent sends gzip-compressed JSON
828+ var smap seriesSlice
829+ gz := getGzipReader (t , metricsBytes )
830+ dec := json .NewDecoder (gz )
831+ if err := dec .Decode (& smap ); err != nil {
832+ return nil , err
833+ }
834+ return smap .Series , nil
835+ }
836+ // zlib magic bytes (0x78 0x??): serializer exporter sends zlib-compressed protobuf
837+ zr , err := zlib .NewReader (bytes .NewReader (metricsBytes ))
838+ if err != nil {
839+ return nil , err
840+ }
841+ pl := new (gogen.MetricPayload )
842+ b , err := io .ReadAll (zr )
843+ if err != nil {
844+ return nil , err
845+ }
846+ if err := pl .Unmarshal (b ); err != nil {
847+ return nil , err
848+ }
849+ var result []series
850+ for _ , s := range pl .GetSeries () {
851+ points := make ([]point , len (s .GetPoints ()))
852+ for i , p := range s .GetPoints () {
853+ points [i ] = point {
854+ Timestamp : int (p .GetTimestamp ()),
855+ Value : p .GetValue (),
856+ }
857+ }
858+ result = append (result , series {
859+ Metric : s .GetMetric (),
860+ Points : points ,
861+ Tags : s .GetTags (),
862+ })
863+ }
864+ return result , nil
865+ }
866+
867+ func missingMetrics (got map [string ]series , expected map [string ]struct {}) []string {
868+ var missing []string
869+ for k := range expected {
870+ if _ , ok := got [k ]; ! ok {
871+ missing = append (missing , k )
821872 }
822873 }
874+ return missing
823875}
824876
825877func TestIntegrationLogsHostMetadata (t * testing.T ) {
0 commit comments