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