Line data Source code
1 : #define _GNU_SOURCE
2 : #include "fd_ipecho_server.h"
3 :
4 : #include "../../util/fd_util.h"
5 : #include "../../util/net/fd_ip4.h"
6 :
7 : #include <errno.h>
8 : #include <unistd.h>
9 : #include <poll.h>
10 : #include <sys/epoll.h>
11 : #include <sys/socket.h>
12 : #include <netinet/in.h>
13 :
14 0 : #define STATE_READING (0)
15 0 : #define STATE_WRITING (1)
16 :
17 0 : #define CLOSE_OK ( 0)
18 0 : #define CLOSE_EXPECTED_EOF (-1)
19 0 : #define CLOSE_PEER_RESET (-2)
20 0 : #define CLOSE_LARGE_REQUEST (-3)
21 0 : #define CLOSE_BAD_HEADER (-4)
22 0 : #define CLOSE_BAD_TRAILER (-5)
23 0 : #define CLOSE_BAD_LENGTH (-6)
24 0 : #define CLOSE_EVICTED (-7)
25 :
26 : struct fd_ipecho_server_connection {
27 : int state;
28 :
29 : uint ipv4;
30 :
31 : ushort parent;
32 :
33 : ulong request_bytes_read;
34 : uchar request_bytes[ 22UL ];
35 : ulong response_bytes_written;
36 : uchar response_bytes[ 27UL ];
37 : };
38 :
39 : typedef struct fd_ipecho_server_connection fd_ipecho_server_connection_t;
40 :
41 : #define POOL_NAME conn_pool
42 0 : #define POOL_T fd_ipecho_server_connection_t
43 : #define POOL_IDX_T ushort
44 0 : #define POOL_NEXT parent
45 : #include "../../util/tmpl/fd_pool.c"
46 :
47 : struct fd_ipecho_server {
48 : int sockfd;
49 :
50 : int epoll_fd;
51 :
52 : ushort shred_version;
53 :
54 : ulong evict_idx;
55 : ulong max_connection_cnt;
56 :
57 : fd_ipecho_server_connection_t * pool;
58 : struct pollfd * pollfds;
59 :
60 : fd_ipecho_server_metrics_t metrics[ 1 ];
61 :
62 : ulong magic;
63 : };
64 :
65 : FD_FN_CONST ulong
66 0 : fd_ipecho_server_align( void ) {
67 0 : return 128UL;
68 0 : }
69 :
70 : FD_FN_CONST ulong
71 0 : fd_ipecho_server_footprint( ulong max_connection_cnt ) {
72 0 : ulong l = FD_LAYOUT_INIT;
73 0 : l = FD_LAYOUT_APPEND( l, fd_ipecho_server_align(), sizeof(fd_ipecho_server_t) );
74 0 : l = FD_LAYOUT_APPEND( l, conn_pool_align(), conn_pool_footprint( max_connection_cnt ) );
75 0 : l = FD_LAYOUT_APPEND( l, alignof(struct pollfd), (1UL+max_connection_cnt)*sizeof(struct pollfd) );
76 0 : return FD_LAYOUT_FINI( l, fd_ipecho_server_align() );
77 0 : }
78 :
79 : void *
80 : fd_ipecho_server_new( void * shmem,
81 0 : ulong max_connection_cnt ) {
82 0 : if( FD_UNLIKELY( !shmem ) ) {
83 0 : FD_LOG_WARNING(( "NULL shmem" ));
84 0 : return NULL;
85 0 : }
86 :
87 0 : if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)shmem, fd_ipecho_server_align() ) ) ) {
88 0 : FD_LOG_WARNING(( "misaligned shmem" ));
89 0 : return NULL;
90 0 : }
91 :
92 0 : FD_SCRATCH_ALLOC_INIT( l, shmem );
93 0 : fd_ipecho_server_t * server = FD_SCRATCH_ALLOC_APPEND( l, fd_ipecho_server_align(), sizeof(fd_ipecho_server_t) );
94 0 : void * pool = FD_SCRATCH_ALLOC_APPEND( l, conn_pool_align(), conn_pool_footprint( max_connection_cnt ) );
95 0 : server->pollfds = FD_SCRATCH_ALLOC_APPEND( l, alignof(struct pollfd), (1UL+max_connection_cnt)*sizeof(struct pollfd) );
96 :
97 0 : server->pool = conn_pool_join( conn_pool_new( pool, max_connection_cnt ) );
98 0 : FD_TEST( server->pool );
99 :
100 0 : server->sockfd = -1;
101 0 : server->epoll_fd = -1;
102 :
103 0 : for( ulong i=0UL; i<max_connection_cnt; i++ ) {
104 0 : server->pollfds[ i ].fd = -1;
105 0 : server->pollfds[ i ].events = POLLIN;
106 0 : }
107 0 : server->pollfds[ max_connection_cnt ].fd = -1;
108 :
109 0 : server->evict_idx = 0UL;
110 0 : server->max_connection_cnt = max_connection_cnt;
111 :
112 0 : memset( &server->metrics, 0, sizeof(server->metrics) );
113 :
114 0 : FD_COMPILER_MFENCE();
115 0 : FD_VOLATILE( server->magic ) = FD_IPECHO_SERVER_MAGIC;
116 0 : FD_COMPILER_MFENCE();
117 :
118 0 : return server;
119 0 : }
120 :
121 : fd_ipecho_server_t *
122 0 : fd_ipecho_server_join( void * shipe ) {
123 0 : if( FD_UNLIKELY( !shipe ) ) {
124 0 : FD_LOG_WARNING(( "NULL shipe" ));
125 0 : return NULL;
126 0 : }
127 :
128 0 : if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)shipe, fd_ipecho_server_align() ) ) ) {
129 0 : FD_LOG_WARNING(( "misaligned shipe" ));
130 0 : return NULL;
131 0 : }
132 :
133 0 : fd_ipecho_server_t * server = (fd_ipecho_server_t *)shipe;
134 :
135 0 : if( FD_UNLIKELY( server->magic!=FD_IPECHO_SERVER_MAGIC ) ) {
136 0 : FD_LOG_WARNING(( "bad magic" ));
137 0 : return NULL;
138 0 : }
139 :
140 0 : return server;
141 0 : }
142 :
143 : static void
144 0 : epoll_maybe_add_listener( fd_ipecho_server_t * server ) {
145 0 : if( FD_UNLIKELY( -1==server->sockfd ) ) return;
146 0 : if( FD_UNLIKELY( server->shred_version==0U ) ) return;
147 0 : struct epoll_event ev = { .events = EPOLLIN, .data.u64 = server->max_connection_cnt };
148 0 : if( FD_UNLIKELY( -1==epoll_ctl( server->epoll_fd, EPOLL_CTL_ADD, server->sockfd, &ev ) ) ) {
149 0 : if( FD_LIKELY( errno==EEXIST ) ) return;
150 0 : FD_LOG_ERR(( "epoll_ctl(ADD) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
151 0 : }
152 0 : }
153 :
154 : void
155 : fd_ipecho_server_init( fd_ipecho_server_t * server,
156 : int epoll_fd,
157 : uint address,
158 : ushort port,
159 0 : ushort shred_version ) {
160 :
161 0 : FD_TEST( -1!=epoll_fd );
162 0 : server->epoll_fd = epoll_fd;
163 :
164 : /* If the shred version is 0 that means that the shred version has not
165 : been set yet. */
166 0 : server->shred_version = shred_version;
167 :
168 0 : server->sockfd = socket( AF_INET, SOCK_STREAM|SOCK_NONBLOCK, 0 );
169 0 : if( FD_UNLIKELY( -1==server->sockfd ) ) FD_LOG_ERR(( "socket() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
170 :
171 0 : int optval = 1;
172 0 : if( FD_UNLIKELY( -1==setsockopt( server->sockfd, SOL_SOCKET, SO_REUSEADDR, &optval, sizeof( optval ) ) ) )
173 0 : FD_LOG_ERR(( "setsockopt failed (%i-%s)", errno, strerror( errno ) ));
174 :
175 0 : struct sockaddr_in addr = {
176 0 : .sin_family = AF_INET,
177 0 : .sin_port = fd_ushort_bswap( port ),
178 0 : .sin_addr.s_addr = address,
179 0 : };
180 :
181 0 : if( FD_UNLIKELY( -1==bind( server->sockfd, fd_type_pun( &addr ), sizeof( addr ) ) ) ) {
182 0 : FD_LOG_ERR(( "bind(%i,AF_INET," FD_IP4_ADDR_FMT ":%u) failed (%i-%s)",
183 0 : server->sockfd, FD_IP4_ADDR_FMT_ARGS( address ), port,
184 0 : errno, fd_io_strerror( errno ) ));
185 0 : }
186 0 : if( FD_UNLIKELY( -1==listen( server->sockfd, (int)server->max_connection_cnt ) ) ) FD_LOG_ERR(( "listen() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
187 :
188 0 : server->pollfds[ server->max_connection_cnt ] = (struct pollfd){ .fd = server->sockfd, .events = POLLIN, .revents = 0 };
189 0 : epoll_maybe_add_listener( server );
190 0 : }
191 :
192 : void
193 0 : fd_ipecho_server_fini( fd_ipecho_server_t * server ) {
194 0 : fd_ipecho_server_close_conns( server );
195 0 : if( FD_UNLIKELY( -1!=server->sockfd ) ) {
196 0 : FD_TEST( -1!=close( server->sockfd ) );
197 0 : server->sockfd = -1;
198 0 : server->pollfds[ server->max_connection_cnt ].fd = -1;
199 0 : }
200 0 : }
201 :
202 : void
203 0 : fd_ipecho_server_close_conns( fd_ipecho_server_t * server ) {
204 0 : for( ulong i=0UL; i<server->max_connection_cnt; i++ ) {
205 0 : if( FD_UNLIKELY( -1!=server->pollfds[ i ].fd ) ) {
206 0 : FD_TEST( -1!=close( server->pollfds[ i ].fd ) );
207 0 : server->pollfds[ i ].fd = -1;
208 0 : conn_pool_ele_release( server->pool, &server->pool[ i ] );
209 0 : }
210 0 : }
211 0 : server->metrics->connection_cnt = 0UL;
212 0 : server->evict_idx = 0UL;
213 0 : }
214 :
215 : void
216 : fd_ipecho_server_set_shred_version( fd_ipecho_server_t * server,
217 0 : ushort shred_version ) {
218 0 : server->shred_version = shred_version;
219 0 : epoll_maybe_add_listener( server );
220 0 : }
221 :
222 : static inline int
223 0 : is_expected_network_error( int err ) {
224 0 : return
225 0 : err==ENETDOWN ||
226 0 : err==EPROTO ||
227 0 : err==ENOPROTOOPT ||
228 0 : err==EHOSTDOWN ||
229 0 : err==ENONET ||
230 0 : err==EHOSTUNREACH ||
231 0 : err==EOPNOTSUPP ||
232 0 : err==ENETUNREACH ||
233 0 : err==ETIMEDOUT ||
234 0 : err==ENETRESET ||
235 0 : err==ECONNABORTED ||
236 0 : err==ECONNRESET ||
237 0 : err==EPIPE ||
238 0 : err==EPERM || /* iptables */
239 0 : err==ENOMEM; /* net stack OOM */
240 0 : }
241 :
242 : static void
243 : close_conn( fd_ipecho_server_t * server,
244 : ulong conn_idx,
245 0 : int reason ) {
246 0 : (void)reason;
247 0 : FD_TEST( server->pollfds[ conn_idx ].fd!=-1 );
248 :
249 0 : if( FD_UNLIKELY( -1==close( server->pollfds[ conn_idx ].fd ) ) ) FD_LOG_ERR(( "close failed (%i-%s)", errno, strerror( errno ) ));
250 0 : server->pollfds[ conn_idx ].fd = -1;
251 0 : conn_pool_ele_release( server->pool, &server->pool[ conn_idx ] );
252 :
253 0 : FD_TEST( server->metrics->connection_cnt );
254 0 : server->metrics->connection_cnt--;
255 0 : if( FD_UNLIKELY( reason==CLOSE_OK ) ) server->metrics->connections_closed_ok++;
256 0 : else server->metrics->connections_closed_error++;
257 0 : }
258 :
259 : static void
260 : epoll_conn_add( fd_ipecho_server_t * server,
261 0 : ulong conn_idx ) {
262 0 : struct epoll_event ev = { .events = EPOLLIN, .data.u64 = conn_idx };
263 0 : if( FD_UNLIKELY( -1==epoll_ctl( server->epoll_fd, EPOLL_CTL_ADD, server->pollfds[ conn_idx ].fd, &ev ) ) )
264 0 : FD_LOG_ERR(( "epoll_ctl(ADD) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
265 0 : server->pollfds[ conn_idx ].events = POLLIN;
266 0 : }
267 :
268 : static void
269 : epoll_update_out( fd_ipecho_server_t * server,
270 0 : ulong conn_idx ) {
271 0 : if( FD_UNLIKELY( -1==server->pollfds[ conn_idx ].fd ) ) return;
272 0 : short events = server->pool[ conn_idx ].state==STATE_WRITING ? POLLOUT : POLLIN;
273 0 : if( FD_LIKELY( server->pollfds[ conn_idx ].events==events ) ) return;
274 0 : struct epoll_event ev = {
275 0 : .events = ((events & POLLIN) ? EPOLLIN : 0U) | ((events & POLLOUT) ? EPOLLOUT : 0U),
276 0 : .data.u64 = conn_idx,
277 0 : };
278 0 : if( FD_UNLIKELY( -1==epoll_ctl( server->epoll_fd, EPOLL_CTL_MOD, server->pollfds[ conn_idx ].fd, &ev ) ) )
279 0 : FD_LOG_ERR(( "epoll_ctl(MOD) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
280 0 : server->pollfds[ conn_idx ].events = events;
281 0 : }
282 :
283 :
284 : static void
285 0 : accept_conns( fd_ipecho_server_t * server ) {
286 0 : for(;;) {
287 0 : struct sockaddr_in addr;
288 0 : socklen_t addr_len = sizeof(addr);
289 0 : int fd = accept4( server->pollfds[ server->max_connection_cnt ].fd, fd_type_pun( &addr ), &addr_len, SOCK_NONBLOCK|SOCK_CLOEXEC );
290 :
291 0 : if( FD_UNLIKELY( -1==fd ) ) {
292 0 : if( FD_LIKELY( EAGAIN==errno ) ) break;
293 0 : else if( FD_LIKELY( is_expected_network_error( errno ) ) ) continue;
294 0 : else FD_LOG_ERR(( "accept4() failed (%i-%s)", errno, strerror( errno ) ));
295 0 : }
296 :
297 0 : if( FD_UNLIKELY( !conn_pool_free( server->pool ) ) ) {
298 0 : close_conn( server, server->evict_idx, CLOSE_EVICTED );
299 0 : server->evict_idx = (server->evict_idx+1UL) % server->max_connection_cnt;
300 0 : }
301 0 : ulong conn_id = conn_pool_idx_acquire( server->pool );
302 :
303 0 : server->pollfds[ conn_id ].fd = fd;
304 0 : server->pollfds[ conn_id ].events = POLLIN;
305 0 : epoll_conn_add( server, conn_id );
306 0 : server->pool[ conn_id ].ipv4 = addr.sin_addr.s_addr;
307 0 : server->pool[ conn_id ].state = STATE_READING;
308 0 : server->pool[ conn_id ].request_bytes_read = 0UL;
309 0 : server->pool[ conn_id ].response_bytes_written = 0UL;
310 :
311 0 : server->metrics->connection_cnt++;
312 0 : }
313 0 : }
314 :
315 : static void
316 : read_conn( fd_ipecho_server_t * server,
317 0 : ulong conn_idx ) {
318 0 : fd_ipecho_server_connection_t * conn = &server->pool[ conn_idx ];
319 :
320 0 : if( FD_UNLIKELY( conn->state!=STATE_READING ) ) {
321 0 : close_conn( server, conn_idx, CLOSE_EXPECTED_EOF );
322 0 : return;
323 0 : }
324 :
325 0 : long sz = read( server->pollfds[ conn_idx ].fd, conn->request_bytes+conn->request_bytes_read, sizeof(conn->request_bytes)-conn->request_bytes_read );
326 0 : if( FD_UNLIKELY( -1==sz && errno==EAGAIN ) ) return; /* No data to read, continue. */
327 0 : else if( -1==sz && is_expected_network_error( errno ) ) {
328 0 : close_conn( server, conn_idx, CLOSE_PEER_RESET );
329 0 : return;
330 0 : }
331 0 : else if( FD_UNLIKELY( -1==sz ) ) FD_LOG_ERR(( "read failed (%i-%s)", errno, strerror( errno ) )); /* Unexpected programmer error, abort */
332 :
333 0 : if( FD_UNLIKELY( !sz && conn->request_bytes_read!=21UL ) ) {
334 0 : close_conn( server, conn_idx, CLOSE_BAD_LENGTH );
335 0 : return;
336 0 : }
337 :
338 : /* New data was read... process it */
339 0 : server->metrics->bytes_read += (ulong)sz;
340 0 : conn->request_bytes_read += (ulong)sz;
341 0 : if( FD_UNLIKELY( conn->request_bytes_read==sizeof(conn->request_bytes) ) ) {
342 0 : close_conn( server, conn_idx, CLOSE_LARGE_REQUEST );
343 0 : return;
344 0 : }
345 :
346 0 : if( FD_UNLIKELY( conn->request_bytes_read<21UL ) ) return;
347 :
348 0 : if( FD_UNLIKELY( memcmp( conn->request_bytes, "\0\0\0\0", 4UL ) ) ) {
349 0 : close_conn( server, conn_idx, CLOSE_BAD_HEADER );
350 0 : return;
351 0 : }
352 :
353 0 : if( FD_UNLIKELY( conn->request_bytes[ 20UL ]!='\n' ) ) {
354 0 : close_conn( server, conn_idx, CLOSE_BAD_TRAILER );
355 0 : return;
356 0 : }
357 :
358 0 : uchar response[ 27UL ] = {
359 0 : 0, 0, 0, 0, /* Magic */
360 0 : 0, 0, 0, 0, /* IP address variant */
361 0 : 0, 0, 0, 0, /* IP address */
362 0 : 1, /* Shred version option variant */
363 0 : 0, 0, /* Shred version */
364 0 : 0, /* [...] 12 bytes of trailing garbage, as in Agave */
365 0 : };
366 :
367 0 : FD_STORE( uint, response+8UL, conn->ipv4 );
368 0 : FD_STORE( ushort, response+13UL, server->shred_version );
369 :
370 : /* Now have a complete request ... buffer response */
371 0 : conn->state = STATE_WRITING;
372 0 : conn->response_bytes_written = 0UL;
373 0 : memcpy( conn->response_bytes, response, sizeof(response) );
374 0 : }
375 :
376 : static void
377 : write_conn( fd_ipecho_server_t * server,
378 0 : ulong conn_idx ) {
379 0 : fd_ipecho_server_connection_t * conn = &server->pool[ conn_idx ];
380 :
381 0 : if( FD_LIKELY( conn->state==STATE_READING ) ) return;
382 :
383 0 : long sz = sendto( server->pollfds[ conn_idx ].fd, conn->response_bytes+conn->response_bytes_written, sizeof(conn->response_bytes)-conn->response_bytes_written, MSG_NOSIGNAL, NULL, 0 );
384 0 : if( FD_UNLIKELY( -1==sz && errno==EAGAIN ) ) return; /* No data was written, continue. */
385 0 : if( FD_UNLIKELY( -1==sz && is_expected_network_error( errno ) ) ) {
386 0 : close_conn( server, conn_idx, CLOSE_PEER_RESET );
387 0 : return;
388 0 : }
389 0 : if( FD_UNLIKELY( -1==sz ) ) FD_LOG_ERR(( "write failed (%i-%s)", errno, strerror( errno ) )); /* Unexpected programmer error, abort */
390 :
391 0 : server->metrics->bytes_written += (ulong)sz;
392 0 : conn->response_bytes_written += (ulong)sz;
393 0 : if( FD_UNLIKELY( conn->response_bytes_written<sizeof(conn->response_bytes) ) ) return;
394 :
395 0 : close_conn( server, conn_idx, CLOSE_OK );
396 0 : }
397 :
398 :
399 : int
400 : fd_ipecho_server_epoll_poll( fd_ipecho_server_t * server,
401 0 : int * charge_busy ) {
402 0 : FD_TEST( -1!=server->epoll_fd );
403 :
404 0 : if( FD_UNLIKELY( server->shred_version==0U ) ) return 0;
405 :
406 0 : struct epoll_event evs[ 64 ];
407 0 : int nfds = epoll_pwait( server->epoll_fd, evs, 64, 0, NULL );
408 0 : if( FD_UNLIKELY( -1==nfds ) ) {
409 0 : if( FD_LIKELY( errno==EINTR ) ) return 0;
410 0 : FD_LOG_ERR(( "epoll_pwait failed (%i-%s)", errno, fd_io_strerror( errno ) ));
411 0 : }
412 0 : if( FD_LIKELY( !nfds ) ) return 0;
413 :
414 0 : *charge_busy = 1;
415 0 : for( int i=0; i<nfds; i++ ) {
416 0 : ulong conn_idx = evs[ i ].data.u64;
417 0 : if( FD_UNLIKELY( conn_idx==server->max_connection_cnt ) ) {
418 0 : accept_conns( server );
419 0 : continue;
420 0 : }
421 0 : if( FD_UNLIKELY( -1==server->pollfds[ conn_idx ].fd ) ) continue;
422 0 : if( FD_LIKELY( evs[ i ].events & (EPOLLIN|EPOLLHUP|EPOLLERR) ) ) read_conn( server, conn_idx );
423 0 : if( FD_UNLIKELY( -1==server->pollfds[ conn_idx ].fd ) ) continue;
424 0 : if( FD_LIKELY( evs[ i ].events & EPOLLOUT ) ) write_conn( server, conn_idx );
425 0 : if( FD_UNLIKELY( -1==server->pollfds[ conn_idx ].fd ) ) continue;
426 0 : epoll_update_out( server, conn_idx );
427 0 : }
428 :
429 0 : return 1;
430 0 : }
431 :
432 : fd_ipecho_server_metrics_t *
433 0 : fd_ipecho_server_metrics( fd_ipecho_server_t * server ) {
434 0 : return server->metrics;
435 0 : }
436 :
437 : int
438 0 : fd_ipecho_server_sockfd( fd_ipecho_server_t * server ) {
439 0 : return server->sockfd;
440 0 : }
|