Line data Source code
1 : #define _GNU_SOURCE
2 : #include "fd_event_client.h"
3 :
4 : #include "../../waltz/resolv/fd_netdb.h"
5 : #include "../../waltz/http/fd_url.h"
6 : #include "../../waltz/grpc/fd_grpc_client.h"
7 : #include "../../waltz/grpc/fd_grpc_client_private.h"
8 : #include "../../ballet/pb/fd_pb_tokenize.h"
9 : #include "../../ballet/pb/fd_pb_encode.h"
10 : #include "../../ballet/hex/fd_hex.h"
11 : #include "../../util/net/fd_ip4.h"
12 : #include "../../util/log/fd_log.h"
13 : #include "../keyguard/fd_keyguard.h"
14 :
15 : #include "../../waltz/tlsrec/fd_tlsrec.h"
16 : #include "../../ballet/ed25519/fd_x25519.h"
17 :
18 : #include <netinet/tcp.h>
19 : #include <unistd.h>
20 : #include <errno.h>
21 : #include <sys/epoll.h>
22 : #include <sys/socket.h>
23 : #include <netinet/in.h>
24 :
25 0 : #define DISCONNECT_REASON_IDENTITY_CHANGED (0)
26 0 : #define DISCONNECT_REASON_CONNECT_FAILED (1)
27 0 : #define DISCONNECT_REASON_DNS_RESOLVE_FAILED (2)
28 0 : #define DISCONNECT_REASON_TIMEOUT (3)
29 0 : #define DISCONNECT_REASON_TRANSPORT_FAILED (4)
30 0 : #define DISCONNECT_REASON_PEER_CLOSED (5)
31 0 : #define DISCONNECT_REASON_INVALID_CURSOR (6)
32 0 : #define DISCONNECT_REASON_AUTH_FAILED (7)
33 0 : #define DISCONNECT_REASON_INVALID_PROTOBUF (8)
34 :
35 0 : #define FD_EVENT_CLIENT_REQ_CTX_AUTHENTICATE (1UL)
36 0 : #define FD_EVENT_CLIENT_REQ_CTX_STREAM_EVENTS (3UL)
37 :
38 0 : #define FD_EVENT_CLIENT_HEARTBEAT_NANOS (15L*(long)1e9)
39 0 : #define FD_EVENT_CLIENT_RESPONSE_TIMEOUT_NANOS (60L*(long)1e9)
40 :
41 0 : #define FD_EVENT_CLIENT_TX_RATE_BPS (2L*1000L*1000L)
42 0 : #define FD_EVENT_CLIENT_TX_BURST (256L<<10)
43 0 : #define FD_EVENT_CLIENT_CREDIT_STALL_NANOS (20L*(long)1e9)
44 :
45 0 : #define FD_EVENT_CLIENT_TOKEN_SZ (217UL)
46 :
47 : struct fd_event_client {
48 : fd_grpc_client_t * grpc_client;
49 : fd_grpc_client_metrics_t grpc_metrics[1];
50 : fd_grpc_h2_stream_t * event_stream;
51 :
52 : char client_version[ 10UL ];
53 : char commit_hash[ 41UL ];
54 : char action[ 16UL ];
55 : uchar identity_pubkey[ 32UL ];
56 :
57 : int has_genesis_hash;
58 : fd_hash_t genesis_hash[1];
59 :
60 : int connect_fail_logged;
61 :
62 : ushort has_shred_version;
63 : ushort shred_version;
64 :
65 : ulong event_id;
66 :
67 : ulong instance_id;
68 : ulong boot_id;
69 : ulong machine_id;
70 :
71 : int defer_disconnect;
72 : ulong consecutive_failure_count;
73 :
74 : long now; /* of the current poll */
75 :
76 : long last_stream_send_ns;
77 : long last_response_ns;
78 :
79 : long tx_tokens;
80 : long tx_tokens_ns;
81 : long stall_since;
82 : ulong stall_rem;
83 : ulong stall_events_sent;
84 :
85 : int auth_send_pending;
86 :
87 : ulong state;
88 : union {
89 : struct {
90 : long reconnect_deadline;
91 : } disconnected;
92 :
93 : struct {
94 : long connect_deadline;
95 : } connecting;
96 :
97 : struct {
98 : long connected_timestamp;
99 : } connected;
100 : };
101 :
102 : int so_sndbuf;
103 : int sockfd;
104 :
105 : int epoll_fd;
106 : int epoll_out_armed;
107 :
108 : int use_tls;
109 : fd_tls_t tls[1];
110 : fd_chacha_rng_t tls_rng[1];
111 : fd_tlsrec_conn_t tls_conn[1];
112 :
113 : /* wallclock deadline for auth handshake, LONG_MAX if not
114 : authenticating. */
115 : long auth_deadline;
116 :
117 : char server_fqdn[ FD_FQDN_BUF_MAX ]; /* cstr */
118 : ulong server_fqdn_len;
119 : uint server_ip4_addr;
120 : ushort server_tcp_port;
121 :
122 : fd_rng_t * rng;
123 : fd_circq_t * circq;
124 : fd_keyguard_client_t * keyguard_client;
125 :
126 : /* Stateless-auth bearer value: "hex(challenge_token).hex(signature)",
127 : built from the challenge token returned by Authenticate and our
128 : ed25519 signature over it. Presented as the `authorization: Bearer
129 : <...>` header on the StreamEvents request. */
130 : char auth_bearer[ 2UL*FD_EVENT_CLIENT_TOKEN_SZ + 1UL + 2UL*64UL + 1UL ];
131 : ulong auth_bearer_len;
132 :
133 : fd_event_client_metrics_t metrics;
134 : };
135 :
136 :
137 : static void
138 : tls_init( fd_event_client_t * client,
139 0 : fd_x509_ca_store_t const * ca_store ) {
140 0 : fd_tls_t * tls = fd_tls_join( fd_tls_new( client->tls ) );
141 0 : FD_TEST( tls );
142 :
143 0 : uchar rng_key[ FD_CHACHA_KEY_SZ ];
144 0 : if( FD_UNLIKELY( !fd_rng_secure( rng_key, sizeof(rng_key) ) ) ) FD_LOG_CRIT(( "fd_rng_secure failed" ));
145 0 : fd_chacha_rng_init( client->tls_rng, rng_key, FD_CHACHA_RNG_ALGO_CHACHA8 );
146 0 : fd_memzero_explicit( rng_key, sizeof(rng_key) );
147 0 : tls->rng = client->tls_rng;
148 :
149 0 : static uchar const alpn[] = { 2, 'h', '2' };
150 0 : fd_memcpy( tls->alpn, alpn, sizeof(alpn) );
151 0 : tls->alpn_sz = sizeof(alpn);
152 :
153 0 : if( FD_UNLIKELY( client->server_fqdn_len>=sizeof(tls->server_name) ) ) {
154 0 : FD_LOG_ERR(( "Server name is %lu bytes, longer than the %lu byte maximum: check [tiles.event.url]",
155 0 : client->server_fqdn_len, sizeof(tls->server_name)-1UL ));
156 0 : }
157 0 : fd_memcpy( tls->server_name, client->server_fqdn, client->server_fqdn_len+1UL );
158 0 : tls->server_name_len = (ushort)client->server_fqdn_len;
159 :
160 0 : tls->ca_store = ca_store;
161 0 : }
162 :
163 : FD_FN_CONST ulong
164 0 : fd_event_client_align( void ) {
165 0 : return alignof( fd_event_client_t );
166 0 : }
167 :
168 : FD_FN_CONST ulong
169 0 : fd_event_client_footprint( ulong buf_max ) {
170 0 : ulong l;
171 0 : l = FD_LAYOUT_INIT;
172 0 : l = FD_LAYOUT_APPEND( l, alignof(fd_event_client_t), sizeof(fd_event_client_t) );
173 0 : l = FD_LAYOUT_APPEND( l, fd_grpc_client_align(), fd_grpc_client_footprint( buf_max ) );
174 0 : return FD_LAYOUT_FINI( l, alignof(fd_event_client_t) );
175 0 : }
176 :
177 : void *
178 : fd_event_client_new( void * shmem,
179 : fd_keyguard_client_t * keyguard_client,
180 : fd_rng_t * rng,
181 : fd_circq_t * circq,
182 : int epoll_fd,
183 : int so_sndbuf,
184 : char const * _url,
185 : uchar const * identity_pubkey,
186 : char const * client_version,
187 : char const * commit_hash,
188 : char const * action,
189 : ulong instance_id,
190 : ulong boot_id,
191 : ulong machine_id,
192 : ulong buf_max,
193 : int use_tls,
194 0 : fd_x509_ca_store_t const * ca_store ) {
195 0 : if( FD_UNLIKELY( !shmem ) ) {
196 0 : FD_LOG_WARNING(( "NULL shmem" ));
197 0 : return NULL;
198 0 : }
199 :
200 0 : if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)shmem, fd_event_client_align() ) ) ) {
201 0 : FD_LOG_WARNING(( "misaligned shmem" ));
202 0 : return NULL;
203 0 : }
204 :
205 0 : FD_SCRATCH_ALLOC_INIT( l, shmem );
206 0 : fd_event_client_t * client = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_event_client_t), sizeof(fd_event_client_t) );
207 0 : void * grpc_client_mem = FD_SCRATCH_ALLOC_APPEND( l, fd_grpc_client_align(), fd_grpc_client_footprint( buf_max ) );
208 :
209 0 : fd_url_t url[1];
210 0 : _Bool _is_ssl = 0;
211 0 : if( FD_UNLIKELY( fd_url_parse_endpoint( url,
212 0 : _url,
213 0 : strlen( _url ),
214 0 : &client->server_tcp_port,
215 0 : &_is_ssl,
216 0 : "[tiles.event.url]" ) ) ) {
217 0 : FD_LOG_ERR(( "Could not parse [tiles.event.url]" ));
218 0 : }
219 0 : if( FD_UNLIKELY( url->host_len>=FD_FQDN_BUF_MAX ) ) {
220 0 : FD_LOG_CRIT(( "Invalid url->host_len" )); /* unreachable */
221 0 : }
222 0 : fd_cstr_fini( fd_cstr_append_text( fd_cstr_init( client->server_fqdn ), url->host, url->host_len ) );
223 0 : client->server_fqdn_len = url->host_len;
224 :
225 0 : fd_memcpy( client->identity_pubkey, identity_pubkey, 32UL );
226 0 : fd_cstr_ncpy( client->client_version, client_version, sizeof( client->client_version ) );
227 0 : fd_cstr_fini( fd_cstr_append_text( fd_cstr_init( client->commit_hash ), commit_hash, fd_ulong_min( strlen( commit_hash ), sizeof( client->commit_hash )-1UL ) ) );
228 0 : fd_cstr_ncpy( client->action, action, sizeof( client->action ) );
229 :
230 0 : client->event_id = 0UL;
231 :
232 0 : client->instance_id = instance_id;
233 0 : client->boot_id = boot_id;
234 0 : client->machine_id = machine_id;
235 :
236 0 : client->has_genesis_hash = 0;
237 0 : client->has_shred_version = 0;
238 0 : client->connect_fail_logged = 0;
239 :
240 0 : client->so_sndbuf = so_sndbuf;
241 0 : client->sockfd = -1;
242 0 : FD_TEST( -1!=epoll_fd );
243 0 : client->epoll_fd = epoll_fd;
244 0 : client->use_tls = use_tls;
245 0 : if( use_tls ) {
246 0 : if( FD_UNLIKELY( !ca_store ) ) FD_LOG_ERR(( "NULL ca_store, but [tiles.event.url] requires TLS" ));
247 0 : tls_init( client, ca_store );
248 0 : }
249 0 : client->auth_deadline = LONG_MAX;
250 0 : client->auth_send_pending = 0;
251 0 : client->state = FD_EVENT_CLIENT_STATE_DISCONNECTED;
252 0 : client->disconnected.reconnect_deadline = 0L;
253 :
254 0 : client->defer_disconnect = INT_MAX;
255 0 : client->consecutive_failure_count = 7UL; /* Start high, so if server is down we don't keep retrying on boot */
256 0 : client->now = 0L;
257 0 : client->last_stream_send_ns = 0L;
258 0 : client->last_response_ns = 0L;
259 0 : client->tx_tokens = FD_EVENT_CLIENT_TX_BURST;
260 0 : client->tx_tokens_ns = 0L;
261 0 : client->stall_since = 0L;
262 :
263 0 : client->circq = circq;
264 0 : client->rng = rng;
265 0 : client->keyguard_client = keyguard_client;
266 :
267 0 : extern fd_grpc_client_callbacks_t fd_event_client_grpc_callbacks;
268 0 : client->grpc_client = fd_grpc_client_new( grpc_client_mem, &fd_event_client_grpc_callbacks, client->grpc_metrics, client, buf_max, fd_rng_ulong( rng ) );
269 0 : FD_TEST( client->grpc_client );
270 :
271 0 : memset( &client->metrics, 0, sizeof(client->metrics) );
272 0 : memset( client->grpc_metrics, 0, sizeof(fd_grpc_client_metrics_t) );
273 :
274 0 : fd_grpc_client_set_version( client->grpc_client, client->client_version, strlen( client->client_version ) );
275 0 : fd_grpc_client_set_authority( client->grpc_client, client->server_fqdn, client->server_fqdn_len, client->server_tcp_port );
276 :
277 0 : return (void *)client;
278 0 : }
279 :
280 : fd_event_client_t *
281 0 : fd_event_client_join( void * shec ) {
282 0 : if( FD_UNLIKELY( !shec ) ) {
283 0 : FD_LOG_WARNING(( "NULL shec" ));
284 0 : return NULL;
285 0 : }
286 :
287 0 : if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)shec, fd_event_client_align() ) ) ) {
288 0 : FD_LOG_WARNING(( "misaligned shec" ));
289 0 : return NULL;
290 0 : }
291 :
292 0 : fd_event_client_t * client = (fd_event_client_t *)shec;
293 :
294 0 : return client;
295 0 : }
296 :
297 : fd_event_client_metrics_t const *
298 0 : fd_event_client_metrics( fd_event_client_t const * client ) {
299 : /* Update bytes from grpc metrics */
300 0 : ((fd_event_client_t *)client)->metrics.bytes_written = client->grpc_metrics->stream_chunks_tx_bytes;
301 0 : ((fd_event_client_t *)client)->metrics.bytes_read = client->grpc_metrics->stream_chunks_rx_bytes;
302 0 : return &client->metrics;
303 0 : }
304 :
305 : ulong
306 0 : fd_event_client_state( fd_event_client_t const * client ) {
307 0 : return client->state;
308 0 : }
309 :
310 : ulong
311 54 : fd_event_client_id_reserve( fd_event_client_t * client ) {
312 54 : return client->event_id++;
313 54 : }
314 :
315 : void
316 : fd_event_client_init_genesis( fd_event_client_t * client,
317 0 : fd_genesis_meta_t const * meta ) {
318 0 : *client->genesis_hash = meta->genesis_hash;
319 0 : client->has_genesis_hash = 1;
320 0 : }
321 :
322 : void
323 : fd_event_client_init_shred_version( fd_event_client_t * client,
324 0 : ushort shred_version ) {
325 0 : client->shred_version = shred_version;
326 0 : client->has_shred_version = 1;
327 0 : }
328 :
329 : static void
330 : backoff( fd_event_client_t * client,
331 0 : long now ) {
332 0 : ulong backoff_base = 1UL << fd_ulong_min( client->consecutive_failure_count, 7UL ); /* max 4 mins */
333 0 : ulong backoff_jitter = fd_rng_ulong_roll( client->rng, backoff_base );
334 0 : client->disconnected.reconnect_deadline = now + (long)( backoff_base + backoff_jitter )*(long)1e9;
335 0 : if( FD_UNLIKELY( client->consecutive_failure_count < 8UL ) ) client->consecutive_failure_count++;
336 0 : }
337 :
338 : static void
339 : disconnect( fd_event_client_t * client,
340 : long now,
341 : int reason,
342 : int err,
343 0 : int _backoff ) {
344 0 : if( FD_LIKELY( -1!=client->sockfd ) ) {
345 0 : if( FD_UNLIKELY( -1==close( client->sockfd ) ) ) FD_LOG_ERR(( "close() failed (%d-%s)", errno, fd_io_strerror( errno ) ));
346 0 : client->sockfd = -1;
347 0 : client->epoll_out_armed = 0;
348 0 : client->state = FD_EVENT_CLIENT_STATE_DISCONNECTED;
349 0 : fd_circq_reset_cursor( client->circq );
350 0 : }
351 :
352 0 : client->event_stream = NULL;
353 0 : client->auth_deadline = LONG_MAX;
354 0 : client->auth_send_pending = 0;
355 0 : client->stall_since = 0L;
356 :
357 0 : client->auth_bearer[ 0 ] = '\0';
358 0 : client->auth_bearer_len = 0UL;
359 :
360 0 : switch( reason ) {
361 0 : case DISCONNECT_REASON_IDENTITY_CHANGED:
362 0 : FD_LOG_INFO(( "disconnected: identity changed" ));
363 0 : break;
364 0 : case DISCONNECT_REASON_CONNECT_FAILED:
365 0 : if( FD_UNLIKELY( !client->connect_fail_logged ) ) FD_LOG_WARNING(( "connecting to telemetry server " FD_IP4_ADDR_FMT ":%u failed %s(%i-%s)%s", FD_IP4_ADDR_FMT_ARGS( client->server_ip4_addr ), client->server_tcp_port, fd_log_style_dim(), errno, fd_io_strerror( errno ), fd_log_style_normal() ));
366 0 : else FD_LOG_INFO(( "connecting to telemetry server " FD_IP4_ADDR_FMT ":%u failed (%i-%s)", FD_IP4_ADDR_FMT_ARGS( client->server_ip4_addr ), client->server_tcp_port, errno, fd_io_strerror( errno ) ));
367 0 : client->connect_fail_logged = 1;
368 0 : client->metrics.transport_fail_cnt++;
369 0 : break;
370 0 : case DISCONNECT_REASON_DNS_RESOLVE_FAILED:
371 0 : if( FD_UNLIKELY( !client->connect_fail_logged ) ) FD_LOG_WARNING(( "failed to resolve telemetry server host %.*s %s(%d-%s)%s", (int)client->server_fqdn_len, client->server_fqdn, fd_log_style_dim(), err, fd_gai_strerror( err ), fd_log_style_normal() ));
372 0 : else FD_LOG_INFO(( "failed to resolve telemetry server host %.*s (%d-%s)", (int)client->server_fqdn_len, client->server_fqdn, err, fd_gai_strerror( err ) ));
373 0 : client->connect_fail_logged = 1;
374 0 : client->metrics.transport_fail_cnt++;
375 0 : break;
376 0 : case DISCONNECT_REASON_TIMEOUT:
377 0 : FD_LOG_INFO(( "connection failed: timeout" ));
378 0 : client->metrics.transport_fail_cnt++;
379 0 : break;
380 0 : case DISCONNECT_REASON_TRANSPORT_FAILED:
381 0 : FD_LOG_WARNING(( "disconnected from telemetry server: transport failed %s(%d-%s)%s", fd_log_style_dim(), err, fd_io_strerror( err ), fd_log_style_normal() ));
382 0 : client->metrics.transport_fail_cnt++;
383 0 : break;
384 0 : case DISCONNECT_REASON_PEER_CLOSED:
385 0 : FD_LOG_WARNING(( "disconnected from telemetry server: peer closed connection" ));
386 0 : client->metrics.transport_fail_cnt++;
387 0 : break;
388 0 : case DISCONNECT_REASON_INVALID_CURSOR:
389 0 : FD_LOG_WARNING(( "disconnected from telemetry server: invalid cursor" ));
390 0 : client->metrics.transport_fail_cnt++;
391 0 : break;
392 0 : case DISCONNECT_REASON_AUTH_FAILED:
393 0 : FD_LOG_WARNING(( "disconnected from telemetry server: authentication failed" ));
394 0 : client->metrics.transport_fail_cnt++;
395 0 : break;
396 0 : case DISCONNECT_REASON_INVALID_PROTOBUF:
397 0 : FD_LOG_WARNING(( "disconnected from telemetry server: invalid protobuf message received" ));
398 0 : client->metrics.transport_fail_cnt++;
399 0 : break;
400 0 : default:
401 0 : FD_LOG_WARNING(( "disconnected from telemetry server: unknown reason %d", reason ));
402 0 : client->metrics.transport_fail_cnt++;
403 0 : break;
404 0 : }
405 :
406 0 : if( FD_LIKELY( _backoff ) ) backoff( client, now );
407 0 : }
408 :
409 : void
410 : fd_event_client_set_identity( fd_event_client_t * client,
411 0 : uchar const * identity_pubkey ) {
412 0 : fd_memcpy( client->identity_pubkey, identity_pubkey, 32UL );
413 0 : disconnect( client, client->now, DISCONNECT_REASON_IDENTITY_CHANGED, 0, 0 );
414 0 : }
415 :
416 : static void
417 : reconnect( fd_event_client_t * client,
418 : long now,
419 0 : int * charge_busy ) {
420 0 : FD_TEST( client->state==FD_EVENT_CLIENT_STATE_DISCONNECTED );
421 :
422 0 : if( FD_UNLIKELY( now<client->disconnected.reconnect_deadline ) ) return;
423 :
424 0 : *charge_busy = 1;
425 0 : client->metrics.connect_attempt_cnt++;
426 :
427 0 : FD_LOG_INFO(( "connecting to event server %s://%.*s:%u", client->use_tls ? "https" : "http", (int)client->server_fqdn_len, client->server_fqdn, client->server_tcp_port ));
428 :
429 : /* FIXME IPv6 support */
430 0 : fd_addrinfo_t hints = {0};
431 0 : hints.ai_family = AF_INET;
432 0 : fd_addrinfo_t * res = NULL;
433 0 : uchar scratch[ 4096 ];
434 0 : void * pscratch = scratch;
435 0 : int err = fd_getaddrinfo( client->server_fqdn, &hints, &res, &pscratch, sizeof(scratch) );
436 0 : if( FD_UNLIKELY( err ) ) {
437 0 : disconnect( client, now, DISCONNECT_REASON_DNS_RESOLVE_FAILED, err, 1 );
438 0 : return;
439 0 : }
440 :
441 0 : if( FD_UNLIKELY( !res || !res->ai_addr ) ) {
442 0 : disconnect( client, now, DISCONNECT_REASON_DNS_RESOLVE_FAILED, 0, 1 );
443 0 : return;
444 0 : }
445 :
446 0 : uint const ip4_addr = ((struct sockaddr_in *)res->ai_addr)->sin_addr.s_addr;
447 0 : client->server_ip4_addr = ip4_addr;
448 :
449 0 : client->sockfd = socket( AF_INET, SOCK_STREAM|SOCK_NONBLOCK, 0 );
450 0 : if( FD_UNLIKELY( -1==client->sockfd ) ) FD_LOG_ERR(( "socket() failed (%d-%s)", errno, fd_io_strerror( errno ) ));
451 :
452 0 : struct epoll_event ev = { .events = EPOLLIN|EPOLLOUT, .data.fd = client->sockfd };
453 0 : if( FD_UNLIKELY( -1==epoll_ctl( client->epoll_fd, EPOLL_CTL_ADD, client->sockfd, &ev ) ) ) FD_LOG_ERR(( "epoll_ctl(ADD,sockfd) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
454 0 : client->epoll_out_armed = 1;
455 :
456 0 : struct sockaddr_in addr;
457 0 : fd_memset( &addr, 0, sizeof( addr ) );
458 0 : addr.sin_family = AF_INET;
459 0 : addr.sin_port = fd_ushort_bswap( client->server_tcp_port );
460 0 : addr.sin_addr.s_addr = ip4_addr;
461 :
462 0 : int tcp_nodelay = 1;
463 0 : if( FD_UNLIKELY( -1==setsockopt( client->sockfd, SOL_TCP, TCP_NODELAY, &tcp_nodelay, sizeof(int) ) ) ) FD_LOG_ERR(( "setsockopt failed (%d-%s)", errno, fd_io_strerror( errno ) ));
464 0 : if( FD_UNLIKELY( -1==setsockopt( client->sockfd, SOL_SOCKET, SO_SNDBUF, &client->so_sndbuf, sizeof(int) ) ) ) FD_LOG_ERR(( "setsockopt(SOL_SOCKET,SO_SNDBUF,%i) failed (%i-%s)", client->so_sndbuf, errno, fd_io_strerror( errno ) ));
465 :
466 0 : if( FD_UNLIKELY( -1==connect( client->sockfd, fd_type_pun_const( &addr ), sizeof(struct sockaddr_in) ) && errno!=EINPROGRESS ) ) {
467 0 : disconnect( client, now, DISCONNECT_REASON_CONNECT_FAILED, errno, 1 );
468 0 : return;
469 0 : }
470 :
471 0 : if( client->use_tls ) {
472 0 : fd_tls_t * tls = client->tls;
473 0 : if( FD_UNLIKELY( !fd_rng_secure( tls->kex_private_key, 32UL ) ) ) FD_LOG_CRIT(( "fd_rng_secure failed" ));
474 0 : fd_x25519_public( tls->kex_public_key, tls->kex_private_key );
475 0 : fd_tlsrec_conn_init( client->tls_conn, tls, 0 );
476 0 : }
477 :
478 0 : fd_grpc_client_reset( client->grpc_client );
479 :
480 0 : client->state = FD_EVENT_CLIENT_STATE_CONNECTING;
481 0 : client->connecting.connect_deadline = now+(long)1L*(long)1e9; /* 1 second to connect */
482 0 : }
483 :
484 : static int
485 : fd_event_client_try_send_authenticate( fd_event_client_t * client,
486 0 : long now ) {
487 0 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( client->grpc_client ) ) ) return 0;
488 0 : if( FD_UNLIKELY( fd_grpc_client_request_stream_busy( client->grpc_client ) ) ) return 0;
489 :
490 0 : fd_pb_encoder_t auth_req[1];
491 0 : uchar buffer[ 256UL ];
492 0 : fd_pb_encoder_init( auth_req, buffer, sizeof(buffer) );
493 :
494 0 : fd_pb_push_bytes( auth_req, 1U, client->identity_pubkey, 32UL );
495 0 : fd_pb_push_string( auth_req, 2U, client->client_version, strlen( client->client_version ) );
496 0 : fd_pb_push_string( auth_req, 3U, client->commit_hash, strlen( client->commit_hash ) );
497 0 : fd_pb_push_bytes( auth_req, 4U, client->genesis_hash, 32UL );
498 0 : fd_pb_push_uint64( auth_req, 5U, client->shred_version );
499 0 : fd_pb_push_uint64( auth_req, 6U, client->instance_id );
500 0 : fd_pb_push_uint64( auth_req, 7U, client->machine_id );
501 0 : fd_pb_push_uint64( auth_req, 8U, client->boot_id );
502 0 : fd_pb_push_string( auth_req, 9U, client->action, strlen( client->action ) );
503 :
504 0 : fd_grpc_h2_stream_t * stream = fd_grpc_client_request_start1(
505 0 : client->grpc_client,
506 0 : "/events.v1.EventService/Authenticate", strlen("/events.v1.EventService/Authenticate"),
507 0 : FD_EVENT_CLIENT_REQ_CTX_AUTHENTICATE,
508 0 : buffer, fd_pb_encoder_out_sz( auth_req ),
509 0 : NULL, 0UL,
510 0 : 0 /* not streaming */ );
511 :
512 0 : if( FD_UNLIKELY( !stream ) ) return 0;
513 :
514 0 : fd_grpc_client_deadline_set( stream, FD_GRPC_DEADLINE_HEADER, now+(long)2e9 );
515 0 : fd_grpc_client_deadline_set( stream, FD_GRPC_DEADLINE_RX_END, now+(long)2e9 );
516 :
517 0 : client->auth_send_pending = 0;
518 0 : FD_LOG_INFO(( "Requesting auth challenge from event server " FD_IP4_ADDR_FMT ":%u (%.*s)",
519 0 : FD_IP4_ADDR_FMT_ARGS( client->server_ip4_addr ), client->server_tcp_port,
520 0 : (int)client->server_fqdn_len, client->server_fqdn ));
521 0 : return 1;
522 0 : }
523 :
524 : static void
525 0 : fd_event_client_grpc_conn_established( void * app_ctx ) {
526 0 : fd_event_client_t * client = app_ctx;
527 :
528 0 : long now = client->now;
529 0 : client->state = FD_EVENT_CLIENT_STATE_AUTHENTICATING;
530 0 : client->auth_deadline = now + (long)2e9;
531 0 : client->auth_send_pending = 1;
532 :
533 0 : fd_event_client_try_send_authenticate( client, now );
534 0 : }
535 :
536 : static void
537 : fd_event_client_handle_auth_challenge_resp( fd_event_client_t * client,
538 : void const * protobuf,
539 0 : ulong protobuf_sz ) {
540 0 : fd_pb_inbuf_t inbuf[1];
541 0 : fd_pb_inbuf_init( inbuf, protobuf, protobuf_sz );
542 :
543 0 : if( FD_UNLIKELY( protobuf_sz==0UL ) ) {
544 0 : FD_LOG_WARNING(( "Empty auth challenge response" ));
545 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
546 0 : return;
547 0 : }
548 :
549 0 : fd_pb_tlv_t challenge_tlv;
550 0 : if( FD_UNLIKELY( !fd_pb_read_tlv( inbuf, &challenge_tlv ) ) ) {
551 0 : FD_LOG_WARNING(( "Failed to parse auth challenge response" ));
552 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
553 0 : return;
554 0 : }
555 :
556 0 : if( FD_UNLIKELY( challenge_tlv.field_id!=1U || challenge_tlv.wire_type!=FD_PB_WIRE_TYPE_LEN ) ) {
557 0 : FD_LOG_WARNING(( "Unexpected field in auth challenge response" ));
558 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
559 0 : return;
560 0 : }
561 :
562 0 : ulong challenge_len = challenge_tlv.len;
563 0 : if( FD_UNLIKELY( challenge_len!=FD_EVENT_CLIENT_TOKEN_SZ ) ) {
564 0 : FD_LOG_WARNING(( "Invalid challenge token size: %lu bytes (expected %lu)", challenge_len, FD_EVENT_CLIENT_TOKEN_SZ ));
565 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
566 0 : return;
567 0 : }
568 :
569 0 : if( FD_UNLIKELY( fd_pb_inbuf_sz( inbuf )<challenge_len ) ) {
570 0 : FD_LOG_WARNING(( "Truncated auth challenge response" ));
571 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
572 0 : return;
573 0 : }
574 :
575 0 : uchar challenge_token[ FD_EVENT_CLIENT_TOKEN_SZ ];
576 0 : memcpy( challenge_token, inbuf->cur, challenge_len );
577 0 : inbuf->cur += challenge_len;
578 :
579 0 : if( FD_UNLIKELY( fd_pb_inbuf_sz( inbuf ) ) ) {
580 0 : FD_LOG_WARNING(( "Trailing data in auth challenge response" ));
581 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
582 0 : return;
583 0 : }
584 :
585 0 : uchar sign_request[ 100UL + FD_EVENT_CLIENT_TOKEN_SZ ];
586 0 : static char const sign_prefix[ 100 ] =
587 0 : " " /* 32 spaces */
588 0 : " " /* 32 spaces */
589 0 : "Firedancer event challenge-response";
590 0 : memcpy( sign_request, sign_prefix, sizeof(sign_prefix) );
591 0 : memcpy( sign_request+100, challenge_token, challenge_len );
592 :
593 0 : uchar signature[ 64UL ];
594 0 : fd_keyguard_client_sign( client->keyguard_client,
595 0 : signature,
596 0 : sign_request, 100UL+challenge_len,
597 0 : FD_KEYGUARD_SIGN_TYPE_ED25519 );
598 :
599 : /* Build "hex(challenge_token).hex(signature)" for the bearer token. */
600 0 : fd_hex_encode( client->auth_bearer, challenge_token, FD_EVENT_CLIENT_TOKEN_SZ );
601 0 : client->auth_bearer[ 2UL*FD_EVENT_CLIENT_TOKEN_SZ ] = '.';
602 0 : fd_hex_encode( client->auth_bearer + 2UL*FD_EVENT_CLIENT_TOKEN_SZ+1UL, signature, 64UL );
603 0 : client->auth_bearer_len = 2UL*FD_EVENT_CLIENT_TOKEN_SZ + 1UL + 2UL*64UL;
604 0 : client->auth_bearer[ client->auth_bearer_len ] = '\0';
605 :
606 0 : client->event_stream = NULL;
607 0 : client->metrics.transport_success_cnt++;
608 0 : client->state = FD_EVENT_CLIENT_STATE_CONNECTED;
609 0 : client->connected.connected_timestamp = client->now;
610 0 : client->connect_fail_logged = 0;
611 0 : FD_LOG_NOTICE(( "connected to telemetry server %s%s://%.*s:%u%s",
612 0 : fd_log_style_bold(), client->use_tls ? "https" : "http",
613 0 : (int)client->server_fqdn_len, client->server_fqdn, client->server_tcp_port,
614 0 : fd_log_style_normal() ));
615 0 : }
616 :
617 : static void
618 : fd_event_client_grpc_conn_dead( void * app_ctx,
619 : uint h2_err,
620 0 : int closed_by ) {
621 0 : fd_event_client_t * client = app_ctx;
622 0 : FD_LOG_WARNING(( "telemetry connection closed %s %s(%u-%s)%s",
623 0 : closed_by ? "by peer" : "due to error",
624 0 : fd_log_style_dim(), h2_err, fd_h2_strerror( h2_err ), fd_log_style_normal() ));
625 0 : client->defer_disconnect = DISCONNECT_REASON_PEER_CLOSED;
626 0 : }
627 :
628 : static void
629 : fd_event_client_grpc_tx_complete( void * app_ctx,
630 0 : ulong request_ctx ) {
631 0 : (void)app_ctx; (void)request_ctx;
632 0 : }
633 :
634 : void
635 : fd_event_client_grpc_rx_start( void * app_ctx,
636 0 : ulong request_ctx ) {
637 0 : (void)app_ctx; (void)request_ctx;
638 0 : }
639 :
640 : static void
641 : fd_event_client_handle_stream_events_resp( fd_event_client_t * client,
642 : void const * protobuf,
643 0 : ulong protobuf_sz ) {
644 0 : fd_pb_inbuf_t inbuf[1];
645 0 : fd_pb_inbuf_init( inbuf, protobuf, protobuf_sz );
646 :
647 0 : ulong nonce_ack = 0UL;
648 0 : if( FD_LIKELY( protobuf_sz ) ) {
649 0 : fd_pb_tlv_t event_id;
650 0 : if( FD_UNLIKELY( !fd_pb_read_tlv( inbuf, &event_id ) ||
651 0 : event_id.field_id!=1U /* event_id */ ||
652 0 : event_id.wire_type!=FD_PB_WIRE_TYPE_VARINT ) ) {
653 0 : FD_LOG_WARNING(( "Event gRPC rx msg: invalid Protobuf" ));
654 0 : client->defer_disconnect = DISCONNECT_REASON_INVALID_PROTOBUF;
655 0 : return;
656 0 : }
657 0 : nonce_ack = event_id.varint;
658 :
659 0 : if( FD_UNLIKELY( fd_pb_inbuf_sz( inbuf ) ) ) {
660 0 : FD_LOG_WARNING(( "Event gRPC rx msg: trailing data in StreamEventsResponse" ));
661 0 : client->defer_disconnect = DISCONNECT_REASON_INVALID_PROTOBUF;
662 0 : return;
663 0 : }
664 0 : }
665 :
666 0 : client->metrics.events_acked++;
667 0 : client->last_response_ns = client->now;
668 0 : if( FD_UNLIKELY( nonce_ack==ULONG_MAX ) ) return;
669 :
670 0 : client->metrics.last_acked_id = nonce_ack;
671 :
672 0 : int err = fd_circq_pop_until( client->circq, nonce_ack );
673 0 : if( FD_UNLIKELY( -1==err ) ) {
674 0 : FD_LOG_WARNING(( "Event gRPC rx msg: invalid cursor ack %lu", nonce_ack ));
675 0 : client->defer_disconnect = DISCONNECT_REASON_INVALID_CURSOR;
676 0 : }
677 0 : }
678 :
679 : void
680 : fd_event_client_grpc_rx_msg( void * app_ctx,
681 : void const * protobuf,
682 : ulong protobuf_sz,
683 0 : ulong request_ctx ) {
684 0 : fd_event_client_t * client = app_ctx;
685 :
686 0 : switch( request_ctx ) {
687 0 : case FD_EVENT_CLIENT_REQ_CTX_AUTHENTICATE:
688 0 : fd_event_client_handle_auth_challenge_resp( client, protobuf, protobuf_sz );
689 0 : break;
690 0 : case FD_EVENT_CLIENT_REQ_CTX_STREAM_EVENTS:
691 0 : fd_event_client_handle_stream_events_resp( client, protobuf, protobuf_sz );
692 0 : break;
693 0 : default:
694 0 : FD_LOG_WARNING(( "Unknown request_ctx: %lu, disconnecting", request_ctx ));
695 0 : client->defer_disconnect = DISCONNECT_REASON_INVALID_PROTOBUF;
696 0 : break;
697 0 : }
698 0 : }
699 :
700 : void
701 : fd_event_client_grpc_rx_end( void * app_ctx,
702 : ulong request_ctx,
703 0 : fd_grpc_resp_hdrs_t * resp ) {
704 0 : fd_event_client_t * client = app_ctx;
705 :
706 0 : if( FD_UNLIKELY( resp->h2_status!=200 ) ) {
707 0 : FD_LOG_WARNING(( "telemetry server request failed %s(HTTP status %u)%s", fd_log_style_dim(), resp->h2_status, fd_log_style_normal() ));
708 0 : client->defer_disconnect = DISCONNECT_REASON_TRANSPORT_FAILED;
709 0 : return;
710 0 : }
711 :
712 0 : resp->grpc_msg_len = (uint)fd_url_unescape( resp->grpc_msg, resp->grpc_msg_len );
713 0 : if( !resp->grpc_msg_len ) {
714 0 : fd_memcpy( resp->grpc_msg, "unknown error", 13 );
715 0 : resp->grpc_msg_len = 13;
716 0 : }
717 :
718 0 : if( FD_UNLIKELY( resp->grpc_status!=FD_GRPC_STATUS_OK ) ) {
719 0 : switch( request_ctx ) {
720 0 : case FD_EVENT_CLIENT_REQ_CTX_AUTHENTICATE:
721 0 : FD_LOG_WARNING(( "telemetry server authentication failed: %.*s %s(%u-%s)%s",
722 0 : (int)resp->grpc_msg_len, resp->grpc_msg,
723 0 : fd_log_style_dim(), resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ), fd_log_style_normal() ));
724 0 : client->defer_disconnect = DISCONNECT_REASON_AUTH_FAILED;
725 0 : return;
726 0 : case FD_EVENT_CLIENT_REQ_CTX_STREAM_EVENTS:
727 0 : FD_LOG_WARNING(( "telemetry server event stream failed: %.*s %s(%u-%s)%s",
728 0 : (int)resp->grpc_msg_len, resp->grpc_msg,
729 0 : fd_log_style_dim(), resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ), fd_log_style_normal() ));
730 0 : client->defer_disconnect = DISCONNECT_REASON_PEER_CLOSED;
731 0 : return;
732 0 : default:
733 0 : FD_LOG_WARNING(( "telemetry server request failed: %.*s %s(%u-%s)%s",
734 0 : (int)resp->grpc_msg_len, resp->grpc_msg,
735 0 : fd_log_style_dim(), resp->grpc_status, fd_grpc_status_cstr( resp->grpc_status ), fd_log_style_normal() ));
736 0 : client->defer_disconnect = DISCONNECT_REASON_TRANSPORT_FAILED;
737 0 : return;
738 0 : }
739 0 : }
740 :
741 0 : if( request_ctx==FD_EVENT_CLIENT_REQ_CTX_STREAM_EVENTS ) {
742 0 : FD_LOG_INFO(( "telemetry server event stream ended gracefully" ));
743 0 : client->defer_disconnect = DISCONNECT_REASON_PEER_CLOSED;
744 0 : }
745 0 : }
746 :
747 : void
748 : fd_event_client_grpc_rx_timeout( void * app_ctx,
749 : ulong request_ctx FD_PARAM_UNUSED,
750 0 : int deadline_kind FD_PARAM_UNUSED ) {
751 0 : FD_LOG_WARNING(( "Event gRPC rx timeout" ));
752 0 : fd_event_client_t * client = (fd_event_client_t *)app_ctx;
753 0 : client->defer_disconnect = DISCONNECT_REASON_TRANSPORT_FAILED;
754 0 : client->event_stream = NULL;
755 0 : }
756 :
757 : static void
758 0 : fd_event_client_grpc_ping_ack( void * app_ctx ) {
759 0 : (void)app_ctx;
760 0 : FD_LOG_WARNING(( "Event gRPC ping ack" ));
761 0 : }
762 :
763 : static long
764 : pace_refill( fd_event_client_t * client,
765 0 : long now ) {
766 0 : long ns_per_byte = (long)1e9/FD_EVENT_CLIENT_TX_RATE_BPS;
767 0 : long dt = fd_long_max( now-client->tx_tokens_ns, 0L );
768 0 : long grant = fd_long_min( dt/ns_per_byte, FD_EVENT_CLIENT_TX_BURST-client->tx_tokens );
769 0 : client->tx_tokens += grant;
770 0 : client->tx_tokens_ns = client->tx_tokens==FD_EVENT_CLIENT_TX_BURST ? now : client->tx_tokens_ns+grant*ns_per_byte;
771 0 : return client->tx_tokens;
772 0 : }
773 :
774 : static void
775 : tx( fd_event_client_t * client,
776 : long now,
777 0 : int * charge_busy ) {
778 0 : FD_TEST( client->state==FD_EVENT_CLIENT_STATE_CONNECTED );
779 :
780 0 : long tokens = pace_refill( client, now );
781 :
782 0 : if( FD_UNLIKELY( client->event_stream && client->grpc_client->request_stream != NULL && client->grpc_client->request_stream!=client->event_stream ) ) return;
783 :
784 0 : if( FD_UNLIKELY( client->event_stream ) ) {
785 0 : if( FD_UNLIKELY( fd_grpc_client_stream_send_is_blocked( client->grpc_client ) ) ) return;
786 0 : } else {
787 0 : if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( client->grpc_client ) ) ) return;
788 0 : }
789 :
790 0 : if( FD_UNLIKELY( !client->event_stream ) ) {
791 0 : client->event_stream = fd_grpc_client_request_start1(
792 0 : client->grpc_client,
793 0 : "/events.v1.EventService/StreamEvents", strlen("/events.v1.EventService/StreamEvents"),
794 0 : FD_EVENT_CLIENT_REQ_CTX_STREAM_EVENTS,
795 0 : NULL, 0UL, /* headers only; first message sent later */
796 0 : client->auth_bearer, client->auth_bearer_len,
797 0 : 1 /* streaming */ );
798 0 : if( FD_UNLIKELY( !client->event_stream ) ) return; /* transient; retry next poll */
799 0 : fd_grpc_client_deadline_set( client->event_stream, FD_GRPC_DEADLINE_HEADER, now+(long)10e9 /* 10s */ );
800 0 : client->last_stream_send_ns = now;
801 0 : client->last_response_ns = now;
802 0 : *charge_busy = 1;
803 0 : return;
804 0 : }
805 :
806 0 : if( FD_UNLIKELY( tokens<=0L ) ) return;
807 :
808 0 : ulong msg_sz;
809 0 : uchar const * msg = fd_circq_cursor_advance( client->circq, &msg_sz );
810 0 : if( FD_LIKELY( !msg ) ) {
811 : /* Nothing to send. If the stream has been quiet long enough that an
812 : intermediary proxy might kill it, send a zero-length
813 : StreamEventsRequest to heartbeat. */
814 0 : if( FD_UNLIKELY( now-client->last_stream_send_ns>FD_EVENT_CLIENT_HEARTBEAT_NANOS ) ) {
815 0 : if( FD_LIKELY( fd_grpc_client_stream_send_msg1( client->grpc_client, client->event_stream, (uchar const *)"", 0UL ) ) ) {
816 0 : client->last_stream_send_ns = now;
817 0 : *charge_busy = 1;
818 0 : }
819 0 : }
820 0 : return;
821 0 : }
822 :
823 0 : int result = fd_grpc_client_stream_send_msg1( client->grpc_client, client->event_stream, msg, msg_sz );
824 0 : if( FD_UNLIKELY( !result ) ) return; /* Only reason for failure is too big message, so just skip it */
825 :
826 0 : client->tx_tokens -= (long)msg_sz;
827 0 : client->metrics.events_sent++;
828 0 : client->last_stream_send_ns = now;
829 0 : *charge_busy = 1;
830 0 : }
831 :
832 : static void
833 : epoll_out( fd_event_client_t * client,
834 0 : int want ) {
835 0 : if( FD_UNLIKELY( client->sockfd<0 ) ) return;
836 0 : if( FD_LIKELY( want==client->epoll_out_armed ) ) return;
837 0 : struct epoll_event ev = { .events = EPOLLIN | (want ? EPOLLOUT : 0U), .data.fd = client->sockfd };
838 0 : if( FD_UNLIKELY( -1==epoll_ctl( client->epoll_fd, EPOLL_CTL_MOD, client->sockfd, &ev ) ) ) FD_LOG_ERR(( "epoll_ctl(MOD,sockfd) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
839 0 : client->epoll_out_armed = want;
840 0 : }
841 :
842 : long
843 : fd_event_client_next_deadline( fd_event_client_t const * client,
844 0 : long now ) {
845 0 : if( FD_UNLIKELY( !client->has_genesis_hash || !client->has_shred_version ) ) return LONG_MAX; /* frag-driven */
846 :
847 0 : switch( client->state ) {
848 0 : case FD_EVENT_CLIENT_STATE_DISCONNECTED:
849 0 : return fd_long_max( client->disconnected.reconnect_deadline, now );
850 0 : case FD_EVENT_CLIENT_STATE_CONNECTING:
851 : /* Completion/failure arrives as an EPOLLOUT/EPOLLERR event. */
852 0 : return client->connecting.connect_deadline;
853 0 : default: break;
854 0 : }
855 :
856 : /* Decrypted TLS bytes buffered in the gRPC client do not make the fd
857 : readable: drain now or they wait for the next unrelated event.
858 : Except while TX is blocked: then they are held because the HTTP/2
859 : rings are full behind the parked send, and only the EPOLLOUT armed
860 : for that send can make progress. */
861 0 : if( FD_UNLIKELY( client->use_tls &&
862 0 : fd_grpc_client_tls_rx_pending( client->grpc_client ) &&
863 0 : !fd_grpc_client_tls_tx_pending( client->grpc_client ) ) ) return now;
864 :
865 0 : long deadline = fd_grpc_client_next_deadline( client->grpc_client );
866 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_AUTHENTICATING ) ) {
867 0 : deadline = fd_long_min( deadline, client->auth_deadline );
868 0 : if( FD_UNLIKELY( client->auth_send_pending ) ) deadline = now;
869 0 : }
870 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_CONNECTED && client->consecutive_failure_count ) ) {
871 0 : deadline = fd_long_min( deadline, client->connected.connected_timestamp+10L*(long)1e9 );
872 0 : }
873 0 : if( FD_LIKELY( client->state==FD_EVENT_CLIENT_STATE_CONNECTED && client->event_stream ) ) {
874 0 : deadline = fd_long_min( deadline, fd_long_min( client->last_response_ns +FD_EVENT_CLIENT_RESPONSE_TIMEOUT_NANOS,
875 0 : client->last_stream_send_ns+FD_EVENT_CLIENT_HEARTBEAT_NANOS ) );
876 0 : if( FD_UNLIKELY( client->tx_tokens<=0L && fd_circq_unsent_cnt( client->circq ) ) ) {
877 0 : deadline = fd_long_min( deadline, client->tx_tokens_ns + (1L-client->tx_tokens)*(long)1e9/FD_EVENT_CLIENT_TX_RATE_BPS );
878 0 : }
879 0 : if( FD_UNLIKELY( client->stall_since ) ) deadline = fd_long_min( deadline, client->stall_since+FD_EVENT_CLIENT_CREDIT_STALL_NANOS );
880 0 : }
881 0 : return deadline;
882 0 : }
883 :
884 : static int
885 : credit_stall_check( fd_event_client_t * client,
886 0 : long now ) {
887 0 : ulong rem = fd_grpc_client_tx_starved( client->grpc_client );
888 0 : if( FD_LIKELY( !rem ) ) {
889 0 : client->stall_since = 0L;
890 0 : return 0;
891 0 : }
892 0 : if( !client->stall_since || rem!=client->stall_rem || client->metrics.events_sent!=client->stall_events_sent ) {
893 0 : client->stall_since = now;
894 0 : client->stall_rem = rem;
895 0 : client->stall_events_sent = client->metrics.events_sent;
896 0 : return 0;
897 0 : }
898 0 : if( FD_LIKELY( now-client->stall_since<=FD_EVENT_CLIENT_CREDIT_STALL_NANOS ) ) return 0;
899 0 : FD_LOG_WARNING(( "telemetry server stopped granting flow-control credit for over %ld seconds, reconnecting", FD_EVENT_CLIENT_CREDIT_STALL_NANOS/(long)1e9 ));
900 0 : client->metrics.credit_stall_cnt++;
901 0 : disconnect( client, now, DISCONNECT_REASON_TIMEOUT, 0, 1 );
902 0 : return 1;
903 0 : }
904 :
905 : static void
906 : poll1( fd_event_client_t * client,
907 : long now,
908 0 : int * charge_busy ) {
909 0 : if( FD_UNLIKELY( !client->has_genesis_hash || !client->has_shred_version ) ) return;
910 :
911 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_DISCONNECTED ) ) reconnect( client, now, charge_busy );
912 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_CONNECTING ) ) {
913 0 : if( FD_UNLIKELY( now>client->connecting.connect_deadline ) ) {
914 0 : disconnect( client, now, DISCONNECT_REASON_TIMEOUT, 0, 1 );
915 0 : return;
916 0 : }
917 0 : }
918 : /* Check auth handshake timeout */
919 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_AUTHENTICATING && now>client->auth_deadline ) ) {
920 0 : FD_LOG_WARNING(( "auth handshake timed out" ));
921 0 : client->metrics.handshake_timeout_cnt++;
922 0 : disconnect( client, now, DISCONNECT_REASON_TIMEOUT, 0, 1 );
923 0 : return;
924 0 : }
925 0 : if( FD_LIKELY( client->state!=FD_EVENT_CLIENT_STATE_DISCONNECTED ) ) {
926 0 : int rxtx_err;
927 0 : if( client->use_tls )
928 0 : rxtx_err = fd_grpc_client_rxtx_tls( client->grpc_client, client->tls_conn, client->sockfd, now, charge_busy );
929 0 : else
930 0 : rxtx_err = fd_grpc_client_rxtx_socket( client->grpc_client, client->sockfd, now, charge_busy );
931 0 : if( FD_UNLIKELY( -1==rxtx_err ) ) {
932 0 : disconnect( client, now, DISCONNECT_REASON_TRANSPORT_FAILED, errno, 1 );
933 0 : return;
934 0 : }
935 0 : }
936 :
937 0 : if( FD_UNLIKELY( client->defer_disconnect!=INT_MAX ) ) {
938 0 : int reason = client->defer_disconnect;
939 0 : client->defer_disconnect = INT_MAX;
940 0 : if( reason==DISCONNECT_REASON_AUTH_FAILED ) client->metrics.auth_fail_cnt++;
941 0 : if( reason==DISCONNECT_REASON_INVALID_PROTOBUF ) client->metrics.invalid_msg_cnt++;
942 0 : disconnect( client, now, reason, 0, 1 );
943 0 : return;
944 0 : }
945 :
946 0 : if( FD_UNLIKELY( client->state==FD_EVENT_CLIENT_STATE_AUTHENTICATING && client->auth_send_pending ) ) {
947 0 : fd_event_client_try_send_authenticate( client, now );
948 0 : }
949 :
950 0 : if( FD_LIKELY( client->state==FD_EVENT_CLIENT_STATE_CONNECTED ) ) {
951 0 : if( FD_UNLIKELY( client->event_stream && now-client->last_response_ns>FD_EVENT_CLIENT_RESPONSE_TIMEOUT_NANOS ) ) {
952 0 : FD_LOG_WARNING(( "no response from telemetry server in over %ld seconds", FD_EVENT_CLIENT_RESPONSE_TIMEOUT_NANOS/(long)1e9 ));
953 0 : disconnect( client, now, DISCONNECT_REASON_TIMEOUT, 0, 1 );
954 0 : return;
955 0 : }
956 0 : if( FD_UNLIKELY( client->consecutive_failure_count && (now-client->connected.connected_timestamp>10L*(long)1e9 ) ) ) client->consecutive_failure_count = 0UL;
957 0 : if( FD_UNLIKELY( credit_stall_check( client, now ) ) ) return;
958 0 : tx( client, now, charge_busy );
959 0 : }
960 0 : }
961 :
962 : void
963 : fd_event_client_poll( fd_event_client_t * client,
964 : long now,
965 0 : int * charge_busy ) {
966 0 : client->now = now;
967 0 : poll1( client, now, charge_busy );
968 :
969 : /* poll1 flushes before tx() enqueues: flush once more so bytes queued
970 : this poll go out now instead of arming EPOLLOUT on a writable
971 : socket (a spurious waker roundtrip per message). */
972 0 : if( FD_UNLIKELY( client->state!=FD_EVENT_CLIENT_STATE_DISCONNECTED && fd_grpc_client_tx_pending( client->grpc_client ) ) ) {
973 0 : int flush_err = client->use_tls ? fd_grpc_client_tls_flush ( client->grpc_client, client->sockfd )
974 0 : : fd_grpc_client_tx_flush_socket( client->grpc_client, client->sockfd );
975 0 : if( FD_UNLIKELY( -1==flush_err ) ) {
976 0 : disconnect( client, now, DISCONNECT_REASON_TRANSPORT_FAILED, errno, 1 );
977 0 : return;
978 0 : }
979 0 : }
980 :
981 : /* Arm EPOLLOUT only on write demand: unsent HTTP/2 bytes, or a TLS
982 : handshake blocked on write (which is also what an in-progress
983 : TCP connect looks like: the handshake write EAGAINed, and connect
984 : completion is the writability event). Anything else leaves an
985 : idle-writable socket keeping the epoll set ready. */
986 0 : int want = fd_grpc_client_tx_pending( client->grpc_client );
987 0 : want |= client->use_tls && fd_grpc_client_tls_tx_pending( client->grpc_client );
988 0 : epoll_out( client, want );
989 0 : }
990 :
991 : fd_grpc_client_callbacks_t fd_event_client_grpc_callbacks = {
992 : .conn_established = fd_event_client_grpc_conn_established,
993 : .conn_dead = fd_event_client_grpc_conn_dead,
994 : .tx_complete = fd_event_client_grpc_tx_complete,
995 : .rx_start = fd_event_client_grpc_rx_start,
996 : .rx_msg = fd_event_client_grpc_rx_msg,
997 : .rx_end = fd_event_client_grpc_rx_end,
998 : .rx_timeout = fd_event_client_grpc_rx_timeout,
999 : .ping_ack = fd_event_client_grpc_ping_ack,
1000 : };
|