From dbcc2b708637b625aa09a1ea7ef5a9d01b7f7ff7 Mon Sep 17 00:00:00 2001 From: Ian Oberst Date: Wed, 16 Sep 2026 21:44:30 -0700 Subject: [PATCH 1/2] Size Datadog buffers to available metrics Bound initial raw and offset capacity by available candidates and series and cursor capacity by remaining source values. Preserve exact payload limits, bounded conversion windows, and checkpoint delivery semantics. --- sink/datadog_api.go | 10 +- sink/datadog_capacity_test.go | 333 ++++++++++++++++++++++++++++++++++ sink/datadog_stream.go | 15 +- 3 files changed, 353 insertions(+), 5 deletions(-) create mode 100644 sink/datadog_capacity_test.go diff --git a/sink/datadog_api.go b/sink/datadog_api.go index bec5a61..93ac6b7 100644 --- a/sink/datadog_api.go +++ b/sink/datadog_api.go @@ -182,16 +182,20 @@ func prepareDatadogPayload(ctx context.Context, series []datadogV2.MetricSeries, } var raw bytes.Buffer + candidates := min(maxSeries, len(series)-start) grow := rawTarget - if estimated := maxSeries * 512; estimated < grow { - grow = estimated + // This is only an initial capacity hint: long series can still grow the + // buffer, and the exact byte checks below remain authoritative. Compare + // before multiplying to keep the estimate bounded without overflow. + if candidates <= rawTarget/512 { + grow = candidates * 512 } if grow > 0 { raw.Grow(grow) } raw.Write(datadogPayloadPrefix) - offsets := make([]int, 0, maxSeries) + offsets := make([]int, 0, candidates) end := start for end < len(series) && len(offsets) < maxSeries { select { diff --git a/sink/datadog_capacity_test.go b/sink/datadog_capacity_test.go new file mode 100644 index 0000000..fd2c23d --- /dev/null +++ b/sink/datadog_capacity_test.go @@ -0,0 +1,333 @@ +package sink + +import ( + "context" + "crypto/sha256" + "errors" + "fmt" + "io" + "net/http" + "strings" + "sync" + "testing" + "time" + + "github.com/DataDog/datadog-api-client-go/v2/api/datadog" + "github.com/cashapp/blip/v2" + "github.com/cashapp/blip/v2/test/mock" + "github.com/stretchr/testify/require" +) + +var datadogCapacityCounts = []int{0, 1, 10, 61, 383, 1000, 9999, 10000, 10001, 100000} + +func capacityMetrics(n int) *blip.Metrics { + m := getBlipMetrics(n, blip.GAUGE, 1, false) + m.Begin = time.Unix(1700000000, 0) + return m +} + +// Exercise the real APIClient.CallAPI body copy without decoding or retaining +// requests in the measured path. HTTP transport costs are measured separately. +func capacityAPISender() *Datadog { + cfg := datadog.NewConfiguration() + cfg.HTTPClient = &http.Client{Transport: &mock.Transport{ + RoundTripFunc: func(r *http.Request) (*http.Response, error) { + _, err := io.Copy(io.Discard, r.Body) + r.Body.Close() + return &http.Response{StatusCode: http.StatusAccepted, + Header: make(http.Header), Body: io.NopCloser(strings.NewReader(`{"errors":[]}`))}, err + }, + }} + return newTestDatadogSender(&datadogAPISubmitter{client: datadog.NewAPIClient(cfg), apiKey: "test"}) +} + +func BenchmarkDatadogCapacity(b *testing.B) { + for _, n := range datadogCapacityCounts { + b.Run(fmt.Sprintf("collect/%d", n), func(b *testing.B) { + m, s := capacityMetrics(n), newTestDatadogSender(nil) + state, _ := newDatadogSendCheckpoint(m, nil) + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + _, _, _, _, err := s.collectDatadogSeries(context.Background(), m, state.domains, datadogMetricCursor{}, datadogMaxSeriesPerPayload) + if err != nil { + b.Fatal(err) + } + } + }) + b.Run(fmt.Sprintf("send/%d", n), func(b *testing.B) { + m, s := capacityMetrics(n), capacityAPISender() + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + checkpoint, err := s.SendWithCheckpoint(context.Background(), m, nil) + if err != nil || checkpoint != nil { + b.Fatalf("send: %v, %v", checkpoint, err) + } + } + }) + } +} + +func BenchmarkDatadogCapacityLongTags(b *testing.B) { + for _, n := range []int{1, 61, 1000} { + for _, compressible := range []bool{false, true} { + b.Run(fmt.Sprintf("%d/compressible=%t", n, compressible), func(b *testing.B) { + m, s := capacityMetrics(n), capacityAPISender() + for i := range m.Values["testdomain"] { + tag := strings.Repeat("x", 4096) + if !compressible { + var text strings.Builder + for j := 0; j < 64; j++ { + fmt.Fprintf(&text, "%x", sha256.Sum256([]byte(fmt.Sprintf("%d/%d", i, j)))) + } + tag = text.String() + } + m.Values["testdomain"][i].Group = map[string]string{"table": tag} + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + checkpoint, err := s.SendWithCheckpoint(context.Background(), m, nil) + if err != nil || checkpoint != nil { + b.Fatalf("send: %v, %v", checkpoint, err) + } + } + }) + } + } +} + +func TestDatadogCapacityWindows(t *testing.T) { + for _, n := range datadogCapacityCounts { + t.Run(fmt.Sprint(n), func(t *testing.T) { + m, s := capacityMetrics(n), newTestDatadogSender(nil) + state, err := newDatadogSendCheckpoint(m, nil) + require.NoError(t, err) + cursor, total := datadogMetricCursor{}, 0 + for { + series, cursors, next, done, err := s.collectDatadogSeries(context.Background(), m, state.domains, cursor, datadogMaxSeriesPerPayload) + require.NoError(t, err) + require.LessOrEqual(t, len(series), datadogMaxSeriesPerPayload) + require.Len(t, cursors, len(series)) + for i, v := range series { + require.Equal(t, fmt.Sprintf("testdomain.testmetric%d", total+i+1), v.Metric) + } + total += len(series) + if done { + break + } + require.NotEqual(t, cursor, next) + cursor = next + } + require.Equal(t, n, total) + }) + } +} + +func TestDatadogCapacitySkippedValuesAndCursor(t *testing.T) { + m := capacityMetrics(6) + v := m.Values["testdomain"] + v[1].Type = blip.UNKNOWN + v[2].Meta = map[string]string{"ts": "invalid"} + m.Values = map[string][]blip.MetricValue{"a": v[:3], "b": nil, "c": v[3:]} + s := newTestDatadogSender(nil) + series, cursors, next, done, err := s.collectDatadogSeries(context.Background(), m, []string{"a", "b", "c"}, datadogMetricCursor{metric: 1}, 2) + require.NoError(t, err) + require.False(t, done) + require.Len(t, series, 2) + require.Equal(t, "c.testmetric4", series[0].Metric) + require.Equal(t, []datadogMetricCursor{{domain: 2, metric: 1}, {domain: 2, metric: 2}}, cursors) + require.Equal(t, cursors[1], next) + series, _, next, done, err = s.collectDatadogSeries(context.Background(), m, []string{"a", "b", "c"}, next, 10000) + require.NoError(t, err) + require.True(t, done) + require.Len(t, series, 1) + require.Equal(t, "c.testmetric6", series[0].Metric) + require.Equal(t, datadogMetricCursor{domain: 3}, next) +} + +func TestDatadogSmallWindowCapacity(t *testing.T) { + m, s := capacityMetrics(61), newTestDatadogSender(nil) + for _, offset := range []int{0, 60, 61} { + series, cursors, _, done, err := s.collectDatadogSeries(context.Background(), m, []string{"testdomain"}, datadogMetricCursor{metric: offset}, datadogMaxSeriesPerPayload) + require.NoError(t, err) + require.True(t, done) + require.Equal(t, 61-offset, cap(series)) + require.Equal(t, 61-offset, cap(cursors)) + } +} + +func TestDatadogCapacityPayloadGrowthAndStart(t *testing.T) { + series := testMetricSeries(3) + series[1].Metric = strings.Repeat("long_name", 1024) + series[1].Tags = []string{strings.Repeat("long_tag", 1024)} + for _, compress := range []bool{false, true} { + payload, end, err := prepareDatadogPayload(context.Background(), series, 1, 10000, compress, defaultDatadogPayloadLimits()) + require.NoError(t, err) + require.Equal(t, 3, end) + decoded, err := decodePreparedMetricPayload(payload) + require.NoError(t, err) + require.Equal(t, series[1:], decoded.Series) + } +} + +type capacitySubmitFunc func(context.Context, preparedDatadogPayload) (datadogSubmitResult, error) + +func (f capacitySubmitFunc) Submit(ctx context.Context, p preparedDatadogPayload) (datadogSubmitResult, error) { + return f(ctx, p) +} + +func TestDatadogCapacityStreaming413AndResume(t *testing.T) { + for _, failure := range []error{errors.New("network failure"), context.DeadlineExceeded, context.Canceled} { + t.Run(failure.Error(), func(t *testing.T) { + var attempted []int + var accepted []string + failed := false + s := newTestDatadogSender(capacitySubmitFunc(func(_ context.Context, p preparedDatadogPayload) (datadogSubmitResult, error) { + attempted = append(attempted, p.seriesCount) + if p.seriesCount > 2 { + return datadogSubmitResult{statusCode: 413}, errors.New("too large") + } + if len(accepted) == 2 && !failed { + failed = true + return datadogSubmitResult{}, failure + } + decoded, err := decodePreparedMetricPayload(p) + require.NoError(t, err) + for _, v := range decoded.Series { + accepted = append(accepted, v.Metric) + } + // Intake errors on a success still acknowledge this chunk. + return datadogSubmitResult{statusCode: 202, errors: []string{"test intake warning"}}, nil + })) + m := capacityMetrics(10) + checkpoint, err := s.SendWithCheckpoint(context.Background(), m, nil) + require.ErrorIs(t, err, failure) + require.NotNil(t, checkpoint) + require.Equal(t, []int{10, 5, 2, 2}, attempted) + require.Equal(t, int64(2), s.maxSeriesPerRequest.Load()) + checkpoint, err = s.SendWithCheckpoint(context.Background(), m, checkpoint) + require.NoError(t, err) + require.Nil(t, checkpoint) + require.Len(t, accepted, 10) + for i, name := range accepted { + require.Equal(t, fmt.Sprintf("testdomain.testmetric%d", i+1), name) + } + attempted = nil + _, err = s.SendWithCheckpoint(context.Background(), capacityMetrics(3), nil) + require.NoError(t, err) + require.Equal(t, []int{2, 1}, attempted) + }) + } +} + +func TestDatadogCapacityStreaming413Exhaustion(t *testing.T) { + for _, n := range []int{1, 1000} { + t.Run(fmt.Sprint(n), func(t *testing.T) { + calls := 0 + s := newTestDatadogSender(capacitySubmitFunc(func(context.Context, preparedDatadogPayload) (datadogSubmitResult, error) { + calls++ + return datadogSubmitResult{statusCode: 413}, errors.New("too large") + })) + checkpoint, err := s.SendWithCheckpoint(context.Background(), capacityMetrics(n), nil) + require.Error(t, err) + state := checkpoint.(*datadogSendCheckpoint) + require.Zero(t, state.sentSeries) + require.Equal(t, datadogMetricCursor{}, state.cursor) + require.Equal(t, min(n, datadogMax413Retries+1), calls) + }) + } +} + +func TestDatadogCapacityExactByteLimits(t *testing.T) { + series := testMetricSeries(1) + for _, compress := range []bool{false, true} { + p, _, err := prepareDatadogPayload(context.Background(), series, 0, 10000, compress, defaultDatadogPayloadLimits()) + require.NoError(t, err) + for _, delta := range []int{-1, 0, 1} { + limits := defaultDatadogPayloadLimits() + limits.maxCompressed = p.compressedBytes + delta + limits.targetCompressed = limits.maxCompressed + limits.maxDecompressed = p.uncompressedBytes + delta + limits.targetDecompressed = limits.maxDecompressed + _, _, err := prepareDatadogPayload(context.Background(), series, 0, 10000, compress, limits) + if delta < 0 { + require.Error(t, err) + } else { + require.NoError(t, err) + } + } + } +} + +func TestDatadogCapacityConcurrentSend(t *testing.T) { + s := capacityAPISender() + m := capacityMetrics(61) + var wg sync.WaitGroup + for i := 0; i < 4; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for j := 0; j < 10; j++ { + checkpoint, err := s.SendWithCheckpoint(context.Background(), m, nil) + if err != nil || checkpoint != nil { + t.Errorf("send: %v, %v", checkpoint, err) + } + } + }() + } + wg.Wait() +} + +func TestDatadogCapacityCheckpointDuringQueueOverflow(t *testing.T) { + blocked, release := make(chan struct{}), make(chan struct{}) + var accepted []string + calls := 0 + s := newTestDatadogSender(capacitySubmitFunc(func(_ context.Context, p preparedDatadogPayload) (datadogSubmitResult, error) { + calls++ + if calls == 2 { + close(blocked) + <-release + return datadogSubmitResult{}, errors.New("in-flight failure") + } + decoded, err := decodePreparedMetricPayload(p) + if err != nil { + return datadogSubmitResult{}, err + } + for _, v := range decoded.Series { + accepted = append(accepted, v.Metric) + } + return datadogSubmitResult{statusCode: 202}, nil + })) + s.maxSeriesPerRequest.Store(2) + retry := NewRetry(RetryArgs{MonitorId: "capacity-overflow", Sink: s, BufferSize: 60, SendTimeout: 10 * time.Second, SendRetryWait: time.Millisecond}) + done := make(chan error, 1) + go func() { done <- retry.Send(context.Background(), capacityMetrics(10)) }() + select { + case <-blocked: + case <-time.After(5 * time.Second): + t.Fatal("send did not reach second chunk") + } + // Enqueuing 61 entries evicts the in-flight item and the first new item. + // Its acknowledged prefix must stay accepted, while Retry keeps its normal + // newest-first overflow behavior for the remaining queue. + for i := 0; i < 61; i++ { + m := capacityMetrics(1) + m.Values["testdomain"][0].Name = fmt.Sprintf("queued_%d", i) + require.NoError(t, retry.Send(context.Background(), m)) + } + close(release) + select { + case err := <-done: + require.NoError(t, err) + case <-time.After(10 * time.Second): + t.Fatal("retry did not drain") + } + require.Len(t, accepted, 62) + require.Equal(t, []string{"testdomain.testmetric1", "testdomain.testmetric2"}, accepted[:2]) + for i := 0; i < 60; i++ { + require.Equal(t, fmt.Sprintf("testdomain.queued_%d", 60-i), accepted[i+2]) + } + require.Equal(t, -1, retry.top) +} diff --git a/sink/datadog_stream.go b/sink/datadog_stream.go index 586d5fa..5fe99d0 100644 --- a/sink/datadog_stream.go +++ b/sink/datadog_stream.go @@ -217,8 +217,19 @@ func newDatadogSendCheckpoint(metrics *blip.Metrics, checkpoint any) (*datadogSe } func (s *Datadog) collectDatadogSeries(ctx context.Context, metrics *blip.Metrics, domains []string, start datadogMetricCursor, limit int) ([]datadogV2.MetricSeries, []datadogMetricCursor, datadogMetricCursor, bool, error) { - series := make([]datadogV2.MetricSeries, 0, limit) - cursors := make([]datadogMetricCursor, 0, limit) + // Count source values only until the window is full. Slice lengths avoid + // a second conversion pass, and starting at the cursor avoids rescanning + // earlier domains for each window. Skipped values can overestimate capacity. + capacity := 0 + for domain := start.domain; domain < len(domains) && capacity < limit; domain++ { + remaining := len(metrics.Values[domains[domain]]) + if domain == start.domain { + remaining -= start.metric + } + capacity += min(remaining, limit-capacity) + } + series := make([]datadogV2.MetricSeries, 0, capacity) + cursors := make([]datadogMetricCursor, 0, capacity) cursor := start for cursor.domain < len(domains) { From 6a99d5c2a8f2066da93ff899a26ffad26b802eb0 Mon Sep 17 00:00:00 2001 From: Ian Oberst Date: Thu, 17 Sep 2026 06:52:44 -0700 Subject: [PATCH 2/2] Refine Datadog raw buffer capacity hints Sample bounded name and tag lengths to reduce repeated buffer growth for long series. Use a median to avoid extrapolating isolated large values, and retain exact payload checks and the small-series capacity floor. --- sink/datadog_api.go | 53 ++++++++++++++--- sink/datadog_capacity_hint_test.go | 96 ++++++++++++++++++++++++++++++ 2 files changed, 142 insertions(+), 7 deletions(-) create mode 100644 sink/datadog_capacity_hint_test.go diff --git a/sink/datadog_api.go b/sink/datadog_api.go index 93ac6b7..0a118ce 100644 --- a/sink/datadog_api.go +++ b/sink/datadog_api.go @@ -10,6 +10,7 @@ import ( "fmt" "io" "net/http" + "sort" "strings" "time" @@ -183,13 +184,7 @@ func prepareDatadogPayload(ctx context.Context, series []datadogV2.MetricSeries, var raw bytes.Buffer candidates := min(maxSeries, len(series)-start) - grow := rawTarget - // This is only an initial capacity hint: long series can still grow the - // buffer, and the exact byte checks below remain authoritative. Compare - // before multiplying to keep the estimate bounded without overflow. - if candidates <= rawTarget/512 { - grow = candidates * 512 - } + grow := datadogRawCapacity(series[start:start+candidates], rawTarget) if grow > 0 { raw.Grow(grow) } @@ -286,6 +281,50 @@ func prepareDatadogPayload(ctx context.Context, series []datadogV2.MetricSeries, return prepared, end, nil } +// datadogRawCapacity uses the median of at most eight evenly spaced series. +// It inspects at most sixteen tags per series without scanning string contents. +// Allow extra room for JSON strings, but keep the existing small-series floor. +// This is only a hint: escaping, other fields, and unsampled outliers can still +// require growth. Exact payload checks remain authoritative. +func datadogRawCapacity(series []datadogV2.MetricSeries, target int) int { + count := len(series) + if count == 0 || target <= 0 { + return 0 + } + if count > target/512 { + return target + } + samples := min(count, 8) + perSeries := target / count + // perSeries >= 512. Cap before doubling or multiplying to avoid overflow. + stringBudget := (perSeries - 128) / 2 + var sizes [8]int + for i := 0; i < samples; i++ { + // Include both ends; quotient/remainder avoids multiplying count by i. + divisor := max(samples-1, 1) + index := (count-1)/divisor*i + (count-1)%divisor*i/divisor + value := &series[index] + size := min(len(value.Metric), stringBudget) + tags := min(len(value.Tags), 16) + for j := 0; j < tags; j++ { + divisor := max(tags-1, 1) + index := (len(value.Tags)-1)/divisor*j + (len(value.Tags)-1)%divisor*j/divisor + size += min(len(value.Tags[index]), stringBudget-size) + } + sizes[i] = max(512, 128+2*size) + if size == stringBudget { + sizes[i] = perSeries + } + } + sort.Ints(sizes[:samples]) + // The upper median avoids extrapolating a single unusually large series. + estimate := sizes[samples/2] + if estimate == perSeries { + return target + } + return estimate * count +} + func payloadPrefixForSeries(raw []byte, end int) []byte { prefix := make([]byte, end+len(datadogPayloadSuffix)) copy(prefix, raw[:end]) diff --git a/sink/datadog_capacity_hint_test.go b/sink/datadog_capacity_hint_test.go new file mode 100644 index 0000000..014b92f --- /dev/null +++ b/sink/datadog_capacity_hint_test.go @@ -0,0 +1,96 @@ +package sink + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "testing" + + "github.com/DataDog/datadog-api-client-go/v2/api/datadogV2" + "github.com/stretchr/testify/require" +) + +func TestDatadogRawCapacityBounds(t *testing.T) { + for _, n := range []int{0, 1, 2, 8, 9, 61, 383, 1000, 10000} { + series := testMetricSeries(n) + for _, target := range []int{0, 1, 511, 512, 513, 4095, 4718592, int(^uint(0) >> 1)} { + capacity := datadogRawCapacity(series, target) + require.Equal(t, min(n*512, target), capacity, "ordinary series retain the existing hint") + if n == 0 { + continue + } + series[n-1].Tags = []string{strings.Repeat("large", 20000)} + capacity = datadogRawCapacity(series, target) + require.GreaterOrEqual(t, capacity, 0) + require.LessOrEqual(t, capacity, target) + series[n-1] = testMetricSeries(1)[0] + } + } +} + +func TestDatadogRawCapacityIsolatedOutlier(t *testing.T) { + for _, position := range []int{0, 499, 999} { + series := testMetricSeries(1000) + series[position].Tags = []string{strings.Repeat("x", 65536)} + require.Equal(t, 512000, datadogRawCapacity(series, datadogTargetDecompressedPayloadSize)) + } +} + +func TestDatadogRawCapacityFindsLateTags(t *testing.T) { + series := testMetricSeries(1000) + for i := range series { + series[i].Tags = make([]string, 64) + series[i].Tags[63] = strings.Repeat("x", 4096) + } + require.Equal(t, datadogTargetDecompressedPayloadSize, datadogRawCapacity(series, datadogTargetDecompressedPayloadSize)) +} + +func TestDatadogCapacityHardLimitSeries(t *testing.T) { + for _, compress := range []bool{false, true} { + limit := datadogMaxCompressedPayloadSize + if compress { + limit = datadogMaxDecompressedPayloadSize + } + for _, delta := range []int{-1, 0, 1} { + t.Run(fmt.Sprintf("gzip=%t/delta=%d", compress, delta), func(t *testing.T) { + series := testMetricSeries(1) + series[0].Tags = []string{""} + encoded, err := json.Marshal(series[0]) + require.NoError(t, err) + series[0].Tags[0] = strings.Repeat("x", limit+delta-len(encoded)-len(datadogPayloadPrefix)-len(datadogPayloadSuffix)) + payload, _, err := prepareDatadogPayload(context.Background(), series, 0, 10000, compress, defaultDatadogPayloadLimits()) + if delta > 0 { + require.Error(t, err) + return + } + require.NoError(t, err) + require.Equal(t, limit+delta, payload.uncompressedBytes) + decoded, err := decodePreparedMetricPayload(payload) + require.NoError(t, err) + require.Equal(t, series, decoded.Series) + }) + } + } +} + +func BenchmarkDatadogRawCapacityHint(b *testing.B) { + for _, n := range []int{61, 383, 1000, 10000} { + for _, tags := range []int{2, 64} { + b.Run(fmt.Sprintf("%d/tags=%d", n, tags), func(b *testing.B) { + series := make([]datadogV2.MetricSeries, n) + for i := range series { + series[i].Tags = make([]string, tags) + for j := range series[i].Tags { + series[i].Tags[j] = "tag:synthetic" + } + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + datadogRawCapacity(series, datadogTargetDecompressedPayloadSize) + } + }) + } + } +}