Line data Source code
1 : #include "fd_rpc_client.h"
2 : #include "fd_rpc_client_private.h"
3 :
4 : #include "../../../third_party/picohttpparser/picohttpparser.h"
5 : #include "../../../waltz/http/fd_http.h"
6 : #include "../../../ballet/base58/fd_base58.h"
7 : #include "../../../ballet/json/fd_jtok.h"
8 :
9 : #include <errno.h>
10 : #include <stdio.h>
11 : #include <unistd.h>
12 : #include <strings.h>
13 : #include <sys/socket.h>
14 : #include <sys/types.h>
15 : #include <netinet/ip.h>
16 :
17 : #define MAX_REQUEST_LEN (1024UL)
18 :
19 : void *
20 : fd_rpc_client_new( void * mem,
21 : uint rpc_addr,
22 3 : ushort rpc_port ) {
23 3 : fd_rpc_client_t * rpc = (fd_rpc_client_t *)mem;
24 3 : rpc->request_id = 0UL;
25 3 : rpc->rpc_addr = rpc_addr;
26 3 : rpc->rpc_port = rpc_port;
27 387 : for( ulong i=0; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
28 384 : rpc->requests[ i ].state = FD_RPC_CLIENT_STATE_NONE;
29 384 : rpc->fds[ i ].fd = -1;
30 384 : rpc->fds[ i ].events = POLLIN | POLLOUT;
31 384 : rpc->fds[ i ].revents = 0;
32 384 : }
33 3 : return (void *)rpc;
34 3 : }
35 :
36 : long
37 : fd_rpc_client_wait_ready( fd_rpc_client_t * rpc,
38 0 : long timeout_ns ) {
39 :
40 :
41 0 : struct sockaddr_in addr = {
42 0 : .sin_family = AF_INET,
43 0 : .sin_port = fd_ushort_bswap( rpc->rpc_port ),
44 0 : .sin_addr = { .s_addr = rpc->rpc_addr }
45 0 : };
46 :
47 0 : struct pollfd pfd = {
48 0 : .fd = 0,
49 0 : .events = POLLOUT,
50 0 : .revents = 0
51 0 : };
52 :
53 0 : long start = fd_log_wallclock();
54 0 : for(;;) {
55 0 : pfd.fd = socket( AF_INET, SOCK_STREAM | SOCK_NONBLOCK, 0 );
56 0 : if( FD_UNLIKELY( pfd.fd<0 ) ) return FD_RPC_CLIENT_ERR_NETWORK;
57 :
58 0 : if( FD_UNLIKELY( -1==connect( pfd.fd, fd_type_pun( &addr ), sizeof(addr) ) && errno!=EINPROGRESS ) ) {
59 0 : if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
60 0 : return FD_RPC_CLIENT_ERR_NETWORK;
61 0 : }
62 :
63 0 : for(;;) {
64 0 : long now = fd_log_wallclock();
65 0 : if( FD_UNLIKELY( now-start>=timeout_ns ) ) {
66 0 : if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
67 0 : return FD_RPC_CLIENT_ERR_NETWORK;
68 0 : }
69 :
70 0 : int nfds = poll( &pfd, 1, (int)((now-start) / 1000000) );
71 0 : if( FD_UNLIKELY( 0==nfds ) ) continue;
72 0 : else if( FD_UNLIKELY( -1==nfds && errno==EINTR ) ) continue;
73 0 : else if( FD_UNLIKELY( -1==nfds ) ) {
74 0 : if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
75 0 : return FD_RPC_CLIENT_ERR_NETWORK;
76 0 : } else if( FD_LIKELY( pfd.revents & (POLLERR | POLLHUP) ) ) {
77 0 : if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
78 0 : break;
79 0 : } else if( FD_LIKELY( pfd.revents & POLLOUT ) ) {
80 0 : if( FD_UNLIKELY( close( pfd.fd )<0 ) ) FD_LOG_ERR(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
81 0 : return FD_RPC_CLIENT_SUCCESS;
82 0 : }
83 0 : }
84 0 : }
85 0 : }
86 :
87 : static ulong
88 9 : fd_rpc_available_slot( fd_rpc_client_t * rpc ) {
89 9 : for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
90 9 : if( FD_LIKELY( rpc->requests[i].state==FD_RPC_CLIENT_STATE_NONE ) ) return i;
91 9 : }
92 0 : return ULONG_MAX;
93 9 : }
94 :
95 : static ulong
96 : fd_rpc_find_request( fd_rpc_client_t * rpc,
97 18 : long request_id ) {
98 18 : for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
99 18 : if( FD_LIKELY( rpc->requests[i].state==FD_RPC_CLIENT_STATE_NONE ) ) continue;
100 18 : if( FD_LIKELY( rpc->requests[i].response.request_id!=request_id ) ) continue;
101 18 : return i;
102 18 : }
103 0 : return ULONG_MAX;
104 18 : }
105 :
106 : static long
107 : fd_rpc_client_request( fd_rpc_client_t * rpc,
108 : ulong method,
109 : long request_id,
110 : char * contents,
111 9 : int contents_len ) {
112 9 : ulong idx = fd_rpc_available_slot( rpc );
113 9 : if( FD_UNLIKELY( idx==ULONG_MAX) ) return FD_RPC_CLIENT_ERR_TOO_MANY;
114 :
115 9 : struct fd_rpc_client_request * request = &rpc->requests[ idx ];
116 :
117 9 : if( FD_UNLIKELY( contents_len<0 ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
118 9 : if( FD_UNLIKELY( (ulong)contents_len>=MAX_REQUEST_LEN ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
119 :
120 9 : int printed = snprintf( request->connected.request_bytes, sizeof(request->connected.request_bytes),
121 9 : "POST / HTTP/1.1\r\n"
122 9 : "Host: localhost:12001\r\n"
123 9 : "Content-Length: %d\r\n"
124 9 : "Content-Type: application/json\r\n\r\n"
125 9 : "%s", contents_len, contents );
126 9 : if( FD_UNLIKELY( printed<0 ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
127 9 : if( FD_UNLIKELY( (ulong)printed>=sizeof(request->connected.request_bytes) ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
128 9 : request->connected.request_bytes_cnt = (ulong)printed;
129 :
130 9 : int fd = socket( AF_INET, SOCK_STREAM | SOCK_NONBLOCK, 0 );
131 9 : if( FD_UNLIKELY( fd<0 ) ) return FD_RPC_CLIENT_ERR_NETWORK;
132 :
133 9 : struct sockaddr_in addr = {
134 9 : .sin_family = AF_INET,
135 9 : .sin_port = fd_ushort_bswap( rpc->rpc_port ),
136 9 : .sin_addr = { .s_addr = rpc->rpc_addr }
137 9 : };
138 :
139 9 : if( FD_UNLIKELY( -1==connect( fd, fd_type_pun( &addr ), sizeof(addr) ) && errno!=EINPROGRESS ) ) {
140 0 : if( FD_UNLIKELY( close( fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
141 0 : return FD_RPC_CLIENT_ERR_NETWORK;
142 0 : }
143 :
144 9 : rpc->request_id = request_id;
145 9 : rpc->fds[ idx ].fd = fd;
146 9 : request->response.method = method;
147 9 : request->response.status = FD_RPC_CLIENT_PENDING;
148 9 : request->response.request_id = rpc->request_id;
149 9 : request->connected.request_bytes_sent = 0UL;
150 9 : request->state = FD_RPC_CLIENT_STATE_CONNECTED;
151 9 : return request->response.request_id;
152 9 : }
153 :
154 : long
155 3 : fd_rpc_client_request_latest_block_hash( fd_rpc_client_t * rpc ) {
156 3 : char contents[ MAX_REQUEST_LEN ];
157 3 : long request_id = fd_long_if( rpc->request_id==LONG_MAX, 0L, rpc->request_id+1L );
158 :
159 3 : int contents_len = snprintf( contents, sizeof(contents),
160 3 : "{\"jsonrpc\":\"2.0\",\"id\":%ld,\"method\":\"getLatestBlockhash\",\"params\":[ { \"commitment\": \"processed\" }]}",
161 3 : request_id );
162 :
163 3 : return fd_rpc_client_request( rpc, FD_RPC_CLIENT_METHOD_LATEST_BLOCK_HASH, request_id, contents, contents_len );
164 3 : }
165 :
166 : long
167 6 : fd_rpc_client_request_transaction_count( fd_rpc_client_t * rpc ) {
168 6 : char contents[ MAX_REQUEST_LEN ];
169 6 : long request_id = fd_long_if( rpc->request_id==LONG_MAX, 0L, rpc->request_id+1L );
170 :
171 6 : int contents_len = snprintf( contents, sizeof(contents),
172 6 : "{\"jsonrpc\":\"2.0\",\"id\":%ld,\"method\":\"getTransactionCount\",\"params\":[ { \"commitment\": \"processed\" } ]}",
173 6 : request_id );
174 :
175 6 : return fd_rpc_client_request( rpc, FD_RPC_CLIENT_METHOD_TRANSACTION_COUNT, request_id, contents, contents_len );
176 6 : }
177 :
178 : static void
179 : fd_rpc_mark_error( fd_rpc_client_t * rpc,
180 : ulong idx,
181 0 : long error ) {
182 0 : if( FD_LIKELY( rpc->fds[ idx ].fd>=0 ) ) {
183 0 : if( FD_UNLIKELY( close( rpc->fds[ idx ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
184 0 : rpc->fds[ idx ].fd = -1;
185 0 : }
186 0 : rpc->requests[ idx ].state = FD_RPC_CLIENT_STATE_FINISHED;
187 0 : rpc->requests[ idx ].response.status = error;
188 0 : }
189 :
190 : static ulong
191 : fd_rpc_phr_content_length( struct phr_header * headers,
192 12 : ulong num_headers ) {
193 24 : for( ulong i=0UL; i<num_headers; i++ ) {
194 24 : if( FD_LIKELY( headers[i].name_len!=14UL ) ) continue;
195 12 : if( FD_LIKELY( strncasecmp( headers[i].name, "Content-Length", 14UL ) ) ) continue;
196 12 : ulong content_length = 0UL;
197 12 : if( FD_UNLIKELY( fd_http_parse_content_len( headers[i].value, (ulong)headers[i].value_len, &content_length ) ) ) return ULONG_MAX;
198 12 : if( FD_UNLIKELY( content_length>UINT_MAX ) ) return ULONG_MAX; /* prevent overflow */
199 12 : return content_length;
200 12 : }
201 0 : return ULONG_MAX;
202 12 : }
203 :
204 : static long
205 : parse_response( char * response,
206 : ulong response_len,
207 12 : fd_rpc_client_response_t * result ) {
208 12 : int minor_version;
209 12 : int status;
210 12 : const char * message;
211 12 : ulong message_len;
212 12 : struct phr_header headers[ 32 ];
213 12 : ulong num_headers = 32UL;
214 : /* last_len 0: a non-zero value after the headers completed makes
215 : picohttpparser scan the body for CRLFCRLF and never finish */
216 12 : int http_len = phr_parse_response( response, response_len,
217 12 : &minor_version, &status, &message, &message_len,
218 12 : headers, &num_headers, 0UL );
219 12 : if( FD_UNLIKELY( -2==http_len ) ) return FD_RPC_CLIENT_PENDING;
220 12 : else if( FD_UNLIKELY( -1==http_len ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
221 :
222 12 : if( FD_UNLIKELY( status!=200 ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
223 :
224 12 : ulong content_length = fd_rpc_phr_content_length( headers, num_headers );
225 12 : if( FD_UNLIKELY( content_length==ULONG_MAX ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
226 12 : if( FD_UNLIKELY( content_length+(ulong)http_len > MAX_REQUEST_LEN ) ) return FD_RPC_CLIENT_ERR_TOO_LARGE;
227 12 : if( FD_LIKELY( content_length+(ulong)http_len>response_len ) ) return FD_RPC_CLIENT_PENDING;
228 :
229 9 : fd_jtok_t j[1];
230 9 : fd_jtok_init( j, response+http_len, content_length );
231 9 : fd_jtok_str_t key;
232 9 : fd_jtok_obj_enter( j );
233 :
234 9 : switch( result->method ) {
235 6 : case FD_RPC_CLIENT_METHOD_TRANSACTION_COUNT: {
236 6 : ulong transaction_count = ULONG_MAX;
237 24 : while( fd_jtok_obj_next( j, &key ) ) {
238 18 : if( fd_jtok_str_eq( &key, "result" ) ) fd_jtok_ulong( j, &transaction_count );
239 18 : }
240 6 : if( FD_UNLIKELY( fd_jtok_fini( j ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
241 6 : if( FD_UNLIKELY( transaction_count==ULONG_MAX ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
242 :
243 6 : result->result.transaction_count.transaction_count = transaction_count;
244 6 : return FD_RPC_CLIENT_SUCCESS;
245 6 : }
246 3 : case FD_RPC_CLIENT_METHOD_LATEST_BLOCK_HASH: {
247 3 : char blockhash[ 45 ] = {0};
248 12 : while( fd_jtok_obj_next( j, &key ) ) {
249 9 : if( !fd_jtok_str_eq( &key, "result" ) ) continue;
250 3 : fd_jtok_obj_enter( j );
251 6 : while( fd_jtok_obj_next( j, &key ) ) {
252 3 : if( !fd_jtok_str_eq( &key, "value" ) ) continue;
253 3 : fd_jtok_obj_enter( j );
254 6 : while( fd_jtok_obj_next( j, &key ) ) {
255 3 : if( fd_jtok_str_eq( &key, "blockhash" ) ) fd_jtok_cstr( j, blockhash, sizeof(blockhash) );
256 3 : }
257 3 : }
258 3 : }
259 3 : if( FD_UNLIKELY( fd_jtok_fini( j ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
260 3 : if( FD_UNLIKELY( !fd_base58_decode_32( blockhash, result->result.latest_block_hash.block_hash ) ) ) return FD_RPC_CLIENT_ERR_MALFORMED;
261 3 : return FD_RPC_CLIENT_SUCCESS;
262 3 : }
263 0 : default:
264 0 : FD_TEST( 0 );
265 9 : }
266 0 : return FD_RPC_CLIENT_ERR_MALFORMED;
267 9 : }
268 :
269 : int
270 : fd_rpc_client_service( fd_rpc_client_t * rpc,
271 184851 : int wait ) {
272 184851 : int timeout = wait ? -1 : 0;
273 184851 : int nfds = poll( rpc->fds, FD_RPC_CLIENT_REQUEST_CNT, timeout );
274 184851 : if( FD_UNLIKELY( 0==nfds ) ) return 0;
275 184851 : else if( FD_UNLIKELY( -1==nfds && errno==EINTR ) ) return 0;
276 184851 : else if( FD_UNLIKELY( -1==nfds ) ) FD_LOG_ERR(( "poll failed (%i-%s)", errno, strerror( errno ) ));
277 :
278 23845779 : for( ulong i=0UL; i<FD_RPC_CLIENT_REQUEST_CNT; i++ ) {
279 23660928 : struct fd_rpc_client_request * request = &rpc->requests[i];
280 :
281 23660928 : if( FD_LIKELY( request->state==FD_RPC_CLIENT_STATE_CONNECTED && ( rpc->fds[ i ].revents & POLLOUT ) ) ) {
282 9 : long sent = send( rpc->fds[ i ].fd, request->connected.request_bytes+request->connected.request_bytes_sent,
283 9 : request->connected.request_bytes_cnt-request->connected.request_bytes_sent, MSG_NOSIGNAL );
284 9 : if( FD_UNLIKELY( -1==sent && errno==EAGAIN ) ) continue;
285 9 : if( FD_UNLIKELY( -1==sent ) ) {
286 0 : fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
287 0 : continue;
288 0 : }
289 :
290 9 : request->connected.request_bytes_sent += (ulong)sent;
291 9 : if( FD_UNLIKELY( request->connected.request_bytes_sent==request->connected.request_bytes_cnt ) ) {
292 9 : request->sent.response_bytes_read = 0UL;
293 9 : request->state = FD_RPC_CLIENT_STATE_SENT;
294 9 : }
295 9 : }
296 :
297 23660928 : if( FD_LIKELY( request->state==FD_RPC_CLIENT_STATE_SENT && ( rpc->fds[ i ].revents & POLLIN ) ) ) {
298 12 : long read = recv( rpc->fds[ i ].fd, request->response_bytes+request->sent.response_bytes_read,
299 12 : sizeof(request->response_bytes)-request->sent.response_bytes_read, 0 );
300 12 : if( FD_UNLIKELY( -1==read && errno==EAGAIN ) ) continue;
301 12 : else if( FD_UNLIKELY( -1==read ) ) {
302 0 : fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
303 0 : continue;
304 12 : } else if( FD_UNLIKELY( !read ) ) {
305 0 : fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_NETWORK );
306 0 : continue;
307 0 : }
308 :
309 12 : request->sent.response_bytes_read += (ulong)read;
310 12 : if( FD_UNLIKELY( request->sent.response_bytes_read==sizeof(request->response_bytes) ) ) {
311 0 : fd_rpc_mark_error( rpc, i, FD_RPC_CLIENT_ERR_TOO_LARGE );
312 0 : continue;
313 0 : }
314 :
315 12 : fd_rpc_client_response_t * response = &request->response;
316 12 : long status = parse_response( request->response_bytes,
317 12 : request->sent.response_bytes_read,
318 12 : response );
319 12 : if( FD_LIKELY( status==FD_RPC_CLIENT_PENDING ) ) continue;
320 9 : else if( FD_UNLIKELY( status==FD_RPC_CLIENT_SUCCESS ) ) {
321 9 : if( FD_UNLIKELY( close( rpc->fds[ i ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
322 9 : rpc->fds[ i ].fd = -1;
323 9 : response->status = FD_RPC_CLIENT_SUCCESS;
324 9 : request->state = FD_RPC_CLIENT_STATE_FINISHED;
325 9 : continue;
326 9 : } else {
327 0 : fd_rpc_mark_error( rpc, i, status );
328 0 : continue;
329 0 : }
330 12 : }
331 23660928 : }
332 :
333 184851 : return 1;
334 184851 : }
335 :
336 : fd_rpc_client_response_t *
337 : fd_rpc_client_status( fd_rpc_client_t * rpc,
338 : long request_id,
339 9 : int wait ) {
340 9 : ulong idx = fd_rpc_find_request( rpc, request_id );
341 9 : if( FD_UNLIKELY( idx==ULONG_MAX ) ) return NULL;
342 :
343 9 : if( FD_LIKELY( !wait ) ) return &rpc->requests[ idx ].response;
344 :
345 184860 : for(;;) {
346 184860 : if( FD_LIKELY( rpc->requests[ idx ].state==FD_RPC_CLIENT_STATE_FINISHED ) ) return &rpc->requests[ idx ].response;
347 184851 : fd_rpc_client_service( rpc, 1 );
348 184851 : }
349 9 : }
350 :
351 : void
352 : fd_rpc_client_close( fd_rpc_client_t * rpc,
353 9 : long request_id ) {
354 9 : ulong idx = fd_rpc_find_request( rpc, request_id );
355 9 : if( FD_UNLIKELY( idx==ULONG_MAX ) ) return;
356 :
357 9 : if( FD_LIKELY( rpc->fds[ idx ].fd>=0 ) ) {
358 0 : if( FD_UNLIKELY( close( rpc->fds[ idx ].fd )<0 ) ) FD_LOG_WARNING(( "close() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
359 0 : rpc->fds[ idx ].fd = -1;
360 0 : }
361 9 : rpc->requests[ idx ].state = FD_RPC_CLIENT_STATE_NONE;
362 9 : }
|