LCOV - code coverage report
Current view: top level - disco/bundle - fd_bundle_tile_private.h (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 12 12 100.0 %
Date: 2026-09-17 04:28:31 Functions: 4 12 33.3 %

          Line data    Source code
       1             : #ifndef HEADER_fd_src_disco_bundle_fd_bundle_tile_private_h
       2             : #define HEADER_fd_src_disco_bundle_fd_bundle_tile_private_h
       3             : 
       4             : #include "fd_bundle_auth.h"
       5             : #include "fd_keepalive.h"
       6             : #include "../fd_clock_tile.h"
       7             : #include "../stem/fd_stem.h"
       8             : #include "../keyguard/fd_keyswitch.h"
       9             : #include "../keyguard/fd_keyguard_client.h"
      10             : #include "../../waltz/grpc/fd_grpc_client.h"
      11             : #include "../../waltz/http/fd_url.h"
      12             : #include "../../waltz/resolv/fd_netdb.h"
      13             : #include "../../waltz/fd_rtt_est.h"
      14             : #include "../../util/hist/fd_histf.h"
      15             : 
      16             : #define FD_BUNDLE_CLIENT_MAX_TXN_PER_BUNDLE (5UL)
      17             : 
      18             : /* Pending transaction buffer.  gRPC callbacks push decoded transactions
      19             :    here.  after_credit drains one bundle per call by writing to dcache
      20             :    and calling fd_stem_publish.
      21             : 
      22             :    Sized to match the bundle_verif output link depth. */
      23             : 
      24             : struct fd_bundle_pending_txn {
      25             :   uchar  payload[ FD_TXN_MTU ];
      26             :   ushort payload_sz;
      27             :   uint   source_ipv4;
      28             :   long   first_seen_nanos;
      29             :   ulong  sig;
      30             :   ulong  bundle_seq;
      31             :   ulong  bundle_txn_cnt;
      32             :   uchar  commission;
      33             :   uchar  commission_pubkey[ 32 ];
      34             : };
      35             : 
      36             : typedef struct fd_bundle_pending_txn fd_bundle_pending_txn_t;
      37             : 
      38             : #define DEQUE_NAME pending_txn
      39         180 : #define DEQUE_T    fd_bundle_pending_txn_t
      40             : #include "../../util/tmpl/fd_deque_dynamic.c"
      41             : 
      42             : /* Returns true if the drain loop should continue after popping an
      43             :    entry.  Bundles drain atomically (all txns with matching bundle_seq).
      44             :    Packets drain up to burst consecutive entries. */
      45             : 
      46             : static inline int
      47             : fd_bundle_drain_continue( fd_bundle_pending_txn_t * txns,
      48             :                           ulong                     drain_sig,
      49             :                           ulong                     drain_seq,
      50             :                           ulong                     drain_cnt,
      51         132 :                           ulong                     burst ) {
      52         132 :   if( pending_txn_empty( txns ) ) return 0;
      53         108 :   if( drain_sig==1UL ) return pending_txn_peek_head( txns )->bundle_seq==drain_seq;
      54          54 :   return drain_cnt<burst && pending_txn_peek_head( txns )->sig==0UL;
      55         108 : }
      56             : 
      57             : #include "../../waltz/tlsrec/fd_tlsrec.h"
      58             : #include "../../ballet/x509/fd_x509_ca_store.h"
      59             : #include "../../ballet/x509/fd_x509_verify.h"
      60             : 
      61             : struct fd_bundle_out_ctx {
      62             :   ulong       idx;
      63             :   fd_wksp_t * mem;
      64             :   ulong       chunk0;
      65             :   ulong       wmark;
      66             :   ulong       chunk;
      67             : };
      68             : 
      69             : typedef struct fd_bundle_out_ctx fd_bundle_out_ctx_t;
      70             : 
      71             : /* fd_bundle_metrics_t contains private metric counters.  These get
      72             :    published to fd_metrics periodically. */
      73             : 
      74             : struct fd_bundle_metrics {
      75             :   ulong txn_received_cnt;
      76             :   ulong bundle_received_cnt;
      77             :   ulong packet_received_cnt;
      78             :   ulong proto_received_bytes;
      79             :   ulong shredstream_heartbeat_cnt;
      80             :   ulong ping_ack_cnt;
      81             : 
      82             :   ulong decode_fail_cnt;
      83             :   ulong transport_fail_cnt;
      84             :   ulong missing_builder_info_fail_cnt;
      85             :   ulong backpressure_drop_cnt;
      86             : 
      87             :   fd_histf_t msg_rx_delay[1];
      88             : };
      89             : 
      90             : typedef struct fd_bundle_metrics fd_bundle_metrics_t;
      91             : 
      92             : /* fd_bundle_tile_t is the context object provided to callbacks from
      93             :    stem, and contains all state needed to progress the tile. */
      94             : 
      95             : struct fd_bundle_tile {
      96             :   /* Key switch */
      97             :   fd_keyswitch_t * keyswitch;
      98             : 
      99             :   /* Key guard */
     100             :   fd_keyguard_client_t keyguard_client[1];
     101             : 
     102             :   uint is_ssl : 1;
     103             :   int  keylog_fd;
     104             : 
     105             :   /* Native TLS */
     106             :   fd_tls_t           tls[1];
     107             :   fd_chacha_rng_t    tls_rng[1];
     108             :   fd_tlsrec_conn_t   tls_conn[1];
     109             :   fd_x509_ca_store_t ca_store[1];
     110             : 
     111             :   /* Config */
     112             :   char   server_fqdn[ FD_FQDN_BUF_MAX ]; /* cstr */
     113             :   ulong  server_fqdn_len;
     114             :   char   server_sni[ FD_SNI_BUF_MAX ]; /* cstr */
     115             :   ulong  server_sni_len;
     116             :   ushort server_tcp_port;
     117             : 
     118             :   /* Resolver */
     119             :   fd_netdb_fds_t netdb_fds[1];
     120             :   uint server_ip4_addr; /* last DNS lookup result */
     121             : 
     122             :   /* TCP socket */
     123             :   int  tcp_sock;
     124             :   int  so_rcvbuf;
     125             :   uint tcp_sock_connected : 1;
     126             :   uint defer_reset : 1;
     127             :   uint sock_in_epoll : 1;
     128             :   uint epoll_out_armed : 1;
     129             :   long cached_ts;
     130             : 
     131             :   ulong   waker_client_idx;
     132             :   ulong * waker_fseq;
     133             :   long    next_step_deadline;
     134             : 
     135             :   /* Keepalive via HTTP/2 PINGs (randomized) */
     136             :   long              keepalive_interval;
     137             :   fd_keepalive_t    keepalive[1];
     138             :   fd_rtt_estimate_t rtt[1];
     139             : 
     140             :   /* gRPC client */
     141             :   void *                   grpc_client_mem;
     142             :   ulong                    grpc_buf_max;
     143             :   fd_grpc_client_t *       grpc_client;
     144             :   fd_grpc_client_metrics_t grpc_metrics[1];
     145             :   ulong                    map_seed;
     146             : 
     147             :   /* Bundle authenticator */
     148             :   fd_bundle_auther_t auther;
     149             : 
     150             :   /* Bundle block builder info */
     151             :   uchar builder_pubkey[ 32 ];
     152             :   uchar builder_commission;  /* in [0,100] (percent) */
     153             :   uchar builder_info_avail : 1;  /* Block builder info available? (potentially stale) */
     154             :   uchar builder_info_wait  : 1;  /* Request already in-flight? */
     155             :   long  builder_info_valid_until;
     156             : 
     157             :   /* Bundle subscriptions */
     158             :   uchar packet_subscription_live : 1;  /* Want to subscribe to a stream? */
     159             :   uchar packet_subscription_wait : 1;  /* Request already in-flight? */
     160             :   uchar bundle_subscription_live : 1;
     161             :   uchar bundle_subscription_wait : 1;
     162             : 
     163             :   /* Bundle state */
     164             :   ulong bundle_seq;
     165             :   ulong bundle_txn_cnt;
     166             : 
     167             :   /* Error backoff */
     168             :   fd_rng_t rng[1];
     169             :   uint     backoff_iter;
     170             :   long     backoff_until;
     171             :   long     backoff_reset;
     172             : 
     173             :   /* Stem publish */
     174             :   fd_stem_context_t *       stem;
     175             :   fd_bundle_out_ctx_t       verify_out;
     176             :   fd_bundle_out_ctx_t       plugin_out;
     177             :   fd_bundle_pending_txn_t * pending_txns;
     178             : 
     179             :   /* App metrics */
     180             :   fd_bundle_metrics_t metrics;
     181             : 
     182             :   /* Check engine light */
     183             :   uchar bundle_status_recent;  /* most recently observed 'check engine light' */
     184             :   uchar bundle_status_plugin;  /* last 'plugin' update written */
     185             :   uchar bundle_status_logged;
     186             :   long  last_bundle_status_log_nanos;
     187             : 
     188             :   ulong next_leader_slot; /* from replay_out reset messages, or ULONG_MAX */
     189             :   ulong reset_slot;       /* from replay_out reset messages, or ULONG_MAX */
     190             :   int   sleep_mode;       /* 1 means sleeping, 0 means connecting/connected */
     191             :   long  sleep_check_ns;   /* next wallclock time to re-evaluate sleeping */
     192             : 
     193             :   fd_clock_tile_t clock[1]; /* fast wallclock-ns source for fd_bundle_now */
     194             :   int   halt_signing;     /* 1 means signing is halted, 0 means signing is not halted */
     195             : 
     196             :   /* Staged values from during_frag, committed in after_frag */
     197             :   ulong next_leader_slot_staged;
     198             :   ulong reset_slot_staged;
     199             : 
     200             :   int   in_kind[ 64 ];
     201             :   struct {
     202             :     fd_wksp_t * mem;
     203             :     ulong       chunk0;
     204             :     ulong       wmark;
     205             :   } replay_in;
     206             : };
     207             : 
     208             : typedef struct fd_bundle_tile fd_bundle_tile_t;
     209             : 
     210             : /* Define 'request_ctx' IDs to identify different types of gRPC calls */
     211             : 
     212         102 : #define FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets            4
     213         138 : #define FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles            5
     214          36 : #define FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo      6
     215             : 
     216             : FD_PROTOTYPES_BEGIN
     217             : 
     218             : /* fd_bundle_now returns wallclock nanoseconds from the tile's
     219             :    fd_clock.  Backed by a weak symbol, allowing tests to override the
     220             :    clock source. */
     221             : 
     222             : long
     223             : fd_bundle_now( fd_bundle_tile_t const * ctx );
     224             : 
     225             : /* fd_bundle_client_grpc_callbacks provides callbacks for grpc_client. */
     226             : 
     227             : extern fd_grpc_client_callbacks_t fd_bundle_client_grpc_callbacks;
     228             : 
     229             : /* fd_bundle_client_step is an all-in-one routine to drive client logic.
     230             :    As long as the tile calls this periodically, the client will
     231             :    reconnect to the bundle server, authenticate, and subscribe to
     232             :    packets and bundles. */
     233             : 
     234             : void
     235             : fd_bundle_client_step( fd_bundle_tile_t * bundle,
     236             :                        int *              charge_busy );
     237             : 
     238             : /* fd_bundle_client_step_reconnect drives the 'reconnect' state machine.
     239             :    Once the HTTP/2 conn is established (SETTINGS exchanged), this
     240             :    function drives the auth logic, requests block builder info, sets up
     241             :    packet and bundle subscriptions, and PINGs. */
     242             : 
     243             : int
     244             : fd_bundle_client_step_reconnect( fd_bundle_tile_t * ctx,
     245             :                                  long               now );
     246             : 
     247             : /* fd_bundle_client_next_deadline returns when
     248             :    fd_bundle_client_step next needs to run absent any fd event: the
     249             :    earliest pending timeout (keepalive, gRPC deadlines, builder info
     250             :    expiry, reconnect backoff).  Returns now if TLS bytes are buffered
     251             :    inside OpenSSL (no fd event will announce them), LONG_MAX while
     252             :    connecting (completion is an EPOLLOUT event). */
     253             : 
     254             : long
     255             : fd_bundle_client_next_deadline( fd_bundle_tile_t const * ctx,
     256             :                                 long                     now );
     257             : 
     258             : /* fd_bundle_tile_backoff is called whenever an error occurs.  Stalls
     259             :    forward progress for a randomized amount of time to prevent error
     260             :    floods. */
     261             : 
     262             : void
     263             : fd_bundle_tile_backoff( fd_bundle_tile_t * ctx,
     264             :                         long               now );
     265             : 
     266             : /* fd_bundle_tile_should_stall returns 1 if forward progress should be
     267             :    temporarily prevented due to an error. */
     268             : 
     269             : FD_FN_PURE static inline int
     270             : fd_bundle_tile_should_stall( fd_bundle_tile_t const * ctx,
     271          60 :                              long                     now ) {
     272          60 :   return now < ctx->backoff_until;
     273          60 : }
     274             : 
     275             : /* fd_bundle_tile_housekeeping runs periodically at a low frequency. */
     276             : 
     277             : void
     278             : fd_bundle_tile_housekeeping( fd_bundle_tile_t * ctx );
     279             : 
     280             : /* fd_bundle_client_grpc_rx_start is the first RX callback of a stream. */
     281             : 
     282             : void
     283             : fd_bundle_client_grpc_rx_start(
     284             :     void * app_ctx,
     285             :     ulong  request_ctx
     286             : ) ;
     287             : 
     288             : /* fd_bundle_client_grpc_rx_msg is called by grpc_client when a gRPC
     289             :    message arrives (unary or server-streaming response). */
     290             : 
     291             : void
     292             : fd_bundle_client_grpc_rx_msg(
     293             :     void *       app_ctx,      /* (fd_bundle_tile_t *) */
     294             :     void const * protobuf,
     295             :     ulong        protobuf_sz,
     296             :     ulong        request_ctx   /* FD_BUNDLE_CLIENT_REQ_{...} */
     297             : );
     298             : 
     299             : /* fd_bundle_client_grpc_rx_end is called by grpc_client when a gRPC
     300             :    server-streaming response finishes. */
     301             : 
     302             : void
     303             : fd_bundle_client_grpc_rx_end(
     304             :     void *                app_ctx,
     305             :     ulong                 request_ctx,
     306             :     fd_grpc_resp_hdrs_t * resp
     307             : );
     308             : 
     309             : /* fd_bundle_client_grpc_rx_timeout is called by grpc_client when a
     310             :    gRPC request deadline gets exceeded. */
     311             : 
     312             : void
     313             : fd_bundle_client_grpc_rx_timeout(
     314             :     void * app_ctx,
     315             :     ulong  request_ctx, /* FD_BUNDLE_CLIENT_REQ_{...} */
     316             :     int    deadline_kind /* FD_GRPC_DEADLINE_{HEADER|RX_END} */
     317             : );
     318             : 
     319             : /* fd_bundle_client_status provides a "check engine light".
     320             : 
     321             :    Returns 0 if the client has recently failed and is currently backing
     322             :    off from a reconnect attempt.
     323             : 
     324             :    Returns 1 if the client is currently reconnecting.
     325             : 
     326             :    Returns 2 if all of the following conditions are met:
     327             :    - TCP socket is alive
     328             :    - SSL session is not in an error state
     329             :    - HTTP/2 connection is established (SETTINGS exchange done)
     330             :    - gRPC bundle and packet subscriptions are live
     331             :    - HTTP/2 PING exchange was done recently
     332             : 
     333             :    Return codes are compatible with FD_PLUGIN_MSG_BLOCK_ENGINE_UPDATE_STATUS_{...}. */
     334             : 
     335             : int
     336             : fd_bundle_client_status( fd_bundle_tile_t const * ctx );
     337             : 
     338             : /* fd_bundle_request_ctx_cstr returns the gRPC method name for a
     339             :    FD_BUNDLE_CLIENT_REQ_* ID.  Returns "unknown" the ID is not
     340             :    recognized. */
     341             : 
     342             : FD_FN_CONST char const *
     343             : fd_bundle_request_ctx_cstr( ulong request_ctx );
     344             : 
     345             : /* fd_bundle_client_reset frees all connection-related resources. */
     346             : 
     347             : void
     348             : fd_bundle_client_reset( fd_bundle_tile_t * ctx );
     349             : 
     350             : /* fd_bundle_client_ping_tx enqueues a PING frame for sending.  Returns
     351             :    1 on success and 0 on failure (occurs when frame_tx buf is full). */
     352             : 
     353             : void
     354             : fd_bundle_client_send_ping( fd_bundle_tile_t * ctx );
     355             : 
     356             : FD_PROTOTYPES_END
     357             : 
     358             : #endif /* HEADER_fd_src_disco_bundle_fd_bundle_tile_private_h */

Generated by: LCOV version 1.14