Line data Source code
1 : #define _GNU_SOURCE /* dup3 */
2 : #include "fd_sock_tile_private.h"
3 : #include "../fd_net_common.h"
4 : #include "../../fd_disco_base.h"
5 : #include "../../../discof/repair/fd_repair.h"
6 : #include "../../topo/fd_topo.h"
7 : #include "../../../util/net/fd_eth.h"
8 : #include "../../../util/net/fd_ip4.h"
9 : #include "../../../util/net/fd_udp.h"
10 :
11 : #include <stdalign.h> /* alignof */
12 : #include <errno.h>
13 : #include <fcntl.h> /* fcntl */
14 : #include <unistd.h> /* dup3, close */
15 : #include <netinet/in.h> /* sockaddr_in */
16 : #include <sys/socket.h> /* socket */
17 : #include "../../metrics/fd_metrics.h"
18 :
19 : #include "generated/fd_sock_tile_seccomp.h"
20 :
21 : /* recv/sendmmsg packet count in batch and tango burst depth
22 : FIXME make configurable in the future?
23 : FIXME keep in sync with fd_net_tile_topo.c */
24 0 : #define STEM_BURST (64UL)
25 :
26 : /* Place RX socket file descriptors in contiguous integer range. */
27 0 : #define RX_SOCK_FD_MIN (128)
28 :
29 : /* Controls max ancillary data size.
30 : Must be aligned by alignof(struct cmsghdr) */
31 0 : #define FD_SOCK_CMSG_MAX (64UL)
32 :
33 : static ulong
34 : populate_allowed_seccomp( fd_topo_t const * topo,
35 : fd_topo_tile_t const * tile,
36 : ulong out_cnt,
37 0 : struct sock_filter * out ) {
38 0 : FD_SCRATCH_ALLOC_INIT( l, fd_topo_obj_laddr( topo, tile->tile_obj_id ) );
39 0 : fd_sock_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_sock_tile_t), sizeof(fd_sock_tile_t) );
40 :
41 0 : populate_sock_filter_policy_fd_sock_tile( out_cnt, out, (uint)fd_log_private_logfile_fd(), (uint)ctx->tx_sock, RX_SOCK_FD_MIN, RX_SOCK_FD_MIN+(uint)ctx->sock_cnt );
42 0 : return sock_filter_policy_fd_sock_tile_instr_cnt;
43 0 : }
44 :
45 : static ulong
46 : populate_allowed_fds( fd_topo_t const * topo,
47 : fd_topo_tile_t const * tile,
48 : ulong out_fds_cnt,
49 0 : int * out_fds ) {
50 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
51 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
52 0 : fd_sock_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_sock_tile_t), sizeof(fd_sock_tile_t) );
53 :
54 0 : ulong sock_cnt = ctx->sock_cnt;
55 0 : if( FD_UNLIKELY( out_fds_cnt<sock_cnt+3UL ) ) {
56 0 : FD_LOG_ERR(( "out_fds_cnt %lu", out_fds_cnt ));
57 0 : }
58 :
59 0 : ulong out_cnt = 0UL;
60 :
61 0 : out_fds[ out_cnt++ ] = 2; /* stderr */
62 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) ) {
63 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
64 0 : }
65 0 : out_fds[ out_cnt++ ] = ctx->tx_sock;
66 0 : for( ulong j=0UL; j<sock_cnt; j++ ) {
67 0 : out_fds[ out_cnt++ ] = ctx->pollfd[ j ].fd;
68 0 : }
69 0 : return out_cnt;
70 0 : }
71 :
72 : FD_FN_CONST static inline ulong
73 0 : tx_scratch_footprint( void ) {
74 0 : return STEM_BURST * fd_ulong_align_up( FD_NET_MTU, FD_CHUNK_ALIGN );
75 0 : }
76 :
77 : FD_FN_CONST static inline ulong
78 0 : scratch_align( void ) {
79 0 : return 4096UL;
80 0 : }
81 :
82 : FD_FN_PURE static inline ulong
83 0 : scratch_footprint( fd_topo_tile_t const * tile FD_PARAM_UNUSED ) {
84 0 : ulong l = FD_LAYOUT_INIT;
85 0 : l = FD_LAYOUT_APPEND( l, alignof(fd_sock_tile_t), sizeof(fd_sock_tile_t) );
86 0 : l = FD_LAYOUT_APPEND( l, alignof(struct iovec), STEM_BURST*sizeof(struct iovec) );
87 0 : l = FD_LAYOUT_APPEND( l, alignof(struct cmsghdr), STEM_BURST*FD_SOCK_CMSG_MAX );
88 0 : l = FD_LAYOUT_APPEND( l, alignof(struct sockaddr_in), STEM_BURST*sizeof(struct sockaddr_in) );
89 0 : l = FD_LAYOUT_APPEND( l, alignof(struct mmsghdr), STEM_BURST*sizeof(struct mmsghdr) );
90 0 : l = FD_LAYOUT_APPEND( l, FD_CHUNK_ALIGN, tx_scratch_footprint() );
91 0 : return FD_LAYOUT_FINI( l, scratch_align() );
92 0 : }
93 :
94 : /* create_udp_socket creates and configures a new UDP socket for the
95 : sock tile at the given file descriptor ID. */
96 :
97 : static void
98 : create_udp_socket( int sock_fd,
99 : uint bind_addr,
100 : ushort udp_port,
101 0 : int so_rcvbuf ) {
102 :
103 0 : if( fcntl( sock_fd, F_GETFD, 0 )!=-1 ) {
104 0 : FD_LOG_ERR(( "file descriptor %d already exists", sock_fd ));
105 0 : } else if( errno!=EBADF ) {
106 0 : FD_LOG_ERR(( "fcntl(F_GETFD) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
107 0 : }
108 :
109 0 : int orig_fd = socket( AF_INET, SOCK_DGRAM, IPPROTO_UDP );
110 0 : if( FD_UNLIKELY( orig_fd<0 ) ) {
111 0 : FD_LOG_ERR(( "socket(AF_INET,SOCK_DGRAM,IPPROTO_UDP) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
112 0 : }
113 :
114 0 : int reuseport = 1;
115 0 : if( FD_UNLIKELY( setsockopt( orig_fd, SOL_SOCKET, SO_REUSEPORT, &reuseport, sizeof(int) )<0 ) ) {
116 0 : FD_LOG_ERR(( "setsockopt(SOL_SOCKET,SO_REUSEPORT,1) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
117 0 : }
118 :
119 0 : int ip_pktinfo = 1;
120 0 : if( FD_UNLIKELY( setsockopt( orig_fd, IPPROTO_IP, IP_PKTINFO, &ip_pktinfo, sizeof(int) )<0 ) ) {
121 0 : FD_LOG_ERR(( "setsockopt(IPPROTO_IP,IP_PKTINFO,1) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
122 0 : }
123 :
124 0 : if( FD_UNLIKELY( 0!=setsockopt( orig_fd, SOL_SOCKET, SO_RCVBUF, &so_rcvbuf, sizeof(int) ) ) ) {
125 0 : FD_LOG_ERR(( "setsockopt(SOL_SOCKET,SO_RCVBUF,%i) failed (%i-%s)", so_rcvbuf, errno, fd_io_strerror( errno ) ));
126 0 : }
127 :
128 0 : struct sockaddr_in saddr = {
129 0 : .sin_family = AF_INET,
130 0 : .sin_addr.s_addr = bind_addr,
131 0 : .sin_port = fd_ushort_bswap( udp_port ),
132 0 : };
133 0 : if( FD_UNLIKELY( 0!=bind( orig_fd, fd_type_pun_const( &saddr ), sizeof(struct sockaddr_in) ) ) ) {
134 0 : FD_LOG_ERR(( "bind(0.0.0.0:%i) failed (%i-%s)", udp_port, errno, fd_io_strerror( errno ) ));
135 0 : }
136 :
137 0 : # if defined(__linux__)
138 0 : int dup_res = dup3( orig_fd, sock_fd, O_CLOEXEC );
139 : # else
140 : int dup_res = dup2( orig_fd, sock_fd );
141 : # endif
142 0 : if( FD_UNLIKELY( dup_res!=sock_fd ) ) {
143 0 : FD_LOG_ERR(( "dup2 returned %i (%i-%s)", sock_fd, errno, fd_io_strerror( errno ) ));
144 0 : }
145 :
146 0 : if( FD_UNLIKELY( 0!=close( orig_fd ) ) ) {
147 0 : FD_LOG_ERR(( "close(%d) failed (%i-%s)", orig_fd, errno, fd_io_strerror( errno ) ));
148 0 : }
149 :
150 0 : }
151 :
152 : static void
153 : privileged_init( fd_topo_t const * topo,
154 0 : fd_topo_tile_t const * tile ) {
155 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
156 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
157 0 : fd_sock_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_sock_tile_t), sizeof(fd_sock_tile_t) );
158 0 : struct iovec * batch_iov = FD_SCRATCH_ALLOC_APPEND( l, alignof(struct iovec), STEM_BURST*sizeof(struct iovec) );
159 0 : void * batch_cmsg = FD_SCRATCH_ALLOC_APPEND( l, alignof(struct cmsghdr), STEM_BURST*FD_SOCK_CMSG_MAX );
160 0 : struct sockaddr_in * batch_sa = FD_SCRATCH_ALLOC_APPEND( l, alignof(struct sockaddr_in), STEM_BURST*sizeof(struct sockaddr_in) );
161 0 : struct mmsghdr * batch_msg = FD_SCRATCH_ALLOC_APPEND( l, alignof(struct mmsghdr), STEM_BURST*sizeof(struct mmsghdr) );
162 0 : uchar * tx_scratch = FD_SCRATCH_ALLOC_APPEND( l, FD_CHUNK_ALIGN, tx_scratch_footprint() );
163 0 : FD_DCHECK_CRIT( scratch==ctx, "invalid layout" );
164 :
165 0 : fd_memset( ctx, 0, sizeof(fd_sock_tile_t) );
166 0 : fd_memset( batch_iov, 0, STEM_BURST*sizeof(struct iovec) );
167 0 : fd_memset( batch_sa, 0, STEM_BURST*sizeof(struct sockaddr_in) );
168 0 : fd_memset( batch_msg, 0, STEM_BURST*sizeof(struct mmsghdr) );
169 :
170 0 : ctx->batch_cnt = 0UL;
171 0 : ctx->batch_iov = batch_iov;
172 0 : ctx->batch_cmsg = batch_cmsg;
173 0 : ctx->batch_sa = batch_sa;
174 0 : ctx->batch_msg = batch_msg;
175 0 : ctx->tx_scratch0 = tx_scratch;
176 0 : ctx->tx_scratch1 = tx_scratch + tx_scratch_footprint();
177 0 : ctx->tx_ptr = tx_scratch;
178 0 : ctx->repair_shred_sock_idx = UINT_MAX;
179 :
180 : /* Create receive sockets. Incrementally assign them to file
181 : descriptors starting at sock_fd_min. */
182 :
183 0 : int sock_fd_min = RX_SOCK_FD_MIN;
184 0 : ushort udp_port_candidates[] = {
185 0 : (ushort)tile->sock.net.legacy_transaction_listen_port,
186 0 : (ushort)tile->sock.net.quic_transaction_listen_port,
187 0 : (ushort)tile->sock.net.shred_listen_port,
188 0 : (ushort)tile->sock.net.gossip_listen_port,
189 0 : (ushort)tile->sock.net.repair_client_listen_port,
190 0 : (ushort)tile->sock.net.repair_serve_listen_port,
191 0 : (ushort)tile->sock.net.txsend_src_port,
192 0 : (ushort)tile->sock.net.votor_quic_client_listen_port,
193 0 : (ushort)tile->sock.net.votor_quic_server_listen_port
194 0 : };
195 0 : static char const * udp_port_links[] = {
196 0 : "net_quic", /* legacy_transaction_listen_port */
197 0 : "net_quic", /* quic_transaction_listen_port */
198 0 : "net_shred", /* shred_listen_port (turbine) */
199 0 : "net_gossvf", /* gossip_listen_port */
200 0 : "net_shred", /* shred_listen_port (repair) */
201 0 : "net_rserve", /* repair_serve_listen_port */
202 0 : "net_txsend", /* txsend_src_port */
203 0 : "net_votor", /* votor_quic_client_listen_port */
204 0 : "net_votor" /* votor_quic_server_listen_port */
205 0 : };
206 0 : static uchar const udp_port_protos[] = {
207 0 : DST_PROTO_TPU_UDP, /* legacy_transaction_listen_port */
208 0 : DST_PROTO_TPU_QUIC, /* quic_transaction_listen_port */
209 0 : DST_PROTO_SHRED, /* shred_listen_port (turbine) */
210 0 : DST_PROTO_GOSSIP, /* gossip_listen_port */
211 0 : DST_PROTO_REPAIR, /* shred_listen_port (repair) */
212 0 : DST_PROTO_RSERVE, /* repair_serve_listen_port */
213 0 : DST_PROTO_SEND, /* send_src_port */
214 0 : DST_PROTO_VOTOR, /* votor_quic_client_listen_port */
215 0 : DST_PROTO_VOTOR /* votor_quic_server_listen_port */
216 0 : };
217 0 : for( uint candidate_idx=0U; candidate_idx<9; candidate_idx++ ) {
218 0 : if( !udp_port_candidates[ candidate_idx ] ) continue;
219 0 : uint sock_idx = ctx->sock_cnt;
220 0 : if( sock_idx>=FD_SOCK_TILE_MAX_SOCKETS ) FD_LOG_ERR(( "too many sockets" ));
221 0 : ushort port = (ushort)udp_port_candidates[ candidate_idx ];
222 :
223 0 : char const * target_link = udp_port_links[ candidate_idx ];
224 0 : ctx->link_rx_map[ sock_idx ] = 0xFF;
225 0 : for( ulong j=0UL; j<(tile->out_cnt); j++ ) {
226 0 : if( 0==strcmp( topo->links[ tile->out_link_id[ j ] ].name, target_link ) ) {
227 0 : ctx->proto_id [ sock_idx ] = (uchar)udp_port_protos[ candidate_idx ];
228 0 : ctx->link_rx_map [ sock_idx ] = (uchar)j;
229 0 : ctx->rx_sock_port[ sock_idx ] = (ushort)port;
230 0 : break;
231 0 : }
232 0 : }
233 0 : if( ctx->link_rx_map[ sock_idx ]==0xFF ) {
234 : /* listen port number has no associated links,
235 : i.e. the repair server is disabled, then no net_rserve link. */
236 0 : continue;
237 0 : }
238 :
239 : /* Record the socket index of the repair intake socket, so repair
240 : ping packets can be routed to the repair tile at runtime. */
241 0 : if( tile->sock.net.repair_client_listen_port &&
242 0 : udp_port_candidates[ candidate_idx ]==tile->sock.net.repair_client_listen_port )
243 0 : ctx->repair_shred_sock_idx = sock_idx;
244 :
245 0 : int sock_fd = sock_fd_min + (int)sock_idx;
246 0 : create_udp_socket( sock_fd, tile->sock.net.bind_address, port, tile->sock.so_rcvbuf );
247 0 : ctx->pollfd[ sock_idx ].fd = sock_fd;
248 0 : ctx->pollfd[ sock_idx ].events = POLLIN;
249 0 : ctx->sock_cnt++;
250 0 : }
251 :
252 : /* Create transmit socket */
253 :
254 0 : int tx_sock = socket( AF_INET, SOCK_RAW|SOCK_CLOEXEC, FD_IP4_HDR_PROTOCOL_UDP );
255 0 : if( FD_UNLIKELY( tx_sock<0 ) ) {
256 0 : FD_LOG_ERR(( "socket(AF_INET,SOCK_RAW|SOCK_CLOEXEC,17) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
257 0 : }
258 :
259 0 : if( FD_UNLIKELY( 0!=setsockopt( tx_sock, SOL_SOCKET, SO_SNDBUF, &tile->sock.so_sndbuf, sizeof(int) ) ) ) {
260 0 : FD_LOG_ERR(( "setsockopt(SOL_SOCKET,SO_SNDBUF,%i) failed (%i-%s)", tile->sock.so_sndbuf, errno, fd_io_strerror( errno ) ));
261 0 : }
262 :
263 0 : uchar mcast_ttl = 64;
264 0 : if( FD_UNLIKELY( 0!=setsockopt( tx_sock, IPPROTO_IP, IP_MULTICAST_TTL, &mcast_ttl, sizeof(mcast_ttl) ) ) ) {
265 0 : FD_LOG_ERR(( "setsockopt(IPPROTO_IP,IP_MULTICAST_TTL,%u) failed (%i-%s)", (uint)mcast_ttl, errno, fd_io_strerror( errno ) ));
266 0 : }
267 :
268 0 : ctx->tx_sock = tx_sock;
269 0 : ctx->bind_address = tile->sock.net.bind_address;
270 0 : }
271 :
272 : static void
273 : unprivileged_init( fd_topo_t const * topo,
274 0 : fd_topo_tile_t const * tile ) {
275 0 : fd_sock_tile_t * ctx = fd_topo_obj_laddr( topo, tile->tile_obj_id );
276 :
277 0 : if( FD_UNLIKELY( tile->out_cnt > MAX_NET_OUTS ) ) {
278 0 : FD_LOG_ERR(( "sock tile has %lu out links which exceeds the max (%lu)", tile->out_cnt, MAX_NET_OUTS ));
279 0 : }
280 :
281 0 : ctx->repair_rx = 0xFF;
282 0 : for( ulong i=0UL; i<(tile->out_cnt); i++ ) {
283 0 : if( 0!=strncmp( topo->links[ tile->out_link_id[ i ] ].name, "net_", 4 ) ) {
284 0 : FD_LOG_ERR(( "out link %lu is not a net RX link", i ));
285 0 : }
286 0 : if( 0==strcmp( topo->links[ tile->out_link_id[ i ] ].name, "net_repair" ) ) {
287 0 : if( FD_UNLIKELY( ctx->repair_rx!=0xFF ) ) FD_LOG_ERR(( "multiple net_repair out links" ));
288 0 : ctx->repair_rx = (uchar)i;
289 0 : }
290 0 : fd_topo_link_t const * link = &topo->links[ tile->out_link_id[ i ] ];
291 0 : ctx->link_rx[ i ].base = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
292 0 : ctx->link_rx[ i ].chunk0 = fd_dcache_compact_chunk0( ctx->link_rx[ i ].base, link->dcache );
293 0 : ctx->link_rx[ i ].wmark = fd_dcache_compact_wmark( ctx->link_rx[ i ].base, link->dcache, link->mtu );
294 0 : ctx->link_rx[ i ].chunk = ctx->link_rx[ i ].chunk0;
295 0 : if( FD_UNLIKELY( link->burst < STEM_BURST ) ) {
296 0 : FD_LOG_ERR(( "link %lu dcache burst is too low (%lu<%lu)",
297 0 : tile->out_link_id[ i ], link->burst, STEM_BURST ));
298 0 : }
299 0 : }
300 :
301 0 : if( FD_UNLIKELY( ctx->repair_shred_sock_idx!=UINT_MAX && ctx->repair_rx==0xFF ) ) {
302 0 : FD_LOG_ERR(( "repair intake socket configured but no net_repair out link was found" ));
303 0 : }
304 :
305 0 : FD_CHECK_ERR( tile->in_cnt<=MAX_NET_INS, "too many input links" );
306 :
307 0 : for( ulong i=0UL; i<(tile->in_cnt); i++ ) {
308 0 : if( !strstr( topo->links[ tile->in_link_id[ i ] ].name, "_net" ) ) {
309 0 : FD_LOG_ERR(( "in link %lu is not a net TX link", i ));
310 0 : }
311 0 : fd_topo_link_t const * link = &topo->links[ tile->in_link_id[ i ] ];
312 0 : ctx->link_tx[ i ].base = topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ].wksp;
313 0 : ctx->link_tx[ i ].chunk0 = fd_dcache_compact_chunk0( ctx->link_tx[ i ].base, link->dcache );
314 0 : ctx->link_tx[ i ].wmark = fd_dcache_compact_wmark( ctx->link_tx[ i ].base, link->dcache, link->mtu );
315 0 : }
316 :
317 0 : }
318 :
319 : /* RX PATH (socket->tango) ********************************************/
320 :
321 : /* FIXME Pace RX polling and interleave it with TX jobs to reduce TX
322 : tail latency */
323 :
324 : /* poll_rx_socket does one recvmmsg batch receive on the given socket
325 : index. Returns the number of packets returned by recvmmsg. */
326 :
327 : static ulong
328 : poll_rx_socket( fd_sock_tile_t * ctx,
329 : fd_stem_context_t * stem,
330 : uint sock_idx,
331 : int sock_fd,
332 0 : ushort proto ) {
333 0 : ulong hdr_sz = sizeof(fd_eth_hdr_t) + sizeof(fd_ip4_hdr_t) + sizeof(fd_udp_hdr_t);
334 0 : ulong payload_max = FD_NET_MTU-hdr_sz;
335 0 : uchar rx_link = ctx->link_rx_map[ sock_idx ];
336 0 : ushort dport = ctx->rx_sock_port[ sock_idx ];
337 :
338 0 : fd_sock_link_rx_t * link = ctx->link_rx + rx_link;
339 0 : void * const base = link->base;
340 0 : ulong const chunk0 = link->chunk0;
341 0 : ulong const wmark = link->wmark;
342 0 : ulong chunk_next = link->chunk;
343 0 : uchar * cmsg_next = ctx->batch_cmsg;
344 :
345 0 : for( ulong j=0UL; j<STEM_BURST; j++ ) {
346 0 : ctx->batch_iov[ j ].iov_base = (uchar *)fd_chunk_to_laddr( base, chunk_next ) + hdr_sz;
347 0 : ctx->batch_iov[ j ].iov_len = payload_max;
348 0 : ctx->batch_msg[ j ].msg_hdr = (struct msghdr) {
349 0 : .msg_iov = ctx->batch_iov+j,
350 0 : .msg_iovlen = 1,
351 0 : .msg_name = ctx->batch_sa+j,
352 0 : .msg_namelen = sizeof(struct sockaddr_in),
353 0 : .msg_control = cmsg_next,
354 0 : .msg_controllen = FD_SOCK_CMSG_MAX,
355 0 : };
356 0 : cmsg_next += FD_SOCK_CMSG_MAX;
357 : /* Speculatively prepare all chunk indexes for a receive.
358 : At function exit, chunks into which a packet was received are
359 : committed, all others are freed. */
360 0 : chunk_next = fd_dcache_compact_next( chunk_next, FD_NET_MTU, chunk0, wmark );
361 0 : }
362 :
363 0 : int msg_cnt = recvmmsg( sock_fd, ctx->batch_msg, STEM_BURST, MSG_DONTWAIT, NULL );
364 0 : if( FD_UNLIKELY( msg_cnt<0 ) ) {
365 0 : if( FD_LIKELY( errno==EAGAIN ) ) return 0UL;
366 : /* unreachable if socket is in a valid state */
367 0 : FD_LOG_ERR(( "recvmmsg failed (%i-%s)", errno, fd_io_strerror( errno ) ));
368 0 : }
369 0 : long ts = fd_tickcount();
370 0 : ctx->metrics.sys_recvmmsg_cnt++;
371 :
372 0 : if( FD_UNLIKELY( msg_cnt==0 ) ) return 0UL;
373 :
374 : /* Track the chunk index of the last frag populated, so we can derive
375 : the chunk indexes for the next poll_rx_socket call.
376 : Guaranteed to be set since msg_cnt>0. */
377 0 : ulong last_chunk;
378 :
379 0 : for( ulong j=0; j<(ulong)msg_cnt; j++ ) {
380 0 : uchar * payload = ctx->batch_iov[ j ].iov_base;
381 0 : ulong payload_sz = ctx->batch_msg[ j ].msg_len;
382 0 : struct sockaddr_in * sa = ctx->batch_msg[ j ].msg_hdr.msg_name;
383 0 : ulong frame_sz = payload_sz + hdr_sz;
384 0 : ctx->metrics.rx_bytes_total += frame_sz;
385 0 : if( FD_UNLIKELY( sa->sin_family!=AF_INET ) ) {
386 : /* unreachable */
387 0 : FD_LOG_ERR(( "Received packet with unexpected sin_family %i", sa->sin_family ));
388 0 : }
389 :
390 0 : long daddr = -1;
391 0 : struct cmsghdr * cmsg = CMSG_FIRSTHDR( &ctx->batch_msg[ j ].msg_hdr );
392 0 : if( FD_LIKELY( cmsg ) ) {
393 0 : do {
394 0 : if( FD_LIKELY( (cmsg->cmsg_level==IPPROTO_IP) &
395 0 : (cmsg->cmsg_type ==IP_PKTINFO) ) ) {
396 0 : struct in_pktinfo const * pi = (struct in_pktinfo const *)CMSG_DATA( cmsg );
397 0 : daddr = pi->ipi_addr.s_addr;
398 0 : }
399 0 : cmsg = CMSG_NXTHDR( &ctx->batch_msg[ j ].msg_hdr, cmsg );
400 0 : } while( FD_UNLIKELY( cmsg ) ); /* optimize for 1 cmsg */
401 0 : }
402 0 : if( FD_UNLIKELY( daddr<0L ) ) {
403 : /* unreachable because IP_PKTINFO was set */
404 0 : FD_LOG_ERR(( "Missing IP_PKTINFO on incoming packet" ));
405 0 : }
406 :
407 0 : fd_eth_hdr_t * eth_hdr = (fd_eth_hdr_t *)( payload-42UL );
408 0 : fd_ip4_hdr_t * ip_hdr = (fd_ip4_hdr_t *)( payload-28UL );
409 0 : fd_udp_hdr_t * udp_hdr = (fd_udp_hdr_t *)( payload- 8UL );
410 0 : memset( eth_hdr->dst, 0, 6 );
411 0 : memset( eth_hdr->src, 0, 6 );
412 0 : eth_hdr->net_type = fd_ushort_bswap( FD_ETH_HDR_TYPE_IP );
413 0 : *ip_hdr = (fd_ip4_hdr_t) {
414 0 : .verihl = FD_IP4_VERIHL( 4, 5 ),
415 0 : .net_tot_len = fd_ushort_bswap( (ushort)( payload_sz+28UL ) ),
416 0 : .ttl = 1,
417 0 : .protocol = FD_IP4_HDR_PROTOCOL_UDP,
418 0 : };
419 0 : uint daddr_ = (uint)(ulong)daddr;
420 0 : memcpy( ip_hdr->saddr_c, &sa->sin_addr.s_addr, 4 );
421 0 : memcpy( ip_hdr->daddr_c, &daddr_, 4 );
422 0 : *udp_hdr = (fd_udp_hdr_t) {
423 0 : .net_sport = sa->sin_port,
424 0 : .net_dport = (ushort)fd_ushort_bswap( (ushort)dport ),
425 0 : .net_len = (ushort)fd_ushort_bswap( (ushort)( payload_sz+8UL ) ),
426 0 : .check = 0
427 0 : };
428 :
429 0 : ctx->metrics.rx_pkt_cnt++;
430 0 : ulong chunk = fd_laddr_to_chunk( base, eth_hdr );
431 0 : ulong sig = fd_disco_netmux_sig( sa->sin_addr.s_addr, fd_ushort_bswap( sa->sin_port ), sa->sin_addr.s_addr, proto, hdr_sz );
432 0 : ulong tspub = fd_frag_meta_ts_comp( ts );
433 :
434 : /* When a message arrives on the repair intake port, it is sent to
435 : the shred tile, unless it is a ping message or an alpenglow
436 : repair response (identified by the frame size), then it is sent
437 : to the repair tile. The repair tile does not own any sockets, so
438 : we look up the net_repair link directly. */
439 0 : if( FD_UNLIKELY( sock_idx==ctx->repair_shred_sock_idx && payload_sz<=AG_REPAIR_RESPONSE_MAX_SZ ) ) {
440 0 : fd_sock_link_rx_t * repair_link = ctx->link_rx + ctx->repair_rx;
441 0 : uchar * repair_buf = fd_chunk_to_laddr( repair_link->base, repair_link->chunk );
442 0 : memcpy( repair_buf, eth_hdr, frame_sz );
443 0 : fd_stem_publish( stem, ctx->repair_rx, sig, repair_link->chunk, frame_sz, 0UL, 0UL, tspub );
444 0 : repair_link->chunk = fd_dcache_compact_next( repair_link->chunk, FD_NET_MTU, repair_link->chunk0, repair_link->wmark );
445 0 : } else {
446 0 : fd_stem_publish( stem, rx_link, sig, chunk, frame_sz, 0UL, 0UL, tspub );
447 0 : }
448 :
449 0 : last_chunk = chunk;
450 0 : }
451 :
452 : /* Rewind the chunk index to the first free index. */
453 0 : link->chunk = fd_dcache_compact_next( last_chunk, FD_NET_MTU, chunk0, wmark );
454 0 : return (ulong)msg_cnt;
455 0 : }
456 :
457 : static ulong
458 : poll_rx( fd_sock_tile_t * ctx,
459 0 : fd_stem_context_t * stem ) {
460 0 : ulong pkt_cnt = 0UL;
461 0 : if( FD_UNLIKELY( ctx->batch_cnt ) ) {
462 0 : FD_LOG_ERR(( "Batch is not clean" ));
463 0 : }
464 0 : ctx->tx_idle_cnt = 0; /* restart TX polling */
465 0 : if( FD_UNLIKELY( fd_syscall_poll( ctx->pollfd, ctx->sock_cnt, 0 )<0 ) ) {
466 0 : FD_LOG_ERR(( "fd_syscall_poll failed (%i-%s)", errno, fd_io_strerror( errno ) ));
467 0 : }
468 0 : for( uint j=0UL; j<ctx->sock_cnt; j++ ) {
469 0 : if( ctx->pollfd[ j ].revents & (POLLIN|POLLERR) ) {
470 0 : pkt_cnt += poll_rx_socket(
471 0 : ctx,
472 0 : stem,
473 0 : j,
474 0 : ctx->pollfd[ j ].fd,
475 0 : ctx->proto_id[ j ]
476 0 : );
477 0 : }
478 0 : ctx->pollfd[ j ].revents = 0;
479 0 : }
480 0 : return pkt_cnt;
481 0 : }
482 :
483 : /* TX PATH (tango->socket) ********************************************/
484 :
485 : static void
486 0 : flush_tx_batch( fd_sock_tile_t * ctx ) {
487 0 : ulong batch_cnt = ctx->batch_cnt;
488 0 : for( int j = 0; j < (int)batch_cnt; /* incremented in loop */ ) {
489 0 : int remain = (int)batch_cnt - j;
490 0 : int send_cnt = sendmmsg( ctx->tx_sock, ctx->batch_msg + j, (uint)remain, MSG_DONTWAIT );
491 0 : if( send_cnt>=0 ) {
492 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_NO_ERROR_IDX ]++;
493 0 : }
494 0 : if( FD_UNLIKELY( send_cnt < remain ) ) {
495 0 : ctx->metrics.tx_drop_cnt++;
496 0 : if( FD_UNLIKELY( send_cnt < 0 ) ) {
497 0 : switch( errno ) {
498 0 : case EAGAIN:
499 0 : case ENOBUFS:
500 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_SLOW_IDX ]++;
501 0 : break;
502 0 : case EPERM:
503 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_PERMISSION_IDX ]++;
504 0 : break;
505 0 : case ENETUNREACH:
506 0 : case EHOSTUNREACH:
507 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_UNREACHABLE_IDX ]++;
508 0 : break;
509 0 : case ENONET:
510 0 : case ENETDOWN:
511 0 : case EHOSTDOWN:
512 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_DOWN_IDX ]++;
513 0 : break;
514 0 : default:
515 0 : ctx->metrics.sys_sendmmsg_cnt[ FD_METRICS_ENUM_SOCKET_ERROR_V_OTHER_IDX ]++;
516 : /* log with NOTICE, since flushing has a significant negative performance impact */
517 0 : FD_LOG_NOTICE(( "sendmmsg failed (%i-%s)", errno, fd_io_strerror( errno ) ));
518 0 : }
519 :
520 : /* first message failed, so skip failing message and continue */
521 0 : j++;
522 0 : } else {
523 : /* send_cnt succeeded, so skip those and also the failing message */
524 0 : j += send_cnt + 1;
525 :
526 : /* add the successful count */
527 0 : ctx->metrics.tx_pkt_cnt += (ulong)send_cnt;
528 0 : }
529 :
530 0 : continue;
531 0 : }
532 :
533 : /* send_cnt == batch_cnt, so we sent everything */
534 0 : ctx->metrics.tx_pkt_cnt += (ulong)send_cnt;
535 0 : break;
536 0 : }
537 :
538 0 : ctx->tx_ptr = ctx->tx_scratch0;
539 0 : ctx->batch_cnt = 0;
540 0 : }
541 :
542 : /* before_frag is called when a new frag has been detected. The sock
543 : tile can do early filtering here in the future. For example, it may
544 : want to install routing logic here to take turns with an XDP tile.
545 : (Fast path with slow fallback) */
546 :
547 : static inline int
548 : before_frag( fd_sock_tile_t * ctx FD_PARAM_UNUSED,
549 : ulong in_idx FD_PARAM_UNUSED,
550 : ulong seq FD_PARAM_UNUSED,
551 0 : ulong sig ) {
552 0 : ulong proto = fd_disco_netmux_sig_proto( sig );
553 0 : if( FD_UNLIKELY( proto!=DST_PROTO_OUTGOING ) ) return 1;
554 0 : return 0; /* continue */
555 0 : }
556 :
557 : /* during_frag is called when a new frag passed early filtering.
558 : Speculatively copies data into a sendmmsg buffer. (If all tiles
559 : respect backpressure could eliminate this copy) */
560 :
561 : static inline void
562 : during_frag( fd_sock_tile_t * ctx,
563 : ulong in_idx,
564 : ulong seq FD_PARAM_UNUSED,
565 : ulong sig,
566 : ulong chunk,
567 : ulong sz,
568 0 : ulong ctl FD_PARAM_UNUSED ) {
569 0 : if( FD_UNLIKELY( chunk<ctx->link_tx[ in_idx ].chunk0 || chunk>ctx->link_tx[ in_idx ].wmark || sz>FD_NET_MTU ) ) {
570 0 : FD_LOG_ERR(( "chunk %lu %lu corrupt, not in range [%lu,%lu]", chunk, sz, ctx->link_tx[ in_idx ].chunk0, ctx->link_tx[ in_idx ].wmark ));
571 0 : }
572 :
573 0 : ctx->parsed.invalid = 0;
574 :
575 0 : ulong const hdr_min = sizeof(fd_eth_hdr_t)+sizeof(fd_ip4_hdr_t)+sizeof(fd_udp_hdr_t);
576 0 : if( FD_UNLIKELY( sz<hdr_min ) ) {
577 : /* FIXME support ICMP messages in the future?
578 : Defer the error to after_frag, where we know we weren't
579 : overrun. */
580 0 : ctx->parsed.invalid = 1;
581 0 : return;
582 0 : }
583 :
584 0 : uchar const * frame = fd_chunk_to_laddr_const( ctx->link_tx[ in_idx ].base, chunk );
585 0 : ulong hdr_sz = fd_disco_netmux_sig_hdr_sz( sig );
586 0 : uchar const * payload = frame+hdr_sz;
587 0 : if( FD_UNLIKELY( hdr_sz>sz || hdr_sz<hdr_min ) ) {
588 0 : FD_LOG_ERR(( "packet from in_idx=%lu corrupt: hdr_sz=%lu total_sz=%lu",
589 0 : in_idx, hdr_sz, sz ));
590 0 : }
591 0 : ulong payload_sz = sz-hdr_sz;
592 :
593 0 : fd_ip4_hdr_t const * ip_hdr = (fd_ip4_hdr_t const *)( frame +sizeof(fd_eth_hdr_t) );
594 0 : fd_udp_hdr_t const * udp_hdr = (fd_udp_hdr_t const *)( payload-sizeof(fd_udp_hdr_t) );
595 0 : ctx->parsed.ip_version = FD_IP4_GET_VERSION( *ip_hdr );
596 0 : ctx->parsed.ip_protocol = ip_hdr->protocol;
597 0 : if( FD_UNLIKELY( ( ctx->parsed.ip_version !=4 ) |
598 0 : ( ctx->parsed.ip_protocol!=FD_IP4_HDR_PROTOCOL_UDP ) ) ) {
599 : /* sock tile only supports IPv4 UDP for now. Defer the error to
600 : after_frag, where we know we weren't overrun. */
601 0 : ctx->parsed.invalid = 1;
602 0 : return;
603 0 : }
604 :
605 0 : ulong msg_sz = sizeof(fd_udp_hdr_t) + payload_sz;
606 :
607 0 : ulong batch_idx = ctx->batch_cnt;
608 0 : FD_DCHECK_CRIT( batch_idx<STEM_BURST, "flow control error" );
609 0 : struct mmsghdr * msg = ctx->batch_msg + batch_idx;
610 0 : struct sockaddr_in * sa = ctx->batch_sa + batch_idx;
611 0 : struct iovec * iov = ctx->batch_iov + batch_idx;
612 0 : struct cmsghdr * cmsg = (void *)( (ulong)ctx->batch_cmsg + batch_idx*FD_SOCK_CMSG_MAX );
613 0 : uchar * buf = ctx->tx_ptr;
614 :
615 0 : *iov = (struct iovec) {
616 0 : .iov_base = buf,
617 0 : .iov_len = msg_sz,
618 0 : };
619 0 : sa->sin_family = AF_INET;
620 0 : sa->sin_addr.s_addr = FD_LOAD( uint, ip_hdr->daddr_c );
621 0 : sa->sin_port = 0; /* ignored */
622 :
623 0 : cmsg->cmsg_level = IPPROTO_IP;
624 0 : cmsg->cmsg_type = IP_PKTINFO;
625 0 : cmsg->cmsg_len = CMSG_LEN( sizeof(struct in_pktinfo) );
626 0 : struct in_pktinfo * pi = (struct in_pktinfo *)CMSG_DATA( cmsg );
627 0 : pi->ipi_ifindex = 0;
628 0 : pi->ipi_addr.s_addr = 0;
629 0 : pi->ipi_spec_dst.s_addr = fd_uint_if( !!ip_hdr->saddr, ip_hdr->saddr, ctx->bind_address );
630 :
631 0 : *msg = (struct mmsghdr) {
632 0 : .msg_hdr = {
633 0 : .msg_name = sa,
634 0 : .msg_namelen = sizeof(struct sockaddr_in),
635 0 : .msg_iov = iov,
636 0 : .msg_iovlen = 1,
637 0 : .msg_control = cmsg,
638 0 : .msg_controllen = CMSG_LEN( sizeof(struct in_pktinfo) )
639 0 : }
640 0 : };
641 :
642 0 : memcpy( buf, udp_hdr, sizeof(fd_udp_hdr_t) );
643 0 : fd_memcpy( buf+sizeof(fd_udp_hdr_t), payload, payload_sz );
644 0 : ctx->metrics.tx_bytes_total += sz;
645 0 : }
646 :
647 : /* after_frag is called when a frag was copied into a sendmmsg buffer. */
648 :
649 : static void
650 : after_frag( fd_sock_tile_t * ctx,
651 : ulong in_idx,
652 : ulong seq FD_PARAM_UNUSED,
653 : ulong sig FD_PARAM_UNUSED,
654 : ulong sz,
655 : ulong tsorig FD_PARAM_UNUSED,
656 : ulong tspub FD_PARAM_UNUSED,
657 0 : fd_stem_context_t * stem FD_PARAM_UNUSED ) {
658 : /* Commit the packet added in during_frag. during_frag defers fatal
659 : errors to here (where we know we weren't overrun) by setting the
660 : invalid flag. */
661 :
662 0 : if( FD_UNLIKELY( ctx->parsed.invalid ) ) {
663 0 : ulong const hdr_min = sizeof(fd_eth_hdr_t)+sizeof(fd_ip4_hdr_t)+sizeof(fd_udp_hdr_t);
664 0 : if( FD_UNLIKELY( sz<hdr_min ) ) {
665 : /* FIXME support ICMP messages in the future? */
666 0 : FD_LOG_ERR(( "packet too small %lu (in_idx=%lu)", sz, in_idx ));
667 0 : }
668 0 : if( FD_UNLIKELY( ( ctx->parsed.ip_version !=4 ) |
669 0 : ( ctx->parsed.ip_protocol!=FD_IP4_HDR_PROTOCOL_UDP ) ) ) {
670 0 : FD_LOG_ERR(( "packet from in_idx=%lu: sock tile only supports IPv4 UDP for now", in_idx ));
671 0 : }
672 0 : }
673 :
674 0 : ctx->tx_idle_cnt = 0;
675 0 : ctx->batch_cnt++;
676 : /* Technically leaves a gap. sz is always larger than the payload
677 : written to tx_ptr because Ethernet & IPv4 headers were stripped. */
678 0 : ctx->tx_ptr += fd_ulong_align_up( sz, FD_CHUNK_ALIGN );
679 :
680 0 : if( ctx->batch_cnt >= STEM_BURST ) {
681 0 : flush_tx_batch( ctx );
682 0 : }
683 0 : }
684 :
685 : /* End TX path ********************************************************/
686 :
687 : /* after_credit is called every stem iteration when there are enough
688 : flow control credits to publish a burst of fragments. */
689 :
690 : static inline void
691 : after_credit( fd_sock_tile_t * ctx,
692 : fd_stem_context_t * stem,
693 : int * poll_in FD_PARAM_UNUSED,
694 0 : int * charge_busy ) {
695 0 : if( ctx->tx_idle_cnt > 512 ) {
696 0 : if( ctx->batch_cnt ) {
697 0 : flush_tx_batch( ctx );
698 0 : }
699 0 : ulong pkt_cnt = poll_rx( ctx, stem );
700 0 : *charge_busy = pkt_cnt!=0;
701 0 : }
702 0 : ctx->tx_idle_cnt++;
703 0 : }
704 :
705 : static void
706 0 : metrics_write( fd_sock_tile_t * ctx ) {
707 0 : FD_MCNT_SET( SOCK, SYSCALL_RX, ctx->metrics.sys_recvmmsg_cnt );
708 0 : FD_MCNT_ENUM_COPY( SOCK, SYSCALL_TX, ctx->metrics.sys_sendmmsg_cnt );
709 0 : FD_MCNT_SET( SOCK, PKT_RX, ctx->metrics.rx_pkt_cnt );
710 0 : FD_MCNT_SET( SOCK, PKT_TX, ctx->metrics.tx_pkt_cnt );
711 0 : FD_MCNT_SET( SOCK, PKT_TX_FAILED, ctx->metrics.tx_drop_cnt );
712 0 : FD_MCNT_SET( SOCK, PKT_TX_BYTES, ctx->metrics.tx_bytes_total );
713 0 : FD_MCNT_SET( SOCK, PKT_RX_BYTES, ctx->metrics.rx_bytes_total );
714 0 : }
715 :
716 : static ulong
717 : rlimit_file_cnt( fd_topo_t const * topo,
718 0 : fd_topo_tile_t const * tile ) {
719 0 : fd_sock_tile_t * ctx = fd_topo_obj_laddr( topo, tile->tile_obj_id );
720 0 : return RX_SOCK_FD_MIN + ctx->sock_cnt;
721 0 : }
722 :
723 0 : #define STEM_CALLBACK_CONTEXT_TYPE fd_sock_tile_t
724 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_sock_tile_t)
725 :
726 0 : #define STEM_LAZY ((long)10e6) /* 10ms */
727 :
728 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
729 0 : #define STEM_CALLBACK_AFTER_CREDIT after_credit
730 0 : #define STEM_CALLBACK_BEFORE_FRAG before_frag
731 0 : #define STEM_CALLBACK_DURING_FRAG during_frag
732 0 : #define STEM_CALLBACK_AFTER_FRAG after_frag
733 :
734 : #include "../../stem/fd_stem.c"
735 :
736 : fd_topo_run_tile_t fd_tile_sock = {
737 : .name = "sock",
738 : .rlimit_file_cnt_fn = rlimit_file_cnt,
739 : .populate_allowed_seccomp = populate_allowed_seccomp,
740 : .populate_allowed_fds = populate_allowed_fds,
741 : .scratch_align = scratch_align,
742 : .scratch_footprint = scratch_footprint,
743 : .privileged_init = privileged_init,
744 : .unprivileged_init = unprivileged_init,
745 : .run = stem_run,
746 : };
|