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 :
|