LCOV - code coverage report
Current view: top level - waltz/grpc - fd_grpc_client.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 340 724 47.0 %
Date: 2026-09-17 04:28:31 Functions: 26 48 54.2 %

          Line data    Source code
       1             : #include "fd_grpc_client.h"
       2             : #include "fd_grpc_client_private.h"
       3             : #include "../../third_party/nanopb/pb_encode.h" /* pb_msgdesc_t */
       4             : #include <sys/socket.h>
       5             : #include <poll.h>
       6             : #include "../h2/fd_h2_rbuf_sock.h"
       7             : #include "../tlsrec/fd_tlsrec.h"
       8             : #include "fd_grpc_codec.h"
       9             : 
      10             : static int
      11             : fd_grpc_client_request_continue( fd_grpc_client_t * client );
      12             : 
      13             : ulong
      14         669 : fd_grpc_client_align( void ) {
      15         669 :   return fd_ulong_max( alignof(fd_grpc_client_t), fd_grpc_h2_stream_pool_align() );
      16         669 : }
      17             : 
      18             : ulong
      19         168 : fd_grpc_client_footprint( ulong buf_max ) {
      20         168 :   ulong l = FD_LAYOUT_INIT;
      21         168 :   l = FD_LAYOUT_APPEND( l, alignof(fd_grpc_client_t), sizeof(fd_grpc_client_t) );
      22         168 :   l = FD_LAYOUT_APPEND( l, 1UL, buf_max ); /* nanopb_tx */
      23         168 :   l = FD_LAYOUT_APPEND( l, 1UL, buf_max ); /* frame_scratch */
      24         168 :   l = FD_LAYOUT_APPEND( l, 1UL, buf_max ); /* frame_rx_buf */
      25         168 :   l = FD_LAYOUT_APPEND( l, 1UL, buf_max ); /* frame_tx_buf */
      26         168 :   l = FD_LAYOUT_APPEND( l, fd_grpc_h2_stream_pool_align(), fd_grpc_h2_stream_pool_footprint( FD_GRPC_CLIENT_MAX_STREAMS ) );
      27         168 :   l = FD_LAYOUT_APPEND( l, 1UL, buf_max*FD_GRPC_CLIENT_MAX_STREAMS );
      28         168 :   return FD_LAYOUT_FINI( l, fd_grpc_client_align() );
      29         168 : }
      30             : 
      31             : static void
      32         102 : fd_grpc_h2_stream_reset( fd_grpc_h2_stream_t * stream ) {
      33         102 :   memset( &stream->s, 0, sizeof(fd_h2_stream_t) );
      34         102 :   stream->request_ctx = 0UL;
      35         102 :   memset( &stream->hdrs, 0, sizeof(fd_grpc_resp_hdrs_t) );
      36         102 :   stream->hdrs.grpc_status    = FD_GRPC_STATUS_UNKNOWN;
      37         102 :   stream->hdrs_received       = 0;
      38         102 :   stream->msg_buf_used        = 0UL;
      39         102 :   stream->msg_sz              = 0UL;
      40         102 :   stream->has_header_deadline = 0;
      41         102 :   stream->has_rx_end_deadline = 0;
      42         102 :   stream->tx_wnd_debt         = 0L;
      43         102 : }
      44             : 
      45             : fd_grpc_client_t *
      46             : fd_grpc_client_new( void *                             mem,
      47             :                     fd_grpc_client_callbacks_t const * callbacks,
      48             :                     fd_grpc_client_metrics_t *         metrics,
      49             :                     void *                             app_ctx,
      50             :                     ulong                              buf_max,
      51          84 :                     ulong                              rng_seed ) {
      52          84 :   if( FD_UNLIKELY( !mem ) ) {
      53           0 :     FD_LOG_WARNING(( "NULL mem" ));
      54           0 :     return NULL;
      55           0 :   }
      56          84 :   if( FD_UNLIKELY( buf_max<4096UL ) ) {
      57           0 :     FD_LOG_WARNING(( "undersz buf_max" ));
      58           0 :     return NULL;
      59           0 :   }
      60          84 :   if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)mem, fd_grpc_client_align() ) ) ) {
      61           0 :     FD_LOG_WARNING(( "unaligned mem" ));
      62           0 :     return NULL;
      63           0 :   }
      64             : 
      65          84 :   FD_SCRATCH_ALLOC_INIT( l, mem );
      66          84 :   void * client_mem      = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_grpc_client_t), sizeof(fd_grpc_client_t) );
      67          84 :   void * nanopb_tx       = FD_SCRATCH_ALLOC_APPEND( l, 1UL, buf_max ); /* nanopb_tx */
      68          84 :   void * frame_scratch   = FD_SCRATCH_ALLOC_APPEND( l, 1UL, buf_max ); /* frame_scratch */
      69          84 :   void * frame_rx_buf    = FD_SCRATCH_ALLOC_APPEND( l, 1UL, buf_max ); /* frame_rx_buf */
      70          84 :   void * frame_tx_buf    = FD_SCRATCH_ALLOC_APPEND( l, 1UL, buf_max ); /* frame_tx_buf */
      71          84 :   void * stream_pool_mem = FD_SCRATCH_ALLOC_APPEND( l, fd_grpc_h2_stream_pool_align(), fd_grpc_h2_stream_pool_footprint( FD_GRPC_CLIENT_MAX_STREAMS ) );
      72          84 :   void * stream_buf_mem  = FD_SCRATCH_ALLOC_APPEND( l, 1UL, buf_max*FD_GRPC_CLIENT_MAX_STREAMS );
      73          84 :   ulong end = FD_SCRATCH_ALLOC_FINI( l, fd_grpc_client_align() );
      74          84 :   FD_TEST( end-(ulong)mem == fd_grpc_client_footprint( buf_max ) );
      75             : 
      76          84 :   fd_grpc_client_t * client = client_mem;
      77             : 
      78          84 :   fd_grpc_h2_stream_t * stream_pool =
      79          84 :     fd_grpc_h2_stream_pool_join( fd_grpc_h2_stream_pool_new( stream_pool_mem, FD_GRPC_CLIENT_MAX_STREAMS ) );
      80          84 :   if( FD_UNLIKELY( !stream_pool ) ) FD_LOG_CRIT(( "Failed to create stream pool" )); /* unreachable */
      81             : 
      82          84 :   *client = (fd_grpc_client_t){
      83          84 :     .callbacks         = callbacks,
      84          84 :     .ctx               = app_ctx,
      85          84 :     .stream_pool       = stream_pool,
      86          84 :     .stream_bufs       = stream_buf_mem,
      87          84 :     .nanopb_tx         = nanopb_tx,
      88          84 :     .nanopb_tx_max     = buf_max,
      89          84 :     .frame_scratch     = frame_scratch,
      90          84 :     .frame_scratch_max = buf_max,
      91          84 :     .frame_rx_buf      = frame_rx_buf,
      92          84 :     .frame_rx_buf_max  = buf_max,
      93          84 :     .frame_tx_buf      = frame_tx_buf,
      94          84 :     .frame_tx_buf_max  = buf_max,
      95          84 :     .metrics           = metrics
      96          84 :   };
      97             : 
      98             :   /* FIXME for performance, cache this? */
      99          84 :   fd_h2_hdr_matcher_init( client->matcher, rng_seed );
     100          84 :   fd_h2_hdr_matcher_insert_literal( client->matcher, FD_GRPC_HDR_STATUS,  "grpc-status"  );
     101          84 :   fd_h2_hdr_matcher_insert_literal( client->matcher, FD_GRPC_HDR_MESSAGE, "grpc-message" );
     102             : 
     103          84 :   client->version_len = 5;
     104          84 :   memcpy( client->version, "0.0.0", 5 );
     105             : 
     106         756 :   for( ulong i=0UL; i<FD_GRPC_CLIENT_MAX_STREAMS; i++ ) {
     107         672 :     fd_grpc_h2_stream_t * stream = &client->stream_pool[ i ];
     108         672 :     stream->msg_buf     = (uchar *)stream_buf_mem + (i*buf_max);
     109         672 :     stream->msg_buf_max = buf_max;
     110         672 :     FD_TEST( (ulong)( stream->msg_buf + stream->msg_buf_max )<=end );
     111         672 :   }
     112          84 :   fd_grpc_client_reset( client );
     113             : 
     114          84 :   return client;
     115          84 : }
     116             : 
     117             : void *
     118           3 : fd_grpc_client_delete( fd_grpc_client_t * client ) {
     119           3 :   return client;
     120           3 : }
     121             : 
     122             : void
     123             : fd_grpc_client_set_version( fd_grpc_client_t * client,
     124             :                             char const *       version,
     125           0 :                             ulong              version_len ) {
     126           0 :   if( FD_UNLIKELY( version_len > FD_GRPC_CLIENT_VERSION_LEN_MAX ) ) {
     127           0 :     FD_LOG_WARNING(( "Version string too long (%lu chars), ignoring", version_len ));
     128           0 :     return;
     129           0 :   }
     130           0 :   client->version_len = (uchar)version_len;
     131           0 :   memcpy( client->version, version, version_len );
     132           0 : }
     133             : 
     134             : void
     135             : fd_grpc_client_set_authority( fd_grpc_client_t * client,
     136             :                               char const *       host,
     137             :                               ulong              host_len,
     138           0 :                               ushort             port ) {
     139           0 :   host_len = fd_ulong_min( host_len, sizeof(client->host)-1 );
     140           0 :   fd_cstr_fini( fd_cstr_append_text( fd_cstr_init( client->host ), host, host_len ) );
     141           0 :   client->host_len = (uchar)host_len;
     142           0 :   client->port     = (ushort)port;
     143           0 : }
     144             : 
     145             : int
     146         114 : fd_grpc_client_stream_acquire_is_safe( fd_grpc_client_t * client ) {
     147             :   /* Sufficient quota to start a stream? */
     148         114 :   if( FD_UNLIKELY( client->conn->stream_active_cnt[1]+1 > client->conn->peer_settings.max_concurrent_streams ) ) {
     149          18 :     return 0;
     150          18 :   }
     151             : 
     152             :   /* Free stream object available? */
     153          96 :   if( FD_UNLIKELY( !fd_grpc_h2_stream_pool_free( client->stream_pool ) ) ) {
     154           0 :     return 0;
     155           0 :   }
     156          96 :   if( FD_UNLIKELY( client->stream_cnt >= FD_GRPC_CLIENT_MAX_STREAMS ) ) {
     157           0 :     return 0;
     158           0 :   }
     159             : 
     160          96 :   return 1;
     161          96 : }
     162             : 
     163             : fd_grpc_h2_stream_t *
     164             : fd_grpc_client_stream_acquire( fd_grpc_client_t * client,
     165         102 :                                ulong              request_ctx ) {
     166         102 :   if( FD_UNLIKELY( client->stream_cnt >= FD_GRPC_CLIENT_MAX_STREAMS ) ) {
     167           0 :     FD_LOG_CRIT(( "stream pool exhausted" ));
     168           0 :   }
     169             : 
     170         102 :   fd_h2_conn_t * conn = client->conn;
     171         102 :   uint const stream_id = client->conn->tx_stream_next;
     172         102 :   conn->tx_stream_next += 2U;
     173             : 
     174         102 :   fd_grpc_h2_stream_t * stream = fd_grpc_h2_stream_pool_ele_acquire( client->stream_pool );
     175         102 :   fd_grpc_h2_stream_reset( stream );
     176         102 :   stream->request_ctx = request_ctx;
     177             : 
     178         102 :   fd_h2_stream_open( fd_h2_stream_init( &stream->s ), conn, stream_id );
     179         102 :   client->request_stream = stream;
     180         102 :   client->stream_ids[ client->stream_cnt ] = stream_id;
     181         102 :   client->streams   [ client->stream_cnt ] = stream;
     182         102 :   client->stream_cnt++;
     183             : 
     184         102 :   return stream;
     185         102 : }
     186             : 
     187             : void
     188             : fd_grpc_client_stream_release( fd_grpc_client_t *    client,
     189          27 :                                fd_grpc_h2_stream_t * stream ) {
     190          27 :   if( FD_UNLIKELY( !client->stream_cnt ) ) FD_LOG_CRIT(( "stream map corrupt" )); /* unreachable */
     191             : 
     192             :   /* Deallocate tx_op */
     193          27 :   if( FD_UNLIKELY( stream == client->request_stream ) ) {
     194          15 :     client->request_stream = NULL;
     195          15 :     *client->request_tx_op = (fd_h2_tx_op_t){0};
     196          15 :   }
     197             : 
     198             :   /* Remove stream from map */
     199          27 :   int map_idx = -1;
     200          84 :   for( uint i=0UL; i<(client->stream_cnt); i++ ) {
     201          57 :     if( client->stream_ids[ i ] == stream->s.stream_id ) {
     202          27 :       map_idx = (int)i;
     203          27 :     }
     204          57 :   }
     205          27 :   if( FD_UNLIKELY( map_idx<0 ) ) FD_LOG_CRIT(( "stream map corrupt" )); /* unreachable */
     206          27 :   if( (ulong)map_idx+1 < client->stream_cnt ) {
     207           0 :     client->stream_ids[ map_idx ] = client->stream_ids[ client->stream_cnt-1 ];
     208           0 :     client->streams   [ map_idx ] = client->streams   [ client->stream_cnt-1 ];
     209           0 :   }
     210          27 :   client->stream_cnt--;
     211             : 
     212          27 :   fd_grpc_h2_stream_pool_ele_release( client->stream_pool, stream );
     213          27 : }
     214             : 
     215             : void
     216          90 : fd_grpc_client_reset( fd_grpc_client_t * client ) {
     217          90 :   fd_h2_rbuf_init( client->frame_rx, client->frame_rx_buf, client->frame_rx_buf_max );
     218          90 :   fd_h2_rbuf_init( client->frame_tx, client->frame_tx_buf, client->frame_tx_buf_max );
     219          90 :   fd_h2_conn_init_client( client->conn );
     220          90 :   client->conn->ctx      = client;
     221          90 :   client->h2_hs_done     = 0;
     222          90 :   client->window_update_pending = 0;
     223          90 :   client->request_stream = NULL;
     224          90 :   fd_tlsrec_sock_init( client->tls_sock );
     225          90 :   *client->request_tx_op = (fd_h2_tx_op_t){0};
     226             : 
     227             :   /* Disable RX flow control */
     228          90 :   client->conn->self_settings.initial_window_size = (1U<<31)-1U;
     229          90 :   client->conn->rx_wnd_max   = (1U<<31)-1U;
     230          90 :   client->conn->rx_wnd_wmark = client->conn->rx_wnd_max - (1U<<20);
     231             : 
     232             :   /* Free all stream objects */
     233         102 :   while( client->stream_cnt ) {
     234          12 :     fd_grpc_h2_stream_t * stream = client->streams[ client->stream_cnt-1 ];
     235          12 :     fd_grpc_client_stream_release( client, stream );
     236          12 :   }
     237          90 : }
     238             : 
     239             : /* fd_grpc_client_send_stream_quota writes a WINDOW_UPDATE frame, which
     240             :    eventually allows the peer to send more data bytes. */
     241             : 
     242             : static void
     243             : fd_grpc_client_send_stream_quota( fd_h2_rbuf_t *        rbuf_tx,
     244             :                                   fd_grpc_h2_stream_t * stream,
     245           0 :                                   uint                  bump ) {
     246           0 :   fd_h2_window_update_t window_update = {
     247           0 :     .hdr = {
     248           0 :       .typlen      = fd_h2_frame_typlen( FD_H2_FRAME_TYPE_WINDOW_UPDATE, 4UL ),
     249           0 :       .r_stream_id = fd_uint_bswap( stream->s.stream_id )
     250           0 :     },
     251           0 :     .increment = fd_uint_bswap( bump )
     252           0 :   };
     253           0 :   fd_h2_rbuf_push( rbuf_tx, &window_update, sizeof(fd_h2_window_update_t) );
     254           0 :   stream->s.rx_wnd += bump;
     255           0 : }
     256             : 
     257             : /* fd_grpc_client_send_timeout is called when a stream timeout triggers.
     258             :    Calls back to the user, writes a RST_STREAM frame, and frees the
     259             :    stream object. */
     260             : 
     261             : static void
     262             : fd_grpc_client_send_timeout( fd_h2_rbuf_t *        rbuf_tx,
     263             :                              fd_grpc_client_t *    client,
     264             :                              fd_grpc_h2_stream_t * stream,
     265           6 :                              int                   deadline_kind ) {
     266           6 :   client->callbacks->rx_timeout( client->ctx, stream->request_ctx, deadline_kind );
     267           6 :   fd_h2_tx_rst_stream( rbuf_tx, stream->s.stream_id, FD_H2_ERR_CANCEL );
     268           6 :   fd_grpc_client_stream_release( client, stream );
     269           6 : }
     270             : 
     271             : long
     272           0 : fd_grpc_client_next_deadline( fd_grpc_client_t const * client ) {
     273           0 :   long deadline = LONG_MAX;
     274           0 :   for( ulong i=0UL; i<(client->stream_cnt); i++ ) {
     275           0 :     fd_grpc_h2_stream_t const * stream = client->streams[ i ];
     276           0 :     if( stream->has_header_deadline ) deadline = fd_long_min( deadline, stream->header_deadline_nanos );
     277           0 :     if( stream->has_rx_end_deadline ) deadline = fd_long_min( deadline, stream->rx_end_deadline_nanos );
     278           0 :   }
     279           0 :   return deadline;
     280           0 : }
     281             : 
     282             : int
     283          60 : fd_grpc_client_tx_pending( fd_grpc_client_t const * client ) {
     284          60 :   return !!fd_h2_rbuf_used_sz( client->frame_tx );
     285          60 : }
     286             : 
     287             : int
     288           0 : fd_grpc_client_tls_rx_pending( fd_grpc_client_t const * client ) {
     289           0 :   return !!fd_tlsrec_sock_rx_avail( client->tls_sock );
     290           0 : }
     291             : 
     292             : int
     293           0 : fd_grpc_client_tls_tx_pending( fd_grpc_client_t const * client ) {
     294           0 :   return fd_tlsrec_sock_tx_pending( client->tls_sock );
     295           0 : }
     296             : 
     297             : ulong
     298           0 : fd_grpc_client_tx_starved( fd_grpc_client_t const * client ) {
     299           0 :   fd_grpc_h2_stream_t const * stream = client->request_stream;
     300           0 :   if( !stream || fd_uint_min( client->conn->tx_wnd, stream->s.tx_wnd ) ) return 0UL;
     301           0 :   return client->request_tx_op->chunk_sz;
     302           0 : }
     303             : 
     304             : void
     305             : fd_grpc_client_service_streams( fd_grpc_client_t * client,
     306          33 :                                 long               ts_nanos ) {
     307          33 :   ulong const meta_frame_max =
     308          33 :     fd_ulong_max( sizeof(fd_h2_window_update_t), sizeof(fd_h2_rst_stream_t) );
     309          33 :   fd_h2_conn_t * conn    = client->conn;
     310          33 :   fd_h2_rbuf_t * rbuf_tx = client->frame_tx;
     311          33 :   if( FD_UNLIKELY( conn->flags ) ) return;
     312          33 :   uint  const wnd_max    = conn->self_settings.initial_window_size;
     313          33 :   uint  const wnd_thres  = wnd_max / 2;
     314         105 :   for( ulong i=0UL; i<(client->stream_cnt); i++ ) {
     315          72 :     if( FD_UNLIKELY( fd_h2_rbuf_free_sz( rbuf_tx )<meta_frame_max ) ) break;
     316          72 :     fd_grpc_h2_stream_t * stream = client->streams[ i ];
     317             : 
     318          72 :     if( FD_UNLIKELY( ( stream->has_header_deadline ) &
     319          72 :                      ( stream->header_deadline_nanos - ts_nanos <= 0L ) ) ) {
     320           3 :       fd_grpc_client_send_timeout( rbuf_tx, client, stream, FD_GRPC_DEADLINE_HEADER );
     321           3 :       i--; /* stream removed */
     322           3 :       continue;
     323           3 :     }
     324             : 
     325          69 :     if( FD_UNLIKELY( ( stream->has_rx_end_deadline ) &
     326          69 :                      ( stream->rx_end_deadline_nanos - ts_nanos <= 0L ) ) ) {
     327           3 :       fd_grpc_client_send_timeout( rbuf_tx, client, stream, FD_GRPC_DEADLINE_RX_END );
     328           3 :       i--; /* stream removed */
     329           3 :       continue;
     330           3 :     }
     331             : 
     332          66 :     if( FD_UNLIKELY( stream->s.rx_wnd < wnd_thres ) ) {
     333           0 :       uint const bump = wnd_max - stream->s.rx_wnd;
     334           0 :       fd_grpc_client_send_stream_quota( rbuf_tx, stream, bump );
     335           0 :     }
     336          66 :   }
     337          33 : }
     338             : 
     339             : #if FD_H2_HAS_SOCKETS
     340             : 
     341             : int
     342             : fd_grpc_client_tls_flush( fd_grpc_client_t * client,
     343           0 :                           int                sock_fd ) {
     344           0 :   int rc = fd_tlsrec_sock_flush( client->tls_sock, sock_fd );
     345           0 :   return rc<0 ? -1 : rc;
     346           0 : }
     347             : 
     348             : int
     349             : fd_grpc_client_rxtx_socket( fd_grpc_client_t * client,
     350             :                             int                sock_fd,
     351             :                             long               now,
     352          27 :                             int *              charge_busy ) {
     353          27 :   fd_h2_conn_t * conn = client->conn;
     354          27 :   ulong const frame_rx_lo_0 = client->frame_rx->lo_off;
     355          27 :   ulong const frame_rx_hi_0 = client->frame_rx->hi_off;
     356          27 :   ulong const frame_tx_lo_1 = client->frame_tx->lo_off;
     357          27 :   ulong const frame_tx_hi_1 = client->frame_tx->hi_off;
     358             : 
     359          27 :   int rx_err = fd_h2_rbuf_recvmsg( client->frame_rx, sock_fd, MSG_NOSIGNAL|MSG_DONTWAIT );
     360          27 :   if( FD_UNLIKELY( rx_err ) ) {
     361           0 :     FD_LOG_INFO(( "Disconnected: recvmsg error (%i-%s)", rx_err, fd_io_strerror( rx_err ) ));
     362           0 :     errno = rx_err;
     363           0 :     return -1;
     364           0 :   }
     365             : 
     366          27 :   if( FD_UNLIKELY( conn->flags ) ) fd_h2_tx_control( conn, client->frame_tx, &fd_grpc_client_h2_callbacks );
     367          27 :   fd_h2_rx( conn, client->frame_rx, client->frame_tx, client->frame_scratch, client->frame_scratch_max, &fd_grpc_client_h2_callbacks );
     368          27 :   if( FD_UNLIKELY( client->window_update_pending || client->request_stream ) ) {
     369          27 :     client->window_update_pending = 0;
     370          27 :     fd_grpc_client_request_continue( client ); /* credit or TX ring space may have freed */
     371          27 :   }
     372          27 :   fd_grpc_client_service_streams( client, now );
     373             : 
     374          27 :   int tx_err = fd_h2_rbuf_sendmsg( client->frame_tx, sock_fd, MSG_NOSIGNAL|MSG_DONTWAIT );
     375          27 :   if( FD_UNLIKELY( tx_err && tx_err!=EAGAIN ) ) {
     376           0 :     FD_LOG_WARNING(( "fd_h2_rbuf_sendmsg failed (%i-%s)", tx_err, fd_io_strerror( tx_err ) ));
     377           0 :     errno = tx_err;
     378           0 :     return -1;
     379           0 :   }
     380             : 
     381          27 :   ulong const frame_rx_lo_1 = client->frame_rx->lo_off;
     382          27 :   ulong const frame_rx_hi_1 = client->frame_rx->hi_off;
     383          27 :   ulong const frame_tx_lo_0 = client->frame_tx->lo_off;
     384          27 :   ulong const frame_tx_hi_0 = client->frame_tx->hi_off;
     385             : 
     386          27 :   client->metrics->stream_chunks_rx_bytes += frame_rx_hi_1 - frame_rx_hi_0;
     387          27 :   client->metrics->stream_chunks_tx_bytes += frame_tx_lo_0 - frame_tx_lo_1;
     388             : 
     389          27 :   if( frame_rx_lo_0!=frame_rx_lo_1 || frame_rx_hi_0!=frame_rx_hi_1 ||
     390          27 :       frame_tx_lo_0!=frame_tx_lo_1 || frame_tx_hi_0!=frame_tx_hi_1 ) {
     391           3 :     *charge_busy = 1;
     392           3 :   }
     393             : 
     394          27 :   return tx_err==EAGAIN ? 1 : 0;
     395          27 : }
     396             : 
     397             : int
     398             : fd_grpc_client_tx_flush_socket( fd_grpc_client_t * client,
     399          12 :                                 int                sock_fd ) {
     400          12 :   ulong const lo_0 = client->frame_tx->lo_off;
     401          12 :   int tx_err = fd_h2_rbuf_sendmsg( client->frame_tx, sock_fd, MSG_NOSIGNAL|MSG_DONTWAIT );
     402          12 :   if( FD_UNLIKELY( tx_err && tx_err!=EAGAIN ) ) {
     403           0 :     FD_LOG_WARNING(( "fd_h2_rbuf_sendmsg failed (%i-%s)", tx_err, fd_io_strerror( tx_err ) ));
     404           0 :     errno = tx_err;
     405           0 :     return -1;
     406           0 :   }
     407          12 :   client->metrics->stream_chunks_tx_bytes += client->frame_tx->lo_off - lo_0;
     408          12 :   return tx_err==EAGAIN ? 1 : 0;
     409          12 : }
     410             : 
     411             : int
     412             : fd_grpc_client_rxtx_tls( fd_grpc_client_t * client,
     413             :                          fd_tlsrec_conn_t * tls_conn,
     414             :                          int                sock_fd,
     415             :                          long               now,
     416           0 :                          int *              charge_busy ) {
     417           0 :   fd_h2_conn_t *     conn = client->conn;
     418           0 :   fd_tlsrec_sock_t * sock = client->tls_sock;
     419             : 
     420             :   /* Busy means progress was made.  A flush that moves nothing because
     421             :      the peer stopped reading is not progress: the caller would spin on
     422             :      send(2)/EAGAIN instead of waiting for the EPOLLOUT it armed. */
     423           0 :   ulong parked_0 = sock->tx_sz - sock->tx_off;
     424           0 :   if( FD_UNLIKELY( fd_tlsrec_sock_flush( sock, sock_fd )<0 ) ) return -1;
     425           0 :   if( parked_0 != sock->tx_sz - sock->tx_off ) *charge_busy = 1;
     426             : 
     427           0 :   ulong tcp_rx_sz;
     428           0 :   int rx_err = fd_tlsrec_sock_rx( sock, tls_conn, sock_fd, &tcp_rx_sz );
     429           0 :   if( FD_UNLIKELY( rx_err ) ) {
     430           0 :     if( rx_err==FD_TLSREC_SOCK_ERR_RECV || rx_err==FD_TLSREC_SOCK_ERR_SEND )
     431           0 :       FD_LOG_INFO(( "Disconnected: %s (%i-%s)", fd_tlsrec_sock_strerror( rx_err ), errno, fd_io_strerror( errno ) ));
     432           0 :     else
     433           0 :       FD_LOG_INFO(( "Disconnected: %s", fd_tlsrec_sock_strerror( rx_err ) ));
     434           0 :     return -1;
     435           0 :   }
     436           0 :   if( tcp_rx_sz ) *charge_busy = 1;
     437             : 
     438           0 :   if( FD_UNLIKELY( !fd_tlsrec_conn_is_ready( tls_conn ) ) ) {
     439           0 :     if( FD_UNLIKELY( fd_tlsrec_conn_is_failed( tls_conn ) ) ) {
     440           0 :       FD_LOG_WARNING(( "TLS handshake failed" ));
     441           0 :       return -1;
     442           0 :     }
     443           0 :     return 0;
     444           0 :   }
     445             : 
     446           0 :   if( FD_UNLIKELY( !client->h2_hs_done && !tls_conn->hs.cli.alpn_negotiated ) ) {
     447           0 :     FD_LOG_WARNING(( "TLS handshake failed: not a gRPC server (no h2 ALPN)" ));
     448           0 :     return -1;
     449           0 :   }
     450             : 
     451           0 :   ulong push_sz = fd_ulong_min( fd_tlsrec_sock_rx_avail( sock ),
     452           0 :                                 fd_h2_rbuf_free_sz( client->frame_rx ) );
     453           0 :   if( FD_LIKELY( push_sz ) ) {
     454           0 :     fd_h2_rbuf_push( client->frame_rx, fd_tlsrec_sock_rx_data( sock ), push_sz );
     455           0 :     fd_tlsrec_sock_rx_consume( sock, push_sz );
     456           0 :     client->metrics->stream_chunks_rx_bytes += push_sz;
     457           0 :     *charge_busy = 1;
     458           0 :   }
     459             : 
     460           0 :   if( FD_UNLIKELY( conn->flags ) ) fd_h2_tx_control( conn, client->frame_tx, &fd_grpc_client_h2_callbacks );
     461           0 :   fd_h2_rx( conn, client->frame_rx, client->frame_tx, client->frame_scratch, client->frame_scratch_max, &fd_grpc_client_h2_callbacks );
     462           0 :   if( FD_UNLIKELY( client->window_update_pending || client->request_stream ) ) {
     463           0 :     client->window_update_pending = 0;
     464           0 :     fd_grpc_client_request_continue( client ); /* credit or TX ring space may have freed */
     465           0 :   }
     466           0 :   fd_grpc_client_service_streams( client, now );
     467             : 
     468             :   /* HTTP/2 bytes wait behind parked ciphertext until EPOLLOUT */
     469           0 :   ulong tx_used = fd_h2_rbuf_used_sz( client->frame_tx );
     470           0 :   if( FD_LIKELY( tx_used && !fd_tlsrec_sock_tx_pending( sock ) ) ) {
     471           0 :     uchar plaintext[ FD_TLSREC_PLAINTEXT_MAX ];
     472           0 :     ulong pop_sz = fd_ulong_min( tx_used, sizeof(plaintext) );
     473           0 :     fd_h2_rbuf_pop_copy( client->frame_tx, plaintext, pop_sz );
     474             : 
     475           0 :     ulong consumed;
     476           0 :     int tx_err = fd_tlsrec_sock_tx( sock, tls_conn, sock_fd, plaintext, pop_sz, &consumed );
     477           0 :     if( FD_UNLIKELY( tx_err ) ) {
     478           0 :       FD_LOG_WARNING(( "fd_tlsrec_sock_tx failed: %s", fd_tlsrec_sock_strerror( tx_err ) ));
     479           0 :       return -1;
     480           0 :     }
     481           0 :     FD_CHECK_CRIT( consumed==pop_sz, "mismatched buffer sizes" );
     482             : 
     483           0 :     client->metrics->stream_chunks_tx_bytes += pop_sz;
     484           0 :     *charge_busy = 1;
     485           0 :   }
     486           0 :   return 0;
     487           0 : }
     488             : 
     489             : #endif /* FD_H2_HAS_SOCKETS */
     490             : 
     491             : /* fd_grpc_client_request continue attempts to write a request data
     492             :    frame. */
     493             : 
     494             : static int
     495           9 : fd_grpc_client_request_continue1( fd_grpc_client_t * client ) {
     496           9 :   fd_grpc_h2_stream_t * stream    = client->request_stream;
     497           9 :   fd_h2_stream_t *      h2_stream = &stream->s;
     498           9 :   fd_h2_tx_op_copy( client->conn, h2_stream, client->frame_tx, client->request_tx_op );
     499           9 :   if( FD_UNLIKELY( client->request_tx_op->chunk_sz ) ) return 0;
     500           9 :   client->request_stream = NULL;
     501           9 :   if( FD_UNLIKELY( h2_stream->state != FD_H2_STREAM_STATE_CLOSING_TX ) ) return 0;
     502           9 :   client->metrics->stream_chunks_tx_cnt++;
     503             :   /* Request finished */
     504           9 :   client->callbacks->tx_complete( client->ctx, stream->request_ctx );
     505           9 :   return 1;
     506           9 : }
     507             : 
     508             : static int
     509          27 : fd_grpc_client_request_continue( fd_grpc_client_t * client ) {
     510          27 :   if( FD_UNLIKELY( client->conn->flags & (FD_H2_CONN_FLAGS_DEAD|FD_H2_CONN_FLAGS_SEND_GOAWAY) ) ) return 0;
     511          27 :   if( FD_UNLIKELY( !client->request_stream ) ) return 0;
     512          27 :   if( FD_UNLIKELY( !client->request_tx_op->chunk_sz ) ) return 0;
     513           0 :   return fd_grpc_client_request_continue1( client );
     514          27 : }
     515             : 
     516             : int
     517         114 : fd_grpc_client_is_connected( fd_grpc_client_t * client ) {
     518         114 :   if( FD_UNLIKELY( !client                                          ) ) return 0;
     519         114 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD      ) ) return 0;
     520         114 :   if( FD_UNLIKELY( !client->h2_hs_done                              ) ) return 0;
     521         111 :   return 1;
     522         114 : }
     523             : 
     524             : int
     525         114 : fd_grpc_client_request_is_blocked( fd_grpc_client_t * client ) {
     526         114 :   if( FD_UNLIKELY( !client                                          ) ) return 1;
     527         114 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD      ) ) return 1;
     528         114 :   if( FD_UNLIKELY( !client->h2_hs_done                              ) ) return 1;
     529         114 :   if( FD_UNLIKELY( !fd_h2_rbuf_is_empty( client->frame_tx )         ) ) return 1;
     530         114 :   if( FD_UNLIKELY( client->request_tx_op->chunk_sz > 0UL            ) ) return 1;
     531         114 :   if( FD_UNLIKELY( !fd_grpc_client_stream_acquire_is_safe( client ) ) ) return 1;
     532          96 :   return 0;
     533         114 : }
     534             : 
     535             : int
     536           0 : fd_grpc_client_stream_send_is_blocked( fd_grpc_client_t * client ) {
     537           0 :   if( FD_UNLIKELY( !client                                     ) ) return 1;
     538           0 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD ) ) return 1;
     539           0 :   if( FD_UNLIKELY( !client->h2_hs_done                         ) ) return 1;
     540           0 :   if( FD_UNLIKELY( !fd_h2_rbuf_is_empty( client->frame_tx )    ) ) return 1;
     541           0 :   if( FD_UNLIKELY( client->request_tx_op->chunk_sz > 0UL       ) ) return 1;
     542           0 :   return 0;
     543           0 : }
     544             : 
     545             : int
     546           0 : fd_grpc_client_request_stream_busy( fd_grpc_client_t * client ) {
     547           0 :   return client->request_stream && client->request_tx_op->chunk_sz;
     548           0 : }
     549             : 
     550             : fd_grpc_h2_stream_t *
     551             : fd_grpc_client_request_start(
     552             :     fd_grpc_client_t *   client,
     553             :     char const *         path,
     554             :     ulong                path_len,
     555             :     ulong                request_ctx,
     556             :     pb_msgdesc_t const * fields,
     557             :     void const *         message,
     558             :     char const *         auth_token,
     559             :     ulong                auth_token_sz,
     560             :     int                  is_streaming
     561           9 : ) {
     562           9 :   if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( client ) ) ) return NULL;
     563             : 
     564             :   /* Encode message */
     565           9 :   FD_TEST( client->nanopb_tx_max > sizeof(fd_grpc_hdr_t) );
     566           9 :   uchar * proto_buf = client->nanopb_tx + sizeof(fd_grpc_hdr_t);
     567           9 :   pb_ostream_t ostream = pb_ostream_from_buffer( proto_buf, client->nanopb_tx_max - sizeof(fd_grpc_hdr_t) );
     568           9 :   if( FD_UNLIKELY( !pb_encode( &ostream, fields, message ) ) ) {
     569           0 :     FD_LOG_WARNING(( "Failed to encode Protobuf message (%.*s). This is a bug (insufficient buffer space?)", (int)path_len, path ));
     570           0 :     return NULL;
     571           0 :   }
     572           9 :   ulong const serialized_sz = ostream.bytes_written;
     573             : 
     574             :   /* Create gRPC length prefix */
     575           9 :   fd_grpc_hdr_t hdr = {
     576           9 :     .compressed=0,
     577           9 :     .msg_sz=fd_uint_bswap( (uint)serialized_sz )
     578           9 :   };
     579           9 :   memcpy( client->nanopb_tx, &hdr, sizeof(fd_grpc_hdr_t) );
     580           9 :   ulong const payload_sz = serialized_sz + sizeof(fd_grpc_hdr_t);
     581             : 
     582             :   /* Allocate stream descriptor */
     583           9 :   fd_grpc_h2_stream_t * stream    = fd_grpc_client_stream_acquire( client, request_ctx );
     584           9 :   uint const            stream_id = stream->s.stream_id;
     585             : 
     586             :   /* Write HTTP/2 request headers */
     587           9 :   fd_h2_tx_prepare( client->conn, client->frame_tx, FD_H2_FRAME_TYPE_HEADERS, FD_H2_FLAG_END_HEADERS, stream_id );
     588           9 :   fd_grpc_req_hdrs_t req_meta = {
     589           9 :     .host     = client->host,
     590           9 :     .host_len = client->host_len,
     591           9 :     .port     = client->port,
     592           9 :     .path     = path,
     593           9 :     .path_len = path_len,
     594           9 :     .https    = 1, /* grpc_client assumes TLS encryption for now */
     595             : 
     596           9 :     .bearer_auth     = auth_token,
     597           9 :     .bearer_auth_len = auth_token_sz
     598           9 :   };
     599           9 :   if( FD_UNLIKELY( !fd_grpc_h2_gen_request_hdrs(
     600           9 :       &req_meta,
     601           9 :       client->frame_tx,
     602           9 :       client->version,
     603           9 :       client->version_len
     604           9 :   ) ) ) {
     605           0 :     FD_LOG_WARNING(( "Failed to generate gRPC request headers (%.*s). This is a bug", (int)path_len, path ));
     606           0 :     fd_grpc_client_stream_release( client, stream );
     607           0 :     return NULL;
     608           0 :   }
     609           9 :   fd_h2_tx_commit( client->conn, client->frame_tx );
     610             : 
     611             :   /* Queue request payload for send
     612             :      (Protobuf message might have to be fragmented into multiple HTTP/2
     613             :      DATA frames if the client gets blocked)
     614             :      For streaming requests, don't set END_STREAM flag yet */
     615           9 :   uint flags = is_streaming ? 0U : FD_H2_FLAG_END_STREAM;
     616           9 :   fd_h2_tx_op_init( client->request_tx_op, client->nanopb_tx, payload_sz, flags );
     617           9 :   fd_grpc_client_request_continue1( client );
     618           9 :   client->metrics->requests_sent++;
     619           9 :   client->metrics->streams_active++;
     620             : 
     621           9 :   FD_LOG_DEBUG(( "gRPC request path=%.*s sz=%lu streaming=%d", (int)path_len, path, serialized_sz, is_streaming ));
     622             : 
     623           9 :   return stream;
     624           9 : }
     625             : 
     626             : fd_grpc_h2_stream_t *
     627             : fd_grpc_client_request_start1(
     628             :     fd_grpc_client_t *   client,
     629             :     char const *         path,
     630             :     ulong                path_len,
     631             :     ulong                request_ctx,
     632             :     uchar const *        protobuf,
     633             :     ulong                protobuf_sz,
     634             :     char const *         auth_token,
     635             :     ulong                auth_token_sz,
     636             :     int                  is_streaming
     637           0 : ) {
     638           0 :   if( FD_UNLIKELY( fd_grpc_client_request_is_blocked( client ) ) ) return NULL;
     639             : 
     640           0 :   int const headers_only = (protobuf==NULL);
     641           0 :   if( FD_UNLIKELY( headers_only && !is_streaming ) ) {
     642           0 :     FD_LOG_WARNING(( "headers-only request requires is_streaming (path %.*s). This is a bug", (int)path_len, path ));
     643           0 :     return NULL;
     644           0 :   }
     645             : 
     646             :   /* Validate protobuf size */
     647           0 :   FD_TEST( client->nanopb_tx_max > sizeof(fd_grpc_hdr_t) );
     648           0 :   ulong const max_proto_sz = client->nanopb_tx_max - sizeof(fd_grpc_hdr_t);
     649           0 :   if( FD_UNLIKELY( protobuf_sz > max_proto_sz ) ) {
     650           0 :     FD_LOG_WARNING(( "Protobuf message too large (%lu bytes) for path (%.*s). Max size is %lu bytes", protobuf_sz, (int)path_len, path, max_proto_sz ));
     651           0 :     return NULL;
     652           0 :   }
     653             : 
     654           0 :   if( FD_LIKELY( !headers_only ) ) {
     655             :     /* Copy protobuf to buffer after gRPC header */
     656           0 :     uchar * proto_buf = client->nanopb_tx + sizeof(fd_grpc_hdr_t);
     657           0 :     memcpy( proto_buf, protobuf, protobuf_sz );
     658             : 
     659             :     /* Create gRPC length prefix */
     660           0 :     fd_grpc_hdr_t hdr = {
     661           0 :       .compressed=0,
     662           0 :       .msg_sz=fd_uint_bswap( (uint)protobuf_sz )
     663           0 :     };
     664           0 :     memcpy( client->nanopb_tx, &hdr, sizeof(fd_grpc_hdr_t) );
     665           0 :   }
     666           0 :   ulong const payload_sz = protobuf_sz + sizeof(fd_grpc_hdr_t);
     667             : 
     668             :   /* Allocate stream descriptor */
     669           0 :   fd_grpc_h2_stream_t * stream    = fd_grpc_client_stream_acquire( client, request_ctx );
     670           0 :   uint const            stream_id = stream->s.stream_id;
     671             : 
     672             :   /* Write HTTP/2 request headers */
     673           0 :   fd_h2_tx_prepare( client->conn, client->frame_tx, FD_H2_FRAME_TYPE_HEADERS, FD_H2_FLAG_END_HEADERS, stream_id );
     674           0 :   fd_grpc_req_hdrs_t req_meta = {
     675           0 :     .host     = client->host,
     676           0 :     .host_len = client->host_len,
     677           0 :     .port     = client->port,
     678           0 :     .path     = path,
     679           0 :     .path_len = path_len,
     680           0 :     .https    = 1, /* grpc_client assumes TLS encryption for now */
     681             : 
     682           0 :     .bearer_auth     = auth_token,
     683           0 :     .bearer_auth_len = auth_token_sz
     684           0 :   };
     685           0 :   if( FD_UNLIKELY( !fd_grpc_h2_gen_request_hdrs(
     686           0 :       &req_meta,
     687           0 :       client->frame_tx,
     688           0 :       client->version,
     689           0 :       client->version_len
     690           0 :   ) ) ) {
     691           0 :     FD_LOG_WARNING(( "Failed to generate gRPC request headers (%.*s). This is a bug", (int)path_len, path ));
     692           0 :     fd_grpc_client_stream_release( client, stream );
     693           0 :     return NULL;
     694           0 :   }
     695           0 :   fd_h2_tx_commit( client->conn, client->frame_tx );
     696             : 
     697           0 :   client->metrics->streams_active++;
     698             : 
     699           0 :   if( FD_UNLIKELY( headers_only ) ) {
     700           0 :     client->request_stream = NULL;
     701           0 :     return stream;
     702           0 :   }
     703             : 
     704             :   /* Queue request payload for send
     705             :      (Protobuf message might have to be fragmented into multiple HTTP/2
     706             :      DATA frames if the client gets blocked)
     707             :      For streaming requests, don't set END_STREAM flag yet */
     708           0 :   uint flags = is_streaming ? 0U : FD_H2_FLAG_END_STREAM;
     709           0 :   fd_h2_tx_op_init( client->request_tx_op, client->nanopb_tx, payload_sz, flags );
     710           0 :   fd_grpc_client_request_continue1( client );
     711           0 :   client->metrics->requests_sent++;
     712             : 
     713           0 :   FD_LOG_DEBUG(( "gRPC request path=%.*s sz=%lu streaming=%d", (int)path_len, path, protobuf_sz, is_streaming ));
     714             : 
     715           0 :   return stream;
     716           0 : }
     717             : 
     718             : int
     719             : fd_grpc_client_stream_send_msg(
     720             :     fd_grpc_client_t *   client,
     721             :     fd_grpc_h2_stream_t * stream,
     722             :     pb_msgdesc_t const * fields,
     723             :     void const *         message
     724           0 : ) {
     725           0 :   if( FD_UNLIKELY( !client || !stream ) ) return 0;
     726           0 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD ) ) return 0;
     727           0 :   if( FD_UNLIKELY( !fd_h2_rbuf_is_empty( client->frame_tx ) ) ) return 0;
     728           0 :   if( FD_UNLIKELY( client->request_tx_op->chunk_sz > 0UL ) ) return 0;
     729           0 :   if( FD_UNLIKELY( client->request_stream != NULL && client->request_stream != stream ) ) return 0; /* Another stream has a request in progress */
     730           0 :   if( FD_UNLIKELY( stream->s.state!=FD_H2_STREAM_STATE_OPEN &&
     731           0 :                    stream->s.state!=FD_H2_STREAM_STATE_CLOSING_RX ) ) return 0;
     732             : 
     733             :   /* Encode message */
     734           0 :   FD_TEST( client->nanopb_tx_max > sizeof(fd_grpc_hdr_t) );
     735           0 :   uchar * proto_buf = client->nanopb_tx + sizeof(fd_grpc_hdr_t);
     736           0 :   pb_ostream_t ostream = pb_ostream_from_buffer( proto_buf, client->nanopb_tx_max - sizeof(fd_grpc_hdr_t) );
     737           0 :   if( FD_UNLIKELY( !pb_encode( &ostream, fields, message ) ) ) {
     738           0 :     FD_LOG_WARNING(( "Failed to encode Protobuf message for stream_id=%u. This is a bug (insufficient buffer space?)", stream->s.stream_id ));
     739           0 :     return 0;
     740           0 :   }
     741           0 :   ulong const serialized_sz = ostream.bytes_written;
     742             : 
     743             :   /* Create gRPC length prefix */
     744           0 :   fd_grpc_hdr_t hdr = {
     745           0 :     .compressed=0,
     746           0 :     .msg_sz=fd_uint_bswap( (uint)serialized_sz )
     747           0 :   };
     748           0 :   memcpy( client->nanopb_tx, &hdr, sizeof(fd_grpc_hdr_t) );
     749           0 :   ulong const payload_sz = serialized_sz + sizeof(fd_grpc_hdr_t);
     750             : 
     751             :   /* Queue message payload for send (without END_STREAM flag) */
     752           0 :   client->request_stream = stream;
     753           0 :   fd_h2_tx_op_init( client->request_tx_op, client->nanopb_tx, payload_sz, 0U );
     754           0 :   fd_grpc_client_request_continue1( client );
     755             : 
     756           0 :   FD_LOG_DEBUG(( "gRPC stream_send_msg stream_id=%u sz=%lu", stream->s.stream_id, serialized_sz ));
     757             : 
     758           0 :   return 1;
     759           0 : }
     760             : 
     761             : int
     762             : fd_grpc_client_stream_send_msg1(
     763             :     fd_grpc_client_t *    client,
     764             :     fd_grpc_h2_stream_t * stream,
     765             :     uchar const *         protobuf,
     766             :     ulong                 protobuf_sz
     767           0 : ) {
     768           0 :   if( FD_UNLIKELY( !client || !stream ) ) return 0;
     769           0 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD ) ) return 0;
     770           0 :   if( FD_UNLIKELY( !fd_h2_rbuf_is_empty( client->frame_tx ) ) ) return 0;
     771           0 :   if( FD_UNLIKELY( client->request_tx_op->chunk_sz > 0UL ) ) return 0;
     772           0 :   if( FD_UNLIKELY( client->request_stream != NULL && client->request_stream != stream ) ) return 0; /* Another stream has a request in progress */
     773           0 :   if( FD_UNLIKELY( stream->s.state!=FD_H2_STREAM_STATE_OPEN &&
     774           0 :                    stream->s.state!=FD_H2_STREAM_STATE_CLOSING_RX ) ) return 0;
     775             : 
     776             :   /* Validate protobuf size */
     777           0 :   FD_TEST( client->nanopb_tx_max > sizeof(fd_grpc_hdr_t) );
     778           0 :   ulong const max_proto_sz = client->nanopb_tx_max - sizeof(fd_grpc_hdr_t);
     779           0 :   if( FD_UNLIKELY( protobuf_sz > max_proto_sz ) ) {
     780           0 :     FD_LOG_WARNING(( "Protobuf message too large (%lu bytes) for stream_id=%u. Max size is %lu bytes", protobuf_sz, stream->s.stream_id, max_proto_sz ));
     781           0 :     return 0;
     782           0 :   }
     783             : 
     784             :   /* Copy protobuf to buffer after gRPC header */
     785           0 :   uchar * proto_buf = client->nanopb_tx + sizeof(fd_grpc_hdr_t);
     786           0 :   memcpy( proto_buf, protobuf, protobuf_sz );
     787             : 
     788             :   /* Create gRPC length prefix */
     789           0 :   fd_grpc_hdr_t hdr = {
     790           0 :     .compressed=0,
     791           0 :     .msg_sz=fd_uint_bswap( (uint)protobuf_sz )
     792           0 :   };
     793           0 :   memcpy( client->nanopb_tx, &hdr, sizeof(fd_grpc_hdr_t) );
     794           0 :   ulong const payload_sz = protobuf_sz + sizeof(fd_grpc_hdr_t);
     795             : 
     796             :   /* Queue message payload for send (without END_STREAM flag) */
     797           0 :   client->request_stream = stream;
     798           0 :   fd_h2_tx_op_init( client->request_tx_op, client->nanopb_tx, payload_sz, 0U );
     799           0 :   fd_grpc_client_request_continue1( client );
     800             : 
     801           0 :   FD_LOG_DEBUG(( "gRPC stream_send_msg stream_id=%u sz=%lu", stream->s.stream_id, protobuf_sz ));
     802             : 
     803           0 :   return 1;
     804           0 : }
     805             : 
     806             : int
     807             : fd_grpc_client_stream_close(
     808             :     fd_grpc_client_t *    client,
     809             :     fd_grpc_h2_stream_t * stream
     810           0 : ) {
     811           0 :   if( FD_UNLIKELY( !client || !stream ) ) return 0;
     812           0 :   if( FD_UNLIKELY( client->conn->flags & FD_H2_CONN_FLAGS_DEAD ) ) return 0;
     813           0 :   if( FD_UNLIKELY( !fd_h2_rbuf_is_empty( client->frame_tx ) ) ) return 0;
     814           0 :   if( FD_UNLIKELY( client->request_stream != NULL ) ) return 0; /* Another request in progress */
     815           0 :   if( FD_UNLIKELY( stream->s.state!=FD_H2_STREAM_STATE_OPEN &&
     816           0 :                    stream->s.state!=FD_H2_STREAM_STATE_CLOSING_RX ) ) return 0;
     817             : 
     818             :   /* Send empty DATA frame with END_STREAM flag to close the stream */
     819           0 :   fd_h2_stream_close_tx( &stream->s, client->conn );
     820           0 :   fd_h2_tx_prepare( client->conn, client->frame_tx, FD_H2_FRAME_TYPE_DATA, FD_H2_FLAG_END_STREAM, stream->s.stream_id );
     821           0 :   fd_h2_tx_commit( client->conn, client->frame_tx );
     822             : 
     823           0 :   FD_LOG_DEBUG(( "gRPC stream_close stream_id=%u", stream->s.stream_id ));
     824             : 
     825           0 :   return 1;
     826           0 : }
     827             : 
     828             : void
     829             : fd_grpc_client_deadline_set( fd_grpc_h2_stream_t * stream,
     830             :                              int                   deadline_kind,
     831           9 :                              long                  ts_nanos ) {
     832           9 :   switch( deadline_kind ) {
     833           6 :   case FD_GRPC_DEADLINE_HEADER:
     834           6 :     stream->header_deadline_nanos = ts_nanos;
     835           6 :     stream->has_header_deadline   = 1;
     836           6 :     break;
     837           3 :   case FD_GRPC_DEADLINE_RX_END:
     838           3 :     stream->rx_end_deadline_nanos = ts_nanos;
     839           3 :     stream->has_rx_end_deadline   = 1;
     840           3 :     break;
     841           9 :   }
     842           9 : }
     843             : 
     844             : /* Lookup stream by ID */
     845             : 
     846             : static fd_h2_stream_t *
     847             : fd_grpc_h2_stream_query( fd_h2_conn_t * conn,
     848           6 :                          uint           stream_id ) {
     849           6 :   fd_grpc_client_t * client = conn->ctx;
     850           6 :   for( ulong i=0UL; i<client->stream_cnt; i++ ) {
     851           3 :     if( client->stream_ids[ i ] == stream_id ) {
     852           3 :       return &client->streams[ i ]->s;
     853           3 :     }
     854           3 :   }
     855           3 :   return NULL;
     856           6 : }
     857             : 
     858             : static void
     859           0 : fd_grpc_h2_conn_established( fd_h2_conn_t * conn ) {
     860           0 :   fd_grpc_client_t * client = conn->ctx;
     861           0 :   client->h2_hs_done = 1;
     862           0 :   client->callbacks->conn_established( client->ctx );
     863           0 : }
     864             : 
     865             : static void
     866             : fd_grpc_h2_conn_final( fd_h2_conn_t * conn,
     867             :                        uint           h2_err,
     868           0 :                        int            closed_by ) {
     869           0 :   fd_grpc_client_t * client = conn->ctx;
     870           0 :   client->callbacks->conn_dead( client->ctx, h2_err, closed_by );
     871           0 : }
     872             : 
     873             : /* React to response data */
     874             : 
     875             : void
     876             : fd_grpc_h2_cb_headers(
     877             :     fd_h2_conn_t *   conn,
     878             :     fd_h2_stream_t * h2_stream,
     879             :     void const *     data,
     880             :     ulong            data_sz,
     881             :     ulong            flags
     882           0 : ) {
     883           0 :   fd_grpc_h2_stream_t * stream = fd_grpc_h2_stream_upcast( h2_stream );
     884           0 :   fd_grpc_client_t * client = conn->ctx;
     885             : 
     886           0 :   int h2_status = fd_grpc_h2_read_response_hdrs( &stream->hdrs, client->matcher, data, data_sz );
     887           0 :   if( FD_UNLIKELY( h2_status!=FD_H2_SUCCESS ) ) {
     888             :     /* Failed to parse HTTP/2 headers */
     889           0 :     fd_h2_stream_error( h2_stream, conn, client->frame_tx, FD_H2_ERR_PROTOCOL );
     890           0 :     client->callbacks->rx_end( client->ctx, stream->request_ctx, &stream->hdrs ); /* invalidates stream->hdrs */
     891           0 :     fd_grpc_client_stream_release( client, stream );
     892           0 :     return;
     893           0 :   }
     894             : 
     895           0 :   if( !stream->hdrs_received && !!( flags & FD_H2_FLAG_END_HEADERS) ) {
     896             :     /* Got initial response header */
     897           0 :     stream->hdrs_received = 1;
     898           0 :     stream->has_header_deadline = 0;
     899           0 :     if( FD_LIKELY( ( stream->hdrs.h2_status==200  ) &
     900           0 :                    ( !!stream->hdrs.is_grpc_proto ) ) ) {
     901           0 :       client->callbacks->rx_start( client->ctx, stream->request_ctx );
     902           0 :     }
     903           0 :   }
     904             : 
     905           0 :   if( ( flags & (FD_H2_FLAG_END_HEADERS|FD_H2_FLAG_END_STREAM) )
     906           0 :               ==(FD_H2_FLAG_END_HEADERS|FD_H2_FLAG_END_STREAM)   ) {
     907           0 :     client->callbacks->rx_end( client->ctx, stream->request_ctx, &stream->hdrs );
     908           0 :     fd_grpc_client_stream_release( client, stream );
     909           0 :     return;
     910           0 :   }
     911           0 : }
     912             : 
     913             : void
     914             : fd_grpc_h2_cb_data(
     915             :     fd_h2_conn_t *   conn,
     916             :     fd_h2_stream_t * h2_stream,
     917             :     void const *     data,
     918             :     ulong            data_sz,
     919             :     ulong            flags
     920           6 : ) {
     921           6 :   fd_grpc_client_t *    client = conn->ctx;
     922           6 :   fd_grpc_h2_stream_t * stream = fd_grpc_h2_stream_upcast( h2_stream );
     923           6 :   if( FD_UNLIKELY( ( stream->hdrs.h2_status!=200 ) |
     924           6 :                    ( !stream->hdrs.is_grpc_proto ) ) ) {
     925           0 :     goto check_end_stream;
     926           0 :   }
     927             : 
     928           6 :   do {
     929             : 
     930             :     /* Read header bytes */
     931           6 :     if( stream->msg_buf_used < sizeof(fd_grpc_hdr_t) ) {
     932           6 :       ulong hdr_frag_sz = fd_ulong_min( sizeof(fd_grpc_hdr_t) - stream->msg_buf_used, data_sz );
     933           6 :       fd_memcpy( stream->msg_buf + stream->msg_buf_used, data, hdr_frag_sz );
     934           6 :       stream->msg_buf_used += hdr_frag_sz;
     935           6 :       data     = (void const *)( (ulong)data + (ulong)hdr_frag_sz );
     936           6 :       data_sz -= hdr_frag_sz;
     937           6 :       if( FD_UNLIKELY( stream->msg_buf_used < sizeof(fd_grpc_hdr_t) ) ) goto check_end_stream;
     938             : 
     939             :       /* Header complete */
     940           6 :       stream->msg_sz = fd_uint_bswap( FD_LOAD( uint, (void *)( (ulong)stream->msg_buf+1 ) ) );
     941           6 :       if( FD_UNLIKELY( sizeof(fd_grpc_hdr_t)  + stream->msg_sz > stream->msg_buf_max ) ) {
     942           6 :         FD_LOG_WARNING(( "Received oversized gRPC message (%lu bytes), killing request", stream->msg_sz ));
     943           6 :         client->callbacks->rx_end( client->ctx, stream->request_ctx, &stream->hdrs );
     944           6 :         fd_h2_stream_error( h2_stream, conn, client->frame_tx, FD_H2_ERR_INTERNAL );
     945           6 :         fd_grpc_client_stream_release( client, stream );
     946           6 :         return;
     947           6 :       }
     948           6 :     }
     949             : 
     950             :     /* Read payload bytes */
     951           0 :     ulong wmark    = sizeof(fd_grpc_hdr_t) + stream->msg_sz;
     952           0 :     ulong chunk_sz = fd_ulong_min( stream->msg_buf_used+data_sz, wmark ) - stream->msg_buf_used;
     953           0 :     if( FD_UNLIKELY( chunk_sz>data_sz ) ) FD_LOG_CRIT(( "integer underflow" )); /* unreachable */
     954           0 :     fd_memcpy( stream->msg_buf + stream->msg_buf_used, data, chunk_sz );
     955           0 :     stream->msg_buf_used += chunk_sz;
     956           0 :     data     = (void const *)( (ulong)data + (ulong)chunk_sz );
     957           0 :     data_sz -= chunk_sz;
     958             : 
     959           0 :     client->metrics->stream_chunks_rx_cnt++;
     960             : 
     961           0 :     if( stream->msg_buf_used >= wmark ) {
     962             :       /* Data complete */
     963           0 :       void const * msg_ptr = stream->msg_buf + sizeof(fd_grpc_hdr_t);
     964           0 :       client->callbacks->rx_msg( client->ctx, msg_ptr, stream->msg_sz, stream->request_ctx );
     965           0 :       stream->msg_buf_used = 0UL;
     966           0 :       stream->msg_sz       = 0UL;
     967           0 :     }
     968             : 
     969           0 :   } while( data_sz );
     970             : 
     971             : 
     972             :   /* Check whether the server might want to end the stream in a DATA frame so no stream is leaked.
     973             :      This shouldn't happen in gRPC as server responses should indicate end of stream in headers */
     974           0 :   check_end_stream:
     975           0 :   if( flags & FD_H2_FLAG_END_STREAM ) {
     976             :     /* FIXME incomplete gRPC message */
     977           0 :     if( FD_UNLIKELY( stream->msg_buf_used ) ) {
     978           0 :       FD_LOG_WARNING(( "Received incomplete gRPC message" ));
     979           0 :     }
     980           0 :     client->callbacks->rx_end( client->ctx, stream->request_ctx, &stream->hdrs );
     981           0 :     fd_grpc_client_stream_release( client, stream );
     982           0 :   }
     983           0 : }
     984             : 
     985             : /* Server might kill our request */
     986             : 
     987             : static void
     988             : fd_grpc_h2_rst_stream( fd_h2_conn_t *   conn,
     989             :                        fd_h2_stream_t * h2_stream,
     990             :                        uint             error_code,
     991           3 :                        int              closed_by ) {
     992           3 :   fd_grpc_client_t * client = conn->ctx;
     993           3 :   if( closed_by==1 ) {
     994           3 :     FD_LOG_WARNING(( "server %s%.*s%s terminated request %s(%u-%s)%s",
     995           3 :                      fd_log_style_bold(), (int)client->host_len, client->host, fd_log_style_normal(),
     996           3 :                      fd_log_style_dim(), error_code, fd_h2_strerror( error_code ), fd_log_style_normal() ));
     997           3 :   } else {
     998           0 :     FD_LOG_WARNING(( "request to server %s%.*s%s failed %s(%u-%s)%s",
     999           0 :                      fd_log_style_bold(), (int)client->host_len, client->host, fd_log_style_normal(),
    1000           0 :                      fd_log_style_dim(), error_code, fd_h2_strerror( error_code ), fd_log_style_normal() ));
    1001           0 :   }
    1002           3 :   fd_grpc_h2_stream_t * stream = fd_grpc_h2_stream_upcast( h2_stream );
    1003           3 :   client->callbacks->rx_end( client->ctx, stream->request_ctx, &stream->hdrs ); /* invalidates stream->hdrs */
    1004           3 :   fd_grpc_client_stream_release( client, stream );
    1005           3 : }
    1006             : 
    1007             : /* A HTTP/2 flow control change might unblock a queued request send op */
    1008             : 
    1009             : void
    1010             : fd_grpc_h2_window_update( fd_h2_conn_t * conn,
    1011           0 :                           uint           increment ) {
    1012           0 :   (void)increment;
    1013             :   /* Defer tx resumption to after fd_h2_rx: filling frame_tx mid-rx can
    1014             :      starve control-frame responses (SETTINGS ACK, PONG). */
    1015           0 :   ((fd_grpc_client_t *)conn->ctx)->window_update_pending = 1;
    1016           0 : }
    1017             : 
    1018             : void
    1019             : fd_grpc_h2_stream_window_update( fd_h2_conn_t *   conn,
    1020             :                                  fd_h2_stream_t * stream,
    1021           0 :                                  uint             increment ) {
    1022           0 :   (void)increment;
    1023             :   /* Repay initial-window-shrink debt before the new credit counts. */
    1024           0 :   fd_grpc_h2_stream_t * st = fd_grpc_h2_stream_upcast( stream );
    1025           0 :   long repay = fd_long_min( st->tx_wnd_debt, (long)stream->tx_wnd );
    1026           0 :   if( FD_UNLIKELY( repay>0L ) ) {
    1027           0 :     stream->tx_wnd  -= (uint)repay;
    1028           0 :     st->tx_wnd_debt -= repay;
    1029           0 :   }
    1030           0 :   ((fd_grpc_client_t *)conn->ctx)->window_update_pending = 1;
    1031           0 : }
    1032             : 
    1033             : void
    1034             : fd_grpc_h2_initial_window_update( fd_h2_conn_t * conn,
    1035           0 :                                   long           delta ) {
    1036           0 :   fd_grpc_client_t * client = conn->ctx;
    1037           0 :   for( uint i=0U; i<client->stream_cnt; i++ ) {
    1038           0 :     fd_grpc_h2_stream_t * stream = client->streams[ i ];
    1039           0 :     fd_h2_stream_t *      s      = &stream->s;
    1040             :     /* The effective window can go negative on a shrink below consumed
    1041             :        credit (RFC 9113 section 6.9.2).  tx_wnd is unsigned, so carry
    1042             :        the deficit in tx_wnd_debt; WINDOW_UPDATEs repay it before
    1043             :        granting new credit. */
    1044           0 :     long wnd = (long)s->tx_wnd - stream->tx_wnd_debt + delta;
    1045           0 :     if( FD_UNLIKELY( wnd > 0x7fffffffL ) ) {
    1046           0 :       fd_h2_conn_error( conn, FD_H2_ERR_FLOW_CONTROL );
    1047           0 :       return;
    1048           0 :     }
    1049           0 :     s->tx_wnd           = (uint)fd_long_max( wnd, 0L );
    1050           0 :     stream->tx_wnd_debt =       fd_long_max( -wnd, 0L );
    1051           0 :   }
    1052             :   /* Do not continue the tx op here: this runs inside fd_h2_rx SETTINGS
    1053             :      processing, and filling frame_tx now can starve the SETTINGS ACK
    1054             :      push (conn death via FD_H2_ERR_INTERNAL).  Defer to after rx. */
    1055           0 :   if( FD_UNLIKELY( delta>0L ) ) client->window_update_pending = 1;
    1056           0 : }
    1057             : 
    1058             : void
    1059           3 : fd_grpc_h2_ping_ack( fd_h2_conn_t * conn ) {
    1060           3 :   fd_grpc_client_t * client = conn->ctx;
    1061           3 :   client->callbacks->ping_ack( client->ctx );
    1062           3 : }
    1063             : 
    1064             : fd_h2_rbuf_t *
    1065           6 : fd_grpc_client_rbuf_tx( fd_grpc_client_t * client ) {
    1066           6 :   return client->frame_tx;
    1067           6 : }
    1068             : 
    1069             : fd_h2_rbuf_t *
    1070           0 : fd_grpc_client_rbuf_rx( fd_grpc_client_t * client ) {
    1071           0 :   return client->frame_rx;
    1072           0 : }
    1073             : 
    1074             : fd_h2_conn_t *
    1075         264 : fd_grpc_client_h2_conn( fd_grpc_client_t * client ) {
    1076         264 :   return client->conn;
    1077         264 : }
    1078             : 
    1079             : /* fd_grpc_client_h2_callbacks specifies h2->grpc_client callbacks.
    1080             :    Stored in .rodata for security.  Must be kept in sync with fd_h2 to
    1081             :    avoid NULL pointers. */
    1082             : 
    1083             : fd_h2_callbacks_t const fd_grpc_client_h2_callbacks = {
    1084             :   .stream_create         = fd_h2_noop_stream_create,
    1085             :   .stream_query          = fd_grpc_h2_stream_query,
    1086             :   .conn_established      = fd_grpc_h2_conn_established,
    1087             :   .conn_final            = fd_grpc_h2_conn_final,
    1088             :   .headers               = fd_grpc_h2_cb_headers,
    1089             :   .data                  = fd_grpc_h2_cb_data,
    1090             :   .rst_stream            = fd_grpc_h2_rst_stream,
    1091             :   .window_update         = fd_grpc_h2_window_update,
    1092             :   .stream_window_update  = fd_grpc_h2_stream_window_update,
    1093             :   .initial_window_update = fd_grpc_h2_initial_window_update,
    1094             :   .ping_ack              = fd_grpc_h2_ping_ack,
    1095             : };

Generated by: LCOV version 1.14