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

          Line data    Source code
       1             : #ifndef HEADER_fd_src_waltz_grpc_fd_grpc_client_h
       2             : #define HEADER_fd_src_waltz_grpc_fd_grpc_client_h
       3             : 
       4             : /* fd_grpc_client.h provides an API for dispatching unary and server-
       5             :    streaming gRPC requests over HTTP/2+TLS. */
       6             : 
       7             : #include "fd_grpc_codec.h"
       8             : #include "../fd_fqdn.h"
       9             : #include "../../third_party/nanopb/pb_firedancer.h" /* pb_msgdesc_t */
      10             : #include "../tlsrec/fd_tlsrec.h"
      11             : 
      12             : struct fd_grpc_client_private;
      13             : typedef struct fd_grpc_client_private fd_grpc_client_t;
      14             : 
      15             : struct fd_grpc_h2_stream;
      16             : typedef struct fd_grpc_h2_stream fd_grpc_h2_stream_t;
      17             : 
      18             : /* FD_GRPC_CLIENT_MAX_STREAMS specifies the max number of inflight
      19             :    unary and server-streaming requests.  Note that grpc_client does
      20             :    not scale well to large numbers due to O(n) algorithms. */
      21             : 
      22         840 : #define FD_GRPC_CLIENT_MAX_STREAMS 8
      23             : 
      24             : /* FD_GRPC_DEADLINE_* identify different types of request deadlines. */
      25             : 
      26          15 : #define FD_GRPC_DEADLINE_HEADER 1 /* deadline by which Response-Headers are received */
      27           9 : #define FD_GRPC_DEADLINE_RX_END 2 /* deadline by which 'end of stream' must have been reached */
      28             : 
      29             : /* fd_grpc_client_metrics_t hold counters that are incremented by a
      30             :    grpc_client. */
      31             : 
      32             : struct fd_grpc_client_metrics {
      33             : 
      34             :   /* wakeup_cnt counts the number of times the gRPC client was polled
      35             :      for I/O. */
      36             :   ulong wakeup_cnt;
      37             : 
      38             :   /* stream_err_cnt counts the number of survivable stream errors.
      39             :      These include out-of-memory conditions and decode failures. */
      40             :   ulong stream_err_cnt;
      41             : 
      42             :   /* conn_err_cnt counts the number of connection errors that resulted
      43             :      in connection termination.  These include protocol and I/O errors. */
      44             :   ulong conn_err_cnt;
      45             : 
      46             :   /* stream_chunks_tx_cnt increments whenever a DATA frame containing
      47             :      request bytes is sent.  stream_chunks_tx_bytes counts the number of
      48             :      stream bytes sent. */
      49             :   ulong stream_chunks_tx_cnt;
      50             :   ulong stream_chunks_tx_bytes;
      51             : 
      52             :   /* stream_chunks_rx_cnt increments whenever a DATA frame containing
      53             :      response bytes is received.  stream_chunks_rx_bytes counts the
      54             :      number of stream bytes received. */
      55             :   ulong stream_chunks_rx_cnt;
      56             :   ulong stream_chunks_rx_bytes;
      57             : 
      58             :   /* requests_sent increments whenever a gRPC request finished sending. */
      59             :   ulong requests_sent;
      60             : 
      61             :   /* streams_active is the number of streams not in 'closed' state. */
      62             :   long streams_active;
      63             : 
      64             :   /* rx_wait_ticks_cum is the cumulative time in ticks that incoming
      65             :      gRPC messages were in a "waiting" state.  The waiting state begins
      66             :      when the first byte of a HTTP/2 frame is received, and ends when
      67             :      all gRPC message bytes are received.
      68             : 
      69             :      This is a rough measure of server-to-client congestion.  On a
      70             :      healthy connection, this value should be close to zero. */
      71             :   long rx_wait_ticks_cum;
      72             : 
      73             :   /* tx_wait_ticks_cum is the cumulative time in ticks that an outgoing
      74             :      message was in a "waiting" state.  The waiting state begins when
      75             :      a message is ready to be sent, and ends when all message bytes were
      76             :      handed to the TCP layer.
      77             : 
      78             :      This is a rough measure of client-to-server congestion, which can
      79             :      be caused by the TCP server receive window, TCP client congestion
      80             :      control, or HTTP/2 server flow control.  On a healthy connection,
      81             :      this value should be close to zero. */
      82             :   long tx_wait_ticks_cum;
      83             : 
      84             : };
      85             : 
      86             : typedef struct fd_grpc_client_metrics fd_grpc_client_metrics_t;
      87             : 
      88             : /* fd_grpc_client_callbacks_t is a virtual function table containing
      89             :    grpc_client->app callbacks. */
      90             : 
      91             : struct fd_grpc_client_callbacks {
      92             : 
      93             :   /* conn_established is called when the initial HTTP/2 SETTINGS
      94             :      exchange concludes.  Technically, requests can be sent before this
      95             :      point, though. */
      96             : 
      97             :   void
      98             :   (* conn_established)( void * app_ctx );
      99             : 
     100             :   /* conn_dead is called when the HTTP/2 connection ends.  To recover
     101             :      from this condition, call fd_grpc_client_reset(). */
     102             : 
     103             :   void
     104             :   (* conn_dead)( void * app_ctx,
     105             :                  uint   h2_err,
     106             :                  int    closed_by );
     107             : 
     108             :   /* tx_complete marks the completion of a tx operation. */
     109             : 
     110             :   void
     111             :   (* tx_complete)( void * app_ctx,
     112             :                    ulong  request_ctx );
     113             : 
     114             :   /* rx_start signals that the server sent back a response header
     115             :      indicating success.  rx_start is always called before the first
     116             :      call to rx_msg for that request_ctx. */
     117             : 
     118             :   void
     119             :   (* rx_start)( void * app_ctx,
     120             :                 ulong  request_ctx );
     121             : 
     122             :   /* rx_msg delivers a gRPC message.  May be called multiple times for
     123             :      the same request (server streaming). */
     124             : 
     125             :   void
     126             :   (* rx_msg)( void *       app_ctx,
     127             :               void const * protobuf,
     128             :               ulong        protobuf_sz,
     129             :               ulong        request_ctx );
     130             : 
     131             :   /* rx_end indicates that no more rx_msg callbacks will be delivered
     132             :      for a request. */
     133             : 
     134             :   void
     135             :   (* rx_end)( void *                app_ctx,
     136             :               ulong                 request_ctx,
     137             :               fd_grpc_resp_hdrs_t * resp );
     138             : 
     139             :   /* rx_timeout indicates that a request deadline was exceeded.
     140             :      deadline_kind indicates which timer fired. */
     141             : 
     142             :   void
     143             :   (* rx_timeout)( void * app_ctx,
     144             :                   ulong  request_ctx,
     145             :                   int    deadline_kind );
     146             : 
     147             :   /* ping_ack delivers an acknowledgement of a PING that was previously
     148             :      sent by fd_h2_tx_ping. */
     149             : 
     150             :   void
     151             :   (* ping_ack)( void * app_ctx );
     152             : 
     153             : };
     154             : 
     155             : typedef struct fd_grpc_client_callbacks fd_grpc_client_callbacks_t;
     156             : 
     157             : FD_PROTOTYPES_BEGIN
     158             : 
     159             : /* Constructors */
     160             : 
     161             : ulong
     162             : fd_grpc_client_align( void );
     163             : 
     164             : ulong
     165             : fd_grpc_client_footprint( ulong buf_max );
     166             : 
     167             : fd_grpc_client_t *
     168             : fd_grpc_client_new( void *                             mem,
     169             :                     fd_grpc_client_callbacks_t const * callbacks,
     170             :                     fd_grpc_client_metrics_t *         metrics,
     171             :                     void *                             app_ctx,
     172             :                     ulong                              buf_max,
     173             :                     ulong                              rng_seed );
     174             : 
     175             : void *
     176             : fd_grpc_client_delete( fd_grpc_client_t * client );
     177             : 
     178             : /* fd_grpc_client_next_deadline returns the earliest stream deadline
     179             :    (header or rx-end) across all inflight requests, or LONG_MAX if no
     180             :    stream has a deadline armed. */
     181             : 
     182             : FD_FN_PURE long
     183             : fd_grpc_client_next_deadline( fd_grpc_client_t const * client );
     184             : 
     185             : /* fd_grpc_client_tx_pending returns 1 if the client has HTTP/2 frame
     186             :    bytes buffered that could not yet be written to the transport. */
     187             : 
     188             : FD_FN_PURE int
     189             : fd_grpc_client_tx_pending( fd_grpc_client_t const * client );
     190             : 
     191             : /* fd_grpc_client_tls_rx_pending returns 1 if the client holds decrypted
     192             :    TLS plaintext not yet pushed into the HTTP/2 layer.  Such bytes do
     193             :    not make the socket readable, so the caller must step the client
     194             :    without waiting for an fd event.  fd_grpc_client_tls_tx_pending
     195             :    returns 1 if encrypted TLS bytes are waiting to be written to the
     196             :    socket (e.g. a handshake flight blocked on EAGAIN).  RX continues
     197             :    while it is set; arm EPOLLOUT to learn when the send can resume. */
     198             : 
     199             : FD_FN_PURE int
     200             : fd_grpc_client_tls_rx_pending( fd_grpc_client_t const * client );
     201             : 
     202             : FD_FN_PURE int
     203             : fd_grpc_client_tls_tx_pending( fd_grpc_client_t const * client );
     204             : 
     205             : /* fd_grpc_client_tx_starved returns the number of request bytes parked
     206             :    mid-send because the peer's HTTP/2 flow-control window (stream or
     207             :    connection) is exhausted, or 0 if not starved.  Only a WINDOW_UPDATE
     208             :    from the peer can resume it. */
     209             : 
     210             : FD_FN_PURE ulong
     211             : fd_grpc_client_tx_starved( fd_grpc_client_t const * client );
     212             : 
     213             : /* fd_grpc_client_reset cancels all inflight requests and abandons the
     214             :    HTTP/2 client connection.  Config params are kept intact (e.g. host,
     215             :    port, version). */
     216             : 
     217             : void
     218             : fd_grpc_client_reset( fd_grpc_client_t * client );
     219             : 
     220             : /* fd_grpc_client_set_version sets the gRPC client's version string
     221             :    (relayed via user-agent header).  No reference to the provided string
     222             :    is kept (the content is copied out to the client object).  version
     223             :    does not have to be null-terminated.  version_len must be
     224             :    FD_GRPC_CLIENT_VERSION_LEN_MAX or less, otherwise a warning is logged
     225             :    and the client's version string remains unchanged. */
     226             : 
     227             : #define FD_GRPC_CLIENT_VERSION_LEN_MAX (63UL)
     228             : 
     229             : void
     230             : fd_grpc_client_set_version( fd_grpc_client_t * client,
     231             :                             char const *       version,
     232             :                             ulong              version_len );
     233             : 
     234             : /* fd_grpc_client_set_authority sets the authority header to the
     235             :    specified hostname and port number.  host_len should be less than
     236             :    FD_FQDN_BUF_MAX, otherwise host is truncated. */
     237             : 
     238             : void
     239             : fd_grpc_client_set_authority( fd_grpc_client_t * client,
     240             :                               char const *       host,
     241             :                               ulong              host_len,
     242             :                               ushort             port );
     243             : 
     244             : /* fd_grpc_client_rxtx_socket drives I/O against a TCP socket.
     245             :    (recvmsg(2) and sendmsg(2)).  Uses MSG_NOSIGNAL|MSG_DONTWAIT flags.
     246             : 
     247             :    Returns -1 if an error was encountered, and errno will be set.
     248             :    Returns 1 if send would block with EAGAIN.  Otherwise, returns 0. */
     249             : 
     250             : int
     251             : fd_grpc_client_rxtx_socket( fd_grpc_client_t * client,
     252             :                             int                sock_fd,
     253             :                             long               now,
     254             :                             int *              charge_busy );
     255             : 
     256             : /* fd_grpc_client_tx_flush_socket attempts one sendmsg(2) of pending
     257             :    HTTP/2 frame bytes to the socket.  Returns -1 (errno set) on a hard
     258             :    send error, 1 if send would block with EAGAIN, and 0 otherwise (a
     259             :    short write may leave bytes pending; check fd_grpc_client_tx_pending). */
     260             : 
     261             : int
     262             : fd_grpc_client_tx_flush_socket( fd_grpc_client_t * client,
     263             :                                 int                sock_fd );
     264             : 
     265             : /* fd_grpc_client_tls_flush writes pending encrypted TLS bytes to the
     266             :    socket, looping until the TLS TX buffer is empty or send blocks.
     267             :    Returns -1 (errno set) on a hard send error, 1 if send would block
     268             :    with EAGAIN (bytes remain, see fd_grpc_client_tls_tx_pending), and 0
     269             :    if the buffer is now empty. */
     270             : 
     271             : int
     272             : fd_grpc_client_tls_flush( fd_grpc_client_t * client,
     273             :                           int                sock_fd );
     274             : 
     275             : /* fd_grpc_client_rxtx_tls drives I/O against a TCP socket with
     276             :    fd_tlsrec encryption.  Handles TLS handshake, encrypts outgoing
     277             :    frames, and decrypts incoming data.
     278             : 
     279             :    Returns 0 on success and -1 if there is an unrecoverable error. */
     280             : 
     281             : int
     282             : fd_grpc_client_rxtx_tls( fd_grpc_client_t * client,
     283             :                          fd_tlsrec_conn_t * tls_conn,
     284             :                          int                sock_fd,
     285             :                          long               now,
     286             :                          int *              charge_busy );
     287             : 
     288             : /* fd_grpc_client_request_start queues a gRPC request for send.  The
     289             :    request includes one Protobuf message (unary request).  The client
     290             :    can only write one request payload at a time, but can have multiple
     291             :    requests pending for responses.
     292             : 
     293             :    path is the HTTP request path which usually follows the pattern
     294             :    '/path.to.package/Service.Function'.  If auth_token_sz is greater
     295             :    than zero, adds a request header 'authorization: Bearer *auth_token'.
     296             : 
     297             :    request_ctx is an arbitrary number used to identify the request.  It
     298             :    echoes in callbacks.
     299             : 
     300             :    fields points to a generated nanopb descriptor.  message points to a
     301             :    generated nanopb struct that the user filled in with info.  Calls
     302             :    pb_encode() internally.
     303             : 
     304             :    auth_token is an optional authorization header.  The header value is
     305             :    prepended with "Bearer ".  auth_token_sz==0 omits the auth header.
     306             : 
     307             :    is_streaming: If 0, this is a unary request and the stream is closed
     308             :    after sending the first message (END_STREAM flag set).  If non-zero,
     309             :    this is a client streaming request and the stream remains open for
     310             :    additional messages via fd_grpc_client_stream_send_msg().  The stream
     311             :    must be explicitly closed with fd_grpc_client_stream_close().
     312             : 
     313             :    Conditions for starting send:
     314             :    - The connection is not dead and the HTTP/2 handshake is complete.
     315             :    - Client has quota to open a new stream (MAX_CONCURRENT_STREAMS)
     316             :    - There is no other request still sending.
     317             :    - The message serialized size does not exceed buf_max (set in
     318             :      fd_grpc_client_new())
     319             :    - rbuf_tx is empty.  (HTTP/2 frames all flushed out to sockets) */
     320             : 
     321             : fd_grpc_h2_stream_t *
     322             : fd_grpc_client_request_start(
     323             :     fd_grpc_client_t *   client,
     324             :     char const *         path,
     325             :     ulong                path_len, /* in [0,128) */
     326             :     ulong                request_ctx,
     327             :     pb_msgdesc_t const * fields,
     328             :     void const *         message,
     329             :     char const *         auth_token,
     330             :     ulong                auth_token_sz,
     331             :     int                  is_streaming
     332             : );
     333             : 
     334             : fd_grpc_h2_stream_t *
     335             : fd_grpc_client_request_start1(
     336             :     fd_grpc_client_t *   client,
     337             :     char const *         path,
     338             :     ulong                path_len, /* in [0,128) */
     339             :     ulong                request_ctx,
     340             :     uchar const *        protobuf,
     341             :     ulong                protobuf_sz,
     342             :     char const *         auth_token,
     343             :     ulong                auth_token_sz,
     344             :     int                  is_streaming );
     345             : 
     346             : /* fd_grpc_client_stream_send_msg sends an additional message on an
     347             :    already-open client streaming request.  This function can only be
     348             :    called after fd_grpc_client_request_start() was called with
     349             :    is_streaming=1.
     350             : 
     351             :    Returns 1 on success, 0 if the operation failed (connection dead,
     352             :    buffers blocked, or encoding failure).
     353             : 
     354             :    Conditions for sending:
     355             :    - The connection is alive
     356             :    - No other send operation is in progress
     357             :    - rbuf_tx is empty
     358             :    - The message serialized size does not exceed buf_max */
     359             : 
     360             : int
     361             : fd_grpc_client_stream_send_msg(
     362             :     fd_grpc_client_t *    client,
     363             :     fd_grpc_h2_stream_t * stream,
     364             :     pb_msgdesc_t const *  fields,
     365             :     void const *          message
     366             : );
     367             : 
     368             : int
     369             : fd_grpc_client_stream_send_msg1(
     370             :     fd_grpc_client_t *    client,
     371             :     fd_grpc_h2_stream_t * stream,
     372             :     uchar const *         protobuf,
     373             :     ulong                 protobuf_sz );
     374             : 
     375             : /* fd_grpc_client_stream_close explicitly closes a client streaming
     376             :    request by sending an empty DATA frame with the END_STREAM flag.
     377             :    This signals to the server that no more messages will be sent.
     378             : 
     379             :    This function should be called after all messages have been sent via
     380             :    fd_grpc_client_stream_send_msg() to complete the client stream.
     381             : 
     382             :    Returns 1 on success, 0 if the operation failed (connection dead or
     383             :    buffers blocked).
     384             : 
     385             :    Conditions for closing:
     386             :    - The connection is alive
     387             :    - No other send operation is in progress
     388             :    - rbuf_tx is empty */
     389             : 
     390             : int
     391             : fd_grpc_client_stream_close(
     392             :     fd_grpc_client_t *    client,
     393             :     fd_grpc_h2_stream_t * stream
     394             : );
     395             : 
     396             : /* fd_grpc_client_deadline_set sets a request deadline (used to
     397             :    configure timeouts).  deadline_kind is FD_GRPC_DEADLINE_*.  Logs
     398             :    error and aborts app if deadline_kind is unsupported.
     399             : 
     400             :    Behavior for different deadline kinds:
     401             :    - HEADER: Deadline by which gRPC Response-Headers must have been
     402             :              received
     403             :    - RX_END: Deadline by which the response stream must have been ended.
     404             :              For unary responses, this is the point at which the message
     405             :              has been fully received.  For server-streaming responses,
     406             :              it is the point at which the last message has been
     407             :              received, and there are no more messages remaining.  (Under
     408             :              the hood, this is indicated by the HTTP/2 END_STREAM flag.) */
     409             : 
     410             : void
     411             : fd_grpc_client_deadline_set( fd_grpc_h2_stream_t * stream,
     412             :                              int                   deadline_kind,
     413             :                              long                  ts_nanos );
     414             : 
     415             : /* fd_grpc_client_is_connected returns 1 if HTTP/2 SETTINGS were
     416             :    exchanged, the TLS handshake is complete (if applicable), and the
     417             :    conn hasn't died.  Otherwise, returns 0. */
     418             : 
     419             : int
     420             : fd_grpc_client_is_connected( fd_grpc_client_t * client );
     421             : 
     422             : /* fd_grpc_client_request_is_blocked returns 1 if a call to
     423             :    fd_grpc_client_request_start would certainly fail.  Reasons include
     424             :    SSL / HTTP/2 handshake not complete, or buffers blocked. */
     425             : 
     426             : int
     427             : fd_grpc_client_request_is_blocked( fd_grpc_client_t * client );
     428             : 
     429             : /* fd_grpc_client_stream_send_is_blocked returns 1 if a message cannot
     430             :    currently be sent on an already open stream (conn dead, handshake
     431             :    incomplete, tx buffer non-empty, or a send already in flight), 0
     432             :    otherwise.  Unlike fd_grpc_client_request_is_blocked it does not
     433             :    require headroom to open a new stream, so a server advertising a low
     434             :    SETTINGS_MAX_CONCURRENT_STREAMS cannot stall sends on the open
     435             :    stream. */
     436             : 
     437             : int
     438             : fd_grpc_client_stream_send_is_blocked( fd_grpc_client_t * client );
     439             : 
     440             : int
     441             : fd_grpc_client_request_stream_busy( fd_grpc_client_t * client );
     442             : 
     443             : /* Pointers to internals for testing */
     444             : 
     445             : fd_h2_rbuf_t *
     446             : fd_grpc_client_rbuf_tx( fd_grpc_client_t * client );
     447             : 
     448             : fd_h2_rbuf_t *
     449             : fd_grpc_client_rbuf_rx( fd_grpc_client_t * client );
     450             : 
     451             : fd_h2_conn_t *
     452             : fd_grpc_client_h2_conn( fd_grpc_client_t * client );
     453             : 
     454             : extern fd_h2_callbacks_t const fd_grpc_client_h2_callbacks;
     455             : 
     456             : FD_PROTOTYPES_END
     457             : 
     458             : #endif /* HEADER_fd_src_waltz_grpc_fd_grpc_client_h */

Generated by: LCOV version 1.14