Line data Source code
1 : #include "fd_ipecho_client.h"
2 : #include "fd_ipecho_server.h"
3 :
4 : #include "../genesis/fd_genesi_tile.h"
5 : #include "../genesis/genesis_hash.h"
6 : #include "../../disco/topo/fd_topo.h"
7 : #include "../../disco/topo/fd_dns_resolve.h"
8 : #include "../../disco/metrics/fd_metrics.h"
9 : #include "../../disco/waker/fd_waker.h"
10 :
11 : #include <netinet/in.h>
12 : #include <sys/socket.h>
13 : #include <time.h> /* CLOCK_REALTIME for seccomp filter */
14 : #include <poll.h>
15 :
16 : #include "generated/fd_ipecho_tile_seccomp.h"
17 :
18 0 : #define FD_IPECHO_MAX_CONNECTION_CNT (1024UL)
19 :
20 : struct fd_ipecho_tile_ctx {
21 : int retrieving;
22 :
23 : fd_ipecho_server_t * server;
24 : fd_ipecho_client_t * client;
25 :
26 : ulong waker_client_idx;
27 : ulong * waker_fseq;
28 :
29 : uint bind_address;
30 : ushort bind_port;
31 :
32 : ushort expected_shred_version;
33 :
34 : fd_ip4_port_t entrypoints[ FD_TOPO_GOSSIP_ENTRYPOINTS_MAX ];
35 : ulong entrypoints_cnt;
36 :
37 : fd_wksp_t * genesi_in_mem;
38 : ulong genesi_in_chunk0;
39 : ulong genesi_in_wmark;
40 : };
41 :
42 : typedef struct fd_ipecho_tile_ctx fd_ipecho_tile_ctx_t;
43 :
44 : FD_FN_CONST static inline ulong
45 0 : scratch_align( void ) {
46 0 : return alignof( fd_ipecho_tile_ctx_t );
47 0 : }
48 :
49 : FD_FN_PURE static inline ulong
50 0 : scratch_footprint( fd_topo_tile_t const * tile ) {
51 0 : (void)tile;
52 :
53 0 : ulong l = FD_LAYOUT_INIT;
54 0 : l = FD_LAYOUT_APPEND( l, alignof(fd_ipecho_tile_ctx_t), sizeof(fd_ipecho_tile_ctx_t) );
55 0 : l = FD_LAYOUT_APPEND( l, fd_ipecho_client_align(), fd_ipecho_client_footprint() );
56 0 : l = FD_LAYOUT_APPEND( l, fd_ipecho_server_align(), fd_ipecho_server_footprint( FD_IPECHO_MAX_CONNECTION_CNT ) );
57 0 : return FD_LAYOUT_FINI( l, scratch_align() );
58 0 : }
59 :
60 : static inline void
61 0 : metrics_write( fd_ipecho_tile_ctx_t * ctx ) {
62 0 : fd_ipecho_server_metrics_t * metrics = fd_ipecho_server_metrics( ctx->server );
63 :
64 0 : FD_MGAUGE_SET( IPECHO, CONN_ACTIVE, metrics->connection_cnt );
65 0 : FD_MCNT_SET( IPECHO, BYTES_READ, metrics->bytes_read );
66 0 : FD_MCNT_SET( IPECHO, BYTES_WRITTEN, metrics->bytes_written );
67 0 : ulong conn_closed[ FD_METRICS_ENUM_CONN_CLOSE_RESULT_CNT ];
68 0 : conn_closed[ FD_METRICS_ENUM_CONN_CLOSE_RESULT_V_OK_IDX ] = metrics->connections_closed_ok;
69 0 : conn_closed[ FD_METRICS_ENUM_CONN_CLOSE_RESULT_V_ERROR_IDX ] = metrics->connections_closed_error;
70 0 : FD_MCNT_ENUM_COPY( IPECHO, CONN_CLOSED, conn_closed );
71 0 : }
72 :
73 : static inline void
74 : poll_client( fd_ipecho_tile_ctx_t * ctx,
75 : fd_stem_context_t * stem,
76 0 : int * charge_busy ) {
77 0 : if( FD_UNLIKELY( !ctx->client ) ) return;
78 :
79 0 : ushort shred_version;
80 0 : int result = fd_ipecho_client_poll( ctx->client, &shred_version, charge_busy );
81 0 : if( FD_UNLIKELY( !result ) ) {
82 0 : if( FD_UNLIKELY( ctx->expected_shred_version && ctx->expected_shred_version!=shred_version ) ) {
83 0 : FD_LOG_ERR(( "Expected shred version %hu but entrypoint returned %hu",
84 0 : ctx->expected_shred_version, shred_version ));
85 0 : }
86 :
87 0 : FD_LOG_INFO(( "retrieved shred version %hu from entrypoint", shred_version ));
88 0 : FD_MGAUGE_SET( IPECHO, CURRENT_SHRED_VERSION, shred_version );
89 0 : fd_stem_publish( stem, 0UL, shred_version, 0UL, 0UL, 0UL, 0UL, 0UL );
90 0 : fd_ipecho_server_set_shred_version( ctx->server, shred_version );
91 0 : ctx->retrieving = 0;
92 0 : return;
93 0 : } else if( FD_UNLIKELY( -1==result ) ) {
94 0 : FD_LOG_ERR(( "Could not determine shred version from entrypoints. Please "
95 0 : "check you can connect to the entrypoints provided." ));
96 0 : }
97 0 : }
98 :
99 : static inline void
100 : after_credit( fd_ipecho_tile_ctx_t * ctx,
101 : fd_stem_context_t * stem,
102 : int * opt_poll_in,
103 0 : int * charge_busy ) {
104 0 : (void)opt_poll_in;
105 :
106 0 : if( FD_UNLIKELY( ctx->retrieving ) ) {
107 0 : poll_client( ctx, stem, charge_busy );
108 0 : return;
109 0 : }
110 :
111 0 : if( FD_UNLIKELY( fd_fseq_query( ctx->waker_fseq )==1UL ) ) {
112 0 : fd_fseq_update( ctx->waker_fseq, 0UL );
113 0 : fd_ipecho_server_epoll_poll( ctx->server, charge_busy ); /* one batch; the rearm re-fires leftovers */
114 0 : fd_waker_client_rearm( ctx->waker_client_idx );
115 0 : } else {
116 0 : fd_log_sleep( (long)1e6 );
117 0 : }
118 0 : }
119 :
120 : static inline int
121 : returnable_frag( fd_ipecho_tile_ctx_t * ctx,
122 : ulong in_idx,
123 : ulong seq,
124 : ulong sig,
125 : ulong chunk,
126 : ulong sz,
127 : ulong ctl,
128 : ulong tsorig,
129 : ulong tspub,
130 0 : fd_stem_context_t * stem ) {
131 0 : (void)in_idx; (void)seq; (void)sig; (void)sz; (void)ctl; (void)tspub;
132 0 : fd_genesis_meta_t const * genesis_meta = fd_chunk_to_laddr( ctx->genesi_in_mem, chunk );
133 :
134 0 : if( FD_UNLIKELY( genesis_meta->bootstrap ) ) {
135 0 : ushort shred_version = compute_shred_version( genesis_meta->genesis_hash.uc, NULL, 0UL );
136 0 : FD_TEST( shred_version );
137 :
138 0 : FD_MGAUGE_SET( IPECHO, CURRENT_SHRED_VERSION, shred_version );
139 0 : fd_stem_publish( stem, 0UL, shred_version, 0UL, 0UL, 0UL, tsorig, fd_frag_meta_ts_comp( fd_tickcount() ) );
140 0 : fd_ipecho_server_set_shred_version( ctx->server, shred_version );
141 0 : ctx->retrieving = 0;
142 0 : }
143 :
144 0 : return 0;
145 0 : }
146 :
147 : static void
148 : privileged_init( fd_topo_t const * topo,
149 0 : fd_topo_tile_t const * tile ) {
150 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
151 :
152 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
153 0 : fd_ipecho_tile_ctx_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof( fd_ipecho_tile_ctx_t ), sizeof( fd_ipecho_tile_ctx_t ) );
154 0 : void * _client = FD_SCRATCH_ALLOC_APPEND( l, fd_ipecho_client_align(), fd_ipecho_client_footprint() );
155 0 : void * _server = FD_SCRATCH_ALLOC_APPEND( l, fd_ipecho_server_align(), fd_ipecho_server_footprint( FD_IPECHO_MAX_CONNECTION_CNT ) );
156 :
157 0 : ctx->bind_address = tile->ipecho.bind_address;
158 0 : ctx->bind_port = tile->ipecho.bind_port;
159 :
160 0 : ctx->expected_shred_version = tile->ipecho.expected_shred_version;
161 :
162 0 : ctx->entrypoints_cnt = tile->ipecho.entrypoints_cnt;
163 0 : fd_dns_resolve_peers( tile->ipecho.entrypoints[ 0 ], sizeof(tile->ipecho.entrypoints[ 0 ]), tile->ipecho.entrypoints_cnt, "gossip.entrypoints", ctx->entrypoints );
164 :
165 0 : ctx->retrieving = 1;
166 0 : if( FD_LIKELY( ctx->entrypoints_cnt ) ) {
167 0 : ctx->client = fd_ipecho_client_join( fd_ipecho_client_new( _client ) );
168 0 : FD_TEST( ctx->client );
169 0 : fd_ipecho_client_init( ctx->client, ctx->entrypoints, ctx->entrypoints_cnt );
170 0 : } else {
171 0 : ctx->client = NULL;
172 0 : }
173 :
174 0 : ctx->server = fd_ipecho_server_join( fd_ipecho_server_new( _server, FD_IPECHO_MAX_CONNECTION_CNT ) );
175 0 : FD_TEST( ctx->server );
176 0 : ctx->waker_client_idx = tile->waker_client_idx;
177 0 : FD_TEST( ctx->waker_client_idx!=ULONG_MAX );
178 0 : fd_ipecho_server_init( ctx->server, FD_WAKER_INNER_FD( ctx->waker_client_idx ), ctx->bind_address, ctx->bind_port, ctx->expected_shred_version );
179 :
180 0 : ulong scratch_top = FD_SCRATCH_ALLOC_FINI( l, scratch_align() );
181 0 : if( FD_UNLIKELY( scratch_top > (ulong)scratch + scratch_footprint( tile ) ) )
182 0 : FD_LOG_ERR(( "scratch overflow %lu %lu %lu", scratch_top - (ulong)scratch - scratch_footprint( tile ), scratch_top, (ulong)scratch + scratch_footprint( tile ) ));
183 0 : }
184 :
185 : static void
186 : unprivileged_init( fd_topo_t const * topo,
187 0 : fd_topo_tile_t const * tile ) {
188 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
189 :
190 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
191 0 : fd_ipecho_tile_ctx_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof( fd_ipecho_tile_ctx_t ), sizeof( fd_ipecho_tile_ctx_t ) );
192 :
193 0 : FD_MGAUGE_SET( IPECHO, CURRENT_SHRED_VERSION, tile->ipecho.expected_shred_version );
194 :
195 0 : FD_TEST( ctx->waker_client_idx!=ULONG_MAX );
196 0 : ctx->waker_fseq = fd_fseq_join( fd_topo_obj_laddr( topo, tile->waker_fseq_obj_id ) );
197 0 : FD_TEST( ctx->waker_fseq );
198 :
199 : /* In some topologies (e.g. firedancer-dev gossip), the ipecho tile
200 : has no input links. Guard against dereferencing a missing
201 : link/dcache. */
202 0 : if( FD_LIKELY( tile->in_cnt>0UL ) ) {
203 0 : ulong link_id = tile->in_link_id[ 0UL ];
204 0 : void * dcache = topo->links[ link_id ].dcache;
205 0 : ctx->genesi_in_mem = topo->workspaces[ topo->objs[ topo->links[ link_id ].dcache_obj_id ].wksp_id ].wksp;
206 0 : ctx->genesi_in_chunk0 = fd_dcache_compact_chunk0( ctx->genesi_in_mem, dcache );
207 0 : ctx->genesi_in_wmark = fd_dcache_compact_wmark ( ctx->genesi_in_mem, dcache, topo->links[ link_id ].mtu );
208 0 : } else {
209 0 : ctx->genesi_in_mem = NULL;
210 0 : ctx->genesi_in_chunk0 = 0UL;
211 0 : ctx->genesi_in_wmark = 0UL;
212 0 : }
213 0 : }
214 :
215 : static ulong
216 : rlimit_file_cnt( fd_topo_t const * topo FD_PARAM_UNUSED,
217 0 : fd_topo_tile_t const * tile ) {
218 : /* pipefd, socket, stderr, logfile, and one spare for
219 : new accept() connections */
220 0 : ulong base = 5UL;
221 0 : return base +
222 0 : tile->ipecho.entrypoints_cnt + /* for the client */
223 0 : FD_IPECHO_MAX_CONNECTION_CNT; /* for the server's connections */
224 0 : }
225 :
226 : static ulong
227 : populate_allowed_seccomp( fd_topo_t const * topo,
228 : fd_topo_tile_t const * tile,
229 : ulong out_cnt,
230 0 : struct sock_filter * out ) {
231 0 : (void)topo;
232 :
233 0 : uint epoll_inner_fd = (uint)FD_WAKER_INNER_FD( tile->waker_client_idx );
234 0 : uint epoll_outer_fd = (uint)FD_WAKER_OUTER_FD;
235 :
236 0 : populate_sock_filter_policy_fd_ipecho_tile( out_cnt, out, (uint)fd_log_private_logfile_fd(), epoll_inner_fd, epoll_outer_fd );
237 0 : return sock_filter_policy_fd_ipecho_tile_instr_cnt;
238 0 : }
239 :
240 : static ulong
241 : populate_allowed_fds( fd_topo_t const * topo,
242 : fd_topo_tile_t const * tile,
243 : ulong out_fds_cnt,
244 0 : int * out_fds ) {
245 :
246 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
247 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
248 0 : fd_ipecho_tile_ctx_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_ipecho_tile_ctx_t), sizeof(fd_ipecho_tile_ctx_t) );
249 :
250 0 : if( FD_UNLIKELY( out_fds_cnt<5UL+tile->ipecho.entrypoints_cnt ) ) FD_LOG_ERR(( "out_fds_cnt %lu", out_fds_cnt ));
251 :
252 0 : ulong out_cnt = 0UL;
253 0 : out_fds[ out_cnt++ ] = 2; /* stderr */
254 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) )
255 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
256 :
257 : /* All of the fds managed by the client. */
258 0 : for( ulong i=0UL; i<tile->ipecho.entrypoints_cnt; i++ ) {
259 0 : int fd = fd_ipecho_client_get_pollfds( ctx->client )[ i ].fd;
260 0 : if( FD_LIKELY( fd!=-1 ) ) out_fds[ out_cnt++ ] = fd;
261 0 : }
262 :
263 : /* The server's socket. */
264 0 : out_fds[ out_cnt++ ] = fd_ipecho_server_sockfd( ctx->server );
265 :
266 0 : out_fds[ out_cnt++ ] = FD_WAKER_OUTER_FD; /* waker outer epoll fd (rearm) */
267 0 : out_fds[ out_cnt++ ] = FD_WAKER_INNER_FD( tile->waker_client_idx ); /* waker inner epoll fd */
268 0 : return out_cnt;
269 0 : }
270 :
271 0 : #define STEM_BURST (2UL)
272 0 : #define STEM_LAZY ((long)10e6) /* 10ms */
273 :
274 0 : #define STEM_CALLBACK_CONTEXT_TYPE fd_ipecho_tile_ctx_t
275 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_ipecho_tile_ctx_t)
276 :
277 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
278 0 : #define STEM_CALLBACK_AFTER_CREDIT after_credit
279 0 : #define STEM_CALLBACK_RETURNABLE_FRAG returnable_frag
280 :
281 : #include "../../disco/stem/fd_stem.c"
282 :
283 : fd_topo_run_tile_t fd_tile_ipecho = {
284 : .name = "ipecho",
285 : .rlimit_file_cnt_fn = rlimit_file_cnt,
286 : .populate_allowed_seccomp = populate_allowed_seccomp,
287 : .populate_allowed_fds = populate_allowed_fds,
288 : .scratch_align = scratch_align,
289 : .scratch_footprint = scratch_footprint,
290 : .privileged_init = privileged_init,
291 : .unprivileged_init = unprivileged_init,
292 : .run = stem_run,
293 : .allow_connect = 1,
294 : .keep_host_networking = 1
295 : };
|