Skip to content

Commit a2776c2

Browse files
jszwedkoblt
andauthored
fix(metrics): Make v2 series obey max per-payload series points (1.3.x backport) (#2081)
## Summary Backports #2055 to `releases/1.3.x`. Repairs a differential between DogStatsD and ADP: the v2 series request builder now obeys both `max_metrics_per_payload` **and** `max_series_points_per_payload` for `/api/v2/series`, splitting oversized payloads the same way DogStatsD does. v3 already handled this correctly. The regression was introduced by the metrics-encoder refactor (splitting into `endpoint.rs` / `v2` / `v3` modules), which 1.3.x carries. Clean cherry-pick of `0025407262` (#2055) — 1.3.x shares main's encoder structure. ## Validation - `cargo check -p saluki-components` passes. - Unit tests pass, including the new regression test `v2_series_builder_enforces_max_series_points_per_payload`. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: blt <brian.troutwine@datadoghq.com> Co-authored-by: jesse.szwedko <jesse.szwedko@datadoghq.com>
1 parent f546aa0 commit a2776c2

3 files changed

Lines changed: 91 additions & 21 deletions

File tree

lib/saluki-components/src/encoders/datadog/metrics/endpoint.rs

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,17 +35,24 @@ impl MetricsEndpoint {
3535

3636
pub struct EndpointConfiguration {
3737
compression_scheme: CompressionScheme,
38+
// `max_metrics_per_payload` controls the number of series/sketches
39+
// per-payload to the endpoint, above which multiple payloads are split.
3840
max_metrics_per_payload: usize,
41+
// `max_series_points_per_payload` controls the number of total points per
42+
// payload, across all series. Sketches are not constrained.
43+
max_series_points_per_payload: usize,
3944
additional_tags: SharedTagSet,
4045
}
4146

4247
impl EndpointConfiguration {
4348
pub fn new(
44-
compression_scheme: CompressionScheme, max_metrics_per_payload: usize, additional_tags: Option<SharedTagSet>,
49+
compression_scheme: CompressionScheme, max_metrics_per_payload: usize, max_series_points_per_payload: usize,
50+
additional_tags: Option<SharedTagSet>,
4551
) -> Self {
4652
Self {
4753
compression_scheme,
4854
max_metrics_per_payload,
55+
max_series_points_per_payload,
4956
additional_tags: additional_tags.unwrap_or_default(),
5057
}
5158
}
@@ -58,6 +65,10 @@ impl EndpointConfiguration {
5865
self.max_metrics_per_payload
5966
}
6067

68+
pub fn max_series_points_per_payload(&self) -> usize {
69+
self.max_series_points_per_payload
70+
}
71+
6172
pub fn additional_tags(&self) -> &SharedTagSet {
6273
&self.additional_tags
6374
}

lib/saluki-components/src/encoders/datadog/metrics/mod.rs

Lines changed: 76 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -446,11 +446,15 @@ impl EncoderBuilder for DatadogMetricsConfiguration {
446446
let v2_endpoint_config = EndpointConfiguration::new(
447447
v2_compression_scheme,
448448
self.max_metrics_per_payload,
449+
self.max_series_points_per_payload,
449450
self.additional_tags.clone(),
450451
);
451452
let endpoint_config = EndpointConfiguration::new(
452453
v3_compression_scheme,
453454
self.max_metrics_per_payload,
455+
// Actually enforced by V3PayloadLimits, required for the
456+
// constructor shared between V1/V2/V3.
457+
usize::MAX,
454458
self.additional_tags.clone(),
455459
);
456460

@@ -2104,7 +2108,7 @@ serializer_experimental_use_v3_api:
21042108
assert!(combined_request.compressed_len > single_request.compressed_len);
21052109

21062110
let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2107-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2111+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
21082112
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
21092113
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
21102114

@@ -2129,7 +2133,7 @@ serializer_experimental_use_v3_api:
21292133
];
21302134
let single_request = create_v3_test_request(&metrics[..1]).await;
21312135
let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2132-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2136+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
21332137
let recorder = TestRecorder::default();
21342138
let _local = metrics::set_default_local_recorder(&recorder);
21352139
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
@@ -2150,7 +2154,7 @@ serializer_experimental_use_v3_api:
21502154
let metrics = vec![Metric::counter("v3.telemetry.item_too_big", 1.0)];
21512155
let request = create_v3_test_request(&metrics).await;
21522156
let limits = V3PayloadLimits::new(request.compressed_len - 1, usize::MAX, 10_000, 10_000);
2153-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2157+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
21542158
let recorder = TestRecorder::default();
21552159
let _local = metrics::set_default_local_recorder(&recorder);
21562160
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
@@ -2178,7 +2182,7 @@ serializer_experimental_use_v3_api:
21782182
assert!(combined_request.uncompressed_len > single_request.uncompressed_len);
21792183

21802184
let limits = V3PayloadLimits::new(usize::MAX, single_request.uncompressed_len, 10_000, 10_000);
2181-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2185+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
21822186
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
21832187
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
21842188

@@ -2200,7 +2204,7 @@ serializer_experimental_use_v3_api:
22002204
Metric::counter("v3.points.split.three", 5.0),
22012205
];
22022206
let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
2203-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2207+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
22042208
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
22052209
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
22062210
let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
@@ -2219,7 +2223,7 @@ serializer_experimental_use_v3_api:
22192223
Metric::counter("v3.telemetry.max_points.two", [(123, 3.0), (124, 4.0)]),
22202224
];
22212225
let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 2);
2222-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2226+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
22232227
let recorder = TestRecorder::default();
22242228
let _local = metrics::set_default_local_recorder(&recorder);
22252229
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
@@ -2248,7 +2252,7 @@ serializer_experimental_use_v3_api:
22482252
Metric::counter("v3.points.oversized.after", 7.0),
22492253
];
22502254
let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
2251-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2255+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
22522256
let recorder = TestRecorder::default();
22532257
let _local = metrics::set_default_local_recorder(&recorder);
22542258
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
@@ -2271,7 +2275,7 @@ serializer_experimental_use_v3_api:
22712275
Metric::counter("v3.points.zero.after", 2.0),
22722276
];
22732277
let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000);
2274-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2278+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
22752279
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
22762280
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
22772281
let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
@@ -2294,7 +2298,7 @@ serializer_experimental_use_v3_api:
22942298
assert!(combined_request.compressed_len > single_request.compressed_len);
22952299

22962300
let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2297-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2301+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
22982302
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
22992303
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
23002304
let batch_id = Uuid::now_v7();
@@ -2360,7 +2364,7 @@ serializer_experimental_use_v3_api:
23602364
assert!(combined_request.compressed_len > single_request.compressed_len);
23612365

23622366
let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2363-
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2367+
let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
23642368
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
23652369
let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
23662370
let batch_id = Uuid::now_v7();
@@ -2564,13 +2568,13 @@ serializer_experimental_use_v3_api:
25642568

25652569
#[tokio::test]
25662570
async fn validation_split_flush_assigns_batch_id_to_carried_metric() {
2567-
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 1, None);
2571+
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 1, usize::MAX, None);
25682572
let v2_series_builder = Some(
25692573
v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
25702574
.await
25712575
.expect("V2 request builder should be created"),
25722576
);
2573-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2577+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
25742578
let metrics_builder = MetricsBuilder::default();
25752579
let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
25762580
let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
@@ -2647,7 +2651,7 @@ serializer_experimental_use_v3_api:
26472651

26482652
#[tokio::test]
26492653
async fn authoritative_v3_flushes_previous_point_limit_batch() {
2650-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2654+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
26512655
let recorder = TestRecorder::default();
26522656
let _local = metrics::set_default_local_recorder(&recorder);
26532657
let metrics_builder = MetricsBuilder::default();
@@ -2728,7 +2732,7 @@ serializer_experimental_use_v3_api:
27282732

27292733
#[tokio::test]
27302734
async fn authoritative_v3_sketches_flush_previous_point_limit_batch() {
2731-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2735+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
27322736
let recorder = TestRecorder::default();
27332737
let _local = metrics::set_default_local_recorder(&recorder);
27342738
let metrics_builder = MetricsBuilder::default();
@@ -2800,13 +2804,13 @@ serializer_experimental_use_v3_api:
28002804

28012805
#[tokio::test]
28022806
async fn authoritative_v3_does_not_flush_on_v2_boundary() {
2803-
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 1, None);
2807+
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 1, usize::MAX, None);
28042808
let v2_series_builder = Some(
28052809
v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
28062810
.await
28072811
.expect("V2 request builder should be created"),
28082812
);
2809-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2813+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
28102814
let metrics_builder = MetricsBuilder::default();
28112815
let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
28122816
let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
@@ -2881,13 +2885,13 @@ serializer_experimental_use_v3_api:
28812885

28822886
#[tokio::test]
28832887
async fn shadow_sampled_series_flush_sends_v2_and_v3_beta_with_same_batch_id() {
2884-
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2888+
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
28852889
let v2_series_builder = Some(
28862890
v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
28872891
.await
28882892
.expect("V2 request builder should be created"),
28892893
);
2890-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2894+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
28912895
let metrics_builder = MetricsBuilder::default();
28922896
let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
28932897
let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
@@ -2961,13 +2965,13 @@ serializer_experimental_use_v3_api:
29612965

29622966
#[tokio::test]
29632967
async fn shadow_sample_rate_zero_sends_only_v2_without_validation_headers() {
2964-
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2968+
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
29652969
let v2_series_builder = Some(
29662970
v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
29672971
.await
29682972
.expect("V2 request builder should be created"),
29692973
);
2970-
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, None);
2974+
let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
29712975
let metrics_builder = MetricsBuilder::default();
29722976
let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
29732977
let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
@@ -3038,6 +3042,58 @@ serializer_experimental_use_v3_api:
30383042
fn tag_set<const N: usize>(tags: [&'static str; N]) -> TagSet {
30393043
tags.into_iter().map(Tag::from_static).collect()
30403044
}
3045+
3046+
// Regression test to ensure the V2 series request builder enforces `max_series_points_per_payload`.
3047+
//
3048+
// The test encodes more total points than the configured limit and asserts the builder splits them across
3049+
// multiple payloads without any payload exceeding the limit.
3050+
#[tokio::test]
3051+
async fn v2_series_builder_enforces_max_series_points_per_payload() {
3052+
let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, 10_000, None);
3053+
let mut builder = v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
3054+
.await
3055+
.expect("V2 request builder should be created");
3056+
builder
3057+
.with_len_limits(usize::MAX, usize::MAX)
3058+
.expect("byte limits should be accepted");
3059+
3060+
let mut total_points = 0;
3061+
for i in 0..6_000u32 {
3062+
let metric = Metric::gauge(
3063+
Context::from_parts(MetaString::from(format!("g{i}")), TagSet::default()),
3064+
[(1u64, 1.0), (2u64, 2.0)],
3065+
);
3066+
let mut pending = Some(metric);
3067+
while let Some(metric) = pending.take() {
3068+
if let Some(returned) = builder.encode(metric).await.expect("encode should not error") {
3069+
for request in builder.flush().await {
3070+
let (_, points, _) = request.expect("request should build");
3071+
assert!(
3072+
points <= 10_000,
3073+
"payload carried {} points, over the 10000 limit",
3074+
points
3075+
);
3076+
total_points += points;
3077+
}
3078+
pending = Some(returned);
3079+
}
3080+
}
3081+
}
3082+
for request in builder.flush().await {
3083+
let (_, points, _) = request.expect("request should build");
3084+
assert!(
3085+
points <= 10_000,
3086+
"final payload carried {} points, over the 10000 limit",
3087+
points
3088+
);
3089+
total_points += points;
3090+
}
3091+
3092+
assert_eq!(
3093+
total_points, 12_000,
3094+
"all points should be emitted across the split payloads"
3095+
);
3096+
}
30413097
}
30423098

30433099
#[cfg(test)]

lib/saluki-components/src/encoders/datadog/metrics/v2/mod.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,9 @@ pub async fn create_v2_request_builder(
3838
let mut request_builder =
3939
RequestBuilder::new(encoder, endpoint_config.compression_scheme(), RB_BUFFER_CHUNK_SIZE).await?;
4040
request_builder.with_max_inputs_per_payload(endpoint_config.max_metrics_per_payload());
41+
if matches!(endpoint, MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2) {
42+
request_builder.with_max_data_points_per_payload(endpoint_config.max_series_points_per_payload());
43+
}
4144

4245
Ok(request_builder)
4346
}

0 commit comments

Comments
 (0)