LCOV - code coverage report
Current view: top level - discof/rotor - fd_rotor_tile.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 530 949 55.8 %
Date: 2026-09-17 04:28:31 Functions: 23 108 21.3 %

          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, &eth, &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             : };

Generated by: LCOV version 1.14