diff --git a/deploy/fake-exporter/go.mod b/deploy/fake-exporter/go.mod index 1f3e5db6..01db0190 100644 --- a/deploy/fake-exporter/go.mod +++ b/deploy/fake-exporter/go.mod @@ -11,7 +11,6 @@ require ( ) require ( - github.com/DataDog/sketches-go v1.4.1 // indirect github.com/ProjectASAP/sketchlib-go v0.0.0-20260328221809-b24e56e64e94 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect @@ -19,7 +18,6 @@ require ( github.com/go-logr/stdr v1.2.2 // indirect github.com/google/uuid v1.6.0 // indirect github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 // indirect - github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect github.com/klauspost/cpuid/v2 v2.2.10 // indirect github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.66.1 // indirect @@ -32,7 +30,6 @@ require ( golang.org/x/net v0.50.0 // indirect golang.org/x/sys v0.41.0 // indirect golang.org/x/text v0.34.0 // indirect - google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57 // indirect google.golang.org/grpc v1.79.1 // indirect google.golang.org/protobuf v1.36.11 // indirect diff --git a/deploy/fake-exporter/go.sum b/deploy/fake-exporter/go.sum index e08b52c8..bb233c2b 100644 --- a/deploy/fake-exporter/go.sum +++ b/deploy/fake-exporter/go.sum @@ -1,10 +1,7 @@ -github.com/DataDog/sketches-go v1.4.1 h1:j5G6as+9FASM2qC36lvpvQAj9qsv/jUs3FtO8CwZNAY= -github.com/DataDog/sketches-go v1.4.1/go.mod h1:xJIXldczJyyjnbDop7ZZcLxJdV3+7Kra7H1KMgpgkLk= github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= @@ -12,30 +9,22 @@ github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= -github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= -github.com/golang/protobuf v1.5.2/go.mod h1:XVQd3VNwM+JqD3oG2Ue2ip4fOMUkwXdXDdiuN0vRsmY= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= -github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= -github.com/google/gofuzz v1.2.0 h1:xRy4A+RhZaiKjJ1bPfwQ8sedCA+YS2YcCHW6ec7JMi0= -github.com/google/gofuzz v1.2.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/gopacket v1.1.19 h1:ves8RnFZPGiFnTS0uPQStjwru6uO6h+nlr9j6fL7kF8= github.com/google/gopacket v1.1.19/go.mod h1:iJ8V8n6KS+z2U1A8pUwu8bW5SyEMkXJB8Yo/Vo+TKTo= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853 h1:cLN4IBkmkYZNnk7EAJ0BHIethd+J6LqxFNw5mSiI2bM= github.com/grafana/regexp v0.0.0-20250905093917-f7b3be9d1853/go.mod h1:+JKpmjMGhpgPL+rXZ5nsZieVzvarn86asRlBg4uNGnk= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 h1:HWRh5R2+9EifMyIHV7ZV+MIZqgz+PMpZ14Jynv3O2Zs= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0/go.mod h1:JfhWUomR1baixubs02l85lZYYOm7LV6om4ceouMv45c= github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= @@ -46,8 +35,6 @@ github.com/prometheus/prometheus v0.307.1 h1:Hh3kRMFn+xpQGLe/bR6qpUfW4GXQO0spuYe github.com/prometheus/prometheus v0.307.1/go.mod h1:/7YQG/jOLg7ktxGritmdkZvezE1fa6aWDj0MGDIZvcY= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= -github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= @@ -56,10 +43,6 @@ github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.41.0 h1:VO3BL6OZXRQ1yQc8W6EVfJzINeJ35BkiHx4MYfoQf44= -go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.41.0/go.mod h1:qRDnJ2nv3CQXMK2HUd9K9VtvedsPAce3S+/4LZHjX/s= -go.opentelemetry.io/proto/otlp v1.9.0 h1:l706jCMITVouPOqEnii2fIAuO3IVGBRPV5ICjceRb/A= -go.opentelemetry.io/proto/otlp v1.9.0/go.mod h1:xE+Cx5E/eEHw+ISFkwPLwCZefwVjY+pqKg1qcK03+/4= go.yaml.in/yaml/v2 v2.4.3 h1:6gvOSjQoTB3vt1l+CU+tSyi/HOjfOjRLJ4YwYZGwRO0= go.yaml.in/yaml/v2 v2.4.3/go.mod h1:zSxWcmIDjOzPXpjlTTbAsKokqkDNAVtZO0WOMiT90s8= golang.org/x/net v0.50.0 h1:ucWh9eiCGyDR3vtzso0WMQinm2Dnt8cFMuQa9K33J60= @@ -68,23 +51,16 @@ golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= -golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= -google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 h1:JLQynH/LBHfCTSbDWl+py8C+Rg/k1OVH3xfcaiANuF0= -google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57/go.mod h1:kSJwQxqmFXeo79zOmbrALdflXQeAYcUbgS7PbpMknCY= google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57 h1:mWPCjDEyshlQYzBpMNHaEof6UX1PmHcaUODUywQ0uac= google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57/go.mod h1:j9x/tPzZkyxcgEFkiKEEGxfvyumM01BEtsW8xzOahRQ= google.golang.org/grpc v1.79.1 h1:zGhSi45ODB9/p3VAawt9a+O/MULLl9dpizzNNpq7flY= google.golang.org/grpc v1.79.1/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= -google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= -google.golang.org/protobuf v1.26.0/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= -google.golang.org/protobuf v1.28.0/go.mod h1:HV8QOd/L58Z+nl8r43ehVNZIU/HEI6OcFqwMG9pJV4I= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/opentelemetry-go-patch/sdk/metric/go.mod b/opentelemetry-go-patch/sdk/metric/go.mod index b00c107d..9ece94b4 100644 --- a/opentelemetry-go-patch/sdk/metric/go.mod +++ b/opentelemetry-go-patch/sdk/metric/go.mod @@ -3,7 +3,6 @@ module go.opentelemetry.io/otel/sdk/metric go 1.24.0 require ( - github.com/DataDog/sketches-go v1.4.1 github.com/ProjectASAP/sketchlib-go v0.0.0-20260328221809-b24e56e64e94 github.com/go-logr/logr v1.4.3 github.com/go-logr/stdr v1.2.2 diff --git a/opentelemetry-go-patch/sdk/metric/go.sum b/opentelemetry-go-patch/sdk/metric/go.sum index c724bac2..8c5a8e1d 100644 --- a/opentelemetry-go-patch/sdk/metric/go.sum +++ b/opentelemetry-go-patch/sdk/metric/go.sum @@ -1,5 +1,3 @@ -github.com/DataDog/sketches-go v1.4.1 h1:j5G6as+9FASM2qC36lvpvQAj9qsv/jUs3FtO8CwZNAY= -github.com/DataDog/sketches-go v1.4.1/go.mod h1:xJIXldczJyyjnbDop7ZZcLxJdV3+7Kra7H1KMgpgkLk= github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5 h1:e1If3BxB2QRDt/57ho0PoxxIwouEvuaojy8m9u21IC0= github.com/ProjectASAP/sketchlib-go v0.0.0-20260321024028-d20a9f9151b5/go.mod h1:VmV0RYT6+rXpb1gVxTG3Z17sttm0nfj7RuesrL1oDTI= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= diff --git a/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch.go b/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch.go index 3cf46813..754d392f 100644 --- a/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch.go +++ b/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch.go @@ -8,8 +8,7 @@ import ( "sync" "time" - "github.com/DataDog/sketches-go/ddsketch" - "github.com/DataDog/sketches-go/ddsketch/pb/sketchpb" + ddsketch "github.com/ProjectASAP/sketchlib-go/sketches/DDSketch" "google.golang.org/protobuf/proto" "go.opentelemetry.io/otel" @@ -62,9 +61,11 @@ type ddSketchValues[N int64 | float64] struct { deltaTransmission bool // deltaThreshold is the minimum absolute bucket count change to include. deltaThreshold uint64 - // snapshots holds the proto-serialized snapshot of the last exported sketch - // per series, keyed by attribute.Distinct. Used only when deltaTransmission=true. - snapshots map[attribute.Distinct][]byte + // snapshots holds a clone of the last-exported sketch per series, keyed + // by attribute.Distinct. Used only when deltaTransmission=true so we can + // invoke sketchlib-go's ComputeDelta(prev, curr, threshold) on the next + // cumulative export. + snapshots map[attribute.Distinct]*ddsketch.DDSketch snapshotsMu sync.Mutex newRes func(attribute.Set) FilteredExemplarReservoir[N] @@ -95,7 +96,7 @@ func newDDSketchValues[N int64 | float64]( noSum: noSum, deltaTransmission: deltaTransmission, deltaThreshold: deltaThreshold, - snapshots: make(map[attribute.Distinct][]byte), + snapshots: make(map[attribute.Distinct]*ddsketch.DDSketch), newRes: r, limit: newLimiter[ddSketchSeries[N]](limit), values: make(map[attribute.Distinct]*ddSketchSeries[N]), @@ -106,16 +107,10 @@ func newDDSketchValues[N int64 | float64]( func (d *ddSketchValues[N]) newSeries(attr attribute.Set, value N) *ddSketchSeries[N] { series := d.seriesPool.Get().(*ddSketchSeries[N]) - if series.sketch != nil { - series.sketch.Clear() // reuse internal bucket storage - } else { - sk, err := ddsketch.NewDefaultDDSketch(d.accuracy) - if err != nil { - otel.Handle(err) - return nil - } - series.sketch = sk - } + // sketchlib-go's DDSketch has no in-place Reset(), so we always allocate a + // fresh sketch for a new series. The pool still amortizes the + // ddSketchSeries header allocation, which is the dominant cost. + series.sketch = ddsketch.NewDDSketch(d.accuracy) series.attrs = attr series.seriesID = 0 series.res = d.newRes(attr) // attr-dependent, always recreate @@ -156,10 +151,12 @@ func (d *ddSketchValues[N]) measure( } } - if err := series.sketch.Add(float64(value)); err != nil { - otel.Handle(err) - return - } + // sketchlib-go's DDSketch silently drops non-positive / NaN / Inf values. + // The DataDog sketch returned an explicit error for these; we mirror the + // permissive sketchlib-go behavior since downstream consumers + // (asap-precompute-{go,rs}, the agent decoder) all share the same + // invariant set. + series.sketch.Update(float64(value)) series.updateStats(value, !d.noMinMax) if !d.noSum { series.sum += value @@ -216,7 +213,12 @@ func (d *ddSketch[N]) delta( if series.count == 0 { continue } - if d.exportDataPoint(series, metricdata.DDSketchEncodingProto, nil, t, &dPts[i]) { + payload, encoding, err := serializeDDSketchFull(series.sketch) + if err != nil { + otel.Handle(err) + continue + } + if d.exportDataPoint(series, encoding, payload, t, &dPts[i]) { i++ } } @@ -228,6 +230,7 @@ func (d *ddSketch[N]) delta( series.attrs = attribute.Set{} series.seriesID = 0 series.res = nil // release exemplar reservoir (attr-dependent) + series.sketch = nil series.count = 0 series.sum = 0 series.min = 0 @@ -296,6 +299,7 @@ func (d *ddSketch[N]) cumulative( series.attrs = attribute.Set{} series.seriesID = 0 series.res = nil + series.sketch = nil series.count = 0 series.sum = 0 series.min = 0 @@ -310,48 +314,48 @@ func (d *ddSketch[N]) cumulative( return len(dPts) } -// payloadFor returns the serialized payload and encoding for a cumulative export. -// If deltaTransmission is enabled and a prior snapshot exists, it returns a -// sparse delta; otherwise it returns the full proto payload. +// payloadFor returns the serialized payload and encoding for a cumulative +// export. If deltaTransmission is enabled and a prior snapshot exists, it +// invokes sketchlib-go's ComputeDelta to emit a sparse delta; otherwise it +// emits a full SketchEnvelope-wrapped portable payload. +// +// The snapshot map holds *ddsketch.DDSketch clones (not their serialized +// bytes) so ComputeDelta can iterate buckets directly without re-decoding — +// this matches the path taken by KLL/HLL/CountSketch/CountMinSketch siblings +// and asap-precompute-go's DDSketchWrapper.ComputeDeltaAgainst. func (d *ddSketchValues[N]) payloadFor(key attribute.Distinct, sketch *ddsketch.DDSketch) ([]byte, metricdata.DDSketchEncoding, error) { - fullPayload, err := serializeDDSketch(sketch) - if err != nil { - return nil, metricdata.DDSketchEncodingProto, err - } - if !d.deltaTransmission { - return fullPayload, metricdata.DDSketchEncodingProto, nil + return serializeDDSketchFull(sketch) } d.snapshotsMu.Lock() - snapPayload, hasSnap := d.snapshots[key] + snap, hasSnap := d.snapshots[key] d.snapshotsMu.Unlock() - var ( - payload []byte - encoding metricdata.DDSketchEncoding - ) - if hasSnap && snapPayload != nil { - deltaPayload, deltaErr := ddSketchDeltaPayload(snapPayload, sketch, d.deltaThreshold) - if deltaErr != nil { - // Fall back to full on error. - payload = fullPayload - encoding = metricdata.DDSketchEncodingProto - } else { - payload = deltaPayload - encoding = metricdata.DDSketchEncodingProtoDelta + if hasSnap && snap != nil { + deltaPayload, deltaErr := ddsketch.ComputeDelta(snap, sketch, d.deltaThreshold) + if deltaErr == nil { + // Update snapshot to a clone of the current sketch; the receiver + // will fold this delta into its prior cumulative state, so on + // the next tick we want to delta against this same baseline. + d.snapshotsMu.Lock() + d.snapshots[key] = sketch.Clone() + d.snapshotsMu.Unlock() + return deltaPayload, metricdata.DDSketchEncodingProtoDelta, nil } - } else { - payload = fullPayload - encoding = metricdata.DDSketchEncodingProto + // Fall through to full on error — keeps the emit path always + // producing a valid payload, mirroring asap-precompute-go's + // DDSketchWrapper.ComputeDeltaAgainst contract. } - // Update snapshot with the current full payload. + full, encoding, err := serializeDDSketchFull(sketch) + if err != nil { + return nil, encoding, err + } d.snapshotsMu.Lock() - d.snapshots[key] = fullPayload + d.snapshots[key] = sketch.Clone() d.snapshotsMu.Unlock() - - return payload, encoding, nil + return full, encoding, nil } func (d *ddSketch[N]) exportDataPoint( @@ -361,16 +365,6 @@ func (d *ddSketch[N]) exportDataPoint( t time.Time, dest *metricdata.DDSketchDataPoint[N], ) bool { - // In delta() path payload is nil — serialize inline. - if payload == nil { - var err error - payload, err = serializeDDSketch(series.sketch) - if err != nil { - otel.Handle(err) - return false - } - } - dp := dest if series.seriesID != 0 { dp.SeriesID = series.seriesID @@ -395,65 +389,28 @@ func (d *ddSketch[N]) exportDataPoint( return true } -func serializeDDSketch(sk *ddsketch.DDSketch) ([]byte, error) { +// serializeDDSketchFull emits the canonical full-state portable wire format: +// a sketchlib-go SketchEnvelope wrapping a DDSketchState. The Producer and +// HashSpec fields are stripped to match what the asap-precompute-{go,rs} +// wrappers emit (see integration/parity/golden_test.go's per-sketch +// overlays); this preserves byte-parity with the cross-language fixtures. +func serializeDDSketchFull(sk *ddsketch.DDSketch) ([]byte, metricdata.DDSketchEncoding, error) { if sk == nil { - return nil, nil - } - return proto.Marshal(sk.ToProto()) -} - -// ddSketchDeltaPayload computes a sparse delta between a proto-serialized -// snapshot and the current sketch. Only buckets whose count changed by at -// least threshold are included. -func ddSketchDeltaPayload(snapPayload []byte, current *ddsketch.DDSketch, threshold uint64) ([]byte, error) { - var snap sketchpb.DDSketch - if err := proto.Unmarshal(snapPayload, &snap); err != nil { - return serializeDDSketch(current) - } - - curr := current.ToProto() - delta := &sketchpb.DDSketch{ - Mapping: curr.Mapping, - ZeroCount: curr.ZeroCount - snap.ZeroCount, - } - delta.PositiveValues = ddStoreDelta(snap.PositiveValues, curr.PositiveValues, float64(threshold)) - delta.NegativeValues = ddStoreDelta(snap.NegativeValues, curr.NegativeValues, float64(threshold)) - return proto.Marshal(delta) -} - -// ddStoreDelta returns a sparse Store with only buckets where |Δcount| ≥ threshold. -func ddStoreDelta(snap, curr *sketchpb.Store, threshold float64) *sketchpb.Store { - if curr == nil { - return nil + return nil, metricdata.DDSketchEncodingProto, nil } - snapCounts := ddStoreToMap(snap) - currCounts := ddStoreToMap(curr) - - out := &sketchpb.Store{BinCounts: make(map[int32]float64)} - for idx, cnt := range currCounts { - d := cnt - snapCounts[idx] - if d >= threshold || d <= -threshold { - out.BinCounts[idx] = d - } - } - if len(out.BinCounts) == 0 { - return nil - } - return out -} - -// ddStoreToMap converts a sketchpb.Store into a flat index→count map. -func ddStoreToMap(s *sketchpb.Store) map[int32]float64 { - m := make(map[int32]float64) - if s == nil { - return m - } - for idx, cnt := range s.BinCounts { - m[idx] += cnt + env, err := sk.SerializePortable() + if err != nil { + return nil, metricdata.DDSketchEncodingProto, err } - for i, cnt := range s.ContiguousBinCounts { - idx := s.ContiguousBinIndexOffset + int32(i) - m[idx] += cnt + // Strip producer / hash_spec so the payload bytes are stable across + // sketchlib-go version bumps and identical to the + // asap-precompute-{go,rs} wrapper outputs (see + // integration/parity/golden_test.go::TestGenerateGoldenFixtures). + env.Producer = nil + env.HashSpec = nil + bytes, err := proto.Marshal(env) + if err != nil { + return nil, metricdata.DDSketchEncodingProto, err } - return m + return bytes, metricdata.DDSketchEncodingProto, nil } diff --git a/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch_test.go b/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch_test.go index 3ed052e9..edf50649 100644 --- a/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch_test.go +++ b/opentelemetry-go-patch/sdk/metric/internal/aggregate/ddsketch_test.go @@ -12,7 +12,11 @@ import ( "testing" "time" + envpb "github.com/ProjectASAP/sketchlib-go/proto/sketch_envelope" + ddsketch "github.com/ProjectASAP/sketchlib-go/sketches/DDSketch" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "google.golang.org/protobuf/proto" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/sdk/metric/metricdata" @@ -445,3 +449,110 @@ func sampleProcessCPUSeconds() float64 { } return 0 } + +// TestDDSketchPayloadIsSketchlibPortableEnvelope confirms the post-#262 +// migration to sketchlib-go: the SDK aggregator's DDSketchDataPoint.Sketch +// payload is now a SketchEnvelope-wrapped DDSketchState (the wire format the +// agent processor and asap-precompute-{go,rs} consumers expect), not the +// pre-migration DataDog sketchpb.DDSketch shape that collided on field 1's +// wire type. The structural mismatch was diagnosed in PR #269. +func TestDDSketchPayloadIsSketchlibPortableEnvelope(t *testing.T) { + c := new(clock) + t.Cleanup(c.Register()) + + ctx := t.Context() + meas, comp := Builder[float64]{ + Temporality: metricdata.DeltaTemporality, + Filter: attrFltr, + AggregationLimit: 4, + }.DDSketch(testDDSketchAccuracy, false, false, false, 0) + + attrs := attribute.NewSet( + attribute.String("service", "checkout"), + attribute.String("region", "us-east-1"), + ) + for _, v := range []float64{1.5, 2.5, 3.5, 100.0, 99.0} { + meas(ctx, v, attrs) + } + + got := new(metricdata.Aggregation) + require.Equal(t, 1, comp(got)) + agg := (*got).(metricdata.DDSketch[float64]) + require.Len(t, agg.DataPoints, 1) + dp := agg.DataPoints[0] + require.Equal(t, metricdata.DDSketchEncodingProto, dp.Encoding) + require.NotEmpty(t, dp.Sketch) + + // Round-trip through sketchlib-go's portable envelope schema. This + // pins the SDK's wire format to the same shape decodeDDSketchEnvelope + // (in opentelemetry-collector-contrib-patch/processor/ddsketchprocessor) + // and asap-precompute-rs's DDSketchWrapper::decode_envelope expect. + var env envpb.SketchEnvelope + require.NoError(t, proto.Unmarshal(dp.Sketch, &env)) + state := env.GetDdsketch() + require.NotNil(t, state, "envelope must carry a DDSketchState variant") + require.InDelta(t, testDDSketchAccuracy, state.Alpha, 1e-12) + require.Equal(t, uint64(5), state.Count) + + // Reconstruct via the canonical NewFromState entrypoint and confirm + // quantiles agree within the accuracy bound (the actual value the + // agent will emit on the consume side). + recovered, err := ddsketch.NewFromState(state) + require.NoError(t, err) + q, ok := recovered.Quantile(0.99) + require.True(t, ok) + require.InDelta(t, 100.0, q, 100.0*testDDSketchAccuracy) +} + +// TestDDSketchDeltaEncodingViaComputeDelta confirms the cumulative-mode +// deltaTransmission path emits a DDSketchEncodingProtoDelta payload on the +// second tick (after a snapshot exists), and that the bytes round-trip +// through sketchlib-go's ApplyDelta. +func TestDDSketchDeltaEncodingViaComputeDelta(t *testing.T) { + c := new(clock) + t.Cleanup(c.Register()) + + ctx := t.Context() + meas, comp := Builder[float64]{ + Temporality: metricdata.CumulativeTemporality, + Filter: attrFltr, + AggregationLimit: 4, + }.DDSketch(testDDSketchAccuracy, false, false, true, 1) + + attrs := attribute.NewSet( + attribute.String("service", "checkout"), + attribute.String("region", "us-east-1"), + ) + for _, v := range []float64{1.5, 2.5, 3.5} { + meas(ctx, v, attrs) + } + + got := new(metricdata.Aggregation) + require.Equal(t, 1, comp(got)) + agg := (*got).(metricdata.DDSketch[float64]) + require.Len(t, agg.DataPoints, 1) + // First export is full state — there's no prior snapshot to delta against. + require.Equal(t, metricdata.DDSketchEncodingProto, agg.DataPoints[0].Encoding) + + // Second tick: add more values. With deltaTransmission=true and a + // snapshot now present, this must emit a delta payload. + for _, v := range []float64{10.0, 20.0, 30.0} { + meas(ctx, v, attrs) + } + require.Equal(t, 1, comp(got)) + agg = (*got).(metricdata.DDSketch[float64]) + require.Len(t, agg.DataPoints, 1) + dp := agg.DataPoints[0] + require.Equal(t, metricdata.DDSketchEncodingProtoDelta, dp.Encoding) + require.NotEmpty(t, dp.Sketch) + + // The delta bytes must apply cleanly onto a sketch holding the prior + // state — this is exactly what the agent's cumulative-state path does + // when it receives a delta envelope from a downstream SDK. + prior := ddsketch.NewDDSketch(testDDSketchAccuracy) + for _, v := range []float64{1.5, 2.5, 3.5} { + prior.Update(v) + } + require.NoError(t, ddsketch.ApplyDelta(prior, dp.Sketch)) + assert.Equal(t, uint64(6), prior.GetCount()) +} diff --git a/opentelemetry-go-patch/sdk/metric/metricdata/data.go b/opentelemetry-go-patch/sdk/metric/metricdata/data.go index c89098d2..bc72f08a 100644 --- a/opentelemetry-go-patch/sdk/metric/metricdata/data.go +++ b/opentelemetry-go-patch/sdk/metric/metricdata/data.go @@ -308,7 +308,11 @@ type QuantileValue struct { Value float64 } -// DDSketch represents distributions encoded as DataDog DDSketch payloads. +// DDSketch represents distributions encoded as DDSketch payloads. The +// SDK aggregator emits the sketchlib-go portable wire format +// (SketchEnvelope wrapping a DDSketchState); this is the format every +// downstream consumer (asap-precompute-{go,rs}, the agent's +// ddsketchprocessor decoder, the ASAPQuery backend) decodes against. type DDSketch[N int64 | float64] struct { // DataPoints are the individual aggregated measurements with unique // attributes. @@ -324,12 +328,20 @@ func (DDSketch[N]) privateAggregation() {} type DDSketchEncoding string const ( - // DDSketchEncodingProto indicates the sketch bytes are encoded as the - // serialization of github.com/DataDog/sketches-go/ddsketch/pb/sketchpb.DDSketch. + // DDSketchEncodingProto indicates the sketch bytes are a full-state + // proto serialization of sketchlib-go's + // `proto/sketch_envelope.SketchEnvelope` carrying a `DDSketchState` + // in its `ddsketch` oneof variant. This is the wire format the SDK + // aggregator emits via DDSketch.SerializePortable + proto.Marshal, + // and the format the agent / backend / asap-precompute decoders + // consume. DDSketchEncodingProto DDSketchEncoding = "ddsketch_proto" - // DDSketchEncodingProtoDelta indicates a sparse delta payload: only buckets - // whose count changed by at least DeltaThreshold since the previous export - // are included. Encoded as sketchpb.DDSketch proto. + // DDSketchEncodingProtoDelta indicates a sparse delta payload: only + // buckets whose count changed by at least DeltaThreshold since the + // previous export are included. Encoded as sketchlib-go's + // `proto/ddsketch.DDSketchDelta` (the bare delta message; not + // envelope-wrapped) — the same shape the receiver applies via + // sketchlib-go's ApplyDelta. DDSketchEncodingProtoDelta DDSketchEncoding = "ddsketch_proto_delta" )