@@ -2659,6 +2659,78 @@ serializer_experimental_use_v3_api:
26592659 . expect ( "request builder should stop cleanly" ) ;
26602660 }
26612661
2662+ #[ tokio:: test]
2663+ async fn authoritative_v3_sketches_flush_previous_point_limit_batch ( ) {
2664+ let v3_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 10_000 , None ) ;
2665+ let recorder = TestRecorder :: default ( ) ;
2666+ let _local = metrics:: set_default_local_recorder ( & recorder) ;
2667+ let metrics_builder = MetricsBuilder :: default ( ) ;
2668+ let telemetry = ComponentTelemetry :: from_builder ( & metrics_builder) ;
2669+ let serializer_telemetry = V3SerializerTelemetry :: from_builder ( & metrics_builder) ;
2670+ let ( events_tx, events_rx) = tokio:: sync:: mpsc:: channel ( 1 ) ;
2671+ let ( payloads_tx, mut payloads_rx) = tokio:: sync:: mpsc:: channel ( 8 ) ;
2672+
2673+ let request_builder_handle = tokio:: spawn ( run_request_builder (
2674+ None ,
2675+ None ,
2676+ MetricsEncoderMode :: V2Only ,
2677+ MetricsEncoderMode :: V3Enabled ,
2678+ V3RuntimeConfig {
2679+ endpoint_config : v3_endpoint_config,
2680+ payload_limits : V3PayloadLimits :: new ( usize:: MAX , usize:: MAX , 10_000 , 3 ) ,
2681+ series_endpoint_uri : V3_SERIES_ENDPOINT_URI . to_string ( ) ,
2682+ shadow_series_endpoint_uri : "/api/intake/metrics/v3beta/series" . to_string ( ) ,
2683+ series_shadow_config : SeriesShadowConfig :: new ( 0.0 ) ,
2684+ serializer_telemetry,
2685+ } ,
2686+ telemetry,
2687+ events_rx,
2688+ payloads_tx,
2689+ Duration :: from_millis ( 250 ) ,
2690+ false ,
2691+ ) ) ;
2692+
2693+ let mut events = EventsBuffer :: default ( ) ;
2694+ assert ! ( events
2695+ . try_push( Event :: Metric ( Metric :: distribution(
2696+ "authoritative.v3.sketch.points.one" ,
2697+ [ ( 123 , 1.0 ) , ( 124 , 2.0 ) ]
2698+ ) ) )
2699+ . is_none( ) ) ;
2700+ assert ! ( events
2701+ . try_push( Event :: Metric ( Metric :: distribution(
2702+ "authoritative.v3.sketch.points.two" ,
2703+ [ ( 123 , 3.0 ) , ( 124 , 4.0 ) ]
2704+ ) ) )
2705+ . is_none( ) ) ;
2706+ events_tx
2707+ . send ( events)
2708+ . await
2709+ . expect ( "events should be sent to request builder" ) ;
2710+
2711+ for expected in [ "point-limit" , "timeout" ] {
2712+ let payload = timeout ( Duration :: from_secs ( 1 ) , payloads_rx. recv ( ) )
2713+ . await
2714+ . unwrap_or_else ( |_| panic ! ( "{expected} sketches payload should arrive before timeout" ) )
2715+ . expect ( "payload channel should remain open" ) ;
2716+ let Payload :: Http ( http_payload) = payload else {
2717+ panic ! ( "expected HTTP payload" ) ;
2718+ } ;
2719+ let ( _, request) = http_payload. into_parts ( ) ;
2720+ assert_eq ! ( V3_SKETCHES_ENDPOINT_URI , request. uri( ) ) ;
2721+ }
2722+ assert_eq ! (
2723+ recorder. counter( ( "serializer.v3_payload_split_reason" , & [ ( "reason" , "max_points" ) ] ) ) ,
2724+ Some ( 1 )
2725+ ) ;
2726+
2727+ drop ( events_tx) ;
2728+ request_builder_handle
2729+ . await
2730+ . expect ( "request builder task should complete" )
2731+ . expect ( "request builder should stop cleanly" ) ;
2732+ }
2733+
26622734 #[ tokio:: test]
26632735 async fn authoritative_v3_does_not_flush_on_v2_boundary ( ) {
26642736 let v2_endpoint_config = EndpointConfiguration :: new ( CompressionScheme :: noop ( ) , 1 , None ) ;
0 commit comments