LCOV - code coverage report
Current view: top level - app/shared_dev/rpc_client - fd_rpc_client.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 192 267 71.9 %
Date: 2026-09-17 04:28:31 Functions: 11 13 84.6 %

          Line data    Source code
       1             : #include "fd_rpc_client.h"
       2             : #include "fd_rpc_client_private.h"
       3             : 
       4             : #include "../../../third_party/picohttpparser/picohttpparser.h"
       5             : #include "../../../waltz/http/fd_http.h"
       6             : #include "../../../ballet/base58/fd_base58.h"
       7             : #include "../../../ballet/json/fd_jtok.h"
       8             : 
       9             : #include <errno.h>
      10             : #include <stdio.h>
      11             : #include <unistd.h>
      12             : #include <strings.h>
      13             : #include <sys/socket.h>
      14             : #include <sys/types.h>
      15             : #include <netinet/ip.h>
      16             : 
      17             : #define MAX_REQUEST_LEN (1024UL)
      18             : 
      19             : void *
      20             : fd_rpc_client_new( void * mem,
      21             :                    uint   rpc_addr,
      22           3 :                    ushort rpc_port ) {
      23           3 :   fd_rpc_client_t * rpc = (fd_rpc_client_t *)mem;
      24           3 :   rpc->request_id = 0UL;
      25           3 :   rpc->rpc_addr = rpc_addr;
      26           3 :   rpc->rpc_port = rpc_port;
      27         387 :   for( ulong i=0; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
      28         384 :     rpc->requests[ i ].state = FD_RPC_CLIENT_STATE_NONE;
      29         384 :     rpc->fds[ i ].fd = -1;
      30         384 :     rpc->fds[ i ].events = POLLIN | POLLOUT;
      31         384 :     rpc->fds[ i ].revents = 0;
      32         384 :   }
      33           3 :   return (void *)rpc;
      34           3 : }
      35             : 
      36             : long
      37             : fd_rpc_client_wait_ready( fd_rpc_client_t * rpc,
      38           0 :                           long              timeout_ns ) {
      39             : 
      40             : 
      41           0 :   struct sockaddr_in addr = {
      42           0 :     .sin_family = AF_INET,
      43           0 :     .sin_port   = fd_ushort_bswap( rpc->rpc_port ),
      44           0 :     .sin_addr   = { .s_addr = rpc->rpc_addr }
      45           0 :   };
      46             : 
      47           0 :   struct pollfd pfd = {
      48           0 :     .fd = 0,
      49           0 :     .events = POLLOUT,
      50           0 :     .revents = 0
      51           0 :   };
      52             : 
      53           0 :   long start = fd_log_wallclock();
      54           0 :   for(;;) {
      55           0 :     pfd.fd = socket( AF_INET, SOCK_STREAM | SOCK_NONBLOCK, 0 );
      56           0 :     if( FD_UNLIKELY( pfd.fd<0 ) ) return FD_RPC_CLIENT_ERR_NETWORK;
      57             : 
      58           0 :     if( FD_UNLIKELY( -1==connect( pfd.fd, fd_type_pun( &addr ), sizeof(addr) ) && errno!=EINPROGRESS ) ) {
      59           0 :       if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
      60           0 :       return FD_RPC_CLIENT_ERR_NETWORK;
      61           0 :     }
      62             : 
      63           0 :     for(;;) {
      64           0 :       long now = fd_log_wallclock();
      65           0 :       if( FD_UNLIKELY( now-start>=timeout_ns ) ) {
      66           0 :         if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
      67           0 :         return FD_RPC_CLIENT_ERR_NETWORK;
      68           0 :       }
      69             : 
      70           0 :       int nfds = poll( &pfd, 1, (int)((now-start) / 1000000) );
      71           0 :       if( FD_UNLIKELY( 0==nfds ) ) continue;
      72           0 :       else if( FD_UNLIKELY( -1==nfds && errno==EINTR ) ) continue;
      73           0 :       else if( FD_UNLIKELY( -1==nfds ) ) {
      74           0 :         if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
      75           0 :         return FD_RPC_CLIENT_ERR_NETWORK;
      76           0 :       } else if( FD_LIKELY( pfd.revents & (POLLERR | POLLHUP) ) ) {
      77           0 :         if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
      78           0 :         break;
      79           0 :       } else if( FD_LIKELY( pfd.revents & POLLOUT ) ) {
      80           0 :         if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
      81           0 :         return FD_RPC_CLIENT_SUCCESS;
      82           0 :       }
      83           0 :     }
      84           0 :   }
      85           0 : }
      86             : 
      87             : static ulong
      88           9 : fd_rpc_available_slot( fd_rpc_client_t * rpc ) {
      89           9 :   for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
      90           9 :     if( FD_LIKELY( rpc->requests[i].state==FD_RPC_CLIENT_STATE_NONE ) ) return i;
      91           9 :   }
      92           0 :   return ULONG_MAX;
      93           9 : }
      94             : 
      95             : static ulong
      96             : fd_rpc_find_request( fd_rpc_client_t * rpc,
      97          18 :                      long              request_id ) {
      98          18 :   for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
      99          18 :     if( FD_LIKELY( rpc->requests[i].state==FD_RPC_CLIENT_STATE_NONE ) ) continue;
     100          18 :     if( FD_LIKELY( rpc->requests[i].response.request_id!=request_id ) ) continue;
     101          18 :     return i;
     102          18 :   }
     103           0 :   return ULONG_MAX;
     104          18 : }
     105             : 
     106             : static long
     107             : fd_rpc_client_request( fd_rpc_client_t * rpc,
     108             :                        ulong             method,
     109             :                        long              request_id,
     110             :                        char *            contents,
     111           9 :                        int               contents_len ) {
     112           9 :   ulong idx = fd_rpc_available_slot( rpc );
     113           9 :   if( FD_UNLIKELY( idx==ULONG_MAX) ) return FD_RPC_CLIENT_ERR_TOO_MANY;
     114             : 
     115           9 :   struct fd_rpc_client_request * request = &rpc->requests[ idx ];
     116             : 
     117           9 :   if( FD_UNLIKELY( contents_len<0 ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
     118           9 :   if( FD_UNLIKELY( (ulong)contents_len>=MAX_REQUEST_LEN ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
     119             : 
     120           9 :   int printed = snprintf( request->connected.request_bytes, sizeof(request->connected.request_bytes),
     121           9 :                           "POST / HTTP/1.1\r\n"
     122           9 :                           "Host: localhost:12001\r\n"
     123           9 :                           "Content-Length: %d\r\n"
     124           9 :                           "Content-Type: application/json\r\n\r\n"
     125           9 :                           "%s", contents_len, contents );
     126           9 :   if( FD_UNLIKELY( printed<0 ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
     127           9 :   if( FD_UNLIKELY( (ulong)printed>=sizeof(request->connected.request_bytes) ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
     128           9 :   request->connected.request_bytes_cnt = (ulong)printed;
     129             : 
     130           9 :   int fd = socket( AF_INET, SOCK_STREAM | SOCK_NONBLOCK, 0 );
     131           9 :   if( FD_UNLIKELY( fd<0 ) ) return FD_RPC_CLIENT_ERR_NETWORK;
     132             : 
     133           9 :   struct sockaddr_in addr = {
     134           9 :     .sin_family = AF_INET,
     135           9 :     .sin_port   = fd_ushort_bswap( rpc->rpc_port ),
     136           9 :     .sin_addr   = { .s_addr = rpc->rpc_addr }
     137           9 :   };
     138             : 
     139           9 :   if( FD_UNLIKELY( -1==connect( fd, fd_type_pun( &addr ), sizeof(addr) ) && errno!=EINPROGRESS ) ) {
     140           0 :     if( FD_UNLIKELY( close( fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
     141           0 :     return FD_RPC_CLIENT_ERR_NETWORK;
     142           0 :   }
     143             : 
     144           9 :   rpc->request_id = request_id;
     145           9 :   rpc->fds[ idx ].fd = fd;
     146           9 :   request->response.method = method;
     147           9 :   request->response.status = FD_RPC_CLIENT_PENDING;
     148           9 :   request->response.request_id = rpc->request_id;
     149           9 :   request->connected.request_bytes_sent = 0UL;
     150           9 :   request->state = FD_RPC_CLIENT_STATE_CONNECTED;
     151           9 :   return request->response.request_id;
     152           9 : }
     153             : 
     154             : long
     155           3 : fd_rpc_client_request_latest_block_hash( fd_rpc_client_t * rpc ) {
     156           3 :   char contents[ MAX_REQUEST_LEN ];
     157           3 :   long request_id = fd_long_if( rpc->request_id==LONG_MAX, 0L, rpc->request_id+1L );
     158             : 
     159           3 :   int contents_len = snprintf( contents, sizeof(contents),
     160           3 :                                "{\"jsonrpc\":\"2.0\",\"id\":%ld,\"method\":\"getLatestBlockhash\",\"params\":[ { \"commitment\": \"processed\" }]}",
     161           3 :                                request_id );
     162             : 
     163           3 :   return fd_rpc_client_request( rpc, FD_RPC_CLIENT_METHOD_LATEST_BLOCK_HASH, request_id, contents, contents_len );
     164           3 : }
     165             : 
     166             : long
     167           6 : fd_rpc_client_request_transaction_count( fd_rpc_client_t * rpc ) {
     168           6 :   char contents[ MAX_REQUEST_LEN ];
     169           6 :   long request_id = fd_long_if( rpc->request_id==LONG_MAX, 0L, rpc->request_id+1L );
     170             : 
     171           6 :   int contents_len = snprintf( contents, sizeof(contents),
     172           6 :                                "{\"jsonrpc\":\"2.0\",\"id\":%ld,\"method\":\"getTransactionCount\",\"params\":[ { \"commitment\": \"processed\" } ]}",
     173           6 :                                request_id );
     174             : 
     175           6 :   return fd_rpc_client_request( rpc, FD_RPC_CLIENT_METHOD_TRANSACTION_COUNT, request_id, contents, contents_len );
     176           6 : }
     177             : 
     178             : static void
     179             : fd_rpc_mark_error( fd_rpc_client_t * rpc,
     180             :                    ulong             idx,
     181           0 :                    long              error ) {
     182           0 :   if( FD_LIKELY( rpc->fds[ idx ].fd>=0 ) ) {
     183           0 :     if( FD_UNLIKELY( close( rpc->fds[ idx ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
     184           0 :     rpc->fds[ idx ].fd = -1;
     185           0 :   }
     186           0 :   rpc->requests[ idx ].state = FD_RPC_CLIENT_STATE_FINISHED;
     187           0 :   rpc->requests[ idx ].response.status = error;
     188           0 : }
     189             : 
     190             : static ulong
     191             : fd_rpc_phr_content_length( struct phr_header * headers,
     192          12 :                            ulong               num_headers ) {
     193          24 :   for( ulong i=0UL; i<num_headers; i++ ) {
     194          24 :     if( FD_LIKELY( headers[i].name_len!=14UL ) ) continue;
     195          12 :     if( FD_LIKELY( strncasecmp( headers[i].name, "Content-Length", 14UL ) ) ) continue;
     196          12 :     ulong content_length = 0UL;
     197          12 :     if( FD_UNLIKELY( fd_http_parse_content_len( headers[i].value, (ulong)headers[i].value_len, &content_length ) ) ) return ULONG_MAX;
     198          12 :     if( FD_UNLIKELY( content_length>UINT_MAX ) ) return ULONG_MAX; /* prevent overflow */
     199          12 :     return content_length;
     200          12 :   }
     201           0 :   return ULONG_MAX;
     202          12 : }
     203             : 
     204             : static long
     205             : parse_response( char *                     response,
     206             :                 ulong                      response_len,
     207          12 :                 fd_rpc_client_response_t * result ) {
     208          12 :   int minor_version;
     209          12 :   int status;
     210          12 :   const char * message;
     211          12 :   ulong message_len;
     212          12 :   struct phr_header headers[ 32 ];
     213          12 :   ulong num_headers = 32UL;
     214             :   /* last_len 0: a non-zero value after the headers completed makes
     215             :      picohttpparser scan the body for CRLFCRLF and never finish */
     216          12 :   int http_len = phr_parse_response( response, response_len,
     217          12 :                                      &minor_version, &status, &message, &message_len,
     218          12 :                                      headers, &num_headers, 0UL );
     219          12 :   if( FD_UNLIKELY( -2==http_len ) ) return FD_RPC_CLIENT_PENDING;
     220          12 :   else if( FD_UNLIKELY( -1==http_len ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     221             : 
     222          12 :   if( FD_UNLIKELY( status!=200 ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     223             : 
     224          12 :   ulong content_length = fd_rpc_phr_content_length( headers, num_headers );
     225          12 :   if( FD_UNLIKELY( content_length==ULONG_MAX ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     226          12 :   if( FD_UNLIKELY( content_length+(ulong)http_len > MAX_REQUEST_LEN ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
     227          12 :   if( FD_LIKELY( content_length+(ulong)http_len>response_len ) ) return FD_RPC_CLIENT_PENDING;
     228             : 
     229           9 :   fd_jtok_t j[1];
     230           9 :   fd_jtok_init( j, response+http_len, content_length );
     231           9 :   fd_jtok_str_t key;
     232           9 :   fd_jtok_obj_enter( j );
     233             : 
     234           9 :   switch( result->method ) {
     235           6 :     case FD_RPC_CLIENT_METHOD_TRANSACTION_COUNT: {
     236           6 :       ulong transaction_count = ULONG_MAX;
     237          24 :       while( fd_jtok_obj_next( j, &key ) ) {
     238          18 :         if( fd_jtok_str_eq( &key, "result" ) ) fd_jtok_ulong( j, &transaction_count );
     239          18 :       }
     240           6 :       if( FD_UNLIKELY( fd_jtok_fini( j ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     241           6 :       if( FD_UNLIKELY( transaction_count==ULONG_MAX ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     242             : 
     243           6 :       result->result.transaction_count.transaction_count = transaction_count;
     244           6 :       return FD_RPC_CLIENT_SUCCESS;
     245           6 :     }
     246           3 :     case FD_RPC_CLIENT_METHOD_LATEST_BLOCK_HASH: {
     247           3 :       char blockhash[ 45 ] = {0};
     248          12 :       while( fd_jtok_obj_next( j, &key ) ) {
     249           9 :         if( !fd_jtok_str_eq( &key, "result" ) ) continue;
     250           3 :         fd_jtok_obj_enter( j );
     251           6 :         while( fd_jtok_obj_next( j, &key ) ) {
     252           3 :           if( !fd_jtok_str_eq( &key, "value" ) ) continue;
     253           3 :           fd_jtok_obj_enter( j );
     254           6 :           while( fd_jtok_obj_next( j, &key ) ) {
     255           3 :             if( fd_jtok_str_eq( &key, "blockhash" ) ) fd_jtok_cstr( j, blockhash, sizeof(blockhash) );
     256           3 :           }
     257           3 :         }
     258           3 :       }
     259           3 :       if( FD_UNLIKELY( fd_jtok_fini( j ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     260           3 :       if( FD_UNLIKELY( !fd_base58_decode_32( blockhash, result->result.latest_block_hash.block_hash ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
     261           3 :       return FD_RPC_CLIENT_SUCCESS;
     262           3 :     }
     263           0 :     default:
     264           0 :       FD_TEST( 0 );
     265           9 :   }
     266           0 :   return FD_RPC_CLIENT_ERR_MALFORMED;
     267           9 : }
     268             : 
     269             : int
     270             : fd_rpc_client_service( fd_rpc_client_t * rpc,
     271      184851 :                        int               wait ) {
     272      184851 :   int timeout = wait ? -1 : 0;
     273      184851 :   int nfds = poll( rpc->fds, FD_RPC_CLIENT_REQUEST_CNT, timeout );
     274      184851 :   if( FD_UNLIKELY( 0==nfds ) ) return 0;
     275      184851 :   else if( FD_UNLIKELY( -1==nfds && errno==EINTR ) ) return 0;
     276      184851 :   else if( FD_UNLIKELY( -1==nfds ) ) FD_LOG_ERR(( "poll failed (%i-%s)", errno, strerror( errno ) ));
     277             : 
     278    23845779 :   for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
     279    23660928 :     struct fd_rpc_client_request * request = &rpc->requests[i];
     280             : 
     281    23660928 :     if( FD_LIKELY( request->state==FD_RPC_CLIENT_STATE_CONNECTED && ( rpc->fds[ i ].revents & POLLOUT ) ) ) {
     282           9 :       long sent = send( rpc->fds[ i ].fd, request->connected.request_bytes+request->connected.request_bytes_sent,
     283           9 :                         request->connected.request_bytes_cnt-request->connected.request_bytes_sent, MSG_NOSIGNAL );
     284           9 :       if( FD_UNLIKELY( -1==sent && errno==EAGAIN ) ) continue;
     285           9 :       if( FD_UNLIKELY( -1==sent ) ) {
     286           0 :         fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
     287           0 :         continue;
     288           0 :       }
     289             : 
     290           9 :       request->connected.request_bytes_sent += (ulong)sent;
     291           9 :       if( FD_UNLIKELY( request->connected.request_bytes_sent==request->connected.request_bytes_cnt ) ) {
     292           9 :         request->sent.response_bytes_read = 0UL;
     293           9 :         request->state = FD_RPC_CLIENT_STATE_SENT;
     294           9 :       }
     295           9 :     }
     296             : 
     297    23660928 :     if( FD_LIKELY( request->state==FD_RPC_CLIENT_STATE_SENT && ( rpc->fds[ i ].revents & POLLIN ) ) ) {
     298          12 :       long read = recv( rpc->fds[ i ].fd, request->response_bytes+request->sent.response_bytes_read,
     299          12 :                         sizeof(request->response_bytes)-request->sent.response_bytes_read, 0 );
     300          12 :       if( FD_UNLIKELY( -1==read && errno==EAGAIN ) ) continue;
     301          12 :       else if( FD_UNLIKELY( -1==read ) ) {
     302           0 :         fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
     303           0 :         continue;
     304          12 :       } else if( FD_UNLIKELY( !read ) ) {
     305           0 :         fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
     306           0 :         continue;
     307           0 :       }
     308             : 
     309          12 :       request->sent.response_bytes_read += (ulong)read;
     310          12 :       if( FD_UNLIKELY( request->sent.response_bytes_read==sizeof(request->response_bytes) ) ) {
     311           0 :         fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_TOO_LARGE );
     312           0 :         continue;
     313           0 :       }
     314             : 
     315          12 :       fd_rpc_client_response_t * response = &request->response;
     316          12 :       long status = parse_response( request->response_bytes,
     317          12 :                                     request->sent.response_bytes_read,
     318          12 :                                     response );
     319          12 :       if( FD_LIKELY( status==FD_RPC_CLIENT_PENDING ) ) continue;
     320           9 :       else if( FD_UNLIKELY( status==FD_RPC_CLIENT_SUCCESS ) ) {
     321           9 :         if( FD_UNLIKELY( close( rpc->fds[ i ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
     322           9 :         rpc->fds[ i ].fd = -1;
     323           9 :         response->status = FD_RPC_CLIENT_SUCCESS;
     324           9 :         request->state = FD_RPC_CLIENT_STATE_FINISHED;
     325           9 :         continue;
     326           9 :       } else {
     327           0 :         fd_rpc_mark_error( rpc, i, status );
     328           0 :         continue;
     329           0 :       }
     330          12 :     }
     331    23660928 :   }
     332             : 
     333      184851 :   return 1;
     334      184851 : }
     335             : 
     336             : fd_rpc_client_response_t *
     337             : fd_rpc_client_status( fd_rpc_client_t * rpc,
     338             :                       long              request_id,
     339           9 :                       int               wait ) {
     340           9 :   ulong idx = fd_rpc_find_request( rpc, request_id );
     341           9 :   if( FD_UNLIKELY( idx==ULONG_MAX ) ) return NULL;
     342             : 
     343           9 :   if( FD_LIKELY( !wait ) ) return &rpc->requests[ idx ].response;
     344             : 
     345      184860 :   for(;;) {
     346      184860 :     if( FD_LIKELY( rpc->requests[ idx ].state==FD_RPC_CLIENT_STATE_FINISHED ) ) return &rpc->requests[ idx ].response;
     347      184851 :     fd_rpc_client_service( rpc, 1 );
     348      184851 :   }
     349           9 : }
     350             : 
     351             : void
     352             : fd_rpc_client_close( fd_rpc_client_t * rpc,
     353           9 :                      long              request_id ) {
     354           9 :   ulong idx = fd_rpc_find_request( rpc, request_id );
     355           9 :   if( FD_UNLIKELY( idx==ULONG_MAX ) ) return;
     356             : 
     357           9 :   if( FD_LIKELY( rpc->fds[ idx ].fd>=0 ) ) {
     358           0 :     if( FD_UNLIKELY( close( rpc->fds[ idx ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
     359           0 :     rpc->fds[ idx ].fd = -1;
     360           0 :   }
     361           9 :   rpc->requests[ idx ].state = FD_RPC_CLIENT_STATE_NONE;
     362           9 : }

Generated by: LCOV version 1.14