Line data Source code
1 : /* The rotor tile is responsible for repairing missing shreds that were
2 : not received via Turbine or missing slots of interest from Votor.
3 : The goal is to ensure that slots we "care" about have their FEC sets
4 : inserted into store.
5 :
6 : Most of rotor is copied over from repair tile, */
7 :
8 : #define _GNU_SOURCE
9 :
10 : #include "../genesis/fd_genesi_tile.h"
11 : #include "../../disco/topo/fd_topo.h"
12 : #include "../../disco/fd_clock_tile.h"
13 : #include "generated/fd_rotor_tile_seccomp.h"
14 : #include "../../disco/keyguard/fd_keyload.h"
15 : #include "../../disco/keyguard/fd_keyguard.h"
16 : #include "../../disco/keyguard/fd_keyswitch.h"
17 : #include "../../disco/metrics/fd_metrics.h"
18 : #include "../../disco/net/fd_net_tile.h"
19 : #include "../../disco/shred/fd_rnonce_ss.h"
20 : #include "../../disco/shred/fd_shred_tile.h"
21 : #include "fd_rotor_tile.h"
22 : #include "fd_rotor_tile_private.h"
23 : #include "../replay/fd_replay_tile.h"
24 : #include "../votor/fd_votor_tile.h"
25 : #include "../../discof/restore/utils/fd_ssmsg.h"
26 : #include "../../util/net/fd_net_headers.h"
27 : #include "../../util/pod/fd_pod_format.h"
28 : #include "../../tango/fd_tango_base.h"
29 :
30 : #include "../repair/fd_repair_metrics.h"
31 : #include "../repair/fd_inflight.h"
32 : #include "../repair/fd_repair.h"
33 : #include "../repair/fd_policy.h"
34 :
35 : #include "../../discof/chainer/fd_chainer.h"
36 : #include "../../disco/store/fd_store.h"
37 : #include "../../flamenco/alpenglow/fd_block_marker_serde.h"
38 :
39 : #define DEBUG_LOGGING 0
40 :
41 : #define IN_KIND_CONTACT (0)
42 186 : #define IN_KIND_NET (1)
43 5520 : #define IN_KIND_SHRED (2)
44 45 : #define IN_KIND_SIGN (3)
45 0 : #define IN_KIND_SNAP (4)
46 0 : #define IN_KIND_GOSSIP (5)
47 0 : #define IN_KIND_GENESIS (6)
48 66 : #define IN_KIND_REPLAY (7)
49 81 : #define IN_KIND_VOTOR (8) /* Alpenglow rooting */
50 :
51 : #define MAX_IN_LINKS (32)
52 : #define MAX_SHRED_TILE_CNT ( 16UL )
53 : #define MAX_SIGN_TILE_CNT ( 16UL )
54 :
55 : /* Max number of pending repair requests recently made to keep track of.
56 : Calculated generally as we estimate around 50k/s/core to sign
57 : requests. Assuming an over-provisioned 4 sign tiles just for repair,
58 : this means we can make up to ~200k requests per second. With a dedup
59 : timeout of 80ms, this means we can make up to ~16k requests within
60 : the dedup timeout window. We round up to the next power of two to
61 : get the dedup cache max. Since we are sizing the dedup cache for a
62 : generous margin, and this number not particularly fragile or
63 : sensitive, we can leave it static. */
64 0 : #define FD_REQLIM_CACHE_MAX (1<<20)
65 :
66 : /* static map from request type to metric array index */
67 : static uint metric_index[AG_REPAIR_KIND_SHRED_FOR_BLOCK_ID + 1] = {
68 : [FD_REPAIR_KIND_PONG] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_PONG_IDX,
69 : [FD_REPAIR_KIND_SHRED] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_NEEDED_WINDOW_IDX,
70 : [FD_REPAIR_KIND_HIGHEST_SHRED] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_NEEDED_HIGHEST_WINDOW_IDX,
71 : [FD_REPAIR_KIND_ORPHAN] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_NEEDED_ORPHAN_IDX,
72 : [AG_REPAIR_KIND_PARENT_FEC_COUNT] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_PARENT_FEC_COUNT_IDX,
73 : [AG_REPAIR_KIND_FEC_ROOT] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_FEC_ROOT_IDX,
74 : [AG_REPAIR_KIND_SHRED_FOR_BLOCK_ID] = FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_V_SHRED_BLOCK_ID_IDX,
75 : };
76 :
77 : struct ctx {
78 : fd_clock_tile_t clock[1];
79 :
80 : ulong repair_seed;
81 :
82 : /* When set (alpenglow only), the repair policy walk emits ONLY
83 : block-id requests (ShredForBlockId, driven by known block_ids and
84 : the event-driven getParentAndFecSetCount/getFecRoot path). All
85 : legacy positional emissions -- HighestShred, window Shred, Orphan,
86 : and the orphan-pass shred-0 -- are suppressed. Used to exercise /
87 : test the block-id repair + catchup path in isolation. */
88 : int block_id_repair_only;
89 :
90 : /* When set, publish_fec_replay re-publishes the entire ancestry path
91 : of FECs -- from the chainer root down to the FEC being delivered,
92 : in root-to-target order -- on every delivery, instead of just the
93 : single delivered FEC. Lets replay reconstruct a fork from root
94 : without relying on incremental delivery. The path is queued onto
95 : deliver_queue and drained one FEC per after_credit. */
96 : int deliver_from_root;
97 : out_ele_t * deliver_queue; /* sized to the chainer's FEC capacity */
98 :
99 : fd_keyswitch_t * keyswitch;
100 : int halt_signing;
101 :
102 : fd_chainer_t * chainer; /* alpenglow chainer */
103 : fd_store_t * store; /* rotor publishes/removes FEC sets to/from the store */
104 : fd_store_map_t store_map[1];
105 : fd_policy_t * policy;
106 : fd_reqlim_t * dedup;
107 : fd_inflights_t * inflights;
108 : fd_repair_t * protocol;
109 :
110 : fd_pubkey_t identity_public_key;
111 :
112 : fd_wksp_t * wksp;
113 :
114 : fd_stem_context_t * stem;
115 :
116 : uchar in_kind[ MAX_IN_LINKS ];
117 : in_ctx_t in_links[ MAX_IN_LINKS ];
118 :
119 : int skip_frag;
120 :
121 : out_ctx_t net_out_ctx[1];
122 : out_ctx_t repair_out_ctx[1];
123 :
124 : /* repair_sign links (to sign tiles 1+) - for round-robin
125 : distribution */
126 : ulong repair_sign_cnt;
127 : out_ctx_t repair_sign_out_ctx[ MAX_SIGN_TILE_CNT ];
128 :
129 : ulong sign_rrobin_idx;
130 :
131 : /* Pending sign requests for async operations */
132 :
133 : uint pending_key_next;
134 : sign_req_t * signs_map; /* contains any request currently in the repair->sign or sign->repair dcache */
135 : sign_pending_t * toss_queue; /* contains any pong or initial warmup request waiting to be dispatched to repair->sign. Size is 2*FD_REPAIR_PEER_MAX */
136 : fd_repair_msg_t * meta_queue; /* contains any alpenglow request waiting to be dispatched to sign->repair. Sized to one block's FEC sets (max_shreds_per_block/FD_FEC_SHRED_CNT) */
137 :
138 : ushort net_id;
139 :
140 : /* Buffers for incoming unreliable frags */
141 : uchar net_buf[ FD_NET_MTU ];
142 : uchar sign_buf[ sizeof(fd_ed25519_sig_t) ];
143 :
144 : /* Store chunk for incoming reliable frags */
145 : ulong chunk;
146 : ulong snap_out_chunk; /* store second to last chunk for snap_out */
147 :
148 : fd_ip4_udp_hdrs_t intake_hdr[1];
149 :
150 : fd_rnonce_ss_t repair_nonce_ss[1];
151 : uint ag_nonce; /* simple incrementing nonce for alpenglow requests */
152 :
153 : ulong manifest_slot;
154 : struct {
155 : ulong send_pkt_cnt;
156 : ulong sent_pkt_types[FD_METRICS_ENUM_REPAIR_SENT_REQUEST_TYPE_CNT];
157 : ulong current_slot;
158 : ulong old_shred;
159 : ulong last_requested_slot;
160 : ulong last_requested_orphan;
161 : ulong sign_tile_unavail;
162 : ulong rerequest;
163 : ulong malformed_ping;
164 : ulong unknown_peer_ping;
165 : ulong fail_sigverify_ping;
166 : fd_histf_t slot_compl_time[ 1 ];
167 : fd_histf_t response_latency[ 1 ];
168 :
169 : ulong failed_shred_block_id_cnt;
170 : ulong failed_fec_root_cnt;
171 : ulong failed_parent_fec_count_cnt;
172 :
173 : ulong fecs_delivered; /* diagnostic: FECs pushed to replay via out_queue */
174 : } metrics[ 1 ];
175 :
176 : /* Slot-level metrics */
177 :
178 : fd_repair_metrics_t * slot_metrics;
179 : ulong turbine_slot0; // catchup considered complete after this slot
180 : };
181 : typedef struct ctx ctx_t;
182 :
183 : FD_FN_CONST static inline ulong
184 0 : scratch_align( void ) {
185 0 : return 128UL;
186 0 : }
187 :
188 : FD_FN_PURE static inline ulong
189 0 : loose_footprint( fd_topo_tile_t const * tile FD_PARAM_UNUSED ) {
190 0 : return 1UL * FD_SHMEM_GIGANTIC_PAGE_SZ;
191 0 : }
192 :
193 : FD_FN_PURE static inline ulong
194 0 : scratch_footprint( fd_topo_tile_t const * tile ) {
195 0 : ulong total_sign_depth = tile->rotor.repair_sign_depth * tile->rotor.repair_sign_cnt;
196 0 : int lg_sign_depth = fd_ulong_find_msb( fd_ulong_pow2_up(total_sign_depth) ) + 1;
197 0 : ulong fec_blk_max = tile->rotor.max_shreds_per_block / FD_FEC_SHRED_CNT;
198 :
199 0 : ulong l = FD_LAYOUT_INIT;
200 0 : l = FD_LAYOUT_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
201 0 : l = FD_LAYOUT_APPEND( l, fd_repair_align(), fd_repair_footprint () );
202 0 : l = FD_LAYOUT_APPEND( l, fd_chainer_align(), fd_chainer_footprint ( tile->rotor.slot_max, tile->rotor.max_shreds_per_block ) );
203 0 : l = FD_LAYOUT_APPEND( l, fd_policy_align(), fd_policy_footprint ( FD_REPAIR_PEER_MAX ) );
204 0 : l = FD_LAYOUT_APPEND( l, fd_reqlim_align(), fd_reqlim_footprint ( FD_REQLIM_CACHE_MAX ) );
205 0 : l = FD_LAYOUT_APPEND( l, fd_inflights_align(), fd_inflights_footprint () );
206 0 : l = FD_LAYOUT_APPEND( l, fd_signs_map_align(), fd_signs_map_footprint ( lg_sign_depth ) );
207 0 : l = FD_LAYOUT_APPEND( l, toss_queue_align(), toss_queue_footprint () );
208 0 : l = FD_LAYOUT_APPEND( l, meta_queue_align(), meta_queue_footprint ( fec_blk_max ) );
209 0 : l = FD_LAYOUT_APPEND( l, fd_repair_metrics_align(), fd_repair_metrics_footprint() );
210 0 : l = FD_LAYOUT_APPEND( l, out_queue_align(), out_queue_footprint ( (ulong)tile->rotor.slot_max * FD_CHAINER_SLOT_VER_MAX * fec_blk_max ) );
211 0 : return FD_LAYOUT_FINI( l, scratch_align() );
212 0 : }
213 :
214 : /* Below functions manage the current pending sign requests. */
215 :
216 : static sign_req_t *
217 : sign_map_insert( ctx_t * ctx,
218 : fd_repair_msg_t const * msg,
219 2151 : pong_data_t const * opt_pong_data ) {
220 2151 : if( FD_UNLIKELY( fd_signs_map_key_cnt( ctx->signs_map )==fd_signs_map_key_max( ctx->signs_map ) ) ) return NULL;
221 :
222 2151 : sign_req_t * pending = fd_signs_map_insert( ctx->signs_map, ctx->pending_key_next++ );
223 2151 : if( FD_UNLIKELY( !pending ) ) return NULL; /* Not possible, unless the same key is used twice. */
224 2151 : pending->msg = *msg;
225 2151 : pending->buflen = fd_repair_sz( msg );
226 2151 : if( FD_UNLIKELY( opt_pong_data ) ) pending->pong_data = *opt_pong_data;
227 2151 : return pending;
228 2151 : }
229 :
230 : static int
231 : sign_map_remove( ctx_t * ctx,
232 2151 : ulong key ) {
233 2151 : sign_req_t * pending = fd_signs_map_query( ctx->signs_map, key, NULL );
234 2151 : if( FD_UNLIKELY( !pending ) ) return -1;
235 2151 : fd_signs_map_remove( ctx->signs_map, pending );
236 2151 : return 0;
237 2151 : }
238 :
239 : static void
240 : send_packet( ctx_t * ctx,
241 : fd_stem_context_t * stem,
242 : uint dst_ip_addr,
243 : ushort dst_port,
244 : uint src_ip_addr,
245 : uchar const * payload,
246 : ulong payload_sz,
247 2151 : ulong tsorig ) {
248 2151 : ctx->metrics->send_pkt_cnt++;
249 2151 : uchar * packet = fd_chunk_to_laddr( ctx->net_out_ctx->mem, ctx->net_out_ctx->chunk );
250 2151 : fd_ip4_udp_hdrs_t * hdr = (fd_ip4_udp_hdrs_t *)packet;
251 2151 : *hdr = *ctx->intake_hdr;
252 :
253 2151 : fd_ip4_hdr_t * ip4 = hdr->ip4;
254 2151 : ip4->saddr = src_ip_addr;
255 2151 : ip4->daddr = dst_ip_addr;
256 2151 : ip4->net_id = fd_ushort_bswap( ctx->net_id++ );
257 2151 : ip4->check = 0U;
258 2151 : ip4->net_tot_len = fd_ushort_bswap( (ushort)(payload_sz + sizeof(fd_ip4_hdr_t)+sizeof(fd_udp_hdr_t)) );
259 2151 : ip4->check = fd_ip4_hdr_check_fast( ip4 );
260 :
261 2151 : fd_udp_hdr_t * udp = hdr->udp;
262 2151 : udp->net_dport = dst_port;
263 2151 : udp->net_len = fd_ushort_bswap( (ushort)(payload_sz + sizeof(fd_udp_hdr_t)) );
264 2151 : fd_memcpy( packet+sizeof(fd_ip4_udp_hdrs_t), payload, payload_sz );
265 2151 : hdr->udp->check = 0U;
266 :
267 2151 : ulong tspub = fd_frag_meta_ts_comp( fd_tickcount() );
268 2151 : ulong sig = fd_disco_netmux_sig( dst_ip_addr, dst_port, dst_ip_addr, DST_PROTO_OUTGOING, sizeof(fd_ip4_udp_hdrs_t) );
269 2151 : ulong packet_sz = payload_sz + sizeof(fd_ip4_udp_hdrs_t);
270 2151 : ulong chunk = ctx->net_out_ctx->chunk;
271 2151 : fd_stem_publish( stem, ctx->net_out_ctx->idx, sig, chunk, packet_sz, 0UL, tsorig, tspub );
272 2151 : ctx->net_out_ctx->chunk = fd_dcache_compact_next( chunk, packet_sz, ctx->net_out_ctx->chunk0, ctx->net_out_ctx->wmark );
273 2151 : }
274 :
275 : /* Returns a sign_out context with max available credits.
276 : If no sign_out context has available credits, returns NULL. */
277 : static out_ctx_t *
278 2400 : sign_avail_credits( ctx_t * ctx ) {
279 2400 : out_ctx_t * sign_out = NULL;
280 2400 : ulong max_credits = 0;
281 4800 : for( uint i=0; i<ctx->repair_sign_cnt; i++ ) {
282 2400 : if( ctx->repair_sign_out_ctx[i].credits > max_credits ) {
283 2400 : max_credits = ctx->repair_sign_out_ctx[i].credits;
284 2400 : sign_out = &ctx->repair_sign_out_ctx[i];
285 2400 : }
286 2400 : }
287 2400 : return sign_out;
288 2400 : }
289 :
290 : /* Prepares the signing preimage and publishes a signing request that
291 : will be signed asynchronously by the sign tile. The signed data will
292 : be returned via dcache as a frag. */
293 : static void
294 : fd_repair_send_sign_request( ctx_t * ctx,
295 : out_ctx_t * sign_out,
296 : fd_repair_msg_t const * msg,
297 2151 : pong_data_t const * opt_pong_data ) {
298 :
299 2151 : if( FD_UNLIKELY( ctx->halt_signing ) ) FD_LOG_CRIT(( "can't dispatch sign requests while halting signing" ));
300 :
301 : /* New sign request */
302 2151 : sign_req_t * pending = sign_map_insert( ctx, msg, opt_pong_data );
303 2151 : if( FD_UNLIKELY( !pending ) ) return;
304 :
305 2151 : ulong sig = 0;
306 2151 : ulong preimage_sz = 0;
307 2151 : uchar * dst = fd_chunk_to_laddr( sign_out->mem, sign_out->chunk );
308 :
309 2151 : if( FD_UNLIKELY( msg->kind == FD_REPAIR_KIND_PONG ) ) {
310 0 : uchar pre_image[FD_REPAIR_PONG_PREIMAGE_SZ];
311 0 : preimage_pong( &opt_pong_data->hash, pre_image );
312 0 : preimage_sz = FD_REPAIR_PONG_PREIMAGE_SZ;
313 0 : fd_memcpy( dst, pre_image, preimage_sz );
314 0 : sig = ((ulong)pending->key << 32) | (uint)FD_KEYGUARD_SIGN_TYPE_SHA256_ED25519;
315 2151 : } else {
316 : /* Sign and prepare the message directly into the pending buffer */
317 2151 : uchar * preimage = preimage_req( &pending->msg, &preimage_sz );
318 2151 : fd_memcpy( dst, preimage, preimage_sz );
319 2151 : sig = ((ulong)pending->key << 32) | (uint)FD_KEYGUARD_SIGN_TYPE_ED25519;
320 2151 : }
321 :
322 2151 : fd_stem_publish( ctx->stem, sign_out->idx, sig, sign_out->chunk, preimage_sz, 0UL, 0UL, 0UL );
323 2151 : sign_out->chunk = fd_dcache_compact_next( sign_out->chunk, preimage_sz, sign_out->chunk0, sign_out->wmark );
324 :
325 2151 : ctx->metrics->sent_pkt_types[metric_index[msg->kind]]++;
326 2151 : sign_out->credits--;
327 2151 : }
328 :
329 : /* meta_inflight_record tracks a metadata request (getParentAndFecSetCount
330 : or getFecSetRoot) in inflights so its response can be matched and it
331 : can be redispatched if none arrives. */
332 : static void
333 : meta_inflight_record( ctx_t * ctx,
334 : fd_repair_msg_t const * msg,
335 171 : long now ) {
336 171 : if( FD_LIKELY( msg->kind==AG_REPAIR_KIND_FEC_ROOT ) ) fd_inflights_meta_insert( ctx->inflights, msg->fec_set_root.nonce, AG_REPAIR_KIND_FEC_ROOT, &msg->header.to, msg->fec_set_root.slot, &msg->fec_set_root.block_id, msg->fec_set_root.fec_set_idx, now );
337 60 : else fd_inflights_meta_insert( ctx->inflights, msg->parent_fec_set_count.nonce, AG_REPAIR_KIND_PARENT_FEC_COUNT, &msg->header.to, msg->parent_fec_set_count.slot, &msg->parent_fec_set_count.block_id, 0U, now );
338 171 : }
339 :
340 : /* meta_queue_push_safe queues a metadata request for dispatch. If the
341 : queue is full the request is instead parked straight in inflights,
342 : never sent, and picked up by the redispatch path once it ages out. */
343 : static void
344 : meta_queue_push_safe( ctx_t * ctx,
345 : fd_repair_msg_t const * msg,
346 147 : long now ) {
347 147 : if( FD_LIKELY( !meta_queue_full( ctx->meta_queue ) ) ) meta_queue_push( ctx->meta_queue, *msg );
348 6 : else meta_inflight_record( ctx, msg, now );
349 147 : }
350 :
351 : static inline int
352 : before_frag( ctx_t * ctx,
353 : ulong in_idx,
354 : ulong seq FD_PARAM_UNUSED,
355 5550 : ulong sig ) {
356 5550 : uint in_kind = ctx->in_kind[ in_idx ];
357 5550 : if( FD_LIKELY ( in_kind==IN_KIND_NET ) ) return fd_disco_netmux_sig_proto( sig )!=DST_PROTO_REPAIR;
358 5550 : if( FD_UNLIKELY( in_kind==IN_KIND_SHRED ) ) return fd_int_if( ctx->chainer->root==ULONG_MAX, -1, 0 ); /* not ready to read frag */
359 75 : if( FD_UNLIKELY( in_kind==IN_KIND_GOSSIP ) ) {
360 0 : return sig!=FD_GOSSIP_UPDATE_TAG_CONTACT_INFO &&
361 0 : sig!=FD_GOSSIP_UPDATE_TAG_CONTACT_INFO_REMOVE;
362 0 : }
363 75 : if( FD_UNLIKELY( in_kind==IN_KIND_REPLAY ) ) return sig!=REPLAY_SIG_MISSING_FEC &&
364 33 : sig!=REPLAY_SIG_ROOT_ADVANCED;
365 42 : if( FD_UNLIKELY( in_kind==IN_KIND_VOTOR ) ) return sig!=FD_VOTOR_SIG_REPAIR;
366 0 : return 0;
367 42 : }
368 :
369 : static void
370 : during_frag( ctx_t * ctx,
371 : ulong in_idx,
372 : ulong seq FD_PARAM_UNUSED,
373 : ulong sig,
374 : ulong chunk,
375 : ulong sz,
376 5532 : ulong ctl ) {
377 5532 : ctx->skip_frag = 0;
378 :
379 5532 : uint in_kind = ctx->in_kind[ in_idx ];
380 5532 : in_ctx_t const * in_ctx = &ctx->in_links[ in_idx ];
381 5532 : ctx->chunk = chunk;
382 :
383 5532 : if( FD_UNLIKELY( in_kind==IN_KIND_NET ) ) {
384 0 : ulong hdr_sz = fd_disco_netmux_sig_hdr_sz( sig );
385 0 : FD_TEST( hdr_sz <= sz ); /* Should be ensured by the net tile */
386 0 : uchar const * dcache_entry = fd_net_rx_translate_frag( &in_ctx->net_rx, chunk, ctl, sz );
387 0 : fd_memcpy( ctx->net_buf, dcache_entry, sz );
388 0 : return;
389 0 : }
390 :
391 5532 : if( FD_UNLIKELY( in_kind==IN_KIND_GENESIS ) ) {
392 0 : FD_TEST( sizeof(fd_genesis_meta_t)<=sig );
393 0 : return;
394 0 : }
395 :
396 5532 : if( FD_UNLIKELY( sz!=0UL && ( chunk<in_ctx->chunk0 || chunk>in_ctx->wmark || sz>in_ctx->mtu ) ) )
397 0 : FD_LOG_ERR(( "chunk %lu %lu corrupt, not in range [%lu,%lu] in kind %u", chunk, sz, in_ctx->chunk0, in_ctx->wmark, in_kind ));
398 :
399 5532 : if( FD_UNLIKELY( in_kind==IN_KIND_SNAP ) ) {
400 0 : if( FD_UNLIKELY( fd_ssmsg_sig_message( sig )!=FD_SSMSG_DONE ) ) ctx->snap_out_chunk = chunk;
401 0 : return;
402 0 : }
403 :
404 5532 : if( FD_UNLIKELY( in_kind==IN_KIND_SIGN ) ) {
405 : /* sign_repair is unreliable, so we copy the frag for convention.
406 : Theoretically impossible to overrun. */
407 0 : uchar const * dcache_entry = fd_chunk_to_laddr_const( in_ctx->mem, chunk );
408 0 : fd_memcpy( ctx->sign_buf, dcache_entry, sz );
409 0 : return;
410 0 : }
411 5532 : }
412 :
413 : static inline void
414 : after_snap( ctx_t * ctx,
415 : ulong sig,
416 0 : uchar const * chunk ) {
417 0 : if( FD_UNLIKELY( fd_ssmsg_sig_message( sig )!=FD_SSMSG_DONE ) ) return;
418 0 : fd_snapshot_manifest_t * manifest = (fd_snapshot_manifest_t *)chunk;
419 :
420 0 : fd_chainer_init( ctx->chainer, manifest->slot, (fd_hash_t *)fd_type_pun( manifest->block_id ) );
421 0 : }
422 :
423 : static inline void
424 0 : after_gossip( ctx_t * ctx, fd_gossip_update_message_t const * msg, ulong sig ) {
425 0 : switch( sig ) {
426 0 : case FD_GOSSIP_UPDATE_TAG_CONTACT_INFO_REMOVE: {
427 0 : fd_policy_peer_remove( ctx->policy, fd_type_pun_const( msg->origin ) );
428 0 : break;
429 0 : }
430 0 : case FD_GOSSIP_UPDATE_TAG_CONTACT_INFO: {
431 0 : fd_gossip_contact_info_t const * contact_info = msg->contact_info->value;
432 0 : fd_ip4_port_t repair_peer;
433 0 : repair_peer.addr = contact_info->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_SERVE_REPAIR ].is_ipv6 ? 0U : contact_info->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_SERVE_REPAIR ].ip4;
434 0 : repair_peer.port = contact_info->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_SERVE_REPAIR ].port;
435 0 : if( FD_UNLIKELY( !repair_peer.addr || !repair_peer.port ) ) return;
436 0 : fd_policy_peer_t const * peer = fd_policy_peer_upsert( ctx->policy, fd_type_pun_const( msg->origin ), &repair_peer );
437 0 : if( FD_LIKELY( peer && !toss_queue_full( ctx->toss_queue ) ) ) {
438 : /* The repair process uses a Ping-Pong protocol that incurs one
439 : round-trip time (RTT) for the initial repair request. To
440 : optimize this, we proactively send a placeholder repair
441 : request as soon as we receive a peer's contact information
442 : for the first time, effectively prepaying the RTT cost. */
443 0 : fd_repair_msg_t * init = fd_repair_shred( ctx->protocol, fd_type_pun_const( msg->origin ), (ulong)fd_log_wallclock()/1000000L, 0, 0, 0 );
444 0 : toss_queue_push( ctx->toss_queue, (sign_pending_t){ .msg = *init } );
445 0 : }
446 0 : break;
447 0 : }
448 0 : default: FD_LOG_ERR(( "bad gossip sig %lu", sig ));
449 0 : }
450 0 : }
451 :
452 : static inline void
453 : after_sign( ctx_t * ctx,
454 : ulong in_idx,
455 : ulong sig,
456 2151 : fd_stem_context_t * stem ) {
457 2151 : ulong pending_key = sig >> 32;
458 : /* Look up the pending request. Since the rotor_sign links are
459 : reliable, the incoming sign_repair fragments represent a complete
460 : set of the previously sent outgoing messages. However, with
461 : multiple sign tiles, the responses may arrive interleaved. */
462 :
463 : /* Find which sign tile sent this response and increment its
464 : credits */
465 2151 : for( uint i=0; i<ctx->repair_sign_cnt; i++ ) {
466 2151 : if( ctx->repair_sign_out_ctx[i].in_idx == in_idx ) {
467 2151 : if( FD_LIKELY( ctx->repair_sign_out_ctx[i].credits < ctx->repair_sign_out_ctx[i].max_credits ) ) ctx->repair_sign_out_ctx[i].credits++;
468 2151 : break;
469 2151 : }
470 2151 : }
471 :
472 2151 : sign_req_t * pending_ = fd_signs_map_query( ctx->signs_map, pending_key, NULL );
473 2151 : if( FD_UNLIKELY( !pending_ ) ) FD_LOG_CRIT(( "No pending request found for key %lu", pending_key )); /* implies either bad programmer error or something happened with sign tile */
474 :
475 2151 : sign_req_t pending[1] = { *pending_ }; /* Make a copy of the pending request so we can sign_map_remove immediately. */
476 2151 : sign_map_remove( ctx, pending_key );
477 :
478 : /* This is a pong message */
479 2151 : if( FD_UNLIKELY( pending->msg.kind == FD_REPAIR_KIND_PONG ) ) {
480 0 : fd_policy_peer_t * peer = fd_policy_peer_query( ctx->policy, &pending->pong_data.key );
481 0 : if( FD_LIKELY( peer && peer->ping ) ) peer->ping--; /* prevent underflow if the peer was removed/readded */
482 :
483 0 : fd_memcpy( pending->msg.pong.sig, ctx->sign_buf, 64UL );
484 0 : send_packet( ctx, stem, pending->pong_data.peer_addr.addr, pending->pong_data.peer_addr.port, pending->pong_data.daddr, pending->buf, fd_repair_sz( &pending->msg ), fd_frag_meta_ts_comp( fd_tickcount() ) );
485 0 : return;
486 0 : }
487 :
488 : /* Inject the signature into the pending request */
489 2151 : fd_memcpy( pending->buf + 4, ctx->sign_buf, 64UL );
490 2151 : uint src_ip4 = 0U;
491 :
492 : /* This is a warmup message */
493 2151 : if( FD_UNLIKELY( pending->msg.kind == FD_REPAIR_KIND_SHRED && pending->msg.shred.slot == 0 ) ) {
494 0 : fd_policy_peer_t * peer = fd_policy_peer_query( ctx->policy, &pending->msg.shred.to );
495 0 : if( FD_UNLIKELY( peer ) ) send_packet( ctx, stem, peer->ip4, peer->port, src_ip4, pending->buf, pending->buflen, fd_frag_meta_ts_comp( fd_tickcount() ) );
496 0 : else { /* This is a warmup request for a peer that is no longer active. There's no reason to pick another peer for a warmup rq, so just drop it. */ }
497 0 : return;
498 0 : }
499 :
500 : /* This is a regular repair shred request
501 :
502 : We need to ensure we always send out any shred requests we have,
503 : because policy_next has no way to revisit a shred. But the fact
504 : that peers can drop out of the peer list makes this complicated.
505 : If the peer is still there (common), it's fine. If the peer is not
506 : there, we can add this request to the inflights table, pretend
507 : we've sent it and let the inflight timeout request it down the
508 : line. */
509 :
510 2151 : fd_policy_peer_t * active = fd_policy_peer_query( ctx->policy, &pending->msg.shred.to );
511 2151 : if( FD_UNLIKELY( !active ) ) {
512 : /* Already added to the inflights table, pretend we've sent it
513 : and let the inflight timeout request it down the line. */
514 0 : return;
515 0 : }
516 : /* Happy path - all is well, our peer didn't drop out from beneath
517 : us. */
518 2151 : if( FD_UNLIKELY( pending->msg.kind == FD_REPAIR_KIND_ORPHAN ) ) ctx->metrics->last_requested_orphan = pending->msg.orphan.slot;
519 2151 : else ctx->metrics->last_requested_slot = pending->msg.shred.slot;
520 :
521 2151 : send_packet( ctx, stem, active->ip4, active->port, src_ip4, pending->buf, pending->buflen, fd_frag_meta_ts_comp( fd_tickcount() ) );
522 2151 : }
523 :
524 : /* takes ping after hdr strip. ip4/udp point at the stripped headers
525 : (needed to address the pong back to the sender). */
526 : static inline void
527 : after_ping( ctx_t * ctx,
528 : uchar const * data,
529 : ulong data_sz,
530 : fd_ip4_hdr_t const * ip4,
531 0 : fd_udp_hdr_t const * udp ) {
532 0 : fd_ip4_port_t peer_addr = { .addr=ip4->saddr, .port=udp->net_sport };
533 :
534 0 : fd_repair_ping_t ping[1];
535 0 : int err = fd_repair_ping_de( ping, data, data_sz );
536 0 : if( FD_UNLIKELY( err ) ) {
537 0 : ctx->metrics->malformed_ping++;
538 0 : return;
539 0 : }
540 :
541 0 : fd_policy_peer_t * peer = fd_policy_peer_query( ctx->policy, &ping->ping.from );
542 0 : if( FD_UNLIKELY( !peer ) ) {
543 0 : ctx->metrics->unknown_peer_ping++;
544 0 : return;
545 0 : }
546 0 : if( FD_UNLIKELY( peer->ping ) ) return;
547 0 : if( FD_UNLIKELY( toss_queue_full( ctx->toss_queue ) ) ) return;
548 :
549 0 : fd_sha512_t sha[1];
550 0 : if( FD_UNLIKELY( FD_ED25519_SUCCESS != fd_ed25519_verify( ping->ping.hash.uc, 32UL, ping->ping.sig, ping->ping.from.uc, sha ) ) ) {
551 0 : ctx->metrics->fail_sigverify_ping++;
552 0 : return;
553 0 : }
554 :
555 : /* Any gossip peer can send a ping, but they are bounded to at most
556 : one ping in the queue so they can't evict others' pings without
557 : multiple gossip identities. */
558 :
559 0 : fd_repair_msg_t * pong = fd_repair_pong( ctx->protocol, &ping->ping.hash );
560 0 : toss_queue_push( ctx->toss_queue, (sign_pending_t){ .msg = *pong, .pong_data = { .peer_addr = peer_addr, .hash = ping->ping.hash, .daddr = ip4->daddr, .key = ping->ping.from } } );
561 0 : peer->ping++;
562 0 : }
563 :
564 : /* This is a response for an Alpenglow repair request type, which
565 : returns metadata about a verified slot - not a shred. Specifically,
566 : responses for parent_and_fec_set_count and fec_set_root are routed
567 : directly to repair tile, as they do not contain a shred. */
568 :
569 : static inline void
570 : after_alpen_meta_repair( ctx_t * ctx,
571 141 : ag_repair_response_t * response ) {
572 141 : uint nonce = response->nonce;
573 :
574 141 : fd_inflight_t request[1];
575 141 : if( FD_UNLIKELY( !fd_inflights_meta_match( ctx->inflights, nonce, request ) ) ) return; /* unsolicited */
576 :
577 138 : uint kind = request->key.kind;
578 138 : ulong slot = request->key.slot;
579 138 : uint fec_set_idx = request->key.idx;
580 138 : fd_hash_t block_id = request->block_id;
581 138 : if( FD_UNLIKELY( slot <= ctx->chainer->root ) ) return; /* rooted in flight: obsolete */
582 :
583 138 : fd_pubkey_t to = {0}; /* peer is chosen when the meta queue drains */
584 138 : long now = fd_clock_tile_now( ctx->clock );
585 138 : ulong now_ms = (ulong)(now/(long)1e6);
586 :
587 : /* A response of the wrong kind for this nonce is rejected and the
588 : original request re-queued under a fresh nonce, so one bad peer
589 : response cannot strand the block. */
590 138 : if( FD_UNLIKELY( ( response->kind==AG_REPAIR_RESPONSE_PARENT_FEC_SET_COUNT && kind!=AG_REPAIR_KIND_PARENT_FEC_COUNT ) ||
591 138 : ( response->kind==AG_REPAIR_RESPONSE_FEC_SET_ROOT && kind!=AG_REPAIR_KIND_FEC_ROOT ) ) ) {
592 3 : meta_queue_push_safe( ctx, ( kind==AG_REPAIR_KIND_PARENT_FEC_COUNT )
593 3 : ? ag_repair_parent_and_fec_set_count( ctx->protocol, &to, now_ms, ctx->ag_nonce++, slot, &block_id )
594 3 : : ag_repair_fec_set_root ( ctx->protocol, &to, now_ms, ctx->ag_nonce++, slot, &block_id, fec_set_idx ), now );
595 3 : return;
596 3 : }
597 :
598 : /* Each case verifies the response's merkle proof chains up to the
599 : block id we requested before handing the metadata to the chainer */
600 135 : switch( response->kind ) {
601 42 : case AG_REPAIR_RESPONSE_PARENT_FEC_SET_COUNT: {
602 42 : ag_parent_fec_count_res_t * parent_fec_set_res = &response->parent_fec_set_res;
603 :
604 42 : if( FD_UNLIKELY( ag_repair_parent_fec_count_verify( parent_fec_set_res, &block_id ) ) ) {
605 3 : ctx->metrics->failed_parent_fec_count_cnt++;
606 3 : meta_queue_push_safe( ctx, ag_repair_parent_and_fec_set_count( ctx->protocol, &to, now_ms, ctx->ag_nonce++, slot, &block_id ), now );
607 3 : return;
608 3 : }
609 :
610 39 : int should_repair = !!fd_chainer_verified_parent_fec_count( ctx->chainer, slot, &block_id, parent_fec_set_res->fec_set_count, parent_fec_set_res->parent_slot, &parent_fec_set_res->parent_block_id );
611 39 : if( FD_UNLIKELY( !should_repair ) ) return;
612 :
613 138 : for( uint i=0; i<parent_fec_set_res->fec_set_count; i++ ) {
614 99 : meta_queue_push_safe( ctx, ag_repair_fec_set_root( ctx->protocol, &to, now_ms, ctx->ag_nonce++, slot, &block_id, i*FD_FEC_SHRED_CNT ), now );
615 99 : }
616 39 : break;
617 39 : }
618 93 : case AG_REPAIR_RESPONSE_FEC_SET_ROOT: {
619 93 : ag_fec_root_res_t * fec_set_root = &response->fec_set_root;
620 :
621 93 : if( FD_UNLIKELY( ag_repair_fec_set_root_verify( fec_set_root, &block_id, fec_set_idx ) ) ) {
622 3 : ctx->metrics->failed_fec_root_cnt++;
623 3 : meta_queue_push_safe( ctx, ag_repair_fec_set_root( ctx->protocol, &to, now_ms, ctx->ag_nonce++, slot, &block_id, fec_set_idx ), now );
624 3 : return;
625 3 : }
626 :
627 : /* The response carries only the 20-byte FEC-set root prefix */
628 90 : fd_hash_t fec_root_mr = {0};
629 90 : memcpy( fec_root_mr.uc, fec_set_root->root, FD_SHRED_MERKLE_NODE_SZ );
630 90 : fd_chainer_verified_hash_insert( ctx->chainer, slot, &block_id, fec_set_idx, &fec_root_mr );
631 90 : break;
632 93 : }
633 135 : }
634 135 : }
635 :
636 : /* ag_parse_parent_marker pulls the block's DECLARED parent out of the
637 : BlockMarker that a batch-opening data shred carries: a BlockComponent
638 : whose entry count is 0 is a marker, and the BlockHeaderV1 (shred 0) /
639 : UpdateParentV1 (a later batch) variants both carry (parent_slot,
640 : parent_block_id). Alpenglow chains on this, never on parent_off,
641 : because the block_id double merkle binds exactly these bytes.
642 :
643 : Both variants are ~54 bytes, so they always fit in the opening shred
644 : and never need cross-shred reassembly. Returns 1 and fills the out
645 : params on a well formed marker, 0 otherwise.
646 : https://github.com/anza-xyz/agave/blob/master/entry/src/block_component.rs */
647 :
648 : static int
649 : ag_parse_parent_marker( fd_shred_t const * shred,
650 : ulong * out_parent_slot,
651 75 : fd_hash_t * out_parent_block_id ) {
652 75 : uchar const * payload = fd_shred_data_payload( shred );
653 75 : ulong sz = fd_shred_payload_sz( shred );
654 :
655 75 : fd_block_marker_t marker[1];
656 75 : int err = fd_block_marker_de( marker, payload, sz );
657 75 : if( FD_UNLIKELY( err ) ) return 0;
658 :
659 75 : if( marker->kind==FD_BLOCK_MARKER_KIND_HEADER ) {
660 75 : memcpy( out_parent_slot, &marker->header.parent_slot, 8UL );
661 75 : memcpy( out_parent_block_id->uc, marker->header.parent_block_id.uc, 32UL );
662 75 : return 1;
663 75 : }
664 :
665 0 : if( marker->kind==FD_BLOCK_MARKER_KIND_UPDATE_PARENT ) {
666 0 : memcpy( out_parent_slot, &marker->update_parent.new_parent_slot, 8UL );
667 0 : memcpy( out_parent_block_id->uc, marker->update_parent.new_parent_block_id.uc, 32UL );
668 0 : return 1;
669 0 : }
670 :
671 0 : return 0; /* a footer or genesis cert carries no parent */
672 0 : }
673 :
674 : static inline void
675 : after_alpen_shred( ctx_t * ctx,
676 : ulong sig,
677 : fd_shred_t const * shred,
678 : ulong nonce,
679 5304 : fd_hash_t const * mr ) {
680 5304 : if( FD_UNLIKELY( shred->slot<=ctx->chainer->root ) ) return;
681 5304 : if( FD_UNLIKELY( fd_shred_is_code( fd_shred_type( shred->variant ) ) ) ) return; /* TODO */
682 :
683 : /* Try the root key first, then the positional one; drop the shred if neither matches. */
684 5301 : if( FD_UNLIKELY( fd_shred_sig_src( sig )==SHRED_SIG_SRC_REPAIR && fd_rnonce_ss_normal_repair( (uint)nonce ) ) ) {
685 1731 : fd_pubkey_t peer;
686 1731 : long now = fd_clock_tile_now( ctx->clock );
687 1731 : long rtt = fd_inflights_shred_match( ctx->inflights, nonce, shred->slot, shred->idx, mr, &peer, NULL, now );
688 1731 : if( FD_UNLIKELY( !rtt ) ) {
689 3 : rtt = fd_inflights_shred_match( ctx->inflights, nonce, shred->slot, shred->idx, NULL, &peer, NULL, now );
690 3 : }
691 1731 : if( FD_UNLIKELY( !rtt ) ) return;
692 1728 : fd_policy_peer_response_update( ctx->policy, &peer, rtt );
693 1728 : fd_histf_sample( ctx->metrics->response_latency, (ulong)rtt );
694 1728 : }
695 :
696 : /* Parent discovery on shred 0 of a slot.
697 :
698 : TODO: mid-block UpdateParent markers (at FEC-set boundaries after a
699 : DATA_COMPLETE) also rebind the double-merkle parent; not handled
700 : here yet. */
701 5298 : ulong parent_slot = AG_UNKNOWN_SLOT;
702 5298 : fd_hash_t parent_block_id = {0};
703 5298 : if( FD_UNLIKELY( shred->idx==0U && !ag_parse_parent_marker( shred, &parent_slot, &parent_block_id ) ) ) {
704 0 : FD_LOG_WARNING(( "invalid block header in slot: %lu, ignoring shred 0", shred->slot ));
705 0 : return;
706 0 : }
707 :
708 5298 : int slot_complete = !!(shred->data.flags & FD_SHRED_DATA_FLAG_SLOT_COMPLETE);
709 5298 : fd_chainer_shred_insert( ctx->chainer, shred->slot, shred->idx, slot_complete, mr, parent_slot, &parent_block_id );
710 5298 : }
711 :
712 : /* fec_completes */
713 : static inline void
714 : after_alpen_fec( ctx_t * ctx,
715 : ulong sig,
716 : fd_shred_t * shred,
717 162 : fd_hash_t * mr ) {
718 162 : if( FD_UNLIKELY( shred->slot <= ctx->chainer->root ) ) {
719 0 : fd_store_remove( ctx->store, ctx->store_map, mr );
720 0 : return;
721 0 : }
722 :
723 162 : int slot_complete = !!(shred->data.flags & FD_SHRED_DATA_FLAG_SLOT_COMPLETE);
724 162 : int data_complete = !!(shred->data.flags & FD_SHRED_DATA_FLAG_DATA_COMPLETE);
725 162 : if( fd_chainer_fec_complete( ctx->chainer, shred->slot, shred->fec_set_idx, slot_complete, data_complete, sig==SHRED_SIG_FEC_COMPLETE_LEADER, mr ) ) {
726 0 : fd_store_remove( ctx->store, ctx->store_map, mr );
727 0 : };
728 162 : }
729 :
730 : static inline void
731 : after_votor_block_repair( ctx_t * ctx,
732 36 : fd_votor_repair_t const * nf ) {
733 36 : fd_chainer_slotv_t * slotv = fd_chainer_slot_version_query( ctx->chainer, nf->slot, &nf->block_id );
734 36 : if( FD_LIKELY( slotv ) ) return; /* we already have this NF version recorded, no need for action */
735 :
736 33 : fd_chainer_verified_block_insert( ctx->chainer, nf->slot, nf->block_id );
737 :
738 33 : uint nonce = ctx->ag_nonce++;
739 33 : fd_pubkey_t const * peer = fd_policy_peer_select( ctx->policy );
740 33 : if( FD_LIKELY( peer ) ) {
741 33 : long now = fd_clock_tile_now( ctx->clock );
742 33 : ulong now_ms = (ulong)(now/(long)1e6);
743 33 : fd_repair_msg_t * msg = ag_repair_parent_and_fec_set_count( ctx->protocol,
744 33 : peer,
745 33 : now_ms,
746 33 : nonce,
747 33 : nf->slot,
748 33 : &nf->block_id );
749 33 : meta_queue_push_safe( ctx, msg, now );
750 33 : }
751 33 : }
752 :
753 : static void
754 : after_frag( ctx_t * ctx,
755 : ulong in_idx,
756 : ulong seq FD_PARAM_UNUSED,
757 : ulong sig,
758 : ulong sz,
759 : ulong tsorig FD_PARAM_UNUSED,
760 : ulong tspub FD_PARAM_UNUSED,
761 5673 : fd_stem_context_t * stem ) {
762 5673 : if( FD_UNLIKELY( ctx->skip_frag ) ) return;
763 :
764 5673 : ctx->stem = stem;
765 5673 : in_ctx_t const * in_ctx = &ctx->in_links[ in_idx ];
766 5673 : uint in_kind = ctx->in_kind[ in_idx ];
767 :
768 5673 : switch( in_kind ) {
769 : /* Unreliable frags */
770 141 : case IN_KIND_NET: {
771 141 : fd_eth_hdr_t * eth; fd_ip4_hdr_t * ip4; fd_udp_hdr_t * udp;
772 141 : uchar * data; ulong data_sz;
773 141 : if( FD_UNLIKELY( !fd_ip4_udp_hdr_strip( ctx->net_buf, sz, &data, &data_sz, ð, &ip4, &udp ) ) ) {
774 0 : ctx->metrics->malformed_ping++; // todo generalize
775 0 : return;
776 0 : }
777 :
778 141 : if( FD_LIKELY( data_sz==sizeof(fd_repair_ping_t) ) ) {
779 0 : after_ping( ctx, data, data_sz, ip4, udp );
780 141 : } else { /* alpen repair response */
781 141 : ag_repair_response_t response[1];
782 141 : if( FD_UNLIKELY( ag_repair_response_de( response, data, data_sz, ctx->chainer->fec_blk_max ) ) ) return; /* malformed */
783 141 : after_alpen_meta_repair( ctx, response );
784 141 : }
785 141 : break;
786 141 : }
787 141 : case IN_KIND_REPLAY: {
788 21 : if( FD_UNLIKELY( sig==REPLAY_SIG_MISSING_FEC ) ) {
789 3 : ctx->deliver_from_root = 1;
790 3 : return;
791 3 : }
792 18 : if( FD_LIKELY( sig==REPLAY_SIG_ROOT_ADVANCED ) ) {
793 18 : fd_replay_root_advanced_t const * root = (fd_replay_root_advanced_t const *)fd_type_pun_const( fd_chunk_to_laddr( in_ctx->mem, ctx->chunk ) );
794 18 : if( FD_LIKELY( root->slot > ctx->chainer->root ) ) fd_chainer_publish( ctx->chainer, root->slot, &root->block_id, ctx->store );
795 18 : }
796 18 : break;
797 21 : }
798 0 : case IN_KIND_SIGN: {
799 0 : after_sign( ctx, in_idx, sig, stem );
800 0 : break;
801 21 : }
802 : /* Reliable frags read directly from dcache */
803 0 : case IN_KIND_SNAP: {
804 0 : after_snap( ctx, sig, fd_chunk_to_laddr( ctx->in_links[ in_idx ].mem, ctx->snap_out_chunk ) );
805 0 : break;
806 21 : }
807 0 : case IN_KIND_GENESIS: {
808 0 : fd_genesis_meta_t const * meta = (fd_genesis_meta_t const *)fd_type_pun_const( fd_chunk_to_laddr( in_ctx->mem, ctx->chunk ) );
809 0 : fd_hash_t block_id = {0};
810 0 : if( meta->bootstrap ) fd_chainer_init( ctx->chainer, 0, &block_id );
811 0 : break;
812 21 : }
813 0 : case IN_KIND_GOSSIP: {
814 0 : fd_gossip_update_message_t const * msg = (fd_gossip_update_message_t const *)fd_type_pun_const( fd_chunk_to_laddr( in_ctx->mem, ctx->chunk ) );
815 0 : after_gossip( ctx, msg, sig );
816 0 : break;
817 21 : }
818 36 : case IN_KIND_VOTOR: {
819 36 : fd_votor_repair_t const * nf = fd_chunk_to_laddr_const( in_ctx->mem, ctx->chunk );
820 36 : if( FD_UNLIKELY( nf->slot <= ctx->chainer->root ) ) return;
821 36 : after_votor_block_repair( ctx, nf );
822 36 : break;
823 36 : }
824 5475 : case IN_KIND_SHRED: {
825 5475 : uint sig_src = fd_shred_sig_src( sig );
826 5475 : int sig_res = fd_shred_sig_res( sig );
827 :
828 5475 : if( FD_UNLIKELY( sig_src==SHRED_SIG_FEC_EVICTED ) ) {
829 6 : fd_fec_evicted_t * evicted = (fd_fec_evicted_t *)fd_type_pun( fd_chunk_to_laddr( in_ctx->mem, ctx->chunk ) );
830 6 : fd_chainer_fec_evicted( ctx->chainer, evicted->slot, evicted->fec_set_idx, &evicted->merkle_root );
831 6 : return;
832 6 : }
833 :
834 5469 : uchar * src = fd_chunk_to_laddr( in_ctx->mem, ctx->chunk );
835 5469 : fd_shred_base_t * shred_msg = (fd_shred_base_t *)fd_type_pun( src );
836 5469 : fd_shred_t * shred = &shred_msg->shred; /* completes & shred messages all have a shred header at the same offset (after merkle root) */
837 :
838 5469 : if( FD_UNLIKELY( shred->slot > ctx->metrics->current_slot && sig_src == SHRED_SIG_SRC_TURBINE ) ) {
839 66 : FD_LOG_INFO(( "[Turbine] slot: %lu, root: %lu", shred->slot, ctx->chainer->root ));
840 66 : ctx->metrics->current_slot = shred->slot;
841 66 : }
842 :
843 5469 : if( FD_UNLIKELY( ctx->turbine_slot0 == ULONG_MAX && sig_src == SHRED_SIG_SRC_TURBINE ) ) {
844 39 : ctx->turbine_slot0 = shred->slot;
845 :
846 39 : ulong slot_delta;
847 39 : int cf = __builtin_usubl_overflow( ctx->turbine_slot0, ctx->chainer->root, &slot_delta );
848 39 : if( FD_UNLIKELY( cf || slot_delta > fd_slotv_pool_max( ctx->chainer->slotv_pool ) ) ) {
849 : /* TODO: It's most optimal to define the catchup target as the
850 : first notarize cert we receive. But we currently dont have
851 : any info in the rotor tile to know if we are unstaked or
852 : not. And if we are unstaked, we will not be getting any
853 : certs from votor. So for now we will just use the first
854 : turbine shred we receive. */
855 :
856 0 : FD_LOG_ERR(( "Catchup slot distance exceeds the repair buffer: target %lu - snapshot slot %lu > %lu. "
857 0 : "Restart with a more recent snapshot or increase config rotor.slot_max", ctx->turbine_slot0, ctx->chainer->root, fd_slotv_pool_max( ctx->chainer->slotv_pool ) ));
858 0 : return;
859 0 : }
860 39 : fd_repair_metrics_set_turbine_slot0( ctx->slot_metrics, shred->slot );
861 39 : fd_policy_set_turbine_slot0( ctx->policy, shred->slot );
862 :
863 : /* Catchup optimizations */
864 39 : ulong root = ctx->chainer->root;
865 39 : if( FD_LIKELY( root != ULONG_MAX && shred->slot > root && !ctx->block_id_repair_only ) ) {
866 36 : ulong capacity = toss_queue_max( ctx->toss_queue ) - toss_queue_cnt( ctx->toss_queue );
867 36 : ulong seed_cnt = fd_ulong_min( shred->slot-root, capacity/2 );
868 36 : long now_ms = fd_log_wallclock()/(long)1e6;
869 72 : for( ulong i=1; i<=seed_cnt; i++ ) {
870 36 : fd_pubkey_t const * peer = fd_policy_peer_select( ctx->policy );
871 36 : if( FD_UNLIKELY( !peer ) ) break;
872 36 : fd_repair_msg_t * msg = fd_repair_shred( ctx->protocol, peer, (ulong)now_ms, 0, root + i, 0 );
873 36 : toss_queue_push( ctx->toss_queue, (sign_pending_t){ .msg = *msg } );
874 36 : }
875 36 : }
876 39 : }
877 :
878 5469 : if( FD_UNLIKELY( sig==SHRED_SIG_FEC_COMPLETE || sig==SHRED_SIG_FEC_COMPLETE_LEADER ) ) {
879 162 : fd_fec_complete_t * complete_msg = (fd_fec_complete_t *)fd_type_pun( src );
880 162 : after_alpen_fec( ctx, sig, &complete_msg->last_shred_hdr, &complete_msg->merkle_root );
881 5307 : } else if( FD_LIKELY( sig_res!=SHRED_SIG_RESULT_EQVOC ) ) {
882 5304 : after_alpen_shred( ctx, sig, shred, shred_msg->rnonce, &shred_msg->merkle_root );
883 5304 : }
884 5469 : return;
885 5469 : }
886 0 : default: FD_LOG_ERR(( "bad in_kind %u", in_kind )); /* Should never reach here since before_frag should have filtered out any unexpected frags. */
887 5673 : }
888 5673 : }
889 :
890 : /* Should be called for any shred request made. block_id is the
891 : ShredForBlockId version being repaired and fec_root the root the
892 : chainer holds for that version at the requested FEC set (the key a
893 : response is matched by), or both NULL for a plain positional shred
894 : request. */
895 :
896 : static void
897 1926 : record_inflight_request( ctx_t * ctx, ulong nonce, fd_pubkey_t const * peer, ulong slot, ulong shred_idx, fd_hash_t const * block_id, fd_hash_t const * fec_root, long now ) {
898 1926 : if( FD_LIKELY( block_id && fd_hash_check_zero( block_id ) ) ) block_id = NULL;
899 1926 : fd_inflights_shred_insert( ctx->inflights, nonce, peer, slot, shred_idx, block_id, fec_root, now );
900 1926 : fd_policy_peer_request_update( ctx->policy, peer );
901 1926 : }
902 :
903 : /* ag_policy_block_id_next is development only. */
904 :
905 : /* redispatch_meta re-issues an aged-out metadata request under a fresh
906 : nonce (so a late response to the old one is simply unmatched), via
907 : the meta queue so a peer is chosen at dispatch. Skipped if the slot
908 : has been rooted meanwhile or the chainer already holds the answer. */
909 :
910 : static void
911 : redispatch_meta( ctx_t * ctx,
912 : fd_inflight_t const * req,
913 12 : long now ) {
914 12 : long now_ms = now/(long)1e6;
915 12 : uint kind = req->key.kind;
916 12 : ulong slot = req->key.slot;
917 12 : uint fec_set_idx = req->key.idx;
918 12 : fd_hash_t block_id = req->block_id;
919 :
920 12 : if( FD_UNLIKELY( slot <= ctx->chainer->root ) ) return; /* rooted while outstanding: drop, don't re-request */
921 12 : fd_chainer_slotv_t * slotv = fd_chainer_slot_version_query( ctx->chainer, slot, &block_id );
922 12 : if( kind==AG_REPAIR_KIND_PARENT_FEC_COUNT && slotv && slotv->complete_idx!=UINT_MAX && slotv->parent_slot!=AG_UNKNOWN_SLOT ) return; /* we already have the ParentFecSetCount */
923 6 : if( kind==AG_REPAIR_KIND_FEC_ROOT && fd_chainer_fec_query( ctx->chainer, slot, fec_set_idx, &block_id ) ) return; /* we already have the FecSetRoot */
924 :
925 6 : fd_pubkey_t to = {0};
926 6 : fd_repair_msg_t * msg = ( kind==AG_REPAIR_KIND_PARENT_FEC_COUNT )
927 6 : ? ag_repair_parent_and_fec_set_count( ctx->protocol, &to, (ulong)now_ms, ctx->ag_nonce++, slot, &block_id )
928 6 : : ag_repair_fec_set_root ( ctx->protocol, &to, (ulong)now_ms, ctx->ag_nonce++, slot, &block_id, fec_set_idx );
929 6 : meta_queue_push_safe( ctx, msg, now );
930 6 : }
931 :
932 : /* ag_policy_next is the standard Alpenglow repair pipeline, driven once
933 : per after_credit (ag_policy_block_id_next is the block-id-only
934 : variant). In priority order it:
935 : (1) redispatches an outstanding metadata request
936 : (ParentAndFecSetCount / FecSetRoot) that has aged past its drain
937 : timeout
938 : (2) redispatches an outstanding shred request (normal or block-id)
939 : (3) gates on inflight capacity, then
940 : (4) walks the chainer's slotv treaps to issue new requests -
941 : orphan/ancestry-discovery pass
942 : (5) followed by a shred-fill pass.
943 :
944 : The shred-fill walk issues a request for each still-missing shred
945 : exactly once, advancing highest_requested only on a shred we already
946 : have or a request we actually send, so a budget cutoff never strands
947 : an index; fd_inflights owns all re-requests so a fully-requested slot
948 : is popped from the treap. At most one request is sent per call.
949 :
950 : block_id_repair_only is a development only flag.
951 : TODO revise & refactor later */
952 : static void
953 2226 : ag_policy_next( ctx_t * ctx, out_ctx_t * sign_out, long now, int * charge_busy ) {
954 2226 : fd_chainer_t * chainer = ctx->chainer;
955 2226 : long now_ms = now/(long)1e6;
956 :
957 2226 : fd_pubkey_t const * peer = fd_policy_peer_select( ctx->policy );
958 2226 : if( FD_UNLIKELY( !peer ) ) return;
959 :
960 : /* 1. Redispatch the oldest aged-out request (metadata or shred)
961 : under a fresh nonce. */
962 :
963 2226 : if( FD_UNLIKELY( fd_inflights_should_drain( ctx->inflights, now ) ) ) {
964 15 : fd_inflight_t req[1];
965 15 : fd_inflights_pop( ctx->inflights, req );
966 15 : *charge_busy = 1;
967 :
968 15 : if( FD_UNLIKELY( req->key.kind!=FD_REPAIR_KIND_SHRED ) ) { redispatch_meta( ctx, req, now ); return; }
969 :
970 3 : ulong nonce = req->key.nonce;
971 3 : ulong slot = req->key.slot;
972 3 : ulong shred_idx = req->key.idx;
973 3 : fd_hash_t block_id = req->block_id;
974 :
975 3 : fd_chainer_slotv_t * slotv = fd_chainer_slot_version_query( ctx->chainer, slot, &block_id );
976 3 : if( FD_UNLIKELY( slot > ctx->chainer->root && slotv && !fd_chainer_shred_test( ctx->chainer, slotv, (uint)shred_idx ) ) ) {
977 :
978 : /* A block-id request needs the prior getFecRoot to still be
979 : present: the request is keyed in inflights by that root, and
980 : the chainer attaches the response to this version through it. */
981 :
982 3 : uint fec_set_idx = (uint)shred_idx & ~( (uint)FD_FEC_SHRED_CNT - 1U );
983 3 : int has_block_id = !fd_hash_check_zero( &block_id );
984 3 : fd_chainer_fec_t * fec = has_block_id ? fd_chainer_fec_query( ctx->chainer, slot, fec_set_idx, &block_id ) : NULL;
985 3 : int fec_present = !has_block_id || !!fec;
986 :
987 3 : if( FD_UNLIKELY( !fec_present || fd_reqlim_next( ctx->dedup, fd_reqlim_key( FD_REPAIR_KIND_SHRED, slot, (uint)shred_idx ), now ) ) ) {
988 : /* If getFecRoot hasnt responded, park the request in inflights. */
989 0 : fd_hash_t hash_zero = { 0 };
990 0 : fd_inflights_shred_insert( ctx->inflights, 0UL, &hash_zero, slot, shred_idx, &block_id, fec ? &fec->merkle_root : NULL, now );
991 3 : } else {
992 3 : ctx->metrics->rerequest++;
993 3 : nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 1, slot, (uint)shred_idx, now );
994 3 : fd_repair_msg_t * msg = has_block_id
995 3 : ? ag_repair_shred_block_id( ctx->protocol, peer, (ulong)now_ms, (uint)nonce, slot, &block_id, (uint)shred_idx )
996 3 : : fd_repair_shred( ctx->protocol, peer, (ulong)now_ms, (uint)nonce, slot, shred_idx );
997 3 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
998 3 : record_inflight_request( ctx, nonce, peer, slot, shred_idx, &block_id, fec ? &fec->merkle_root : NULL, now );
999 3 : return;
1000 3 : }
1001 3 : }
1002 3 : }
1003 :
1004 : /* 2. No new shred requests allowed if inflights is near capacity. */
1005 :
1006 2211 : if( FD_UNLIKELY( fd_inflights_outstanding_free( ctx->inflights ) <= fd_signs_map_key_cnt( ctx->signs_map ) ) ) return;
1007 :
1008 : /* 4. Orphan (ancestry) pass
1009 : - Parent slot unknown: request our own shred 0 (contents names the
1010 : parent).
1011 : - Parent slot known but the parent slotv is absent: the slotv
1012 : stays in the treap and we fire an Orphan request for it. */
1013 :
1014 2211 : ulong onext;
1015 2274 : for( ulong oit=fd_chainer_orphan_iter_init( chainer ); !fd_chainer_work_iter_done( oit ); oit=onext ) {
1016 90 : onext = fd_chainer_orphan_iter_next( chainer, oit );
1017 90 : fd_chainer_slotv_t * o = fd_chainer_work_iter_ele( chainer, oit );
1018 :
1019 90 : if( FD_UNLIKELY( ctx->block_id_repair_only && fd_hash_check_zero( &o->block_id ) ) ) continue;
1020 :
1021 84 : if( ctx->block_id_repair_only || ( o->parent_slot==AG_UNKNOWN_SLOT && !fd_hash_check_zero( &o->block_id ) ) ) {
1022 75 : if( !fd_reqlim_next( ctx->dedup, fd_reqlim_key( AG_REPAIR_KIND_PARENT_FEC_COUNT, o->slot, 0 ), now ) ) {
1023 24 : uint nonce = ctx->ag_nonce++;
1024 24 : fd_repair_msg_t * msg = ag_repair_parent_and_fec_set_count( ctx->protocol, peer, (ulong)now_ms, nonce, o->slot, &o->block_id );
1025 24 : *charge_busy = 1;
1026 24 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1027 24 : meta_inflight_record( ctx, msg, now );
1028 24 : return;
1029 24 : }
1030 75 : } else if( o->parent_slot==AG_UNKNOWN_SLOT && !fd_reqlim_next( ctx->dedup, fd_reqlim_key( FD_REPAIR_KIND_SHRED, o->slot, 0 ), now ) ) {
1031 3 : uint nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 1, o->slot, 0U, now );
1032 3 : fd_repair_msg_t * msg = fd_repair_shred( ctx->protocol, peer, (ulong)now_ms, (uint)nonce, o->slot, 0 );
1033 3 : *charge_busy = 1;
1034 3 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1035 3 : record_inflight_request( ctx, nonce, peer, o->slot, 0UL, NULL, NULL, now );
1036 3 : return;
1037 6 : } else if( o->parent_slot!=AG_UNKNOWN_SLOT && !fd_reqlim_next( ctx->dedup, fd_reqlim_key( FD_REPAIR_KIND_ORPHAN, o->slot, UINT_MAX ), now ) ) {
1038 0 : uint nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 0, o->slot, 0U, now );
1039 0 : fd_repair_msg_t * msg = fd_repair_orphan( ctx->protocol, peer, (ulong)now_ms, (uint)nonce, o->slot );
1040 0 : *charge_busy = 1;
1041 0 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1042 0 : return;
1043 0 : }
1044 84 : }
1045 :
1046 : /* 5. Shred-fill pass */
1047 :
1048 2331 : for( ulong it=fd_chainer_repair_iter_init( chainer ); !fd_chainer_work_iter_done( it ); ) {
1049 2094 : fd_chainer_slotv_t * e = fd_chainer_work_iter_ele( chainer, it );
1050 2094 : it = fd_chainer_repair_iter_next( chainer, it );
1051 :
1052 2094 : if( FD_UNLIKELY( ctx->block_id_repair_only && ( fd_hash_check_zero( &e->block_id ) || e->complete_idx==UINT_MAX ) ) ) continue;
1053 :
1054 2085 : ulong slot = e->slot;
1055 :
1056 2085 : if( e->buffered_idx!=UINT_MAX && ( e->highest_requested==UINT_MAX || e->buffered_idx > e->highest_requested ) )
1057 33 : e->highest_requested = e->buffered_idx;
1058 :
1059 2085 : if( FD_UNLIKELY( e->complete_idx==UINT_MAX ) ) {
1060 75 : if( !fd_reqlim_next( ctx->dedup, fd_reqlim_key( FD_REPAIR_KIND_HIGHEST_SHRED, slot, UINT_MAX ), now ) ) {
1061 27 : uint nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 0, slot, 0U, now );
1062 27 : fd_repair_msg_t * msg = fd_repair_highest_shred( ctx->protocol, peer, (ulong)now_ms, nonce, slot, 0 );
1063 27 : *charge_busy = 1;
1064 27 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1065 27 : return;
1066 27 : }
1067 48 : continue;
1068 75 : }
1069 :
1070 : /* make individual shred request */
1071 :
1072 2010 : uint idx = e->highest_requested+1;
1073 2202 : while( FD_UNLIKELY( idx<chainer->fec_blk_max*FD_FEC_SHRED_CNT && fd_chainer_shred_test( chainer, e, idx ) ) ) idx++;
1074 :
1075 2010 : if( FD_UNLIKELY( idx > e->complete_idx ) ) {
1076 39 : e->highest_requested = idx;
1077 39 : fd_chainer_repair_remove( chainer, e );
1078 39 : continue;
1079 1971 : };
1080 :
1081 1971 : if( FD_UNLIKELY( fd_hash_check_zero( &e->block_id ) ) ) {
1082 : /* TODO eager repair time gate */
1083 0 : ulong nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 1, slot, idx, now );
1084 0 : fd_repair_msg_t * msg = fd_repair_shred( ctx->protocol, peer, (ulong)now/(ulong)1e6, (uint)nonce, slot, idx );
1085 0 : *charge_busy = 1;
1086 0 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1087 0 : record_inflight_request( ctx, nonce, peer, slot, idx, NULL, NULL, now );
1088 1971 : } else {
1089 : /* block_id known -> Alpenglow block-id repair. Gate on prior
1090 : getFecRoot being present. Retry this idx on a later walk, once
1091 : the root lands. */
1092 1971 : uint fec_set_idx = idx & ~( (uint)FD_FEC_SHRED_CNT - 1U );
1093 1971 : fd_chainer_fec_t * fec = fd_chainer_fec_query( chainer, slot, fec_set_idx, &e->block_id );
1094 1971 : if( FD_UNLIKELY( !fec ) ) continue;
1095 :
1096 1920 : ulong nonce = fd_rnonce_ss_compute( ctx->repair_nonce_ss, 1, slot, idx, now );
1097 1920 : fd_repair_msg_t * msg = ag_repair_shred_block_id( ctx->protocol, peer, (ulong)now_ms, (uint)nonce, slot, &e->block_id, idx );
1098 1920 : *charge_busy = 1;
1099 1920 : fd_repair_send_sign_request( ctx, sign_out, msg, NULL );
1100 1920 : record_inflight_request( ctx, nonce, peer, slot, idx, &e->block_id, &fec->merkle_root, now );
1101 1920 : }
1102 1920 : e->highest_requested = idx;
1103 1920 : return;
1104 1971 : }
1105 2184 : }
1106 :
1107 : /* publish_fec builds and publishes a single ROTOR_SIG_FEC_REPLAY
1108 : message to replay for the FEC fec owned by version slotv.
1109 : When from_root is set (the deliver_from_root recovery redelivery),
1110 : the block_id is populated on every FEC whose version has a known
1111 : block_id -- even turbine FECs that the normal path leaves zero -- so
1112 : replay can dedup the redelivered path by (slot, block_id) and skip
1113 : blocks it has already replayed. known_id is unaffected by
1114 : redelivery, see below. */
1115 : static void
1116 : publish_fec( ctx_t * ctx,
1117 : fd_stem_context_t * stem,
1118 : fd_chainer_slotv_t * slotv,
1119 : fd_chainer_fec_t * fec,
1120 234 : int from_root ) {
1121 234 : fd_hash_t null_hash = {0};
1122 :
1123 234 : fd_rotor_replay_fec_t * msg = fd_chunk_to_laddr( ctx->repair_out_ctx->mem, ctx->repair_out_ctx->chunk );
1124 234 : msg->slot = fec->slot;
1125 234 : msg->fec_set_idx = fec->fec_set_idx;
1126 234 : msg->mr = fec->merkle_root;
1127 234 : msg->parent_slot = slotv->parent_slot;
1128 234 : msg->parent_block_id = slotv->parent_block_id;
1129 234 : msg->slot_complete = fec->slot_complete;
1130 234 : msg->data_complete = fec->data_complete;
1131 234 : msg->is_leader = fec->is_leader;
1132 :
1133 : /* TODO rename? Unfortunately known_id flag is tightly coupled with
1134 : replay behavior, so worth revisiting. known_id marks the
1135 : votor-driven (cert) versions. A turbine version is never marked
1136 : known, not even when redelivered from root, because a redelivered
1137 : copy must key the same way as the live turbine FECs of the same
1138 : block. block_id is still populated whenever the chainer knows it
1139 : (the slot-complete FEC, or any redelivered FEC) so replay can dedup
1140 : a redelivered block it already fully replayed by {slot, block_id}.
1141 : */
1142 234 : int block_id_known = !fd_hash_check_zero( &slotv->block_id );
1143 234 : msg->known_id = !slotv->turbine;
1144 234 : msg->block_id = ( block_id_known && ( msg->known_id || from_root || fec->slot_complete ) ) ? slotv->block_id : null_hash;
1145 :
1146 234 : if( FD_UNLIKELY( fec->slot_complete ) ) {
1147 108 : FD_BASE58_ENCODE_32_BYTES( slotv->block_id.uc, block_id );
1148 108 : FD_LOG_INFO(( "[%s] slot is complete %lu. num_data_shreds: %u. block_id: %s, parent_slot: %lu, turbine: %d",
1149 108 : __func__,
1150 108 : slotv->slot,
1151 108 : slotv->complete_idx + 1,
1152 108 : block_id,
1153 108 : slotv->parent_slot,
1154 108 : slotv->turbine ));
1155 108 : }
1156 :
1157 234 : fd_stem_publish( stem, ctx->repair_out_ctx->idx, ROTOR_SIG_FEC_REPLAY, ctx->repair_out_ctx->chunk, sizeof(fd_rotor_replay_fec_t), 0UL, 0UL, fd_frag_meta_ts_comp( fd_tickcount() ) );
1158 234 : ctx->repair_out_ctx->chunk = fd_dcache_compact_next( ctx->repair_out_ctx->chunk, sizeof(fd_rotor_replay_fec_t), ctx->repair_out_ctx->chunk0, ctx->repair_out_ctx->wmark );
1159 234 : ctx->metrics->fecs_delivered++;
1160 234 : }
1161 :
1162 : /* full_fec_path_queue queues every FEC from the chainer root down to
1163 : (target_slotv, target_fec), inclusive, onto ctx->deliver_queue in
1164 : root-to-target order. */
1165 : static void
1166 : full_fec_path_queue( ctx_t * ctx,
1167 : fd_chainer_slotv_t * target_slotv,
1168 21 : fd_chainer_fec_t * target_fec ) {
1169 21 : fd_chainer_t * chainer = ctx->chainer;
1170 :
1171 21 : for( fd_chainer_slotv_t * slotv = target_slotv;
1172 63 : slotv && slotv->slot > chainer->root;
1173 42 : slotv = fd_chainer_slot_version_query( chainer, slotv->parent_slot, &slotv->parent_block_id ) ) {
1174 42 : uint slotv_idx = (uint)fd_slotv_pool_idx( chainer->slotv_pool, slotv );
1175 :
1176 42 : uint kmax;
1177 42 : if( FD_LIKELY( slotv==target_slotv ) ) {
1178 21 : kmax = target_fec->fec_set_idx / (uint)FD_FEC_SHRED_CNT; /* target: up to the delivered FEC */
1179 21 : } else {
1180 21 : if( FD_UNLIKELY( slotv->buffered_fec_idx==UINT_MAX ) ) continue; /* ancestor with no complete FEC buffered */
1181 21 : kmax = slotv->buffered_fec_idx / (uint)FD_FEC_SHRED_CNT; /* ancestor: all buffered FECs */
1182 21 : }
1183 :
1184 42 : uint const * fecs = fd_chainer_slotv_fecs( chainer, slotv );
1185 111 : for( int k=(int)kmax; k>=0; k-- ) {
1186 69 : if( FD_UNLIKELY( out_queue_full( ctx->deliver_queue ) ) ) FD_LOG_ERR(( "deliver_from_root queue full" ));
1187 69 : out_queue_push_head( ctx->deliver_queue, (out_ele_t){ .slotv_idx = slotv_idx, .fec_idx = fecs[ k ] } );
1188 69 : }
1189 42 : }
1190 21 : }
1191 :
1192 : /* publish_fec_replay pops one delivered FEC off the chainer's out_queue
1193 : and publishes it to replay on repair_out with ROTOR_SIG_FEC_REPLAY.
1194 : When ctx->deliver_from_root is set the entire ancestry path from the
1195 : chainer root to that FEC is instead queued onto ctx->deliver_queue,
1196 : which after_credit drains one FEC per call. Returns 1 if a delivered
1197 : FEC was consumed. */
1198 : static int
1199 2589 : publish_fec_replay( ctx_t * ctx, fd_stem_context_t * stem ) {
1200 2589 : out_ele_t * out_queue = ctx->chainer->out_queue;
1201 2589 : if( FD_LIKELY( out_queue_empty( out_queue ) ) ) return 0;
1202 :
1203 189 : out_ele_t out_ele = out_queue_pop_head( out_queue );
1204 189 : if( FD_UNLIKELY( out_ele.slotv_idx == UINT_MAX ) ) return 1;
1205 :
1206 186 : fd_chainer_fec_t * fec = fd_fec_pool_ele( ctx->chainer->fec_pool, out_ele.fec_idx );
1207 186 : fd_chainer_slotv_t * slotv = fd_slotv_pool_ele( ctx->chainer->slotv_pool, out_ele.slotv_idx );
1208 :
1209 186 : if( FD_UNLIKELY( ctx->deliver_from_root ) ) {
1210 21 : full_fec_path_queue( ctx, slotv, fec );
1211 21 : ctx->deliver_from_root = 0;
1212 21 : }
1213 165 : else publish_fec( ctx, stem, slotv, fec, 0 /* from_root */ );
1214 :
1215 186 : return 1;
1216 189 : }
1217 :
1218 : static inline void
1219 : after_credit( ctx_t * ctx,
1220 : fd_stem_context_t * stem,
1221 : int * opt_poll_in FD_PARAM_UNUSED,
1222 2658 : int * charge_busy ) {
1223 2658 : long now = fd_clock_tile_now( ctx->clock );
1224 :
1225 : /* deliver_queue has FECs when replay has signaled a bank eviction,
1226 : and we added the full path of FECs from root up until the next FEC
1227 : we need to deliver. */
1228 2658 : if( FD_UNLIKELY( !out_queue_empty( ctx->deliver_queue ) ) ) {
1229 69 : out_ele_t e = out_queue_pop_head( ctx->deliver_queue );
1230 69 : fd_chainer_slotv_t * slotv = fd_slotv_pool_ele( ctx->chainer->slotv_pool, e.slotv_idx );
1231 69 : fd_chainer_fec_t * fec = fd_fec_pool_ele ( ctx->chainer->fec_pool, e.fec_idx );
1232 69 : publish_fec( ctx, stem, slotv, fec, 1 /* from_root: always populate block_id */ );
1233 69 : *charge_busy = 1;
1234 69 : *opt_poll_in = 0;
1235 69 : return;
1236 69 : }
1237 :
1238 : /* Publish any FECs the chainer has delivered for replay. */
1239 2589 : if( publish_fec_replay( ctx, stem ) ) {
1240 189 : *charge_busy = 1;
1241 189 : *opt_poll_in = 0;
1242 189 : return;
1243 189 : }
1244 :
1245 2400 : if( FD_UNLIKELY( ctx->halt_signing ) ) {
1246 0 : *charge_busy = 1;
1247 0 : return;
1248 0 : }
1249 :
1250 : /* Verify that there is at least one sign tile with available credits.
1251 : If not, we can't send any requests and leave early. */
1252 2400 : out_ctx_t * sign_out = sign_avail_credits( ctx );
1253 2400 : if( FD_UNLIKELY( !sign_out ) ) {
1254 0 : ctx->metrics->sign_tile_unavail++;
1255 0 : return;
1256 0 : }
1257 :
1258 : /* If inflights is at capacity, then the only thing we can send is:
1259 : pongs, initial highest window index requests, or resend things that
1260 : are already inflight. Any new requests that would cause an
1261 : inflight to be added to the queue must be deferred. */
1262 :
1263 2400 : if( FD_UNLIKELY( !toss_queue_empty( ctx->toss_queue ) ) ) {
1264 33 : sign_pending_t signable = toss_queue_pop( ctx->toss_queue );
1265 33 : fd_repair_send_sign_request( ctx, sign_out, &signable.msg, signable.msg.kind == FD_REPAIR_KIND_PONG ? &signable.pong_data : NULL );
1266 33 : *charge_busy = 1;
1267 33 : return;
1268 33 : }
1269 :
1270 2367 : if( FD_UNLIKELY( !meta_queue_empty( ctx->meta_queue ) ) ) {
1271 141 : fd_repair_msg_t msg = meta_queue_pop( ctx->meta_queue );
1272 141 : fd_pubkey_t const * peer = fd_policy_peer_select( ctx->policy );
1273 141 : fd_pubkey_t to = {0};
1274 141 : if( FD_LIKELY( peer ) ) {
1275 141 : msg.header.to = *peer;
1276 141 : fd_repair_send_sign_request( ctx, sign_out, &msg, NULL );
1277 141 : } else {
1278 0 : peer = &to;
1279 0 : }
1280 141 : meta_inflight_record( ctx, &msg, now );
1281 141 : *charge_busy = 1;
1282 141 : *opt_poll_in = 0;
1283 141 : return;
1284 141 : }
1285 :
1286 2226 : ag_policy_next( ctx, sign_out, now, charge_busy );
1287 2226 : return;
1288 2367 : }
1289 :
1290 : static void
1291 0 : signs_queue_update_identity( ctx_t * ctx ) {
1292 0 : ulong queue_cnt = toss_queue_cnt( ctx->toss_queue );
1293 0 : for( ulong i=0UL; i<queue_cnt; i++ ) {
1294 0 : sign_pending_t signable = toss_queue_pop( ctx->toss_queue );
1295 0 : switch( signable.msg.kind ) {
1296 0 : case FD_REPAIR_KIND_PONG:
1297 0 : memcpy( signable.msg.pong.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1298 0 : break;
1299 0 : case FD_REPAIR_KIND_SHRED:
1300 0 : memcpy( signable.msg.shred.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1301 0 : break;
1302 0 : case FD_REPAIR_KIND_HIGHEST_SHRED:
1303 0 : memcpy( signable.msg.highest_shred.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1304 0 : break;
1305 0 : case FD_REPAIR_KIND_ORPHAN:
1306 0 : memcpy( signable.msg.orphan.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1307 0 : break;
1308 0 : case AG_REPAIR_KIND_SHRED_FOR_BLOCK_ID:
1309 0 : memcpy( signable.msg.shred_block_id.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1310 0 : break;
1311 0 : default:
1312 0 : FD_LOG_CRIT(( "Unhandled repair kind %u", signable.msg.kind ));
1313 0 : break;
1314 0 : }
1315 0 : toss_queue_push( ctx->toss_queue, signable );
1316 0 : }
1317 0 : queue_cnt = meta_queue_cnt( ctx->meta_queue );
1318 0 : for( ulong i=0UL; i<queue_cnt; i++ ) {
1319 0 : fd_repair_msg_t msg = meta_queue_pop( ctx->meta_queue );
1320 0 : switch( msg.kind ) {
1321 0 : case AG_REPAIR_KIND_FEC_ROOT:
1322 0 : memcpy( msg.fec_set_root.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1323 0 : break;
1324 0 : case AG_REPAIR_KIND_PARENT_FEC_COUNT:
1325 0 : memcpy( msg.parent_fec_set_count.from.uc, ctx->identity_public_key.uc, sizeof(fd_pubkey_t) );
1326 0 : break;
1327 0 : default:
1328 0 : FD_LOG_CRIT(( "Unhandled repair kind %u", msg.kind ));
1329 0 : break;
1330 0 : }
1331 0 : meta_queue_push( ctx->meta_queue, msg );
1332 0 : }
1333 0 : }
1334 :
1335 : static inline void
1336 0 : during_housekeeping( ctx_t * ctx ) {
1337 0 : if( FD_UNLIKELY( fd_clock_tile_recal_due( ctx->clock ) ) ) fd_clock_tile_recal( ctx->clock );
1338 :
1339 0 : if( FD_UNLIKELY( fd_keyswitch_state_query( ctx->keyswitch )==FD_KEYSWITCH_STATE_UNHALT_PENDING ) ) {
1340 0 : FD_LOG_DEBUG(( "keyswitch: unhalting" ));
1341 0 : FD_CHECK_CRIT( ctx->halt_signing, "state machine corruption" );
1342 0 : ctx->halt_signing = 0;
1343 0 : fd_keyswitch_state( ctx->keyswitch, FD_KEYSWITCH_STATE_COMPLETED );
1344 0 : }
1345 :
1346 0 : if( FD_UNLIKELY( fd_keyswitch_state_query( ctx->keyswitch )==FD_KEYSWITCH_STATE_SWITCH_PENDING ) ) {
1347 :
1348 0 : if( !ctx->halt_signing ) {
1349 : /* At this point, stop sending new sign requests to the sign tile
1350 : and wait for all outstanding sign requests to be received back
1351 : from the sign tile. We also need to update any pending
1352 : outgoing sign requests with the new identity key. */
1353 0 : FD_LOG_DEBUG(( "keyswitch: halting signing" ));
1354 0 : ctx->halt_signing = 1;
1355 0 : memcpy( ctx->identity_public_key.uc, ctx->keyswitch->bytes, 32UL );
1356 0 : ctx->protocol->identity_key = ctx->identity_public_key;
1357 0 : signs_queue_update_identity( ctx );
1358 0 : }
1359 :
1360 0 : if( fd_signs_map_key_cnt( ctx->signs_map )==0UL ) {
1361 : /* Once there are no more in flight sign requests, we are ready to
1362 : say that the keyswitch is completed. */
1363 0 : FD_LOG_DEBUG(( "keyswitch: completed, no more outstanding stale sign requests" ));
1364 0 : fd_keyswitch_state( ctx->keyswitch, FD_KEYSWITCH_STATE_COMPLETED );
1365 0 : }
1366 0 : }
1367 0 : }
1368 :
1369 : static void
1370 : privileged_init( fd_topo_t const * topo,
1371 0 : fd_topo_tile_t const * tile ) {
1372 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
1373 :
1374 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
1375 0 : ctx_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
1376 0 : fd_memset( ctx, 0, sizeof(ctx_t) );
1377 :
1378 0 : uchar const * identity_key = fd_keyload_load( tile->rotor.identity_key_path, /* pubkey only: */ 1 );
1379 0 : fd_memcpy( ctx->identity_public_key.uc, identity_key, sizeof(fd_pubkey_t) );
1380 :
1381 0 : FD_TEST( fd_rng_secure( &ctx->repair_seed, sizeof(ulong) ) );
1382 :
1383 0 : ulong rnonce_ss_id = fd_pod_queryf_ulong( topo->props, ULONG_MAX, "rnonce_ss" );
1384 0 : FD_TEST( rnonce_ss_id!=ULONG_MAX );
1385 0 : memcpy( ctx->repair_nonce_ss, fd_topo_obj_laddr( topo, rnonce_ss_id ), sizeof(fd_rnonce_ss_t) );
1386 0 : }
1387 :
1388 : static void
1389 : unprivileged_init( fd_topo_t const * topo,
1390 0 : fd_topo_tile_t const * tile ) {
1391 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
1392 :
1393 0 : ulong total_sign_depth = tile->rotor.repair_sign_depth * tile->rotor.repair_sign_cnt;
1394 0 : int lg_sign_depth = fd_ulong_find_msb( fd_ulong_pow2_up(total_sign_depth) ) + 1;
1395 0 : ulong fec_blk_max = tile->rotor.max_shreds_per_block / FD_FEC_SHRED_CNT;
1396 :
1397 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
1398 0 : ctx_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
1399 0 : ctx->protocol = FD_SCRATCH_ALLOC_APPEND( l, fd_repair_align(), fd_repair_footprint() );
1400 0 : ctx->chainer = FD_SCRATCH_ALLOC_APPEND( l, fd_chainer_align(), fd_chainer_footprint( tile->rotor.slot_max, tile->rotor.max_shreds_per_block ) );
1401 0 : ctx->policy = FD_SCRATCH_ALLOC_APPEND( l, fd_policy_align(), fd_policy_footprint( FD_REPAIR_PEER_MAX ) );
1402 0 : ctx->dedup = FD_SCRATCH_ALLOC_APPEND( l, fd_reqlim_align(), fd_reqlim_footprint( FD_REQLIM_CACHE_MAX ) );
1403 0 : ctx->inflights = FD_SCRATCH_ALLOC_APPEND( l, fd_inflights_align(), fd_inflights_footprint() );
1404 0 : ctx->signs_map = FD_SCRATCH_ALLOC_APPEND( l, fd_signs_map_align(), fd_signs_map_footprint( lg_sign_depth ) );
1405 0 : ctx->toss_queue = FD_SCRATCH_ALLOC_APPEND( l, toss_queue_align(), toss_queue_footprint() );
1406 0 : ctx->meta_queue = FD_SCRATCH_ALLOC_APPEND( l, meta_queue_align(), meta_queue_footprint( fec_blk_max ) );
1407 0 : ctx->slot_metrics = FD_SCRATCH_ALLOC_APPEND( l, fd_repair_metrics_align(), fd_repair_metrics_footprint() );
1408 0 : ctx->deliver_queue = FD_SCRATCH_ALLOC_APPEND( l, out_queue_align(), out_queue_footprint( (ulong)tile->rotor.slot_max * FD_CHAINER_SLOT_VER_MAX * fec_blk_max ) );
1409 0 : ulong scratch_top = FD_SCRATCH_ALLOC_FINI( l, scratch_align() );
1410 0 : if( FD_UNLIKELY( scratch_top > (ulong)scratch + scratch_footprint( tile ) ) )
1411 0 : FD_LOG_ERR(( "scratch overflow %lu %lu %lu", scratch_top - (ulong)scratch - scratch_footprint( tile ), scratch_top, (ulong)scratch + scratch_footprint( tile ) ));
1412 :
1413 0 : ctx->chainer = fd_chainer_join ( fd_chainer_new ( ctx->chainer, tile->rotor.slot_max, tile->rotor.max_shreds_per_block, ctx->repair_seed ) );
1414 0 : ctx->deliver_queue = out_queue_join ( out_queue_new ( ctx->deliver_queue, (ulong)tile->rotor.slot_max * FD_CHAINER_SLOT_VER_MAX * fec_blk_max ) );
1415 :
1416 0 : ctx->protocol = fd_repair_join ( fd_repair_new ( ctx->protocol, &ctx->identity_public_key ) );
1417 0 : ctx->policy = fd_policy_join ( fd_policy_new ( ctx->policy, FD_REPAIR_PEER_MAX, ctx->repair_seed, ctx->repair_nonce_ss ) );
1418 0 : ctx->dedup = fd_reqlim_join ( fd_reqlim_new ( ctx->dedup, FD_REQLIM_CACHE_MAX, ctx->repair_seed ) );
1419 0 : ctx->inflights = fd_inflights_join ( fd_inflights_new ( ctx->inflights, ctx->repair_seed+1234UL ) );
1420 0 : ctx->signs_map = fd_signs_map_join ( fd_signs_map_new ( ctx->signs_map, lg_sign_depth, 0UL ) );
1421 0 : ctx->toss_queue = toss_queue_join ( toss_queue_new ( ctx->toss_queue ) );
1422 0 : ctx->meta_queue = meta_queue_join ( meta_queue_new ( ctx->meta_queue, fec_blk_max ) );
1423 0 : ctx->slot_metrics = fd_repair_metrics_join( fd_repair_metrics_new( ctx->slot_metrics ) );
1424 0 : ctx->keyswitch = fd_keyswitch_join( fd_topo_obj_laddr( topo, tile->id_keyswitch_obj_id ) );
1425 0 : FD_TEST( ctx->keyswitch );
1426 :
1427 0 : ulong store_obj_id = fd_pod_query_ulong( topo->props, "store", ULONG_MAX );
1428 0 : FD_TEST( store_obj_id!=ULONG_MAX );
1429 0 : ctx->store = fd_store_join( fd_topo_obj_laddr( topo, store_obj_id ) );
1430 0 : FD_TEST( ctx->store );
1431 0 : FD_TEST( fd_store_map_ljoin( ctx->store, ctx->store_map ) );
1432 :
1433 0 : ctx->halt_signing = 0;
1434 :
1435 : /* Flip to 1 to exercise block-id-only repair/catchup */
1436 0 : ctx->block_id_repair_only = 0;
1437 0 : ctx->deliver_from_root = 0;
1438 :
1439 : /* Process in links */
1440 :
1441 0 : if( FD_UNLIKELY( tile->in_cnt > MAX_IN_LINKS ) ) FD_LOG_ERR(( "repair tile has too many input links" ));
1442 :
1443 0 : uint sign_repair_in_idx[ MAX_SIGN_TILE_CNT ] = {0};
1444 0 : uint sign_repair_idx = 0;
1445 0 : ulong sign_link_depth = 0;
1446 :
1447 0 : for( uint in_idx=0U; in_idx<(tile->in_cnt); in_idx++ ) {
1448 0 : fd_topo_link_t const * link = &topo->links[ tile->in_link_id[ in_idx ] ];
1449 0 : if( 0==strcmp( link->name, "net_repair" ) ) {
1450 0 : ctx->in_kind[ in_idx ] = IN_KIND_NET;
1451 0 : fd_net_rx_bounds_init( &ctx->in_links[ in_idx ].net_rx, link->dcache );
1452 0 : continue;
1453 0 : } else if( 0==strcmp( link->name, "sign_repair" ) ) {
1454 0 : ctx->in_kind[ in_idx ] = IN_KIND_SIGN;
1455 0 : sign_repair_in_idx[ sign_repair_idx++ ] = in_idx;
1456 0 : sign_link_depth = link->depth;
1457 0 : }
1458 0 : else if( 0==strcmp( link->name, "gossip_out" ) ) ctx->in_kind[ in_idx ] = IN_KIND_GOSSIP;
1459 0 : else if( 0==strcmp( link->name, "shred_out" ) ) ctx->in_kind[ in_idx ] = IN_KIND_SHRED;
1460 0 : else if( 0==strcmp( link->name, "snapin_manif" ) ) ctx->in_kind[ in_idx ] = IN_KIND_SNAP;
1461 0 : else if( 0==strcmp( link->name, "genesi_out" ) ) ctx->in_kind[ in_idx ] = IN_KIND_GENESIS;
1462 0 : else if( 0==strcmp( link->name, "replay_out" ) ) ctx->in_kind[ in_idx ] = IN_KIND_REPLAY;
1463 0 : else if( 0==strcmp( link->name, "votor_out" ) ) ctx->in_kind[ in_idx ] = IN_KIND_VOTOR;
1464 0 : else FD_LOG_ERR(( "repair tile has unexpected input link %s", link->name ));
1465 :
1466 0 : ctx->in_links[ in_idx ].mem = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
1467 0 : ctx->in_links[ in_idx ].chunk0 = fd_dcache_compact_chunk0( ctx->in_links[ in_idx ].mem, link->dcache );
1468 0 : ctx->in_links[ in_idx ].wmark = fd_dcache_compact_wmark ( ctx->in_links[ in_idx ].mem, link->dcache, link->mtu );
1469 0 : ctx->in_links[ in_idx ].mtu = link->mtu;
1470 :
1471 0 : FD_TEST( fd_dcache_compact_is_safe( ctx->in_links[in_idx].mem, link->dcache, link->mtu, link->depth ) );
1472 0 : }
1473 :
1474 0 : ctx->net_out_ctx->idx = UINT_MAX;
1475 0 : ctx->repair_out_ctx->idx = UINT_MAX;
1476 0 : ctx->repair_sign_cnt = 0;
1477 0 : ctx->sign_rrobin_idx = 0;
1478 :
1479 0 : for( uint out_idx=0U; out_idx<(tile->out_cnt); out_idx++ ) {
1480 0 : fd_topo_link_t const * link = &topo->links[ tile->out_link_id[ out_idx ] ];
1481 :
1482 0 : if( 0==strcmp( link->name, "repair_net" ) ) {
1483 :
1484 0 : if( ctx->net_out_ctx->idx!=UINT_MAX ) continue; /* only use first net link */
1485 0 : ctx->net_out_ctx->idx = out_idx;
1486 0 : ctx->net_out_ctx->mem = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
1487 0 : ctx->net_out_ctx->chunk0 = fd_dcache_compact_chunk0( ctx->net_out_ctx->mem, link->dcache );
1488 0 : ctx->net_out_ctx->wmark = fd_dcache_compact_wmark( ctx->net_out_ctx->mem, link->dcache, link->mtu );
1489 0 : ctx->net_out_ctx->chunk = ctx->net_out_ctx->chunk0;
1490 :
1491 0 : } else if( 0==strcmp( link->name, "repair_out" ) ) {
1492 :
1493 0 : out_ctx_t * replay_out = ctx->repair_out_ctx;
1494 0 : replay_out->idx = out_idx;
1495 0 : replay_out->mem = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
1496 0 : replay_out->chunk0 = fd_dcache_compact_chunk0( replay_out->mem, link->dcache );
1497 0 : replay_out->wmark = fd_dcache_compact_wmark( replay_out->mem, link->dcache, link->mtu );
1498 0 : replay_out->chunk = replay_out->chunk0;
1499 :
1500 0 : } else if( 0==strcmp( link->name, "repair_sign" ) ) {
1501 :
1502 0 : out_ctx_t * repair_sign_out = &ctx->repair_sign_out_ctx[ ctx->repair_sign_cnt ];
1503 0 : repair_sign_out->idx = out_idx;
1504 0 : repair_sign_out->mem = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
1505 0 : repair_sign_out->chunk0 = fd_dcache_compact_chunk0( repair_sign_out->mem, link->dcache );
1506 0 : repair_sign_out->wmark = fd_dcache_compact_wmark( repair_sign_out->mem, link->dcache, link->mtu );
1507 0 : repair_sign_out->chunk = repair_sign_out->chunk0;
1508 0 : repair_sign_out->in_idx = sign_repair_in_idx[ ctx->repair_sign_cnt++ ]; /* match to the sign_repair input link */
1509 0 : repair_sign_out->max_credits = sign_link_depth;
1510 0 : repair_sign_out->credits = sign_link_depth;
1511 :
1512 0 : } else {
1513 0 : FD_LOG_ERR(( "repair tile has unexpected output link %s", link->name ));
1514 0 : }
1515 0 : }
1516 0 : FD_TEST( ctx->net_out_ctx->idx!=UINT_MAX );
1517 0 : FD_TEST( ctx->repair_out_ctx->idx!=UINT_MAX );
1518 0 : if( FD_UNLIKELY( ctx->repair_sign_cnt!=sign_repair_idx ) ) {
1519 0 : FD_LOG_ERR(( "Mismatch between repair_sign output links (%lu) and sign_repair input links (%u)", ctx->repair_sign_cnt, sign_repair_idx ));
1520 0 : }
1521 0 : if( FD_UNLIKELY( fd_signs_map_key_max( ctx->signs_map ) < tile->rotor.repair_sign_depth * tile->rotor.repair_sign_cnt ) ) {
1522 0 : FD_LOG_ERR(( "Repair pending signs tracking map is too small: %lu < %lu.", fd_signs_map_key_max( ctx->signs_map ), tile->rotor.repair_sign_depth * tile->rotor.repair_sign_cnt ));
1523 0 : }
1524 :
1525 0 : ctx->wksp = topo->workspaces[ topo->objs[ tile->tile_obj_id ].wksp_id ].wksp;
1526 :
1527 : /* TODO clean these up */
1528 0 : ctx->net_id = (ushort)0;
1529 0 : fd_ip4_udp_hdr_init( ctx->intake_hdr, 0, 0, tile->rotor.repair_client_listen_port );
1530 :
1531 : /* Repair set up */
1532 :
1533 0 : ctx->turbine_slot0 = ULONG_MAX;
1534 :
1535 0 : memset( ctx->metrics, 0, sizeof(ctx->metrics) );
1536 :
1537 0 : fd_histf_join( fd_histf_new( ctx->metrics->slot_compl_time, FD_MHIST_SECONDS_MIN( REPAIR, SLOT_COMPLETE_DURATION_SECONDS ),
1538 0 : FD_MHIST_SECONDS_MAX( REPAIR, SLOT_COMPLETE_DURATION_SECONDS ) ) );
1539 0 : fd_histf_join( fd_histf_new( ctx->metrics->response_latency, FD_MHIST_MIN( REPAIR, RESPONSE_LATENCY_NANOS ),
1540 0 : FD_MHIST_MAX( REPAIR, RESPONSE_LATENCY_NANOS ) ) );
1541 :
1542 0 : fd_clock_tile_init( ctx->clock );
1543 0 : ctx->pending_key_next = 0;
1544 0 : ctx->ag_nonce = 0;
1545 0 : }
1546 :
1547 : static ulong
1548 : populate_allowed_seccomp( fd_topo_t const * topo FD_PARAM_UNUSED,
1549 : fd_topo_tile_t const * tile FD_PARAM_UNUSED,
1550 : ulong out_cnt,
1551 0 : struct sock_filter * out ) {
1552 0 : populate_sock_filter_policy_fd_rotor_tile( out_cnt, out, (uint)fd_log_private_logfile_fd() );
1553 0 : return sock_filter_policy_fd_rotor_tile_instr_cnt;
1554 0 : }
1555 :
1556 : static ulong
1557 : populate_allowed_fds( fd_topo_t const * topo FD_PARAM_UNUSED,
1558 : fd_topo_tile_t const * tile FD_PARAM_UNUSED,
1559 : ulong out_fds_cnt,
1560 0 : int * out_fds ) {
1561 0 : if( FD_UNLIKELY( out_fds_cnt<2UL ) ) FD_LOG_ERR(( "out_fds_cnt %lu", out_fds_cnt ));
1562 :
1563 0 : ulong out_cnt = 0UL;
1564 0 : out_fds[ out_cnt++ ] = 2; /* stderr */
1565 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) )
1566 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
1567 0 : return out_cnt;
1568 0 : }
1569 :
1570 : static inline void
1571 0 : metrics_write( ctx_t * ctx ) {
1572 0 : FD_MGAUGE_SET( ROTOR, SLOT_CURRENT, ctx->metrics->current_slot );
1573 0 : FD_MGAUGE_SET( ROTOR, SLOT_HIGHEST_REPAIRED, ctx->chainer->highest_repaired ); //fd_forest_highest_repaired_slot( ctx->forest ) );
1574 0 : FD_MCNT_SET( ROTOR, SHRED_OLD, ctx->metrics->old_shred );
1575 0 : FD_MCNT_SET( ROTOR, PEER_REQUESTED, fd_policy_peer_pool_used( ctx->policy->peers.pool ) );
1576 0 : FD_MCNT_SET( ROTOR, SHRED_REREQUESTED, ctx->metrics->rerequest );
1577 :
1578 0 : FD_MGAUGE_SET( ROTOR, SLOT_LAST_REQUESTED, ctx->metrics->last_requested_slot );
1579 0 : FD_MGAUGE_SET( ROTOR, ORPHAN_LAST_REQUESTED, ctx->metrics->last_requested_orphan );
1580 0 : FD_MGAUGE_SET( ROTOR, REQUEST_INFLIGHT, fd_inflights_outstanding_cnt( ctx->inflights ) );
1581 :
1582 0 : FD_MCNT_SET ( ROTOR, PKT_TX, ctx->metrics->send_pkt_cnt );
1583 0 : FD_MCNT_ENUM_COPY( ROTOR, REQUEST_TX, ctx->metrics->sent_pkt_types );
1584 :
1585 0 : FD_MHIST_COPY( ROTOR, RESPONSE_LATENCY_NANOS, ctx->metrics->response_latency );
1586 :
1587 0 : FD_MCNT_SET( ROTOR, PING_UNKNOWN_PEER, ctx->metrics->unknown_peer_ping );
1588 0 : FD_MCNT_SET( ROTOR, PING_MALFORMED, ctx->metrics->malformed_ping );
1589 0 : FD_MCNT_SET( ROTOR, PING_SIGNATURE_FAILED, ctx->metrics->fail_sigverify_ping );
1590 :
1591 0 : FD_MCNT_SET( ROTOR, SHRED_BLOCK_ID_FAILED, ctx->metrics->failed_shred_block_id_cnt );
1592 0 : FD_MCNT_SET( ROTOR, FEC_ROOT_FAILED, ctx->metrics->failed_fec_root_cnt );
1593 0 : FD_MCNT_SET( ROTOR, PARENT_FEC_COUNT_FAILED, ctx->metrics->failed_parent_fec_count_cnt );
1594 0 : }
1595 :
1596 : #undef DEBUG_LOGGING
1597 :
1598 : /* At most one sign request is made in after_credit. Then at most one
1599 : message is published in after_frag. */
1600 0 : #define STEM_BURST (3UL)
1601 :
1602 : /* Set LAZY to a reasonable value that keeps housekeeping time low.
1603 : Repair tile's only reliable consumer is replay. */
1604 0 : #define STEM_LAZY (64000)
1605 :
1606 0 : #define STEM_CALLBACK_CONTEXT_TYPE ctx_t
1607 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(ctx_t)
1608 :
1609 0 : #define STEM_CALLBACK_AFTER_CREDIT after_credit
1610 0 : #define STEM_CALLBACK_BEFORE_FRAG before_frag
1611 0 : #define STEM_CALLBACK_DURING_FRAG during_frag
1612 0 : #define STEM_CALLBACK_AFTER_FRAG after_frag
1613 0 : #define STEM_CALLBACK_DURING_HOUSEKEEPING during_housekeeping
1614 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
1615 :
1616 : #include "../../disco/stem/fd_stem.c"
1617 :
1618 : fd_topo_run_tile_t fd_tile_rotor = {
1619 : .name = "rotor",
1620 : .loose_footprint = loose_footprint,
1621 : .populate_allowed_seccomp = populate_allowed_seccomp,
1622 : .populate_allowed_fds = populate_allowed_fds,
1623 : .scratch_align = scratch_align,
1624 : .scratch_footprint = scratch_footprint,
1625 : .unprivileged_init = unprivileged_init,
1626 : .privileged_init = privileged_init,
1627 : .run = stem_run,
1628 : };
|