LCOV - code coverage report
Current view: top level - waltz/grpc - fd_grpc_client_private.h (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 3 3 100.0 %
Date: 2026-09-17 04:28:31 Functions: 1 11 9.1 %

          Line data    Source code
       1             : #ifndef HEADER_fd_src_waltz_grpc_fd_grpc_client_private_h
       2             : #define HEADER_fd_src_waltz_grpc_fd_grpc_client_private_h
       3             : 
       4             : #include "fd_grpc_client.h"
       5             : #include "../grpc/fd_grpc_codec.h"
       6             : #include "../h2/fd_h2.h"
       7             : #include "../tlsrec/fd_tlsrec_sock.h"
       8             : 
       9             : /* fd_grpc_h2_stream_t holds the state of a gRPC request. */
      10             : 
      11             : struct fd_grpc_h2_stream {
      12             :   fd_h2_stream_t s;
      13             : 
      14             :   ulong request_ctx;
      15             :   uint  next;
      16             : 
      17             :   /* Buffer response headers */
      18             :   fd_grpc_resp_hdrs_t hdrs;
      19             : 
      20             :   /* Buffer an incoming gRPC message */
      21             :   uchar * msg_buf;
      22             :   ulong   msg_buf_max;
      23             :   uint    hdrs_received : 1;
      24             :   ulong   msg_buf_used; /* including header */
      25             :   ulong   msg_sz;       /* size of next message */
      26             : 
      27             :   long header_deadline_nanos;  /* deadline to first resp header bit */
      28             :   long rx_end_deadline_nanos;  /* deadline to end of stream signal */
      29             :   long tx_wnd_debt;            /* send-window deficit from an initial-window shrink below consumed credit; repaid from WINDOW_UPDATEs before new credit counts */
      30             :   uint has_header_deadline : 1;
      31             :   uint has_rx_end_deadline : 1;
      32             : };
      33             : 
      34             : /* Declare a pool of stream objects.
      35             : 
      36             :    While only one stream is used to write requests out to the wire at a
      37             :    time, a gRPC client might be waiting for multiple responses. */
      38             : 
      39             : #define POOL_NAME fd_grpc_h2_stream_pool
      40         168 : #define POOL_T fd_grpc_h2_stream_t
      41             : #define POOL_IDX_T uint
      42             : #include "../../util/tmpl/fd_pool.c"
      43             : 
      44             : static inline fd_grpc_h2_stream_t *
      45           9 : fd_grpc_h2_stream_upcast( fd_h2_stream_t * stream ) {
      46             :   return (fd_grpc_h2_stream_t *)( (ulong)stream - offsetof(fd_grpc_h2_stream_t, s) );
      47           9 : }
      48             : 
      49             : /* I/O paths
      50             : 
      51             :    RX path (TLS)
      52             :    - fd_grpc_client_rxtx_tls
      53             :    - calls fd_tlsrec_sock_rx to read and decrypt TCP ciphertext
      54             :    - pushes plaintext into fd_h2_rbuf
      55             : 
      56             :    TX path (TLS)
      57             :    - fd_grpc_client_rxtx_tls
      58             :    - pops plaintext from fd_h2_rbuf
      59             :    - calls fd_tlsrec_sock_tx to encrypt and write TCP ciphertext
      60             : 
      61             :    RX/TX path (plaintext sockets)
      62             :    - fd_grpc_client_rxtx_socket
      63             :    - calls fd_h2_rbuf_recvmsg / fd_h2_rbuf_sendmsg */
      64             : 
      65             : /* gRPC client internal state.  Quick overview:
      66             : 
      67             :    - The client maintains exactly one gRPC connection.
      68             :    - This conn includes a TCP socket, fd_tlsrec handle, and a fd_h2
      69             :      conn (and accompanying buffers).
      70             :    - The client object dies when the connection dies.
      71             : 
      72             :    - The client manages a small pool of stream objects.
      73             :    - Each stream has one of 3 states:
      74             :      - IDLE: marked free in stream_pool
      75             :      - OPEN: sending request data. marked used in stream_pool, present
      76             :           in stream_ids/streams, referred to by request_stream, has
      77             :           associated tx_op object.
      78             :      - CLOSE_TX: request sent, waiting for response. marked used in
      79             :           stream_pool,  present in stream_ids/streams.
      80             :     - Only 1 stream can be in OPEN state.
      81             : 
      82             :    Regular state transitions:
      83             : 
      84             :    - IDLE->OPEN: Client acquires a stream object and starts a tx_op
      85             :      See fd_grpc_client_request_start
      86             :    - OPEN->CLOSE_TX: tx_op finished writing request data and is now
      87             :      waiting for the response.  tx_op object finalized.
      88             :      See fd_grpc_client_request_continue1
      89             :    - CLOSE_TX->IDLE: All response data arrived.  Stream object
      90             :      deallocated.
      91             : 
      92             :    Irregular state transitions:
      93             : 
      94             :    - CLOSE_TX->IDLE: Server aborts stream before request is fully sent.
      95             :    - OPEN->IDLE: Server aborts stream before response is received. */
      96             : 
      97             : struct fd_grpc_client_private {
      98             :   fd_grpc_client_callbacks_t const * callbacks;
      99             :   void *                             ctx;
     100             : 
     101             :   fd_h2_hdr_matcher_t matcher[1];
     102             : 
     103             :   /* HTTP/2 connection */
     104             :   fd_h2_conn_t conn[1];
     105             :   fd_h2_rbuf_t frame_rx[1]; /* unencrypted HTTP/2 RX frame buffer */
     106             :   fd_h2_rbuf_t frame_tx[1]; /* unencrypted HTTP/2 TX frame buffer */
     107             : 
     108             :   /* HTTP/2 authority */
     109             :   char   host[ FD_FQDN_BUF_MAX ];
     110             :   ushort port; /* <=65535 */
     111             :   uchar  host_len; /* <FD_FQDN_BUF_MAX */
     112             : 
     113             :   uint  h2_hs_done : 1;
     114             :   uint  window_update_pending : 1; /* SETTINGS window delta arrived mid-rx; continue tx after fd_h2_rx returns */
     115             : 
     116             :   /* Inflight request
     117             :      Non-NULL until a gRPC request is fully written out. */
     118             :   fd_grpc_h2_stream_t * request_stream;
     119             :   fd_h2_tx_op_t         request_tx_op[1];
     120             : 
     121             :   /* Stream pool */
     122             :   fd_grpc_h2_stream_t * stream_pool;
     123             :   uchar *               stream_bufs;
     124             : 
     125             :   /* Stream map */
     126             :   /* FIXME pull this into a fd_map_tiny.c? */
     127             :   uint                  stream_ids[ FD_GRPC_CLIENT_MAX_STREAMS ];
     128             :   fd_grpc_h2_stream_t * streams   [ FD_GRPC_CLIENT_MAX_STREAMS ];
     129             :   ulong                 stream_cnt;
     130             : 
     131             :   /* Buffers */
     132             :   uchar * nanopb_tx;
     133             :   ulong   nanopb_tx_max;
     134             :   uchar * frame_scratch;
     135             :   ulong   frame_scratch_max;
     136             : 
     137             :   /* Frame buffers */
     138             :   uchar * frame_rx_buf;
     139             :   ulong   frame_rx_buf_max;
     140             :   uchar * frame_tx_buf;
     141             :   ulong   frame_tx_buf_max;
     142             : 
     143             :   fd_tlsrec_sock_t tls_sock[1];
     144             : 
     145             :   /* Version string */
     146             :   uchar version_len;
     147             :   char  version[ FD_GRPC_CLIENT_VERSION_LEN_MAX ];
     148             : 
     149             :   fd_grpc_client_metrics_t * metrics;
     150             : };
     151             : 
     152             : FD_PROTOTYPES_BEGIN
     153             : 
     154             : /* fd_grpc_client_stream_acquire grabs a new stream ID and a stream
     155             :    object. */
     156             : 
     157             : int
     158             : fd_grpc_client_stream_acquire_is_safe( fd_grpc_client_t * client );
     159             : 
     160             : 
     161             : fd_grpc_h2_stream_t *
     162             : fd_grpc_client_stream_acquire( fd_grpc_client_t * client,
     163             :                                ulong              request_ctx );
     164             : 
     165             : void
     166             : fd_grpc_client_stream_release( fd_grpc_client_t *    client,
     167             :                                fd_grpc_h2_stream_t * stream );
     168             : 
     169             : /* fd_grpc_client_service_streams checks all streams for timeouts and
     170             :    optionally generates receive window updates.
     171             :    FIXME poor algorithmic inefficiency (O(n)).  Consider using a service
     172             :          queue/heap */
     173             : 
     174             : void
     175             : fd_grpc_client_service_streams( fd_grpc_client_t * client,
     176             :                                 long               ts_nanos );
     177             : 
     178             : void
     179             : fd_grpc_h2_cb_headers(
     180             :     fd_h2_conn_t *   conn,
     181             :     fd_h2_stream_t * h2_stream,
     182             :     void const *     data,
     183             :     ulong            data_sz,
     184             :     ulong            flags
     185             : );
     186             : 
     187             : void
     188             : fd_grpc_h2_cb_data(
     189             :     fd_h2_conn_t *   conn,
     190             :     fd_h2_stream_t * h2_stream,
     191             :     void const *     data,
     192             :     ulong            data_sz,
     193             :     ulong            flags
     194             : );
     195             : 
     196             : void
     197             : fd_grpc_h2_stream_window_update( fd_h2_conn_t *   conn,
     198             :                                  fd_h2_stream_t * stream,
     199             :                                  uint             increment );
     200             : 
     201             : void
     202             : fd_grpc_h2_initial_window_update( fd_h2_conn_t * conn,
     203             :                                   long           delta );
     204             : 
     205             : FD_PROTOTYPES_END
     206             : 
     207             : #endif /* HEADER_fd_src_waltz_grpc_fd_grpc_client_private_h */

Generated by: LCOV version 1.14