@@ -962,6 +962,95 @@ fn work_queue_take_respects_servable_range_contiguity_and_max() {
962962 ) ;
963963}
964964
965+ #[ test]
966+ fn work_queue_budgeted_take_respects_count_cap ( ) {
967+ let queue = work_queue_with (
968+ 0 ,
969+ ( 1 ..=4 ) . map ( |height| needed ( height, BlockSizeEstimate :: Advertised ( 100 ) ) ) ,
970+ ) ;
971+
972+ let taken = queue. take_in_range_budgeted ( block:: Height ( 1 ) , block:: Height ( 4 ) , 2 , u64:: MAX ) ;
973+
974+ assert_eq ! (
975+ taken. iter( ) . map( |( height, _) | height. 0 ) . collect:: <Vec <_>>( ) ,
976+ vec![ 1 , 2 ]
977+ ) ;
978+ }
979+
980+ #[ test]
981+ fn work_queue_budgeted_take_respects_estimated_byte_cap ( ) {
982+ let queue = work_queue_with (
983+ 0 ,
984+ [
985+ needed ( 1 , BlockSizeEstimate :: Advertised ( 100 ) ) ,
986+ needed ( 2 , BlockSizeEstimate :: Advertised ( 150 ) ) ,
987+ needed ( 3 , BlockSizeEstimate :: Advertised ( 1 ) ) ,
988+ ] ,
989+ ) ;
990+
991+ let taken = queue. take_in_range_budgeted ( block:: Height ( 1 ) , block:: Height ( 3 ) , 3 , 250 ) ;
992+
993+ assert_eq ! (
994+ taken. iter( ) . map( |( height, _) | height. 0 ) . collect:: <Vec <_>>( ) ,
995+ vec![ 1 , 2 ]
996+ ) ;
997+ assert ! ( queue. pending_contains( block:: Height ( 3 ) ) ) ;
998+ }
999+
1000+ #[ test]
1001+ fn work_queue_budgeted_take_stops_at_gaps ( ) {
1002+ let queue = work_queue_with (
1003+ 0 ,
1004+ [
1005+ needed ( 10 , BlockSizeEstimate :: Advertised ( 100 ) ) ,
1006+ needed ( 11 , BlockSizeEstimate :: Advertised ( 100 ) ) ,
1007+ needed ( 13 , BlockSizeEstimate :: Advertised ( 100 ) ) ,
1008+ ] ,
1009+ ) ;
1010+
1011+ let taken = queue. take_in_range_budgeted ( block:: Height ( 10 ) , block:: Height ( 13 ) , 3 , u64:: MAX ) ;
1012+
1013+ assert_eq ! (
1014+ taken. iter( ) . map( |( height, _) | height. 0 ) . collect:: <Vec <_>>( ) ,
1015+ vec![ 10 , 11 ]
1016+ ) ;
1017+ assert ! ( queue. pending_contains( block:: Height ( 13 ) ) ) ;
1018+ }
1019+
1020+ #[ test]
1021+ fn work_queue_budgeted_take_takes_one_oversized_first_item_for_progress ( ) {
1022+ let queue = work_queue_with (
1023+ 0 ,
1024+ [
1025+ needed ( 1 , BlockSizeEstimate :: Advertised ( 500 ) ) ,
1026+ needed ( 2 , BlockSizeEstimate :: Advertised ( 1 ) ) ,
1027+ ] ,
1028+ ) ;
1029+
1030+ let taken = queue. take_in_range_budgeted ( block:: Height ( 1 ) , block:: Height ( 2 ) , 2 , 100 ) ;
1031+
1032+ assert_eq ! (
1033+ taken. iter( ) . map( |( height, _) | height. 0 ) . collect:: <Vec <_>>( ) ,
1034+ vec![ 1 ]
1035+ ) ;
1036+ assert_eq ! ( taken[ 0 ] . 1 . estimated_bytes, 500 ) ;
1037+ assert ! ( queue. pending_contains( block:: Height ( 2 ) ) ) ;
1038+ }
1039+
1040+ #[ test]
1041+ fn work_queue_budgeted_take_preserves_estimates_through_take_and_return ( ) {
1042+ let queue = work_queue_with ( 0 , [ needed ( 10 , BlockSizeEstimate :: Advertised ( 12_345 ) ) ] ) ;
1043+
1044+ let taken = queue. take_in_range_budgeted ( block:: Height ( 10 ) , block:: Height ( 10 ) , 1 , 1 ) ;
1045+ assert_eq ! ( taken. len( ) , 1 ) ;
1046+ assert_eq ! ( taken[ 0 ] . 1 . estimated_bytes, 12_345 ) ;
1047+
1048+ queue. return_items ( [ block:: Height ( 10 ) ] ) ;
1049+ let retaken = queue. take_in_range_budgeted ( block:: Height ( 10 ) , block:: Height ( 10 ) , 1 , 1 ) ;
1050+ assert_eq ! ( retaken. len( ) , 1 ) ;
1051+ assert_eq ! ( retaken[ 0 ] . 1 . estimated_bytes, 12_345 ) ;
1052+ }
1053+
9651054#[ test]
9661055fn work_queue_take_does_not_clamp_high_to_floor ( ) {
9671056 // The committed floor is NOT an upper bound on a take: a peer fetches as far
@@ -1221,7 +1310,7 @@ async fn reactor_suppresses_needed_block_query_when_work_already_covers_tip() {
12211310 ( 1 ..=4 )
12221311 . map ( |height| BlockSyncBlockMeta {
12231312 height : block:: Height ( height) ,
1224- hash : block:: Hash ( [ height as u8 ; 32 ] ) ,
1313+ hash : block:: Hash ( [ u8 :: try_from ( height) . expect ( "test height fits u8" ) ; 32 ] ) ,
12251314 size : BlockSizeEstimate :: Advertised ( 1_000 ) ,
12261315 } )
12271316 . collect ( ) ,
@@ -4442,6 +4531,129 @@ async fn reactor_reserves_worst_case_per_block_not_size_hint() {
44424531 reactor_task. abort ( ) ;
44434532}
44444533
4534+ #[ tokio:: test]
4535+ async fn reactor_packs_small_estimates_under_peer_response_byte_cap ( ) {
4536+ let config = ZakuraBlockSyncConfig {
4537+ max_inflight_block_bytes : 4 * BS_PER_BLOCK_WORST_CASE_BYTES ,
4538+ max_blocks_per_response : 4 ,
4539+ ..ZakuraBlockSyncConfig :: default ( )
4540+ } ;
4541+
4542+ let ( _tip_tx, tip_rx) = watch:: channel ( ( block:: Height ( 0 ) , block:: Hash ( [ 0 ; 32 ] ) ) ) ;
4543+ let startup = BlockSyncStartup :: new (
4544+ BlockSyncFrontiers {
4545+ finalized_height : block:: Height ( 0 ) ,
4546+ verified_block_tip : block:: Height ( 0 ) ,
4547+ verified_block_hash : block:: Hash ( [ 0 ; 32 ] ) ,
4548+ } ,
4549+ ( block:: Height ( 0 ) , block:: Hash ( [ 0 ; 32 ] ) ) ,
4550+ tip_rx,
4551+ config. clone ( ) ,
4552+ ) ;
4553+ let ( handle, mut actions, reactor_task) = spawn_block_sync_reactor ( startup) ;
4554+ let service = BlockSyncService :: new_with_handle_for_test ( config, handle. clone ( ) ) ;
4555+
4556+ let one_worst_case_response =
4557+ u32:: try_from ( BS_PER_BLOCK_WORST_CASE_BYTES ) . expect ( "worst-case block size fits u32" ) ;
4558+ let ( _peer_id, _inbound, mut outbound) = connect_peer_with_status_message (
4559+ & service,
4560+ & mut actions,
4561+ 52 ,
4562+ BlockSyncStatus {
4563+ servable_low : block:: Height ( 1 ) ,
4564+ servable_high : block:: Height ( 4 ) ,
4565+ tip_hash : block:: Hash ( [ 4 ; 32 ] ) ,
4566+ max_blocks_per_response : 4 ,
4567+ max_inflight_requests : 1 ,
4568+ max_response_bytes : one_worst_case_response,
4569+ } ,
4570+ )
4571+ . await ;
4572+
4573+ handle
4574+ . send ( BlockSyncEvent :: NeededBlocks (
4575+ ( 1 ..=4 )
4576+ . map ( |height| BlockSyncBlockMeta {
4577+ height : block:: Height ( height) ,
4578+ hash : block:: Hash ( [ height as u8 ; 32 ] ) ,
4579+ size : BlockSizeEstimate :: Advertised ( 1_000 ) ,
4580+ } )
4581+ . collect ( ) ,
4582+ ) )
4583+ . await
4584+ . expect ( "needed metadata queues" ) ;
4585+
4586+ let ( start_height, count) = wait_for_outbound_getblocks ( & mut outbound) . await ;
4587+ assert_eq ! ( start_height, block:: Height ( 1 ) ) ;
4588+ assert_eq ! (
4589+ count, 4 ,
4590+ "small advertised estimates should pack more than the one worst-case block \
4591+ allowed by the peer response byte cap",
4592+ ) ;
4593+
4594+ reactor_task. abort ( ) ;
4595+ }
4596+
4597+ #[ tokio:: test]
4598+ async fn reactor_tiny_estimates_do_not_exceed_one_worst_case_budget_block ( ) {
4599+ let config = ZakuraBlockSyncConfig {
4600+ max_inflight_block_bytes : BS_PER_BLOCK_WORST_CASE_BYTES ,
4601+ max_blocks_per_response : 4 ,
4602+ ..ZakuraBlockSyncConfig :: default ( )
4603+ } ;
4604+
4605+ let ( _tip_tx, tip_rx) = watch:: channel ( ( block:: Height ( 0 ) , block:: Hash ( [ 0 ; 32 ] ) ) ) ;
4606+ let startup = BlockSyncStartup :: new (
4607+ BlockSyncFrontiers {
4608+ finalized_height : block:: Height ( 0 ) ,
4609+ verified_block_tip : block:: Height ( 0 ) ,
4610+ verified_block_hash : block:: Hash ( [ 0 ; 32 ] ) ,
4611+ } ,
4612+ ( block:: Height ( 0 ) , block:: Hash ( [ 0 ; 32 ] ) ) ,
4613+ tip_rx,
4614+ config. clone ( ) ,
4615+ ) ;
4616+ let ( handle, mut actions, reactor_task) = spawn_block_sync_reactor ( startup) ;
4617+ let service = BlockSyncService :: new_with_handle_for_test ( config, handle. clone ( ) ) ;
4618+
4619+ let ( _peer_id, _inbound, mut outbound) = connect_peer_with_status_message (
4620+ & service,
4621+ & mut actions,
4622+ 53 ,
4623+ BlockSyncStatus {
4624+ servable_low : block:: Height ( 1 ) ,
4625+ servable_high : block:: Height ( 4 ) ,
4626+ tip_hash : block:: Hash ( [ 4 ; 32 ] ) ,
4627+ max_blocks_per_response : 4 ,
4628+ max_inflight_requests : 1 ,
4629+ max_response_bytes : MAX_BS_RESPONSE_BYTES ,
4630+ } ,
4631+ )
4632+ . await ;
4633+
4634+ handle
4635+ . send ( BlockSyncEvent :: NeededBlocks (
4636+ ( 1 ..=4 )
4637+ . map ( |height| BlockSyncBlockMeta {
4638+ height : block:: Height ( height) ,
4639+ hash : block:: Hash ( [ u8:: try_from ( height) . expect ( "test height fits u8" ) ; 32 ] ) ,
4640+ size : BlockSizeEstimate :: Advertised ( 1 ) ,
4641+ } )
4642+ . collect ( ) ,
4643+ ) )
4644+ . await
4645+ . expect ( "needed metadata queues" ) ;
4646+
4647+ let ( start_height, count) = wait_for_outbound_getblocks ( & mut outbound) . await ;
4648+ assert_eq ! ( start_height, block:: Height ( 1 ) ) ;
4649+ assert_eq ! (
4650+ count, 1 ,
4651+ "only one worst-case block of global budget is available, regardless of tiny estimates" ,
4652+ ) ;
4653+
4654+ reactor_task. abort ( ) ;
4655+ }
4656+
44454657#[ tokio:: test]
44464658async fn reactor_zero_pause_threshold_preserves_lag_one_downloads ( ) {
44474659 let config = immediate_body_download_config ( ) ;
0 commit comments