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 : };
|