LCOV - code coverage report
Current view: top level - discof/repair - fd_policy.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 142 302 47.0 %
Date: 2026-09-17 04:28:31 Functions: 9 15 60.0 %

          Line data    Source code
       1             : #include "fd_policy.h"
       2             : #include "../../disco/metrics/fd_metrics.h"
       3             : 
       4             : #define NONCE_NULL        (UINT_MAX)
       5           0 : #define DEFER_REPAIR_MS   (200UL)
       6           0 : #define TARGET_TICK_PER_SLOT (64.0)
       7           0 : #define MS_PER_TICK          (200.0 / TARGET_TICK_PER_SLOT)
       8             : 
       9             : void *
      10          45 : fd_policy_new( void * shmem, ulong peer_max, ulong seed, fd_rnonce_ss_t const * rnonce_ss ) {
      11             : 
      12          45 :   if( FD_UNLIKELY( !shmem ) ) {
      13           0 :     FD_LOG_WARNING(( "NULL mem" ));
      14           0 :     return NULL;
      15           0 :   }
      16             : 
      17          45 :   if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)shmem, fd_policy_align() ) ) ) {
      18           0 :     FD_LOG_WARNING(( "misaligned mem" ));
      19           0 :     return NULL;
      20           0 :   }
      21             : 
      22          45 :   ulong footprint = fd_policy_footprint( peer_max );
      23          45 :   fd_memset( shmem, 0, footprint );
      24             : 
      25          45 :   ulong peer_chain_cnt = fd_policy_peer_map_chain_cnt_est( peer_max );
      26          45 :   FD_SCRATCH_ALLOC_INIT( l, shmem );
      27          45 :   fd_policy_t * policy     = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_align(),            sizeof(fd_policy_t)                            );
      28          45 :   void *        peers      = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_peer_map_align(),   fd_policy_peer_map_footprint( peer_chain_cnt ) );
      29          45 :   void *        peers_pool = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_peer_pool_align(),  fd_policy_peer_pool_footprint( peer_max )      );
      30          45 :   void *        peers_fast = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_peer_dlist_align(), fd_policy_peer_dlist_footprint()               );
      31          45 :   void *        peers_slow = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_peer_dlist_align(), fd_policy_peer_dlist_footprint()               );
      32          45 :   FD_TEST( FD_SCRATCH_ALLOC_FINI( l, fd_policy_align() ) == (ulong)shmem + footprint );
      33             : 
      34          45 :   policy->peers.map     = fd_policy_peer_map_new  ( peers,      peer_chain_cnt, seed );
      35          45 :   policy->peers.pool    = fd_policy_peer_pool_new ( peers_pool, peer_max             );
      36          45 :   policy->peers.fast    = fd_policy_peer_dlist_new( peers_fast                       );
      37          45 :   policy->peers.slow    = fd_policy_peer_dlist_new( peers_slow                       );
      38          45 :   policy->turbine_slot0 = ULONG_MAX;
      39          45 :   policy->rnonce_ss[0]  = *rnonce_ss;
      40             : 
      41          45 :   return shmem;
      42          45 : }
      43             : 
      44             : fd_policy_t *
      45          45 : fd_policy_join( void * shpolicy ) {
      46          45 :   fd_policy_t * policy = (fd_policy_t *)shpolicy;
      47             : 
      48          45 :   if( FD_UNLIKELY( !policy ) ) {
      49           0 :     FD_LOG_WARNING(( "NULL policy" ));
      50           0 :     return NULL;
      51           0 :   }
      52             : 
      53          45 :   if( FD_UNLIKELY( !fd_ulong_is_aligned((ulong)policy, fd_policy_align() ) ) ) {
      54           0 :     FD_LOG_WARNING(( "misaligned policy" ));
      55           0 :     return NULL;
      56           0 :   }
      57             : 
      58          45 :   fd_wksp_t * wksp = fd_wksp_containing( policy );
      59          45 :   if( FD_UNLIKELY( !wksp ) ) {
      60           0 :     FD_LOG_WARNING(( "policy must be part of a workspace" ));
      61           0 :     return NULL;
      62           0 :   }
      63             : 
      64          45 :   policy->peers.map  = fd_policy_peer_map_join  ( policy->peers.map  );
      65          45 :   policy->peers.pool = fd_policy_peer_pool_join ( policy->peers.pool );
      66          45 :   policy->peers.fast = fd_policy_peer_dlist_join( policy->peers.fast );
      67          45 :   policy->peers.slow = fd_policy_peer_dlist_join( policy->peers.slow );
      68             : 
      69          45 :   policy->peers.select.fast_iter = fd_policy_peer_dlist_iter_fwd_init( policy->peers.fast, policy->peers.pool );
      70          45 :   policy->peers.select.slow_iter = fd_policy_peer_dlist_iter_fwd_init( policy->peers.slow, policy->peers.pool );
      71          45 :   policy->peers.select.cnt       = 0;
      72             : 
      73          45 :   return policy;
      74          45 : }
      75             : 
      76             : void *
      77           0 : fd_policy_leave( fd_policy_t const * policy ) {
      78             : 
      79           0 :   if( FD_UNLIKELY( !policy ) ) {
      80           0 :     FD_LOG_WARNING(( "NULL policy" ));
      81           0 :     return NULL;
      82           0 :   }
      83             : 
      84           0 :   return (void *)policy;
      85           0 : }
      86             : 
      87             : void *
      88           0 : fd_policy_delete( void * policy ) {
      89             : 
      90           0 :   if( FD_UNLIKELY( !policy ) ) {
      91           0 :     FD_LOG_WARNING(( "NULL policy" ));
      92           0 :     return NULL;
      93           0 :   }
      94             : 
      95           0 :   if( FD_UNLIKELY( !fd_ulong_is_aligned((ulong)policy, fd_policy_align() ) ) ) {
      96           0 :     FD_LOG_WARNING(( "misaligned policy" ));
      97           0 :     return NULL;
      98           0 :   }
      99             : 
     100           0 :   return policy;
     101           0 : }
     102             : 
     103           0 : static ulong ts_ms( long wallclock ) {
     104           0 :   return (ulong)wallclock / (ulong)1e6;
     105           0 : }
     106             : 
     107             : /* throttle_remaining_ns returns how long until the block passes the
     108             :    eager repair threshold, 0 if it already passes.  Essentially is
     109             :    checking if current duration of block ( from the first shred
     110             :    received until now ) is greater than the highest tick received +
     111             :    200ms. */
     112             : 
     113             : static long
     114           0 : throttle_remaining_ns( fd_policy_t * policy, fd_forest_blk_t * ele ) {
     115           0 :   if( FD_UNLIKELY( ele->slot < policy->turbine_slot0 ) ) return 0L;
     116           0 :   double current_duration = (double)(fd_tickcount() - ele->first_shred_ts) / fd_tempo_tick_per_ns(NULL);
     117           0 :   double tick_plus_buffer = (ele->est_buffered_tick_recv * MS_PER_TICK + DEFER_REPAIR_MS) * 1e6;
     118             : 
     119           0 :   if( current_duration >= tick_plus_buffer ){
     120           0 :     FD_MCNT_INC( REPAIR, EAGER_THRESHOLD_EXCEEDED, 1 );
     121           0 :     return 0L;
     122           0 :   }
     123           0 :   return (long)(tick_plus_buffer - current_duration);
     124           0 : }
     125             : 
     126             : static inline fd_policy_peer_dlist_iter_t
     127             : peer_iter_advance( fd_policy_peer_dlist_iter_t iter,
     128             :                    fd_policy_peer_dlist_t *    dlist,
     129        2436 :                    fd_policy_peer_t *          pool ) {
     130        2436 :   iter = fd_policy_peer_dlist_iter_fwd_next( iter, dlist, pool );
     131        2436 :   if( FD_UNLIKELY( fd_policy_peer_dlist_iter_done( iter, dlist, pool ) ) ) {
     132        1203 :     iter = fd_policy_peer_dlist_iter_fwd_init( dlist, pool );
     133        1203 :   }
     134        2436 :   return iter;
     135        2436 : }
     136             : 
     137             : fd_pubkey_t const *
     138        2436 : fd_policy_peer_select( fd_policy_t * policy ) {
     139        2436 :   fd_policy_peer_dlist_t * fast = policy->peers.fast;
     140        2436 :   fd_policy_peer_dlist_t * slow = policy->peers.slow;
     141        2436 :   fd_policy_peer_t       * pool = policy->peers.pool;
     142             : 
     143        2436 :   if( FD_UNLIKELY( fd_policy_peer_pool_used( pool ) == 0 ) ) return NULL;
     144             : 
     145             :   /* reinit stale iterators.  happens when peers are inserted into a
     146             :      previously-empty list after the iterator was initialized. */
     147        2436 :   int fast_empty = fd_policy_peer_dlist_iter_done( fd_policy_peer_dlist_iter_fwd_init( fast, pool ), fast, pool );
     148        2436 :   int slow_empty = fd_policy_peer_dlist_iter_done( fd_policy_peer_dlist_iter_fwd_init( slow, pool ), slow, pool );
     149             : 
     150        2436 :   if( FD_UNLIKELY( !fast_empty && fd_policy_peer_dlist_iter_done( policy->peers.select.fast_iter, fast, pool ) ) ) {
     151          12 :     policy->peers.select.fast_iter = fd_policy_peer_dlist_iter_fwd_init( fast, pool );
     152          12 :   }
     153        2436 :   if( FD_UNLIKELY( !slow_empty && fd_policy_peer_dlist_iter_done( policy->peers.select.slow_iter, slow, pool ) ) ) {
     154          45 :     policy->peers.select.slow_iter = fd_policy_peer_dlist_iter_fwd_init( slow, pool );
     155          45 :   }
     156             : 
     157        2436 :   fd_policy_peer_t * select;
     158             : 
     159             :   /* select will be set to current iterator status. Then iterator should
     160             :      be advanced for the following peer_select call. */
     161             : 
     162        2436 :   if( FD_UNLIKELY( fast_empty ) ) {
     163         987 :     select = fd_policy_peer_dlist_iter_ele( policy->peers.select.slow_iter, slow, pool );
     164         987 :     policy->peers.select.slow_iter = peer_iter_advance( policy->peers.select.slow_iter, slow, pool );
     165         987 :     return &select->key;
     166         987 :   }
     167             : 
     168        1449 :   if( FD_UNLIKELY( slow_empty ) ) {
     169        1446 :     select = fd_policy_peer_dlist_iter_ele( policy->peers.select.fast_iter, fast, pool );
     170        1446 :     policy->peers.select.fast_iter = peer_iter_advance( policy->peers.select.fast_iter, fast, pool );
     171        1446 :     return &select->key;
     172        1446 :   }
     173             : 
     174             :   /* interleave FD_POLICY_FAST_PER_SLOW fast, 1 slow. */
     175           3 :   if( FD_LIKELY( policy->peers.select.cnt < FD_POLICY_FAST_PER_SLOW ) ) {
     176           3 :     select = fd_policy_peer_dlist_iter_ele( policy->peers.select.fast_iter, fast, pool );
     177           3 :     policy->peers.select.fast_iter = peer_iter_advance( policy->peers.select.fast_iter, fast, pool );
     178           3 :     policy->peers.select.cnt++;
     179           3 :     return &select->key;
     180           3 :   }
     181             : 
     182           0 :   select = fd_policy_peer_dlist_iter_ele( policy->peers.select.slow_iter, slow, pool );
     183           0 :   policy->peers.select.slow_iter = peer_iter_advance( policy->peers.select.slow_iter, slow, pool );
     184           0 :   policy->peers.select.cnt = 0;
     185           0 :   return &select->key;
     186           3 : }
     187             : 
     188             : fd_repair_msg_t const *
     189           0 : fd_policy_next( fd_policy_t * policy, fd_reqlim_t * dedup, fd_forest_t * forest, fd_repair_t * repair, long now, ulong highest_known_slot, int * charge_busy ) {
     190           0 :   fd_forest_blk_t * pool = fd_forest_pool( forest );
     191           0 :   *charge_busy = 0;
     192             : 
     193           0 :   if( FD_UNLIKELY( forest->root == ULONG_MAX ) ) return NULL;
     194           0 :   if( FD_UNLIKELY( fd_policy_peer_pool_used( policy->peers.pool ) == 0 ) ) return NULL;
     195             : 
     196           0 :   fd_repair_msg_t * out = NULL;
     197           0 :   ulong now_ms = ts_ms( now );
     198             : 
     199           0 :   fd_forest_orphan_ent_t * orphanq = fd_forest_orphanq( forest );
     200           0 :   ulong budget = 64UL;
     201           0 :   while( budget-- && fd_forest_orphanq_cnt( orphanq ) && orphanq[ 0 ].due<=now ) {
     202           0 :     *charge_busy = 1;
     203           0 :     fd_forest_orphan_ent_t ent = orphanq[ 0 ];
     204           0 :     fd_forest_orphanq_remove_min( orphanq );
     205           0 :     fd_forest_blk_t * orphan = fd_forest_subtrees_ele_query( fd_forest_subtrees( forest ), &ent.slot, NULL, pool );
     206           0 :     if( FD_UNLIKELY( !orphan || orphan->orphan_seq!=ent.seq ) ) continue;
     207           0 :     ulong key     = fd_reqlim_key( FD_REPAIR_KIND_ORPHAN, ent.slot, UINT_MAX );
     208           0 :     int   deduped = fd_reqlim_next( dedup, key, now );
     209           0 :     ulong seq = forest->orphan_seq_next++;
     210           0 :     orphan->orphan_seq = seq;
     211           0 :     fd_forest_orphan_ent_t nxt = { .due = fd_reqlim_next_due( dedup, key, now ), .slot = ent.slot, .seq = seq };
     212           0 :     fd_forest_orphanq_insert( orphanq, &nxt );
     213           0 :     if( FD_LIKELY( !deduped ) ) {
     214           0 :       uint nonce = fd_rnonce_ss_compute( policy->rnonce_ss, 0, ent.slot, 0U, now );
     215           0 :       out = fd_repair_orphan( repair, fd_policy_peer_select( policy ), now_ms, nonce, ent.slot );
     216           0 :       orphan->req_orphan_cnt++;
     217           0 :       return out;
     218           0 :     }
     219           0 :   }
     220             : 
     221             :   /* Select a slot to operate on 🔪. Advance either the orphan iter or
     222             :      regular iter. */
     223           0 :   fd_forest_iter_t * iter = NULL;
     224           0 :   if( FD_UNLIKELY( fd_forest_reqslist_is_empty( fd_forest_reqslist( forest ), fd_forest_reqspool( forest ) ) ) ) {
     225             :     /* If the main tree has nothing to iterate at the moment, we can
     226             :        request down the ORPHAN trees on slots we know about. */
     227           0 :     iter = &forest->orphiter;
     228           0 :   } else {
     229           0 :     iter = &forest->iter;
     230           0 :   }
     231             : 
     232           0 :   fd_forest_iter_next( iter, forest );
     233           0 :   if( FD_UNLIKELY( fd_forest_iter_done( iter, forest ) ) ) {
     234             :     // This happens when we have already requested all the shreds we know about.
     235           0 :     return NULL;
     236           0 :   }
     237             : 
     238           0 :   fd_forest_blk_t * ele = fd_forest_pool_ele( pool, iter->ele_idx );
     239             : 
     240             :   /* The next request this call would produce.  If it was recently
     241             :      declined and nothing about it changed, skip the turn. */
     242           0 :   uint cand_idx = iter->shred_idx==UINT_MAX ? ele->buffered_idx+1U : iter->shred_idx;
     243           0 :   if( FD_UNLIKELY( ( policy->skip.throttled || iter->shred_idx==UINT_MAX ) &&
     244           0 :                    ele->slot==highest_known_slot &&
     245           0 :                    policy->skip.slot==ele->slot &&
     246           0 :                    policy->skip.idx==cand_idx &&
     247           0 :                    now<policy->skip.until ) ) {
     248           0 :     iter->shred_idx = UINT_MAX;
     249           0 :     return NULL;
     250           0 :   }
     251             : 
     252           0 :   long throttle_ns = ele->slot==highest_known_slot ? throttle_remaining_ns( policy, ele ) : 0L;
     253           0 :   if( FD_UNLIKELY( throttle_ns ) ) {
     254             :     /* When we are at the head of the turbine, we should give turbine the
     255             :        chance to complete the shreds.  Agave waits 200ms from the
     256             :        estimated "correct time" of the highest shred received to repair.
     257             :        i.e. if we've received the first 200 shreds, the 200th has a tick
     258             :        of x. Translate that to millis, and we should wait to request shred
     259             :        201 until x + 200ms.  If we have a hole, i.e. first 200 shreds
     260             :        receive except shred 100, and the 101th shred has a tick of y, we
     261             :        should wait until y + 200ms to request shred 100.
     262             : 
     263             :        Here we did not pass the timeout threshold, so we are not ready
     264             :        to repair this slot yet.  But it's possible we have another fork
     265             :        that we need to repair... so we just should skip to the next SLOT
     266             :        in the main tree iterator.  The likelihood that this ele is the
     267             :        head of turbine is high, which means that the shred_idx of the
     268             :        iterf is likely to be UINT_MAX, which means calling
     269             :        fd_forest_iter_next will advance the iterf to the next slot. */
     270           0 :     iter->shred_idx = UINT_MAX;
     271             :     /* TODO: Heinous... but the easiest way to ensure this slot gets
     272             :        added back to the requests deque is if we set the shred_idx to
     273             :        UINT_MAX, but maybe there should be an explicit API for it. */
     274             : 
     275           0 :     policy->skip.slot      = ele->slot;
     276           0 :     policy->skip.idx       = cand_idx;
     277           0 :     policy->skip.throttled = 1;
     278             :     /* Cap at 1ms: the deadline is derived from tick estimates that can
     279             :        move as more shreds land without changing the candidate. */
     280           0 :     policy->skip.until     = now + fd_long_min( throttle_ns, (long)1e6 );
     281           0 :     return NULL;
     282           0 :   }
     283             : 
     284           0 :   *charge_busy = 1;
     285             : 
     286           0 :   if( FD_UNLIKELY( iter->shred_idx == UINT_MAX ) ) {
     287             :     // We'll never know the the highest shred for the current turbine slot, so there's no point in requesting it.
     288           0 :     if( FD_UNLIKELY( ele->slot < highest_known_slot && !fd_reqlim_next( dedup, fd_reqlim_key( FD_REPAIR_KIND_HIGHEST_SHRED, ele->slot, UINT_MAX ), now ) ) ) {
     289           0 :       uint nonce = fd_rnonce_ss_compute( policy->rnonce_ss, 0, ele->slot, 0U, now );
     290           0 :       out = fd_repair_highest_shred( repair, fd_policy_peer_select( policy ), now_ms, nonce, ele->slot, 0 );
     291           0 :       ele->req_highest_cnt++;
     292           0 :     } else if( FD_LIKELY( ele->slot == highest_known_slot && (ulong)cand_idx < forest->shred_max ) ) {
     293           0 :       ulong key = fd_reqlim_key( FD_REPAIR_KIND_SHRED, ele->slot, cand_idx );
     294           0 :       if( FD_UNLIKELY( fd_reqlim_query( dedup, key, now ) ) ) {
     295           0 :         policy->skip.slot      = ele->slot;
     296           0 :         policy->skip.idx       = cand_idx;
     297           0 :         policy->skip.throttled = 0;
     298           0 :         policy->skip.until     = fd_reqlim_next_due( dedup, key, now );
     299           0 :         *charge_busy = 0;
     300           0 :         return NULL;
     301           0 :       }
     302           0 :       uint nonce = fd_rnonce_ss_compute( policy->rnonce_ss, 1, ele->slot, ele->buffered_idx + 1, now );
     303           0 :       out = fd_repair_shred( repair, fd_policy_peer_select( policy ), now_ms, nonce, ele->slot, ele->buffered_idx + 1 );
     304           0 :       ele->req_window_cnt++;
     305           0 :     }
     306           0 :   } else {
     307             :     /* Regular repair requests are not deduped.  Any potential regular
     308             :        shred request that will be made needs to be handled at the repair
     309             :        tile level to allow repair tile to re-request the same shred if
     310             :        it gets deduped. */
     311           0 :     uint nonce = fd_rnonce_ss_compute( policy->rnonce_ss, 1, ele->slot, iter->shred_idx, now );
     312           0 :     out = fd_repair_shred( repair, fd_policy_peer_select( policy ), now_ms, nonce, ele->slot, iter->shred_idx );
     313           0 :     ele->req_window_cnt++;
     314           0 :     if( FD_UNLIKELY( ele->first_req_ts == 0 ) ) ele->first_req_ts = fd_tickcount();
     315           0 :   }
     316           0 :   return out;
     317           0 : }
     318             : 
     319             : fd_policy_peer_t const *
     320          90 : fd_policy_peer_upsert( fd_policy_t * policy, fd_pubkey_t const * key, fd_ip4_port_t const * addr ) {
     321          90 :   fd_policy_peer_map_t * peer_map = policy->peers.map;
     322          90 :   fd_policy_peer_t * pool = policy->peers.pool;
     323          90 :   fd_policy_peer_t * peer = fd_policy_peer_map_ele_query( peer_map, key, NULL, pool );
     324          90 :   if( FD_UNLIKELY( !peer && fd_policy_peer_pool_free( pool ) ) ) {
     325          90 :     peer = fd_policy_peer_pool_ele_acquire( pool );
     326          90 :     peer->key  = *key;
     327          90 :     peer->ip4  = addr->addr;
     328          90 :     peer->port = addr->port;
     329          90 :     peer->req_cnt       = 0;
     330          90 :     peer->res_cnt       = 0;
     331          90 :     peer->first_req_ts  = 0;
     332          90 :     peer->last_req_ts   = 0;
     333          90 :     peer->first_resp_ts = 0;
     334          90 :     peer->last_resp_ts  = 0;
     335          90 :     peer->total_lat     = 0;
     336          90 :     peer->ewma_lat      = 0;
     337          90 :     peer->stake         = 0;
     338          90 :     peer->unanswered    = 0;
     339          90 :     peer->ping          = 0;
     340             : 
     341          90 :     fd_policy_peer_map_ele_insert( peer_map, peer, pool );
     342          90 :     fd_policy_peer_dlist_ele_push_tail( policy->peers.slow, peer, pool );
     343          90 :     return peer;
     344          90 :   }
     345           0 :   if( FD_LIKELY( peer ) ) {
     346           0 :     peer->ip4  = addr->addr;
     347           0 :     peer->port = addr->port;
     348           0 :   }
     349           0 :   return NULL;
     350          90 : }
     351             : 
     352             : fd_policy_peer_t *
     353        5805 : fd_policy_peer_query( fd_policy_t * policy, fd_pubkey_t const * key ) {
     354        5805 :   if( FD_UNLIKELY( memcmp( key->key, null_pubkey.key, 32UL ) == 0 ) ) return NULL;
     355        5805 :   fd_policy_peer_t * pool = policy->peers.pool;
     356        5805 :   return fd_policy_peer_map_ele_query( policy->peers.map, key, NULL, pool );
     357        5805 : }
     358             : 
     359             : int
     360           0 : fd_policy_peer_remove( fd_policy_t * policy, fd_pubkey_t const * key ) {
     361           0 :   fd_policy_peer_t * pool = policy->peers.pool;
     362           0 :   fd_policy_peer_t * peer = fd_policy_peer_map_ele_query( policy->peers.map, key, NULL, pool );
     363           0 :   if( FD_UNLIKELY( !peer ) ) return 0;
     364             : 
     365           0 :   ulong peer_idx = fd_policy_peer_pool_idx( pool, peer );
     366           0 :   fd_policy_peer_dlist_t * bucket = fd_policy_peer_latency_bucket( policy, peer->ewma_lat, peer->res_cnt );
     367             : 
     368             :   /* Advance iterators past the peer being removed while the dlist links
     369             :      are still intact, so iter_fwd_next can follow the forward pointer. */
     370           0 :   if( FD_UNLIKELY( policy->peers.select.fast_iter == peer_idx ) ) {
     371           0 :     policy->peers.select.fast_iter = fd_policy_peer_dlist_iter_fwd_next( policy->peers.select.fast_iter, bucket, pool );
     372           0 :   }
     373           0 :   if( FD_UNLIKELY( policy->peers.select.slow_iter == peer_idx ) ) {
     374           0 :     policy->peers.select.slow_iter = fd_policy_peer_dlist_iter_fwd_next( policy->peers.select.slow_iter, bucket, pool );
     375           0 :   }
     376             : 
     377           0 :   fd_policy_peer_dlist_ele_remove( bucket, peer, pool );
     378           0 :   fd_policy_peer_map_ele_remove  ( policy->peers.map, key, NULL, pool );
     379           0 :   fd_policy_peer_pool_ele_release( pool,   peer );
     380           0 :   return 1;
     381           0 : }
     382             : 
     383             : void
     384        1926 : fd_policy_peer_request_update( fd_policy_t * policy, fd_pubkey_t const * to ) {
     385        1926 :   fd_policy_peer_t * active = fd_policy_peer_query( policy, to );
     386        1926 :   if( FD_LIKELY( active ) ) {
     387        1926 :     active->req_cnt++;
     388        1926 :     active->unanswered++;
     389        1926 :     active->last_req_ts = fd_tickcount();
     390        1926 :     if( FD_UNLIKELY( active->first_req_ts == 0 ) ) active->first_req_ts = active->last_req_ts;
     391        1926 :   }
     392        1926 : }
     393             : 
     394             : void
     395        1728 : fd_policy_peer_response_update( fd_policy_t * policy, fd_pubkey_t const * to, long rtt /* ns */ ) {
     396        1728 :   fd_policy_peer_t * peer = fd_policy_peer_query( policy, to );
     397        1728 :   if( FD_LIKELY( peer ) ) {
     398        1728 :     long now = fd_tickcount();
     399        1728 :     fd_policy_peer_dlist_t * prev_bucket = fd_policy_peer_latency_bucket( policy, peer->ewma_lat, peer->res_cnt );
     400        1728 :     peer->res_cnt++;
     401        1728 :     peer->unanswered = 0;
     402        1728 :     if( FD_UNLIKELY( peer->first_resp_ts == 0 ) ) peer->first_resp_ts = now;
     403        1728 :     peer->last_resp_ts = now;
     404        1728 :     peer->total_lat   += rtt;
     405             : 
     406        1728 :     if( FD_UNLIKELY( peer->res_cnt == 1 ) ) {
     407          21 :       peer->ewma_lat = rtt;
     408        1707 :     } else {
     409        1707 :       peer->ewma_lat = peer->ewma_lat - peer->ewma_lat / (long)FD_POLICY_EWMA_ALPHA_DENOM
     410        1707 :                       + rtt / (long)FD_POLICY_EWMA_ALPHA_DENOM;
     411        1707 :     }
     412        1728 :     fd_policy_peer_dlist_t * new_bucket = fd_policy_peer_latency_bucket( policy, peer->ewma_lat, peer->res_cnt );
     413        1728 :     if( prev_bucket != new_bucket ) {
     414             :       /* Advance stale iterators */
     415          21 :       ulong peer_idx = fd_policy_peer_pool_idx( policy->peers.pool, peer );
     416          21 :       if( FD_UNLIKELY( policy->peers.select.fast_iter == peer_idx ) ) policy->peers.select.fast_iter = fd_policy_peer_dlist_iter_fwd_next( policy->peers.select.fast_iter, policy->peers.fast, policy->peers.pool );
     417          21 :       if( FD_UNLIKELY( policy->peers.select.slow_iter == peer_idx ) ) policy->peers.select.slow_iter = fd_policy_peer_dlist_iter_fwd_next( policy->peers.select.slow_iter, policy->peers.slow, policy->peers.pool );
     418             : 
     419          21 :       fd_policy_peer_dlist_ele_remove   ( prev_bucket, peer, policy->peers.pool );
     420          21 :       fd_policy_peer_dlist_ele_push_tail( new_bucket,  peer, policy->peers.pool );
     421          21 :     }
     422        1728 :   }
     423        1728 : }
     424             : 
     425             : void
     426          39 : fd_policy_set_turbine_slot0( fd_policy_t * policy, ulong slot ) {
     427          39 :   policy->turbine_slot0 = slot;
     428          39 : }
     429             : 

Generated by: LCOV version 1.14