@@ -813,6 +813,7 @@ async fn run_request_builder(
813813 // Keep a clone so we can add it to the next V3 shadow batch if that next batch samples in.
814814 let metric_for_next_shadow_batch = ( is_series
815815 && matches!( endpoint_mode, MetricsEncoderMode :: V2Only )
816+ && series_shadow_config. is_enabled( )
816817 && !series_shadow_active)
817818 . then( || metric. clone( ) ) ;
818819 let mut v2_encoded = false ;
@@ -852,14 +853,17 @@ async fn run_request_builder(
852853 let should_flush_v3 = match endpoint_mode {
853854 MetricsEncoderMode :: V2Only => series_shadow_active && v2_flushed,
854855 MetricsEncoderMode :: V3Enabled => {
855- if v3_flush_context. payload_limits. point_count_exceeds_limit( * v3_points)
856+ if v3_flush_context. payload_limits. point_count_fits( metric_point_count)
857+ && v3_flush_context. payload_limits. point_count_exceeds_limit( * v3_points)
856858 && v3_metrics. len( ) > 1
857859 {
858860 v3_flush_context
859861 . serializer_telemetry
860862 . record_split_reason( V3PayloadSplitReason :: MaxPoints ) ;
861- split_metric = v3_metrics. pop( ) ;
862- * v3_points = v3_points. saturating_sub( metric_point_count) ;
863+ if let Some ( metric) = v3_metrics. pop( ) {
864+ * v3_points = v3_points. saturating_sub( metric. values( ) . len( ) ) ;
865+ split_metric = Some ( metric) ;
866+ }
863867 true
864868 } else {
865869 v3_flush_context. payload_limits. should_flush_point_count_limit( * v3_points)
@@ -874,9 +878,11 @@ async fn run_request_builder(
874878 if !matches!( endpoint_mode, MetricsEncoderMode :: V3Enabled ) {
875879 // V2 flushes the previous batch without the current metric. Pop it
876880 // from V3 before flushing so both batches cover the same set of metrics.
877- split_metric = if v2_flushed { v3_metrics. pop( ) } else { None } ;
878- if split_metric. is_some( ) {
879- * v3_points = v3_points. saturating_sub( metric_point_count) ;
881+ if v2_flushed {
882+ if let Some ( metric) = v3_metrics. pop( ) {
883+ * v3_points = v3_points. saturating_sub( metric. values( ) . len( ) ) ;
884+ split_metric = Some ( metric) ;
885+ }
880886 }
881887 }
882888 let flush_context = if series_shadow_active {
@@ -2167,6 +2173,32 @@ serializer_experimental_use_v3_api:
21672173 ) ;
21682174 }
21692175
2176+ #[ test]
2177+ fn v3_metric_ranges_drop_oversized_metric_after_previous_range ( ) {
2178+ let metrics = vec ! [
2179+ Metric :: counter( "v3.points.oversized.before" , [ ( 123 , 1.0 ) , ( 124 , 2.0 ) ] ) ,
2180+ Metric :: counter(
2181+ "v3.points.oversized.too_big" ,
2182+ [ ( 123 , 3.0 ) , ( 124 , 4.0 ) , ( 125 , 5.0 ) , ( 126 , 6.0 ) ] ,
2183+ ) ,
2184+ Metric :: counter( "v3.points.oversized.after" , 7.0 ) ,
2185+ ] ;
2186+ let limits = V3PayloadLimits :: new ( usize:: MAX , usize:: MAX , 10_000 , 3 ) ;
2187+ let ep_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 10_000 , None ) ;
2188+ let recorder = TestRecorder :: default ( ) ;
2189+ let _local = metrics:: set_default_local_recorder ( & recorder) ;
2190+ let telemetry = ComponentTelemetry :: from_builder ( & MetricsBuilder :: default ( ) ) ;
2191+ let serializer_telemetry = V3SerializerTelemetry :: from_builder ( & MetricsBuilder :: default ( ) ) ;
2192+ let context = test_v3_flush_context ( & ep_config, limits, & serializer_telemetry, & telemetry) ;
2193+
2194+ let ranges = split_v3_metric_ranges_by_point_limit ( & metrics, context, "series" )
2195+ . into_iter ( )
2196+ . collect :: < Vec < _ > > ( ) ;
2197+
2198+ assert_eq ! ( vec![ 0 ..1 , 2 ..3 ] , ranges) ;
2199+ assert_eq ! ( recorder. counter( "serializer.v3_item_too_big" ) , Some ( 1 ) ) ;
2200+ }
2201+
21702202 #[ test]
21712203 fn v3_metric_ranges_skip_zero_point_metrics ( ) {
21722204 let metrics = vec ! [
@@ -2546,6 +2578,168 @@ serializer_experimental_use_v3_api:
25462578 . expect ( "request builder should stop cleanly" ) ;
25472579 }
25482580
2581+ #[ tokio:: test]
2582+ async fn authoritative_v3_flushes_previous_point_limit_batch ( ) {
2583+ let v3_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 10_000 , None ) ;
2584+ let recorder = TestRecorder :: default ( ) ;
2585+ let _local = metrics:: set_default_local_recorder ( & recorder) ;
2586+ let metrics_builder = MetricsBuilder :: default ( ) ;
2587+ let telemetry = ComponentTelemetry :: from_builder ( & metrics_builder) ;
2588+ let serializer_telemetry = V3SerializerTelemetry :: from_builder ( & metrics_builder) ;
2589+ let ( events_tx, events_rx) = tokio:: sync:: mpsc:: channel ( 1 ) ;
2590+ let ( payloads_tx, mut payloads_rx) = tokio:: sync:: mpsc:: channel ( 8 ) ;
2591+
2592+ let request_builder_handle = tokio:: spawn ( run_request_builder (
2593+ None ,
2594+ None ,
2595+ MetricsEncoderMode :: V3Enabled ,
2596+ MetricsEncoderMode :: V2Only ,
2597+ V3RuntimeConfig {
2598+ endpoint_config : v3_endpoint_config,
2599+ payload_limits : V3PayloadLimits :: new ( usize:: MAX , usize:: MAX , 10_000 , 3 ) ,
2600+ series_endpoint_uri : V3_SERIES_ENDPOINT_URI . to_string ( ) ,
2601+ shadow_series_endpoint_uri : "/api/intake/metrics/v3beta/series" . to_string ( ) ,
2602+ series_shadow_config : SeriesShadowConfig :: new ( 0.0 ) ,
2603+ serializer_telemetry,
2604+ } ,
2605+ telemetry,
2606+ events_rx,
2607+ payloads_tx,
2608+ Duration :: from_millis ( 250 ) ,
2609+ false ,
2610+ ) ) ;
2611+
2612+ let mut events = EventsBuffer :: default ( ) ;
2613+ assert ! ( events
2614+ . try_push( Event :: Metric ( Metric :: counter(
2615+ "authoritative.v3.points.one" ,
2616+ [ ( 123 , 1.0 ) , ( 124 , 2.0 ) ]
2617+ ) ) )
2618+ . is_none( ) ) ;
2619+ assert ! ( events
2620+ . try_push( Event :: Metric ( Metric :: counter(
2621+ "authoritative.v3.points.two" ,
2622+ [ ( 123 , 3.0 ) , ( 124 , 4.0 ) ]
2623+ ) ) )
2624+ . is_none( ) ) ;
2625+ events_tx
2626+ . send ( events)
2627+ . await
2628+ . expect ( "events should be sent to request builder" ) ;
2629+
2630+ let payload = timeout ( Duration :: from_secs ( 1 ) , payloads_rx. recv ( ) )
2631+ . await
2632+ . expect ( "point-limit payload should arrive before timeout" )
2633+ . expect ( "payload channel should remain open" ) ;
2634+ let Payload :: Http ( http_payload) = payload else {
2635+ panic ! ( "expected HTTP payload" ) ;
2636+ } ;
2637+ let ( _, request) = http_payload. into_parts ( ) ;
2638+ assert_eq ! ( V3_SERIES_ENDPOINT_URI , request. uri( ) ) ;
2639+ assert ! ( timeout( Duration :: from_millis( 50 ) , payloads_rx. recv( ) ) . await . is_err( ) ) ;
2640+
2641+ let payload = timeout ( Duration :: from_secs ( 1 ) , payloads_rx. recv ( ) )
2642+ . await
2643+ . expect ( "timeout payload should arrive before timeout" )
2644+ . expect ( "payload channel should remain open" ) ;
2645+ let Payload :: Http ( http_payload) = payload else {
2646+ panic ! ( "expected HTTP payload" ) ;
2647+ } ;
2648+ let ( _, request) = http_payload. into_parts ( ) ;
2649+ assert_eq ! ( V3_SERIES_ENDPOINT_URI , request. uri( ) ) ;
2650+ assert_eq ! (
2651+ recorder. counter( ( "serializer.v3_payload_split_reason" , & [ ( "reason" , "max_points" ) ] ) ) ,
2652+ Some ( 1 )
2653+ ) ;
2654+
2655+ drop ( events_tx) ;
2656+ request_builder_handle
2657+ . await
2658+ . expect ( "request builder task should complete" )
2659+ . expect ( "request builder should stop cleanly" ) ;
2660+ }
2661+
2662+ #[ tokio:: test]
2663+ async fn authoritative_v3_does_not_flush_on_v2_boundary ( ) {
2664+ let v2_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 1 , None ) ;
2665+ let v2_series_builder = Some (
2666+ v2:: create_v2_request_builder ( MetricsEndpoint :: SeriesV2 , & v2_endpoint_config)
2667+ . await
2668+ . expect ( "V2 request builder should be created" ) ,
2669+ ) ;
2670+ let v3_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 10_000 , None ) ;
2671+ let metrics_builder = MetricsBuilder :: default ( ) ;
2672+ let telemetry = ComponentTelemetry :: from_builder ( & metrics_builder) ;
2673+ let serializer_telemetry = V3SerializerTelemetry :: from_builder ( & metrics_builder) ;
2674+ let ( events_tx, events_rx) = tokio:: sync:: mpsc:: channel ( 1 ) ;
2675+ let ( payloads_tx, mut payloads_rx) = tokio:: sync:: mpsc:: channel ( 8 ) ;
2676+
2677+ let request_builder_handle = tokio:: spawn ( run_request_builder (
2678+ v2_series_builder,
2679+ None ,
2680+ MetricsEncoderMode :: V3Enabled ,
2681+ MetricsEncoderMode :: V2Only ,
2682+ V3RuntimeConfig {
2683+ endpoint_config : v3_endpoint_config,
2684+ payload_limits : V3PayloadLimits :: new ( usize:: MAX , usize:: MAX , 10_000 , 10_000 ) ,
2685+ series_endpoint_uri : V3_SERIES_ENDPOINT_URI . to_string ( ) ,
2686+ shadow_series_endpoint_uri : "/api/intake/metrics/v3beta/series" . to_string ( ) ,
2687+ series_shadow_config : SeriesShadowConfig :: new ( 0.0 ) ,
2688+ serializer_telemetry,
2689+ } ,
2690+ telemetry,
2691+ events_rx,
2692+ payloads_tx,
2693+ Duration :: from_millis ( 250 ) ,
2694+ false ,
2695+ ) ) ;
2696+
2697+ let mut events = EventsBuffer :: default ( ) ;
2698+ assert ! ( events
2699+ . try_push( Event :: Metric ( Metric :: counter( "authoritative.v3.decouple.one" , 1.0 ) ) )
2700+ . is_none( ) ) ;
2701+ assert ! ( events
2702+ . try_push( Event :: Metric ( Metric :: counter( "authoritative.v3.decouple.two" , 2.0 ) ) )
2703+ . is_none( ) ) ;
2704+ events_tx
2705+ . send ( events)
2706+ . await
2707+ . expect ( "events should be sent to request builder" ) ;
2708+
2709+ let payload = timeout ( Duration :: from_secs ( 1 ) , payloads_rx. recv ( ) )
2710+ . await
2711+ . expect ( "V2 split payload should arrive before timeout" )
2712+ . expect ( "payload channel should remain open" ) ;
2713+ let Payload :: Http ( http_payload) = payload else {
2714+ panic ! ( "expected HTTP payload" ) ;
2715+ } ;
2716+ let ( _, request) = http_payload. into_parts ( ) ;
2717+ assert_eq ! ( "/api/v2/series" , request. uri( ) ) ;
2718+ assert ! ( timeout( Duration :: from_millis( 50 ) , payloads_rx. recv( ) ) . await . is_err( ) ) ;
2719+
2720+ let mut timeout_flush_uris = Vec :: new ( ) ;
2721+ for _ in 0 ..2 {
2722+ let payload = timeout ( Duration :: from_secs ( 1 ) , payloads_rx. recv ( ) )
2723+ . await
2724+ . expect ( "timeout payload should arrive before timeout" )
2725+ . expect ( "payload channel should remain open" ) ;
2726+ let Payload :: Http ( http_payload) = payload else {
2727+ panic ! ( "expected HTTP payload" ) ;
2728+ } ;
2729+ let ( _, request) = http_payload. into_parts ( ) ;
2730+ timeout_flush_uris. push ( request. uri ( ) . to_string ( ) ) ;
2731+ }
2732+ assert_eq ! ( 2 , timeout_flush_uris. len( ) ) ;
2733+ assert ! ( timeout_flush_uris. iter( ) . any( |uri| uri == "/api/v2/series" ) ) ;
2734+ assert ! ( timeout_flush_uris. iter( ) . any( |uri| uri == V3_SERIES_ENDPOINT_URI ) ) ;
2735+
2736+ drop ( events_tx) ;
2737+ request_builder_handle
2738+ . await
2739+ . expect ( "request builder task should complete" )
2740+ . expect ( "request builder should stop cleanly" ) ;
2741+ }
2742+
25492743 #[ tokio:: test]
25502744 async fn shadow_sampled_series_flush_sends_v2_and_v3_beta_with_same_batch_id ( ) {
25512745 let v2_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 10_000 , None ) ;
0 commit comments