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 */