Line data Source code
1 : /* fd_bundle_client.c steps gRPC related tasks. */
2 :
3 : #define _GNU_SOURCE /* SOL_TCP */
4 : #include "fd_bundle_auth.h"
5 : #include "fd_bundle_tile_private.h"
6 : #include "fd_bundle_tile.h"
7 : #include "proto/block_engine.pb.h"
8 : #include "proto/bundle.pb.h"
9 : #include "proto/packet.pb.h"
10 : #include "../waker/fd_waker.h"
11 : #include "../fd_txn_m.h"
12 : #include "../../waltz/h2/fd_h2_conn.h"
13 : #include "../../waltz/http/fd_url.h" /* fd_url_unescape */
14 : #include "../../ballet/base58/fd_base58.h"
15 : #include "../../ballet/ed25519/fd_x25519.h"
16 : #include "../../third_party/nanopb/pb_decode.h"
17 : #include "../../util/net/fd_ip4.h"
18 :
19 : #include <fcntl.h>
20 : #include <errno.h>
21 : #include <unistd.h> /* close */
22 : #include <poll.h> /* poll */
23 : #include <sys/epoll.h>
24 : #include <sys/socket.h> /* socket */
25 : #include <netinet/in.h>
26 : #include <netinet/ip.h>
27 : #include <netinet/tcp.h>
28 :
29 9 : #define FD_BUNDLE_CLIENT_REQUEST_TIMEOUT ((long)8e9) /* 8 seconds */
30 :
31 :
32 : __attribute__((weak)) long
33 1158 : fd_bundle_now( fd_bundle_tile_t const * ctx ) {
34 1158 : return fd_clock_tile_now( ctx->clock );
35 1158 : }
36 :
37 : void
38 6 : fd_bundle_client_reset( fd_bundle_tile_t * ctx ) {
39 6 : if( FD_UNLIKELY( ctx->tcp_sock >= 0 ) ) {
40 6 : if( FD_UNLIKELY( 0!=close( ctx->tcp_sock ) ) ) {
41 0 : FD_LOG_ERR(( "close(tcp_sock=%i) failed (%i-%s)", ctx->tcp_sock, errno, fd_io_strerror( errno ) ));
42 0 : }
43 6 : ctx->tcp_sock = -1;
44 6 : ctx->tcp_sock_connected = 0;
45 6 : ctx->sock_in_epoll = 0;
46 6 : ctx->epoll_out_armed = 0;
47 6 : }
48 6 : ctx->defer_reset = 0;
49 :
50 6 : ctx->builder_info_avail = 0;
51 6 : ctx->builder_info_wait = 0;
52 6 : ctx->packet_subscription_live = 0;
53 6 : ctx->packet_subscription_wait = 0;
54 6 : ctx->bundle_subscription_live = 0;
55 6 : ctx->bundle_subscription_wait = 0;
56 :
57 6 : fd_memset( ctx->rtt, 0, sizeof(fd_rtt_estimate_t) );
58 :
59 6 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
60 :
61 6 : fd_bundle_auther_reset( &ctx->auther );
62 6 : fd_grpc_client_reset( ctx->grpc_client );
63 6 : }
64 :
65 : static int
66 : fd_bundle_client_do_connect( fd_bundle_tile_t const * ctx,
67 0 : uint ip4_addr ) {
68 0 : struct sockaddr_in addr = {
69 0 : .sin_family = AF_INET,
70 0 : .sin_addr.s_addr = ip4_addr,
71 0 : .sin_port = fd_ushort_bswap( ctx->server_tcp_port )
72 0 : };
73 0 : int err = connect( ctx->tcp_sock, fd_type_pun_const( &addr ), sizeof(struct sockaddr_in) );
74 : /* FD_LIKELY is used here as EINPROGRESS is expected even to local tcp ports */
75 0 : if( FD_LIKELY( err==-1 ) ) {
76 0 : return errno;
77 0 : }
78 0 : return 0;
79 0 : }
80 :
81 : static int
82 0 : fd_bundle_client_get_connect_result( fd_bundle_tile_t const * ctx ) {
83 0 : int so_err = 0;
84 0 : socklen_t so_err_sz = sizeof(so_err);
85 0 : if( FD_UNLIKELY( getsockopt( ctx->tcp_sock, SOL_SOCKET, SO_ERROR, &so_err, &so_err_sz )==-1 ) ) {
86 0 : return errno;
87 0 : }
88 0 : return so_err;
89 0 : }
90 :
91 : static void
92 0 : fd_bundle_client_create_conn( fd_bundle_tile_t * ctx ) {
93 0 : fd_bundle_client_reset( ctx );
94 :
95 : /* FIXME IPv6 support */
96 0 : fd_addrinfo_t hints = {0};
97 0 : hints.ai_family = AF_INET;
98 0 : fd_addrinfo_t * res = NULL;
99 0 : uchar scratch[ 4096 ];
100 0 : void * pscratch = scratch;
101 0 : int err = fd_getaddrinfo( ctx->server_fqdn, &hints, &res, &pscratch, sizeof(scratch) );
102 0 : if( FD_UNLIKELY( err ) ) {
103 0 : FD_LOG_WARNING(( "fd_getaddrinfo `%s` failed (%d-%s)", ctx->server_fqdn, err, fd_gai_strerror( err ) ));
104 0 : fd_bundle_client_reset( ctx );
105 0 : ctx->metrics.transport_fail_cnt++;
106 0 : return;
107 0 : }
108 0 : uint const ip4_addr = ((struct sockaddr_in *)res->ai_addr)->sin_addr.s_addr;
109 0 : ctx->server_ip4_addr = ip4_addr;
110 :
111 0 : int tcp_sock = socket( AF_INET, SOCK_STREAM|SOCK_CLOEXEC, 0 );
112 0 : if( FD_UNLIKELY( tcp_sock<0 ) ) {
113 0 : FD_LOG_ERR(( "socket(AF_INET,SOCK_STREAM|SOCK_CLOEXEC,0) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
114 0 : }
115 0 : ctx->tcp_sock = tcp_sock;
116 :
117 0 : struct epoll_event ev = { .events = EPOLLIN|EPOLLOUT, .data.fd = tcp_sock };
118 0 : if( FD_UNLIKELY( -1==epoll_ctl( FD_WAKER_INNER_FD( ctx->waker_client_idx ), EPOLL_CTL_ADD, tcp_sock, &ev ) ) ) {
119 0 : FD_LOG_ERR(( "epoll_ctl(ADD,tcp_sock) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
120 0 : }
121 0 : ctx->sock_in_epoll = 1;
122 0 : ctx->epoll_out_armed = 1;
123 :
124 0 : if( FD_UNLIKELY( 0!=setsockopt( tcp_sock, SOL_SOCKET, SO_RCVBUF, &ctx->so_rcvbuf, sizeof(int) ) ) ) {
125 0 : FD_LOG_ERR(( "setsockopt(SOL_SOCKET,SO_RCVBUF,%i) failed (%i-%s)", ctx->so_rcvbuf, errno, fd_io_strerror( errno ) ));
126 0 : }
127 :
128 0 : int tcp_nodelay = 1;
129 0 : if( FD_UNLIKELY( 0!=setsockopt( tcp_sock, SOL_TCP, TCP_NODELAY, &tcp_nodelay, sizeof(int) ) ) ) {
130 0 : FD_LOG_ERR(( "setsockopt failed (%d-%s)", errno, fd_io_strerror( errno ) ));
131 0 : }
132 :
133 0 : if( FD_UNLIKELY( fcntl( tcp_sock, F_SETFL, O_NONBLOCK )==-1 ) ) {
134 0 : FD_LOG_ERR(( "fcntl(tcp_sock,F_SETFL,O_NONBLOCK) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
135 0 : }
136 :
137 0 : char const * scheme = ctx->is_ssl ? "https" : "http";
138 :
139 0 : FD_LOG_INFO(( "Connecting to %s://" FD_IP4_ADDR_FMT ":%hu (%.*s)",
140 0 : scheme,
141 0 : FD_IP4_ADDR_FMT_ARGS( ip4_addr ), ctx->server_tcp_port,
142 0 : (int)ctx->server_sni_len, ctx->server_sni ));
143 :
144 0 : int connect_err = fd_bundle_client_do_connect( ctx, ip4_addr );
145 : /* FD_LIKELY as EINPROGRESS is expected */
146 0 : if( FD_LIKELY( connect_err ) ) {
147 0 : if( FD_UNLIKELY( connect_err!=EINPROGRESS ) ) {
148 0 : FD_LOG_WARNING(( "connect(tcp_sock," FD_IP4_ADDR_FMT ":%u) failed (%i-%s)",
149 0 : FD_IP4_ADDR_FMT_ARGS( ip4_addr ), ctx->server_tcp_port,
150 0 : connect_err, fd_io_strerror( connect_err ) ));
151 0 : fd_bundle_client_reset( ctx );
152 0 : ctx->metrics.transport_fail_cnt++;
153 0 : return;
154 0 : }
155 0 : }
156 :
157 0 : if( ctx->is_ssl ) {
158 0 : fd_tls_t * tls = ctx->tls;
159 0 : ulong sni_len = strlen( ctx->server_sni );
160 0 : fd_memcpy( tls->server_name, ctx->server_sni, sni_len );
161 0 : tls->server_name[ sni_len ] = '\0';
162 0 : tls->server_name_len = (ushort)sni_len;
163 0 : if( FD_UNLIKELY( !fd_rng_secure( tls->kex_private_key, 32UL ) ) ) FD_LOG_CRIT(( "fd_rng_secure failed" ));
164 0 : fd_x25519_public( tls->kex_public_key, tls->kex_private_key );
165 0 : fd_tlsrec_conn_init( ctx->tls_conn, tls, 0 );
166 0 : }
167 :
168 0 : fd_grpc_client_reset( ctx->grpc_client );
169 0 : fd_keepalive_init( ctx->keepalive, ctx->rng, ctx->keepalive_interval, ctx->keepalive_interval, fd_bundle_now( ctx ) );
170 0 : }
171 :
172 : static int
173 : fd_bundle_client_drive_io( fd_bundle_tile_t * ctx,
174 : long now,
175 27 : int * charge_busy ) {
176 27 : if( ctx->is_ssl ) {
177 0 : return fd_grpc_client_rxtx_tls( ctx->grpc_client, ctx->tls_conn, ctx->tcp_sock, now, charge_busy );
178 0 : }
179 :
180 27 : return fd_grpc_client_rxtx_socket( ctx->grpc_client, ctx->tcp_sock, now, charge_busy );
181 27 : }
182 :
183 : static void
184 3 : fd_bundle_client_request_builder_info( fd_bundle_tile_t * ctx ) {
185 3 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( ctx->grpc_client ) ) ) return;
186 :
187 3 : block_engine_BlockBuilderFeeInfoRequest req = block_engine_BlockBuilderFeeInfoRequest_init_default;
188 3 : static char const path[] = "/block_engine.BlockEngineValidator/GetBlockBuilderFeeInfo";
189 3 : fd_grpc_h2_stream_t * request = fd_grpc_client_request_start(
190 3 : ctx->grpc_client,
191 3 : path, sizeof(path)-1,
192 3 : FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo,
193 3 : &block_engine_BlockBuilderFeeInfoRequest_msg, &req,
194 3 : ctx->auther.access_token, ctx->auther.access_token_sz,
195 3 : 0 /* is_streaming */
196 3 : );
197 3 : if( FD_UNLIKELY( !request ) ) return;
198 3 : fd_grpc_client_deadline_set(
199 3 : request,
200 3 : FD_GRPC_DEADLINE_RX_END,
201 3 : fd_log_wallclock() + FD_BUNDLE_CLIENT_REQUEST_TIMEOUT );
202 :
203 3 : ctx->builder_info_wait = 1;
204 3 : }
205 :
206 : static void
207 3 : fd_bundle_client_subscribe_packets( fd_bundle_tile_t * ctx ) {
208 3 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( ctx->grpc_client ) ) ) return;
209 :
210 3 : block_engine_SubscribePacketsRequest req = block_engine_SubscribePacketsRequest_init_default;
211 3 : static char const path[] = "/block_engine.BlockEngineValidator/SubscribePackets";
212 3 : fd_grpc_h2_stream_t * request = fd_grpc_client_request_start(
213 3 : ctx->grpc_client,
214 3 : path, sizeof(path)-1,
215 3 : FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets,
216 3 : &block_engine_SubscribePacketsRequest_msg, &req,
217 3 : ctx->auther.access_token, ctx->auther.access_token_sz,
218 3 : 0 /* is_streaming */
219 3 : );
220 3 : if( FD_UNLIKELY( !request ) ) return;
221 3 : fd_grpc_client_deadline_set(
222 3 : request,
223 3 : FD_GRPC_DEADLINE_HEADER,
224 3 : fd_log_wallclock() + FD_BUNDLE_CLIENT_REQUEST_TIMEOUT );
225 :
226 3 : ctx->packet_subscription_wait = 1;
227 3 : }
228 :
229 : static void
230 3 : fd_bundle_client_subscribe_bundles( fd_bundle_tile_t * ctx ) {
231 3 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( ctx->grpc_client ) ) ) return;
232 :
233 3 : block_engine_SubscribeBundlesRequest req = block_engine_SubscribeBundlesRequest_init_default;
234 3 : static char const path[] = "/block_engine.BlockEngineValidator/SubscribeBundles";
235 3 : fd_grpc_h2_stream_t * request = fd_grpc_client_request_start(
236 3 : ctx->grpc_client,
237 3 : path, sizeof(path)-1,
238 3 : FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles,
239 3 : &block_engine_SubscribeBundlesRequest_msg, &req,
240 3 : ctx->auther.access_token, ctx->auther.access_token_sz,
241 3 : 0 /* is_streaming */
242 3 : );
243 3 : if( FD_UNLIKELY( !request ) ) return;
244 3 : fd_grpc_client_deadline_set(
245 3 : request,
246 3 : FD_GRPC_DEADLINE_HEADER,
247 3 : fd_log_wallclock() + FD_BUNDLE_CLIENT_REQUEST_TIMEOUT );
248 :
249 3 : ctx->bundle_subscription_wait = 1;
250 3 : }
251 :
252 : void
253 3 : fd_bundle_client_send_ping( fd_bundle_tile_t * ctx ) {
254 3 : if( FD_UNLIKELY( !ctx->grpc_client ) ) return; /* no client */
255 3 : fd_h2_conn_t * conn = fd_grpc_client_h2_conn( ctx->grpc_client );
256 3 : if( FD_UNLIKELY( !conn ) ) return; /* no conn */
257 3 : if( FD_UNLIKELY( conn->flags ) ) return; /* conn busy */
258 3 : fd_h2_rbuf_t * rbuf_tx = fd_grpc_client_rbuf_tx( ctx->grpc_client );
259 :
260 3 : if( FD_LIKELY( fd_h2_tx_ping( conn, rbuf_tx ) ) ) {
261 3 : long now = fd_bundle_now( ctx );
262 3 : fd_keepalive_tx( ctx->keepalive, ctx->rng, now );
263 3 : FD_LOG_DEBUG(( "Keepalive TX (deadline=+%gs)", (double)( ctx->keepalive->ts_deadline-now )/1e9 ));
264 3 : }
265 3 : }
266 :
267 : static void
268 : fd_bundle_client_epoll_out( fd_bundle_tile_t * ctx,
269 30 : int want ) {
270 30 : if( FD_UNLIKELY( !ctx->sock_in_epoll ) ) return;
271 0 : if( FD_LIKELY( want==(int)ctx->epoll_out_armed ) ) return;
272 0 : struct epoll_event ev = { .events = EPOLLIN | (want ? EPOLLOUT : 0U), .data.fd = ctx->tcp_sock };
273 0 : if( FD_UNLIKELY( -1==epoll_ctl( FD_WAKER_INNER_FD( ctx->waker_client_idx ), EPOLL_CTL_MOD, ctx->tcp_sock, &ev ) ) )
274 0 : FD_LOG_ERR(( "epoll_ctl(MOD,tcp_sock) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
275 0 : ctx->epoll_out_armed = !!want;
276 0 : }
277 :
278 : long
279 : fd_bundle_client_next_deadline( fd_bundle_tile_t const * ctx,
280 0 : long now ) {
281 : /* Connecting: completion/failure arrives as an EPOLLOUT/EPOLLERR
282 : event, nothing is time-driven yet. */
283 0 : if( FD_UNLIKELY( ctx->tcp_sock>=0 && !ctx->tcp_sock_connected ) ) return LONG_MAX;
284 :
285 : /* Disconnected: next action is the reconnect attempt. */
286 0 : if( FD_UNLIKELY( ctx->tcp_sock<0 ) ) return fd_long_max( ctx->backoff_until, now );
287 :
288 : /* Decrypted TLS bytes buffered in the gRPC client do not make the fd
289 : readable: drain now or they wait for the next unrelated event.
290 : Except while TX is blocked: then they are held because the HTTP/2
291 : rings are full behind the parked send, and only the EPOLLOUT armed
292 : for that send can make progress. */
293 0 : if( FD_UNLIKELY( ctx->is_ssl &&
294 0 : fd_grpc_client_tls_rx_pending( ctx->grpc_client ) &&
295 0 : !fd_grpc_client_tls_tx_pending( ctx->grpc_client ) ) ) return now;
296 :
297 0 : long deadline = fd_grpc_client_next_deadline( ctx->grpc_client );
298 0 : if( FD_LIKELY( ctx->keepalive->interval ) )
299 0 : deadline = fd_long_min( deadline, ctx->keepalive->inflight ? ctx->keepalive->ts_deadline
300 0 : : ctx->keepalive->ts_next_tx );
301 0 : if( FD_LIKELY( ctx->builder_info_avail & !ctx->builder_info_wait ) )
302 0 : deadline = fd_long_min( deadline, ctx->builder_info_valid_until );
303 0 : if( FD_UNLIKELY( ctx->backoff_until>now ) )
304 0 : deadline = fd_long_min( deadline, ctx->backoff_until );
305 0 : return deadline;
306 0 : }
307 :
308 : int
309 : fd_bundle_client_step_reconnect( fd_bundle_tile_t * ctx,
310 18 : long now ) {
311 : /* Drive auth */
312 18 : if( FD_UNLIKELY( ctx->auther.needs_poll ) ) {
313 0 : fd_bundle_auther_poll( &ctx->auther, ctx->grpc_client, ctx->keyguard_client );
314 0 : return 1;
315 0 : }
316 18 : if( FD_UNLIKELY( ctx->auther.state!=FD_BUNDLE_AUTH_STATE_DONE_WAIT ) ) return 0;
317 :
318 : /* Request block builder info */
319 18 : int const builder_info_expired = ( ctx->builder_info_valid_until - now )<0;
320 18 : if( FD_UNLIKELY( ( ( !ctx->builder_info_avail ) |
321 18 : ( builder_info_expired ) ) &
322 18 : ( !ctx->builder_info_wait ) ) ) {
323 3 : fd_bundle_client_request_builder_info( ctx );
324 3 : return 1;
325 3 : }
326 :
327 : /* Subscribe to packets */
328 15 : if( FD_UNLIKELY( !ctx->packet_subscription_live && !ctx->packet_subscription_wait ) ) {
329 3 : fd_bundle_client_subscribe_packets( ctx );
330 3 : return 1;
331 3 : }
332 :
333 : /* Subscribe to bundles */
334 12 : if( FD_UNLIKELY( !ctx->bundle_subscription_live && !ctx->bundle_subscription_wait ) ) {
335 3 : fd_bundle_client_subscribe_bundles( ctx );
336 3 : return 1;
337 3 : }
338 :
339 : /* Send a PING */
340 9 : if( FD_UNLIKELY( fd_keepalive_should_tx( ctx->keepalive, now ) ) ) {
341 3 : fd_bundle_client_send_ping( ctx );
342 3 : return 1;
343 3 : }
344 :
345 6 : return 0;
346 9 : }
347 :
348 : static void
349 : fd_bundle_client_step1( fd_bundle_tile_t * ctx,
350 30 : int * charge_busy ) {
351 :
352 : /* Wait for TCP socket to connect */
353 30 : if( FD_UNLIKELY( !ctx->tcp_sock_connected ) ) {
354 0 : if( FD_UNLIKELY( ctx->tcp_sock < 0 ) ) goto reconnect;
355 :
356 0 : struct pollfd pfds[1] = {
357 0 : { .fd = ctx->tcp_sock, .events = POLLOUT }
358 0 : };
359 0 : int poll_res = fd_syscall_poll( pfds, 1, 0 );
360 0 : if( FD_UNLIKELY( poll_res<0 ) ) {
361 0 : FD_LOG_ERR(( "fd_syscall_poll(tcp_sock) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
362 0 : }
363 0 : if( poll_res==0 ) return;
364 :
365 0 : int connect_result = 0;
366 0 : if( pfds[0].revents & (POLLERR|POLLHUP) ) {
367 0 : connect_result = fd_bundle_client_get_connect_result( ctx );
368 0 : connect_failed:
369 0 : FD_LOG_INFO(( "Bundle gRPC connect attempt failed (%i-%s)", connect_result, fd_io_strerror( connect_result ) ));
370 0 : fd_bundle_client_reset( ctx );
371 0 : ctx->metrics.transport_fail_cnt++;
372 0 : *charge_busy = 1;
373 0 : return;
374 0 : }
375 0 : if( FD_UNLIKELY( !(pfds[0].revents & POLLOUT) ) ) return;
376 0 : connect_result = fd_bundle_client_get_connect_result( ctx );
377 0 : if( FD_UNLIKELY( connect_result!=0 ) ) {
378 0 : goto connect_failed;
379 0 : }
380 0 : FD_LOG_DEBUG(( "Bundle TCP socket connected" ));
381 0 : ctx->tcp_sock_connected = 1;
382 0 : fd_bundle_client_epoll_out( ctx, 0 );
383 0 : *charge_busy = 1;
384 0 : }
385 :
386 : /* gRPC conn died? */
387 30 : if( FD_UNLIKELY( !ctx->grpc_client ) ) {
388 0 : long sleep_start;
389 0 : reconnect:
390 0 : sleep_start = fd_bundle_now( ctx );
391 0 : if( FD_UNLIKELY( fd_bundle_tile_should_stall( ctx, sleep_start ) ) ) {
392 0 : long wait_dur = ctx->backoff_until - sleep_start;
393 0 : fd_log_sleep( fd_long_min( wait_dur, 1e6 ) );
394 0 : return;
395 0 : }
396 0 : fd_bundle_client_create_conn( ctx );
397 0 : *charge_busy = 1;
398 0 : return;
399 0 : }
400 :
401 : /* Did a HTTP/2 PING time out */
402 30 : long check_ts = ctx->cached_ts = fd_bundle_now( ctx );
403 30 : if( FD_UNLIKELY( fd_keepalive_is_timeout( ctx->keepalive, check_ts ) ) ) {
404 3 : FD_LOG_WARNING(( "Bundle gRPC timed out (HTTP/2 PING went unanswered for %.2f seconds)",
405 3 : (double)( check_ts - ctx->keepalive->ts_last_tx )/1e9 ));
406 3 : ctx->keepalive->inflight = 0;
407 3 : ctx->defer_reset = 1;
408 3 : *charge_busy = 1;
409 3 : return;
410 3 : }
411 :
412 : /* Drive I/O, SSL handshake, and any inflight requests */
413 27 : if( FD_UNLIKELY( -1==fd_bundle_client_drive_io( ctx, check_ts, charge_busy ) || ctx->defer_reset /* new error? */ ) ) {
414 0 : fd_bundle_client_reset( ctx );
415 0 : ctx->metrics.transport_fail_cnt++;
416 0 : *charge_busy = 1;
417 0 : return;
418 0 : }
419 :
420 : /* Are we ready to issue a new request? */
421 27 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( ctx->grpc_client ) ) ) return;
422 18 : long io_ts = fd_bundle_now( ctx );
423 18 : if( FD_UNLIKELY( fd_bundle_tile_should_stall( ctx, io_ts ) ) ) return;
424 :
425 18 : *charge_busy |= fd_bundle_client_step_reconnect( ctx, io_ts );
426 18 : }
427 :
428 : static void
429 30 : fd_bundle_client_log_status( fd_bundle_tile_t * ctx ) {
430 30 : int status = fd_bundle_client_status( ctx );
431 :
432 30 : int const connected_now = ( status==FD_BUNDLE_STATE_CONNECTED );
433 30 : int const connected_before = ( ctx->bundle_status_logged==FD_BUNDLE_STATE_CONNECTED );
434 :
435 30 : if( FD_UNLIKELY( connected_now!=connected_before ) ) {
436 3 : long ts = fd_log_wallclock();
437 3 : if( FD_LIKELY( ts-(ctx->last_bundle_status_log_nanos) >= (long)1e6 ) ) {
438 3 : if( connected_now ) FD_LOG_INFO(( "Connected to bundle server" ));
439 0 : else FD_LOG_INFO(( "Disconnected from bundle server" ));
440 3 : ctx->last_bundle_status_log_nanos = ts;
441 3 : ctx->bundle_status_logged = (uchar)status;
442 3 : }
443 3 : }
444 30 : }
445 :
446 : void
447 : fd_bundle_client_step( fd_bundle_tile_t * ctx,
448 30 : int * charge_busy ) {
449 : /* Edge-trigger logging with rate limiting */
450 30 : fd_bundle_client_step1( ctx, charge_busy );
451 30 : fd_bundle_client_log_status( ctx );
452 :
453 30 : if( FD_UNLIKELY( ctx->tcp_sock_connected && fd_grpc_client_tx_pending( ctx->grpc_client ) ) ) {
454 12 : int flush_err = ctx->is_ssl ? fd_grpc_client_tls_flush ( ctx->grpc_client, ctx->tcp_sock )
455 12 : : fd_grpc_client_tx_flush_socket( ctx->grpc_client, ctx->tcp_sock );
456 12 : if( FD_UNLIKELY( -1==flush_err ) ) {
457 0 : fd_bundle_client_reset( ctx );
458 0 : ctx->metrics.transport_fail_cnt++;
459 0 : *charge_busy = 1;
460 0 : return;
461 0 : }
462 12 : }
463 :
464 : /* Arm EPOLLOUT only on write demand: unsent HTTP/2 bytes, or a TLS
465 : handshake blocked on write (a permanently armed idle-writable
466 : socket would keep the epoll set ready forever). */
467 30 : if( FD_LIKELY( ctx->tcp_sock_connected && ctx->grpc_client ) ) {
468 30 : int want = fd_grpc_client_tx_pending( ctx->grpc_client );
469 30 : want |= ctx->is_ssl && fd_grpc_client_tls_tx_pending( ctx->grpc_client );
470 30 : fd_bundle_client_epoll_out( ctx, want );
471 30 : }
472 30 : }
473 :
474 : void
475 : fd_bundle_tile_backoff( fd_bundle_tile_t * ctx,
476 27 : long now ) {
477 27 : uint iter = ctx->backoff_iter;
478 27 : if( now < ctx->backoff_reset ) iter = 0U;
479 27 : iter++;
480 :
481 : /* FIXME proper backoff */
482 27 : long wait_ns = (long)2e9;
483 27 : wait_ns = (long)( fd_rng_ulong( ctx->rng ) & ( (1UL<<fd_ulong_find_msb_w_default( (ulong)wait_ns, 0 ))-1UL ) );
484 :
485 27 : ctx->backoff_until = now + wait_ns;
486 27 : ctx->backoff_reset = now + 2*wait_ns;
487 :
488 27 : ctx->backoff_iter = iter;
489 27 : }
490 :
491 : static void
492 0 : fd_bundle_client_grpc_conn_established( void * app_ctx ) {
493 0 : (void)app_ctx;
494 0 : FD_LOG_INFO(( "Bundle gRPC connection established" ));
495 0 : }
496 :
497 : static void
498 : fd_bundle_client_grpc_conn_dead( void * app_ctx,
499 : uint h2_err,
500 0 : int closed_by ) {
501 0 : fd_bundle_tile_t * ctx = app_ctx;
502 0 : FD_LOG_INFO(( "Bundle gRPC connection closed %s (%u-%s)",
503 0 : closed_by ? "by peer" : "due to error",
504 0 : h2_err, fd_h2_strerror( h2_err ) ));
505 0 : ctx->defer_reset = 1;
506 0 : }
507 :
508 : /* Buffers a bundle transaction for deferred publishing by after_credit. */
509 :
510 : static void
511 : fd_bundle_tile_publish_bundle_txn(
512 : fd_bundle_tile_t * ctx,
513 : void const * txn,
514 : ulong txn_sz, /* <=FD_TXN_MTU */
515 : ulong bundle_txn_cnt,
516 : uint source_ipv4
517 72 : ) {
518 72 : if( FD_UNLIKELY( !ctx->builder_info_avail ) ) {
519 0 : ctx->metrics.missing_builder_info_fail_cnt++; /* unreachable */
520 0 : return;
521 0 : }
522 :
523 72 : if( FD_UNLIKELY( pending_txn_full( ctx->pending_txns ) ) ) {
524 0 : ctx->metrics.backpressure_drop_cnt++;
525 0 : return;
526 0 : }
527 :
528 72 : fd_bundle_pending_txn_t * entry = pending_txn_push_tail_nocopy( ctx->pending_txns );
529 72 : fd_memcpy( entry->payload, txn, txn_sz );
530 72 : entry->payload_sz = (ushort)txn_sz;
531 72 : entry->source_ipv4 = source_ipv4;
532 72 : entry->first_seen_nanos = fd_bundle_now( ctx );
533 72 : entry->sig = 1UL;
534 72 : entry->bundle_seq = ctx->bundle_seq;
535 72 : entry->bundle_txn_cnt = bundle_txn_cnt;
536 72 : entry->commission = (uchar)ctx->builder_commission;
537 72 : fd_memcpy( entry->commission_pubkey, ctx->builder_pubkey, 32UL );
538 72 : ctx->metrics.txn_received_cnt++;
539 72 : }
540 :
541 : /* Buffers a regular transaction for deferred publishing by after_credit. */
542 :
543 : static void
544 : fd_bundle_tile_publish_txn(
545 : fd_bundle_tile_t * ctx,
546 : void const * txn,
547 : ulong txn_sz, /* <=FD_TXN_MTU */
548 : uint source_ipv4
549 42 : ) {
550 42 : if( FD_UNLIKELY( pending_txn_full( ctx->pending_txns ) ) ) {
551 3 : ctx->metrics.backpressure_drop_cnt++;
552 3 : return;
553 3 : }
554 :
555 39 : fd_bundle_pending_txn_t * entry = pending_txn_push_tail_nocopy( ctx->pending_txns );
556 39 : fd_memcpy( entry->payload, txn, txn_sz );
557 39 : entry->payload_sz = (ushort)txn_sz;
558 39 : entry->source_ipv4 = source_ipv4;
559 39 : entry->first_seen_nanos = fd_bundle_now( ctx );
560 39 : entry->sig = 0UL;
561 39 : entry->bundle_seq = 0UL;
562 39 : entry->bundle_txn_cnt = 1UL;
563 39 : entry->commission = 0U;
564 39 : fd_memset( entry->commission_pubkey, 0, 32UL );
565 39 : ctx->metrics.txn_received_cnt++;
566 39 : }
567 :
568 : /* Called for each transaction in a bundle. Simply counts up
569 : bundle_txn_cnt, but does not publish anything. */
570 :
571 : static bool
572 : fd_bundle_client_visit_pb_bundle_txn_preflight(
573 : pb_istream_t * istream,
574 : pb_field_t const * field,
575 : void ** arg
576 90 : ) {
577 90 : (void)istream; (void)field;
578 90 : fd_bundle_tile_t * ctx = *arg;
579 90 : ctx->bundle_txn_cnt++;
580 90 : return true;
581 90 : }
582 :
583 : /* Called for each transaction in a bundle. Publishes each transaction
584 : to the tango message bus. */
585 :
586 : static bool
587 : fd_bundle_client_visit_pb_bundle_txn(
588 : pb_istream_t * istream,
589 : pb_field_t const * field,
590 : void ** arg
591 72 : ) {
592 72 : (void)field;
593 72 : fd_bundle_tile_t * ctx = *arg;
594 :
595 72 : packet_Packet packet = packet_Packet_init_default;
596 72 : if( FD_UNLIKELY( !pb_decode( istream, &packet_Packet_msg, &packet ) ) ) {
597 0 : ctx->metrics.decode_fail_cnt++;
598 0 : FD_LOG_WARNING(( "Protobuf decode of (packet.Packet) failed" ));
599 0 : return false;
600 0 : }
601 :
602 72 : if( FD_UNLIKELY( packet.data.size == 0 ) ) {
603 0 : FD_LOG_INFO(( "Bundle server delivered an empty packet, ignoring" ));
604 0 : return true;
605 0 : }
606 :
607 72 : if( FD_UNLIKELY( packet.data.size > FD_TXN_MTU ) ) {
608 0 : FD_LOG_WARNING(( "Bundle server delivered an oversize transaction, ignoring" ));
609 0 : return true;
610 0 : }
611 :
612 72 : uint _ip4; uint ip4 = fd_uint_if( packet.has_meta, fd_cstr_to_ip4_addr( packet.meta.addr, &_ip4 ) ? _ip4 : ctx->server_ip4_addr, ctx->server_ip4_addr );
613 72 : fd_bundle_tile_publish_bundle_txn(
614 72 : ctx,
615 72 : packet.data.bytes, packet.data.size,
616 72 : ctx->bundle_txn_cnt,
617 72 : ip4
618 72 : );
619 :
620 72 : return true;
621 72 : }
622 :
623 : static void
624 : fd_bundle_client_sample_rx_delay(
625 : fd_bundle_tile_t * ctx,
626 : google_protobuf_Timestamp const * ts
627 45 : ) {
628 45 : ulong tsorig = (ulong)ts->seconds*(ulong)1e9 + (ulong)ts->nanos;
629 45 : fd_histf_sample( ctx->metrics.msg_rx_delay, fd_ulong_sat_sub( (ulong)ctx->cached_ts, tsorig ) );
630 45 : }
631 :
632 : /* Called for each BundleUuid in a SubscribeBundlesResponse. */
633 :
634 : static bool
635 : fd_bundle_client_visit_pb_bundle_uuid(
636 : pb_istream_t * istream,
637 : pb_field_t const * field,
638 : void ** arg
639 24 : ) {
640 24 : (void)field;
641 24 : fd_bundle_tile_t * ctx = *arg;
642 :
643 : /* Reset bundle state */
644 :
645 24 : ctx->bundle_txn_cnt = 0UL;
646 :
647 : /* Do two decode passes. This is required because we need to know the
648 : number of transactions in a bundle ahead of time. However, due to
649 : the Protobuf wire encoding, we don't know the number of txns that
650 : will come until we've parsed everything.
651 :
652 : First pass: Count number of bundles. */
653 :
654 24 : pb_istream_t peek = *istream;
655 24 : bundle_BundleUuid bundle = bundle_BundleUuid_init_default;
656 24 : bundle.bundle.packets = (pb_callback_t) {
657 24 : .funcs.decode = fd_bundle_client_visit_pb_bundle_txn_preflight,
658 24 : .arg = ctx
659 24 : };
660 24 : if( FD_UNLIKELY( !pb_decode( &peek, &bundle_BundleUuid_msg, &bundle ) ) ) {
661 0 : ctx->metrics.decode_fail_cnt++;
662 0 : FD_LOG_WARNING(( "Protobuf decode of (bundle.BundleUuid) failed: %s", peek.errmsg ));
663 0 : return false;
664 0 : }
665 :
666 : /* At this opint, ctx->bundle_txn_cnt is correctly set. Too many txns
667 : is treated as a NOP.
668 :
669 : Second pass: Actually publish bundle packets */
670 :
671 24 : if( FD_UNLIKELY( ctx->bundle_txn_cnt>FD_BUNDLE_CLIENT_MAX_TXN_PER_BUNDLE ) ) return true;
672 :
673 21 : if( FD_UNLIKELY( pending_txn_avail( ctx->pending_txns )<ctx->bundle_txn_cnt ) ) {
674 0 : ctx->metrics.backpressure_drop_cnt += ctx->bundle_txn_cnt;
675 0 : return true;
676 0 : }
677 :
678 21 : ctx->bundle_seq++;
679 21 : bundle = (bundle_BundleUuid)bundle_BundleUuid_init_default;
680 21 : bundle.bundle.packets = (pb_callback_t) {
681 21 : .funcs.decode = fd_bundle_client_visit_pb_bundle_txn,
682 21 : .arg = ctx
683 21 : };
684 :
685 21 : ctx->metrics.bundle_received_cnt++;
686 :
687 21 : if( FD_UNLIKELY( !pb_decode( istream, &bundle_BundleUuid_msg, &bundle ) ) ) {
688 0 : ctx->metrics.decode_fail_cnt++;
689 0 : FD_LOG_WARNING(( "Protobuf decode of (bundle.BundleUuid) failed (internal error): %s", istream->errmsg ));
690 0 : return false;
691 0 : }
692 :
693 21 : fd_bundle_client_sample_rx_delay( ctx, &bundle.bundle.header.ts );
694 :
695 21 : return true;
696 21 : }
697 :
698 : /* Handle a SubscribeBundlesResponse from a SubscribeBundles gRPC call. */
699 :
700 : static void
701 : fd_bundle_client_handle_bundle_batch(
702 : fd_bundle_tile_t * ctx,
703 : pb_istream_t * istream
704 24 : ) {
705 24 : if( FD_UNLIKELY( !ctx->builder_info_avail ) ) {
706 3 : ctx->metrics.missing_builder_info_fail_cnt++; /* unreachable */
707 3 : return;
708 3 : }
709 :
710 21 : block_engine_SubscribeBundlesResponse res = block_engine_SubscribeBundlesResponse_init_default;
711 21 : res.bundles = (pb_callback_t) {
712 21 : .funcs.decode = fd_bundle_client_visit_pb_bundle_uuid,
713 21 : .arg = ctx
714 21 : };
715 21 : if( FD_UNLIKELY( !pb_decode( istream, &block_engine_SubscribeBundlesResponse_msg, &res ) ) ) {
716 0 : ctx->metrics.decode_fail_cnt++;
717 0 : FD_LOG_WARNING(( "Protobuf decode of (block_engine.SubscribeBundlesResponse) failed: %s", istream->errmsg ));
718 0 : return;
719 0 : }
720 21 : }
721 :
722 : /* Called for each 'Packet' (a regular transaction) of a
723 : SubscribePacketsResponse. */
724 :
725 : static bool
726 : fd_bundle_client_visit_pb_packet(
727 : pb_istream_t * istream,
728 : pb_field_t const * field,
729 : void ** arg
730 42 : ) {
731 42 : (void)field;
732 42 : fd_bundle_tile_t * ctx = *arg;
733 :
734 42 : packet_Packet packet = packet_Packet_init_default;
735 42 : if( FD_UNLIKELY( !pb_decode( istream, &packet_Packet_msg, &packet ) ) ) {
736 0 : ctx->metrics.decode_fail_cnt++;
737 0 : FD_LOG_WARNING(( "Protobuf decode of (packet.Packet) failed" ));
738 0 : return false;
739 0 : }
740 :
741 42 : if( FD_UNLIKELY( packet.data.size == 0 ) ) {
742 0 : FD_LOG_WARNING(( "Bundle server delivered an empty packet, ignoring" ));
743 0 : return true;
744 0 : }
745 :
746 42 : if( FD_UNLIKELY( packet.data.size > FD_TXN_MTU ) ) {
747 0 : FD_LOG_WARNING(( "Bundle server delivered an oversize transaction, ignoring" ));
748 0 : return true;
749 0 : }
750 :
751 :
752 42 : uint _ip4; uint ip4 = fd_uint_if( packet.has_meta, fd_cstr_to_ip4_addr( packet.meta.addr, &_ip4 ) ? _ip4 : 0U, 0U );
753 42 : fd_bundle_tile_publish_txn( ctx, packet.data.bytes, packet.data.size, ip4 );
754 42 : ctx->metrics.packet_received_cnt++;
755 :
756 42 : return true;
757 42 : }
758 :
759 : /* Handle a SubscribePacketsResponse from a SubscribePackets gRPC call. */
760 :
761 : static void
762 : fd_bundle_client_handle_packet_batch(
763 : fd_bundle_tile_t * ctx,
764 : pb_istream_t * istream
765 24 : ) {
766 24 : block_engine_SubscribePacketsResponse res = block_engine_SubscribePacketsResponse_init_default;
767 24 : res.batch.packets = (pb_callback_t) {
768 24 : .funcs.decode = fd_bundle_client_visit_pb_packet,
769 24 : .arg = ctx
770 24 : };
771 24 : if( FD_UNLIKELY( !pb_decode( istream, &block_engine_SubscribePacketsResponse_msg, &res ) ) ) {
772 0 : ctx->metrics.decode_fail_cnt++;
773 0 : FD_LOG_WARNING(( "Protobuf decode of (block_engine.SubscribePacketsResponse) failed" ));
774 0 : return;
775 0 : }
776 :
777 24 : fd_bundle_client_sample_rx_delay( ctx, &res.header.ts );
778 24 : }
779 :
780 : /* Handle a BlockBuilderFeeInfoResponse from a GetBlockBuilderFeeInfo
781 : gRPC call. */
782 :
783 : static void
784 : fd_bundle_client_handle_builder_fee_info(
785 : fd_bundle_tile_t * ctx,
786 : pb_istream_t * istream
787 9 : ) {
788 9 : block_engine_BlockBuilderFeeInfoResponse res = block_engine_BlockBuilderFeeInfoResponse_init_default;
789 9 : if( FD_UNLIKELY( !pb_decode( istream, &block_engine_BlockBuilderFeeInfoResponse_msg, &res ) ) ) {
790 0 : ctx->metrics.decode_fail_cnt++;
791 0 : FD_LOG_WARNING(( "Protobuf decode of (block_engine.BlockBuilderFeeInfoResponse) failed" ));
792 0 : return;
793 0 : }
794 9 : if( FD_UNLIKELY( res.commission > 100 ) ) {
795 3 : ctx->metrics.decode_fail_cnt++;
796 3 : FD_LOG_WARNING(( "BlockBuilderFeeInfoResponse commission out of range (0-100): %lu", res.commission ));
797 3 : return;
798 3 : }
799 :
800 6 : uchar decoded_builder_pubkey[ 32 ];
801 6 : if( FD_UNLIKELY( !fd_base58_decode_32( res.pubkey, decoded_builder_pubkey ) ) ) {
802 3 : FD_LOG_HEXDUMP_WARNING(( "Invalid pubkey in BlockBuilderFeeInfoResponse", res.pubkey, strnlen( res.pubkey, sizeof(res.pubkey) ) ));
803 3 : return;
804 3 : }
805 :
806 3 : ctx->builder_commission = (uchar)res.commission; /* Apply update atomically */
807 3 : fd_memcpy( ctx->builder_pubkey, decoded_builder_pubkey, sizeof(ctx->builder_pubkey) );
808 :
809 3 : long validity_duration_ns = (long)( 60e9 * 5. ); /* 5 minutes */
810 3 : ctx->builder_info_avail = 1;
811 3 : ctx->builder_info_valid_until = fd_bundle_now( ctx ) + validity_duration_ns;
812 3 : }
813 :
814 : static void
815 : fd_bundle_client_grpc_tx_complete(
816 : void * app_ctx,
817 : ulong request_ctx
818 9 : ) {
819 9 : (void)app_ctx; (void)request_ctx;
820 9 : }
821 :
822 : void
823 : fd_bundle_client_grpc_rx_start(
824 : void * app_ctx,
825 : ulong request_ctx
826 9 : ) {
827 9 : fd_bundle_tile_t * ctx = app_ctx;
828 9 : switch( request_ctx ) {
829 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets:
830 3 : ctx->packet_subscription_live = 1;
831 3 : ctx->packet_subscription_wait = 0;
832 3 : break;
833 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles:
834 3 : ctx->bundle_subscription_live = 1;
835 3 : ctx->bundle_subscription_wait = 0;
836 3 : break;
837 9 : }
838 9 : }
839 :
840 : void
841 : fd_bundle_client_grpc_rx_msg(
842 : void * app_ctx,
843 : void const * protobuf,
844 : ulong protobuf_sz,
845 : ulong request_ctx
846 57 : ) {
847 57 : fd_bundle_tile_t * ctx = app_ctx;
848 57 : ctx->metrics.proto_received_bytes += protobuf_sz;
849 57 : pb_istream_t istream = pb_istream_from_buffer( protobuf, protobuf_sz );
850 57 : switch( request_ctx ) {
851 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthChallenge:
852 0 : if( FD_UNLIKELY( !fd_bundle_auther_handle_challenge_resp( &ctx->auther, protobuf, protobuf_sz ) ) ) {
853 0 : ctx->metrics.decode_fail_cnt++;
854 0 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
855 0 : }
856 0 : break;
857 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthTokens:
858 0 : if( FD_UNLIKELY( !fd_bundle_auther_handle_tokens_resp( &ctx->auther, protobuf, protobuf_sz ) ) ) {
859 0 : ctx->metrics.decode_fail_cnt++;
860 0 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
861 0 : }
862 0 : break;
863 24 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles:
864 24 : fd_bundle_client_handle_bundle_batch( ctx, &istream );
865 24 : break;
866 24 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets:
867 24 : fd_bundle_client_handle_packet_batch( ctx, &istream );
868 24 : break;
869 9 : case FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo:
870 9 : fd_bundle_client_handle_builder_fee_info( ctx, &istream );
871 9 : break;
872 0 : default:
873 0 : FD_LOG_ERR(( "Received unexpected gRPC message (request_ctx=%lu)", request_ctx ));
874 57 : }
875 57 : }
876 :
877 : static void
878 : fd_bundle_client_request_failed( fd_bundle_tile_t * ctx,
879 9 : ulong request_ctx ) {
880 9 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
881 9 : switch( request_ctx ) {
882 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthChallenge:
883 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthTokens:
884 0 : fd_bundle_auther_handle_request_fail( &ctx->auther );
885 0 : break;
886 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo:
887 3 : ctx->builder_info_wait = 0;
888 3 : break;
889 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets:
890 3 : ctx->packet_subscription_live = 0;
891 3 : ctx->packet_subscription_wait = 0;
892 3 : break;
893 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles:
894 3 : ctx->bundle_subscription_live = 0;
895 3 : ctx->bundle_subscription_wait = 0;
896 3 : break;
897 9 : }
898 9 : }
899 :
900 : void
901 : fd_bundle_client_grpc_rx_end(
902 : void * app_ctx,
903 : ulong request_ctx,
904 : fd_grpc_resp_hdrs_t * resp
905 24 : ) {
906 24 : fd_bundle_tile_t * ctx = app_ctx;
907 24 : if( FD_UNLIKELY( resp->h2_status!=200 ) ) {
908 9 : FD_LOG_WARNING(( "gRPC request failed (HTTP status %u)", resp->h2_status ));
909 9 : fd_bundle_client_request_failed( ctx, request_ctx );
910 9 : fd_bundle_auther_reset( &ctx->auther );
911 9 : ctx->defer_reset = 1;
912 9 : return;
913 9 : }
914 :
915 15 : resp->grpc_msg_len = (uint)fd_url_unescape( resp->grpc_msg, resp->grpc_msg_len );
916 15 : if( !resp->grpc_msg_len ) {
917 15 : fd_memcpy( resp->grpc_msg, "unknown error", 13 );
918 15 : resp->grpc_msg_len = 13;
919 15 : }
920 :
921 15 : switch( request_ctx ) {
922 0 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets:
923 0 : ctx->packet_subscription_live = 0;
924 0 : ctx->packet_subscription_wait = 0;
925 0 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
926 0 : ctx->defer_reset = 1;
927 0 : FD_LOG_INFO(( "SubscribePackets stream failed (gRPC status %u-%s). Reconnecting ...",
928 0 : resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ) ));
929 0 : return;
930 12 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles:
931 12 : ctx->bundle_subscription_live = 0;
932 12 : ctx->bundle_subscription_wait = 0;
933 12 : fd_bundle_tile_backoff( ctx, fd_bundle_now( ctx ) );
934 12 : ctx->defer_reset = 1;
935 12 : FD_LOG_INFO(( "SubscribeBundles stream failed (gRPC status %u-%s). Reconnecting ...",
936 12 : resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ) ));
937 12 : return;
938 3 : case FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo:
939 3 : ctx->builder_info_wait = 0;
940 3 : break;
941 0 : default:
942 0 : break;
943 15 : }
944 :
945 3 : if( FD_UNLIKELY( resp->grpc_status!=FD_GRPC_STATUS_OK ) ) {
946 0 : FD_LOG_INFO(( "gRPC request failed (gRPC status %u-%s): %.*s",
947 0 : resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ),
948 0 : (int)resp->grpc_msg_len, resp->grpc_msg ));
949 0 : fd_bundle_client_request_failed( ctx, request_ctx );
950 0 : if( resp->grpc_status==FD_GRPC_STATUS_UNAUTHENTICATED ||
951 0 : resp->grpc_status==FD_GRPC_STATUS_PERMISSION_DENIED ) {
952 0 : fd_bundle_auther_reset( &ctx->auther );
953 0 : }
954 0 : return;
955 0 : }
956 3 : }
957 :
958 : void
959 : fd_bundle_client_grpc_rx_timeout(
960 : void * app_ctx,
961 : ulong request_ctx, /* FD_BUNDLE_CLIENT_REQ_{...} */
962 : int deadline_kind /* FD_GRPC_DEADLINE_{HEADER|RX_END} */
963 6 : ) {
964 6 : (void)deadline_kind;
965 6 : FD_LOG_WARNING(( "Request timed out: %s", fd_bundle_request_ctx_cstr( request_ctx ) ));
966 6 : fd_bundle_tile_t * ctx = app_ctx;
967 6 : ctx->defer_reset = 1;
968 6 : }
969 :
970 : static void
971 3 : fd_bundle_client_grpc_ping_ack( void * app_ctx ) {
972 3 : fd_bundle_tile_t * ctx = app_ctx;
973 3 : long rtt_sample = fd_keepalive_rx( ctx->keepalive, fd_bundle_now( ctx ) );
974 3 : if( FD_LIKELY( rtt_sample ) ) {
975 3 : fd_rtt_sample( ctx->rtt, (float)rtt_sample, 0 );
976 3 : FD_LOG_DEBUG(( "Keepalive ACK" ));
977 3 : }
978 3 : ctx->metrics.ping_ack_cnt++;
979 3 : }
980 :
981 : fd_grpc_client_callbacks_t fd_bundle_client_grpc_callbacks = {
982 : .conn_established = fd_bundle_client_grpc_conn_established,
983 : .conn_dead = fd_bundle_client_grpc_conn_dead,
984 : .tx_complete = fd_bundle_client_grpc_tx_complete,
985 : .rx_start = fd_bundle_client_grpc_rx_start,
986 : .rx_msg = fd_bundle_client_grpc_rx_msg,
987 : .rx_end = fd_bundle_client_grpc_rx_end,
988 : .rx_timeout = fd_bundle_client_grpc_rx_timeout,
989 : .ping_ack = fd_bundle_client_grpc_ping_ack,
990 : };
991 :
992 : int
993 180 : fd_bundle_client_status( fd_bundle_tile_t const * ctx ) {
994 180 : if( FD_UNLIKELY( ( !ctx->tcp_sock_connected ) |
995 180 : ( !ctx->grpc_client ) ) ) {
996 6 : return FD_BUNDLE_STATE_DISCONNECTED;
997 6 : }
998 :
999 174 : fd_h2_conn_t * conn = fd_grpc_client_h2_conn( ctx->grpc_client );
1000 174 : if( FD_UNLIKELY( !conn ) ) {
1001 0 : return FD_BUNDLE_STATE_DISCONNECTED; /* no conn */
1002 0 : }
1003 174 : if( FD_UNLIKELY( conn->flags &
1004 174 : ( FD_H2_CONN_FLAGS_DEAD |
1005 174 : FD_H2_CONN_FLAGS_SEND_GOAWAY ) ) ) {
1006 6 : return FD_BUNDLE_STATE_DISCONNECTED;
1007 6 : }
1008 :
1009 168 : if( FD_UNLIKELY( conn->flags &
1010 168 : ( FD_H2_CONN_FLAGS_CLIENT_INITIAL |
1011 168 : FD_H2_CONN_FLAGS_WAIT_SETTINGS_ACK_0 |
1012 168 : FD_H2_CONN_FLAGS_WAIT_SETTINGS_0 |
1013 168 : FD_H2_CONN_FLAGS_SERVER_INITIAL ) ) ) {
1014 12 : return FD_BUNDLE_STATE_CONNECTING; /* connection is not ready */
1015 12 : }
1016 :
1017 156 : if( FD_UNLIKELY( ctx->auther.state != FD_BUNDLE_AUTH_STATE_DONE_WAIT ) ) {
1018 12 : return FD_BUNDLE_STATE_CONNECTING; /* not authenticated */
1019 12 : }
1020 :
1021 144 : if( FD_UNLIKELY( ( !ctx->builder_info_avail ) |
1022 144 : ( !ctx->packet_subscription_live ) |
1023 144 : ( !ctx->bundle_subscription_live ) ) ) {
1024 27 : return FD_BUNDLE_STATE_CONNECTING; /* not fully connected */
1025 27 : }
1026 :
1027 117 : if( FD_UNLIKELY( fd_keepalive_is_timeout( ctx->keepalive, fd_bundle_now( ctx ) ) ) ) {
1028 3 : return FD_BUNDLE_STATE_DISCONNECTED; /* possible timeout */
1029 3 : }
1030 :
1031 114 : if( FD_UNLIKELY( !fd_grpc_client_is_connected( ctx->grpc_client ) ) ) {
1032 3 : return FD_BUNDLE_STATE_CONNECTING;
1033 3 : }
1034 :
1035 : /* As far as we know, the bundle connection is alive and well. */
1036 111 : return FD_BUNDLE_STATE_CONNECTED;
1037 114 : }
1038 :
1039 : #undef DISCONNECTED
1040 : #undef CONNECTING
1041 : #undef CONNECTED
1042 :
1043 : FD_FN_CONST char const *
1044 6 : fd_bundle_request_ctx_cstr( ulong request_ctx ) {
1045 6 : switch( request_ctx ) {
1046 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthChallenge:
1047 0 : return "GenerateAuthChallenge";
1048 0 : case FD_BUNDLE_CLIENT_REQ_Auth_GenerateAuthTokens:
1049 0 : return "GenerateAuthTokens";
1050 0 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribePackets:
1051 0 : return "SubscribePackets";
1052 6 : case FD_BUNDLE_CLIENT_REQ_Bundle_SubscribeBundles:
1053 6 : return "SubscribeBundles";
1054 0 : case FD_BUNDLE_CLIENT_REQ_Bundle_GetBlockBuilderFeeInfo:
1055 0 : return "GetBlockBuilderFeeInfo";
1056 0 : default:
1057 0 : return "unknown";
1058 6 : }
1059 6 : }
|