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 */