Line data Source code
1 : /* The txsend tile relays new transactions to the current leader.
2 :
3 : To recap, every ~1.6 seconds, a new validator is elected by the
4 : protocol to receive all transactions, and pack them into blocks.
5 : To submit a transaction to Solana, one has to track the "leader
6 : schedule" and get the timing right to submit a transaction early
7 : enough.
8 :
9 : The txsend tile supports the TPU-UDP (connectionless) and TPU-QUIC
10 : transports. This tile has the following jobs:
11 : - decide who to connect to (continually monitor gossip, leader schedule)
12 : - manage a pool of QUIC connections
13 : - actually send transactions
14 :
15 : An important quirk is that txsend tolerates temporary gossip outages
16 : by indefinitely caching old endpoint info (until the gossip feed is
17 : revived). */
18 :
19 : #include "fd_txsend_tile.h"
20 :
21 : #include "../../disco/topo/fd_topo.h"
22 : #include "../../disco/fd_txn_m.h"
23 : #include "../../disco/metrics/fd_metrics.h"
24 : #include "../../disco/events/generated/fd_event_gen.h"
25 : #include "../../disco/events/fd_event_report.h"
26 : #include "../../choreo/tower/fd_tower_serdes.h"
27 : #include "../../flamenco/runtime/fd_system_ids.h"
28 : #include "../../disco/keyguard/fd_keyguard.h"
29 : #include "../../disco/keyguard/fd_keyload.h"
30 : #include "../fd_startup.h"
31 : #include "../tower/fd_tower_tile.h"
32 : #include "../../util/net/fd_net_headers.h"
33 : #include "../../waltz/quic/fd_quic.h"
34 :
35 : #include <time.h>
36 : #include "generated/fd_txsend_tile_seccomp.h"
37 :
38 0 : #define IN_KIND_SIGN (0UL)
39 0 : #define IN_KIND_GOSSIP (1UL)
40 0 : #define IN_KIND_EPOCH (2UL)
41 0 : #define IN_KIND_TOWER (3UL)
42 0 : #define IN_KIND_NET (4UL)
43 :
44 : fd_quic_limits_t quic_limits = {
45 : .conn_cnt = 128UL,
46 : .handshake_cnt = 128UL,
47 : .conn_id_cnt = FD_QUIC_MIN_CONN_ID_CNT,
48 : .inflight_frame_cnt = 16UL * 128UL,
49 : .min_inflight_frame_cnt_conn = 4UL,
50 : .stream_id_cnt = 64UL,
51 : .tx_buf_sz = FD_TXN_MTU,
52 : .stream_pool_cnt = 128UL,
53 : };
54 :
55 : #define MAP_NAME peer_map
56 0 : #define MAP_KEY pubkey
57 : #define MAP_ELE_T peer_entry_t
58 : #define MAP_KEY_T fd_pubkey_t
59 0 : #define MAP_NEXT map.next
60 : #define MAP_KEY_EQ(k0,k1) fd_pubkey_eq( k0, k1 )
61 : #define MAP_KEY_HASH(key,seed) fd_progcache_rec_key_hash1( (key)->uc, (seed) )
62 : #define MAP_IMPL_STYLE 2
63 : #include "../../util/tmpl/fd_map_chain.c"
64 :
65 : FD_FN_CONST static inline ulong
66 0 : scratch_align( void ) {
67 0 : return fd_ulong_max( 128UL, fd_quic_align() );
68 0 : }
69 :
70 : FD_FN_PURE static inline ulong
71 0 : scratch_footprint( fd_topo_tile_t const * tile FD_PARAM_UNUSED ) {
72 0 : ulong l = FD_LAYOUT_INIT;
73 0 : l = FD_LAYOUT_APPEND( l, alignof(fd_txsend_tile_t), sizeof(fd_txsend_tile_t) );
74 0 : l = FD_LAYOUT_APPEND( l, fd_quic_align(), fd_quic_footprint( &quic_limits ) );
75 0 : l = FD_LAYOUT_APPEND( l, peer_map_align(), peer_map_footprint( 2UL*FD_CONTACT_INFO_TABLE_SIZE ) );
76 0 : return FD_LAYOUT_FINI( l, scratch_align() );
77 0 : }
78 :
79 : static void
80 0 : during_housekeeping( fd_txsend_tile_t * ctx ) {
81 0 : if( FD_UNLIKELY( fd_keyswitch_state_query( ctx->keyswitch )==FD_KEYSWITCH_STATE_UNHALT_PENDING ) ) {
82 0 : FD_LOG_DEBUG(( "keyswitch: unhalting" ));
83 0 : ctx->halt_net_frags = 0;
84 0 : fd_keyswitch_state( ctx->keyswitch, FD_KEYSWITCH_STATE_COMPLETED );
85 0 : }
86 :
87 0 : if( FD_UNLIKELY( fd_keyswitch_state_query( ctx->keyswitch )==FD_KEYSWITCH_STATE_SWITCH_PENDING ) ) {
88 0 : FD_LOG_DEBUG(( "keyswitch: switching identity" ));
89 0 : ulong seq_must_complete = ctx->keyswitch->param;
90 0 : if( FD_UNLIKELY( fd_seq_lt( ctx->tower_in_expect_seq, seq_must_complete ) ) ) {
91 : /* See fd_keyswitch.h, we need to flush any in-flight shreds from
92 : the leader pipeline before switching key. */
93 0 : FD_LOG_WARNING(( "Flushing in-flight unpublished votes from tower, must reach seq %lu, currently at %lu ...", seq_must_complete, ctx->tower_in_expect_seq ));
94 0 : return;
95 0 : }
96 :
97 : /* Halt net frags to avoid potential quic callback in after_frag */
98 0 : ctx->halt_net_frags = 1;
99 :
100 0 : fd_quic_set_identity_public_key( ctx->quic, ctx->keyswitch->bytes );
101 :
102 0 : memcpy( ctx->identity_key, ctx->keyswitch->bytes, 32UL );
103 0 : fd_keyswitch_state( ctx->keyswitch, FD_KEYSWITCH_STATE_COMPLETED );
104 0 : }
105 :
106 0 : if( FD_UNLIKELY( ctx->av_keyswitch && fd_keyswitch_state_query( ctx->av_keyswitch )==FD_KEYSWITCH_STATE_SWITCH_PENDING ) ) {
107 0 : ulong seq_must_complete = fd_keyswitch_param_query( ctx->av_keyswitch );
108 0 : if( FD_UNLIKELY( fd_seq_lt( ctx->tower_in_expect_seq, seq_must_complete ) ) ) {
109 0 : FD_LOG_WARNING(( "Flushing in-flight votes from tower, must reach seq %lu, currently at %lu ...", seq_must_complete, ctx->tower_in_expect_seq ));
110 0 : return;
111 0 : }
112 0 : fd_keyswitch_state( ctx->av_keyswitch, FD_KEYSWITCH_STATE_COMPLETED );
113 0 : }
114 0 : }
115 :
116 : static void
117 0 : metrics_write( fd_txsend_tile_t * ctx ) {
118 0 : FD_MCNT_SET( TXSEND, PKT_RX_BYTES, ctx->quic->metrics.net_rx_byte_cnt );
119 0 : FD_MCNT_ENUM_COPY( TXSEND, FRAME_RX, ctx->quic->metrics.frame_rx_cnt );
120 0 : FD_MCNT_SET( TXSEND, PKT_RX, ctx->quic->metrics.net_rx_pkt_cnt );
121 0 : FD_MCNT_SET( TXSEND, STREAM_RX_BYTES, ctx->quic->metrics.stream_rx_byte_cnt );
122 0 : FD_MCNT_SET( TXSEND, STREAM_RX, ctx->quic->metrics.stream_rx_event_cnt );
123 :
124 0 : FD_MCNT_SET( TXSEND, PKT_TX, ctx->quic->metrics.net_tx_pkt_cnt );
125 0 : FD_MCNT_SET( TXSEND, PKT_TX_BYTES, ctx->quic->metrics.net_tx_byte_cnt );
126 0 : FD_MCNT_SET( TXSEND, PKT_TX_RETRY, ctx->quic->metrics.retry_tx_cnt );
127 0 : FD_MCNT_ENUM_COPY( TXSEND, ACK_TX, ctx->quic->metrics.ack_tx );
128 :
129 0 : FD_MGAUGE_ENUM_COPY( TXSEND, CONN_STATE, ctx->quic->metrics.conn_state_cnt );
130 0 : FD_MGAUGE_SET( TXSEND, CONN_IN_USE, ctx->quic->metrics.conn_alloc_cnt );
131 0 : FD_MCNT_SET( TXSEND, CONN_CREATED, ctx->quic->metrics.conn_created_cnt );
132 0 : FD_MCNT_SET( TXSEND, CONN_CLOSED, ctx->quic->metrics.conn_closed_cnt );
133 0 : FD_MCNT_SET( TXSEND, CONN_ABORTED, ctx->quic->metrics.conn_aborted_cnt );
134 0 : FD_MCNT_SET( TXSEND, CONN_TIMED_OUT, ctx->quic->metrics.conn_timeout_cnt );
135 0 : FD_MCNT_SET( TXSEND, CONN_RETRIED, ctx->quic->metrics.conn_retry_cnt );
136 0 : FD_MCNT_SET( TXSEND, CONN_ERROR_NO_SLOTS, ctx->quic->metrics.conn_err_no_slots_cnt );
137 0 : FD_MCNT_SET( TXSEND, CONN_ERROR_RETRY_FAILED, ctx->quic->metrics.conn_err_retry_fail_cnt );
138 :
139 0 : FD_MCNT_ENUM_COPY( TXSEND, PKT_CRYPTO_FAILED, ctx->quic->metrics.pkt_decrypt_fail_cnt );
140 0 : FD_MCNT_ENUM_COPY( TXSEND, PKT_NO_KEY, ctx->quic->metrics.pkt_no_key_cnt );
141 0 : FD_MCNT_ENUM_COPY( TXSEND, PKT_NO_CONN, ctx->quic->metrics.pkt_no_conn_cnt );
142 0 : FD_MCNT_SET( TXSEND, PKT_SRC_INVALID, ctx->quic->metrics.pkt_wrong_src_cnt );
143 0 : FD_MCNT_ENUM_COPY( TXSEND, FRAME_META_ACQUIRED, ctx->quic->metrics.frame_tx_alloc_cnt );
144 0 : FD_MCNT_SET( TXSEND, PKT_NET_HEADER_INVALID, ctx->quic->metrics.pkt_net_hdr_err_cnt );
145 0 : FD_MCNT_SET( TXSEND, PKT_HEADER_INVALID, ctx->quic->metrics.pkt_quic_hdr_err_cnt );
146 0 : FD_MCNT_SET( TXSEND, PKT_UNDERSIZE, ctx->quic->metrics.pkt_undersz_cnt );
147 0 : FD_MCNT_SET( TXSEND, PKT_OVERSIZE, ctx->quic->metrics.pkt_oversz_cnt );
148 0 : FD_MCNT_SET( TXSEND, PKT_RX_VERSION_NEGOTIATION, ctx->quic->metrics.pkt_verneg_cnt );
149 0 : FD_MCNT_ENUM_COPY( TXSEND, PKT_TX_RETRANSMITTED, ctx->quic->metrics.pkt_retransmissions_cnt );
150 :
151 0 : FD_MCNT_SET( TXSEND, HANDSHAKE_CREATED, ctx->quic->metrics.hs_created_cnt );
152 0 : FD_MCNT_SET( TXSEND, HANDSHAKE_ERROR_ALLOC_FAIL, ctx->quic->metrics.hs_err_alloc_fail_cnt );
153 0 : FD_MCNT_SET( TXSEND, HANDSHAKE_EVICTED, ctx->quic->metrics.hs_evicted_cnt );
154 :
155 0 : FD_MCNT_SET( TXSEND, FRAME_PARSE_FAILED, ctx->quic->metrics.frame_rx_err_cnt );
156 :
157 0 : FD_MHIST_COPY( TXSEND, SERVICE_DURATION_SECONDS, ctx->quic->metrics.service_duration );
158 0 : FD_MHIST_COPY( TXSEND, RX_DURATION_SECONDS, ctx->quic->metrics.receive_duration );
159 0 : }
160 :
161 : static void
162 : quic_tls_cv_sign( void * signer_ctx,
163 : uchar signature[ static 64 ],
164 0 : uchar const payload[ static 130 ] ) {
165 0 : fd_txsend_tile_t * ctx = signer_ctx;
166 :
167 0 : fd_keyguard_client_sign( ctx->keyguard_client, signature, payload, 130UL, FD_KEYGUARD_SIGN_TYPE_ED25519 );
168 0 : }
169 :
170 : static void
171 : send_to_net( fd_txsend_tile_t * ctx,
172 : fd_ip4_hdr_t const * ip4_hdr,
173 : fd_udp_hdr_t const * udp_hdr,
174 : uchar const * payload,
175 : ulong payload_sz,
176 0 : long now ) {
177 0 : uint const ip_dst = FD_LOAD( uint, ip4_hdr->daddr_c );
178 0 : ulong const ip_sz = FD_IP4_GET_LEN( *ip4_hdr );
179 :
180 0 : ulong const hdr_sz = sizeof(fd_eth_hdr_t) + ip_sz + sizeof(fd_udp_hdr_t);
181 0 : if( FD_UNLIKELY( payload_sz>FD_ETH_PAYLOAD_MAX-hdr_sz ) ) return;
182 0 : ulong const sz_l2 = hdr_sz + payload_sz;
183 :
184 0 : fd_txsend_out_t * net_out_link = ctx->net_out;
185 0 : uchar * packet_l2 = fd_chunk_to_laddr( net_out_link->mem, net_out_link->chunk );
186 0 : uchar * packet_l3 = packet_l2 + sizeof(fd_eth_hdr_t);
187 0 : uchar * packet_l4 = packet_l3 + ip_sz;
188 0 : uchar * packet_l5 = packet_l4 + sizeof(fd_udp_hdr_t);
189 :
190 0 : fd_memcpy( packet_l2, ctx->packet_hdr->eth, sizeof(fd_eth_hdr_t) );
191 0 : fd_memcpy( packet_l3, ip4_hdr, ip_sz );
192 0 : fd_memcpy( packet_l4, udp_hdr, sizeof(fd_udp_hdr_t) );
193 0 : fd_memcpy( packet_l5, payload, payload_sz );
194 :
195 0 : ulong sig = fd_disco_netmux_sig( ip_dst, 0U, ip_dst, DST_PROTO_OUTGOING, FD_NETMUX_SIG_MIN_HDR_SZ );
196 :
197 0 : ulong tspub = (ulong)fd_frag_meta_ts_comp( now );
198 0 : fd_stem_publish( ctx->stem, net_out_link->idx, sig, net_out_link->chunk, sz_l2, 0UL, 0, tspub );
199 0 : net_out_link->chunk = fd_dcache_compact_next( net_out_link->chunk, sz_l2, net_out_link->chunk0, net_out_link->wmark );
200 0 : }
201 :
202 : static int
203 : quic_tx_aio_send( void * _ctx,
204 : fd_aio_pkt_info_t const * batch,
205 : ulong batch_cnt,
206 : ulong * opt_batch_idx,
207 0 : int flush FD_PARAM_UNUSED ) {
208 0 : fd_txsend_tile_t * ctx = _ctx;
209 :
210 0 : long now = fd_log_wallclock();
211 :
212 0 : for( ulong i=0; i<batch_cnt; i++ ) {
213 0 : if( FD_UNLIKELY( batch[ i ].buf_sz<FD_NETMUX_SIG_MIN_HDR_SZ ) ) continue;
214 0 : uchar * buf = batch[ i ].buf;
215 0 : fd_ip4_hdr_t * ip4_hdr = fd_type_pun( buf );
216 0 : ulong const ip4_len = FD_IP4_GET_LEN( *ip4_hdr );
217 0 : fd_udp_hdr_t * udp_hdr = fd_type_pun( buf + ip4_len );
218 0 : uchar * payload = buf + ip4_len + sizeof(fd_udp_hdr_t);
219 0 : FD_TEST( batch[ i ].buf_sz >= ip4_len + sizeof(fd_udp_hdr_t) );
220 0 : ulong payload_sz = batch[ i ].buf_sz - ip4_len - sizeof(fd_udp_hdr_t);
221 0 : send_to_net( ctx, ip4_hdr, udp_hdr, payload, payload_sz, now );
222 0 : }
223 :
224 0 : if( FD_LIKELY( opt_batch_idx ) ) {
225 0 : *opt_batch_idx = batch_cnt;
226 0 : }
227 :
228 0 : return FD_AIO_SUCCESS;
229 0 : }
230 :
231 : /* quic_conn_deregister removes references to a quic_conn object from
232 : txsend state. */
233 :
234 : static void
235 : quic_conn_deregister( fd_txsend_tile_t * tile,
236 0 : fd_quic_conn_t * conn ) {
237 0 : for( ulong i=0UL; i<tile->conns_len; i++ ) {
238 0 : if( FD_LIKELY( tile->conns[ i ].conn!=conn ) ) continue;
239 0 : peer_entry_t * peer = peer_map_ele_query( tile->peer_map, &tile->conns[ i ].pubkey, NULL, tile->peers );
240 0 : if( FD_LIKELY( peer ) ) {
241 0 : for( ulong j=0UL; j<2UL; j++ ) {
242 0 : if( peer->quic_conns[ j ].quic_conn==conn ) {
243 0 : peer->quic_conns[ j ].quic_conn = NULL;
244 0 : }
245 0 : }
246 0 : }
247 0 : if( FD_UNLIKELY( i!=tile->conns_len-1UL ) ) tile->conns[ i ] = tile->conns[ tile->conns_len-1UL ];
248 0 : tile->conns_len--;
249 0 : return;
250 0 : }
251 0 : }
252 :
253 : /* quic_conn_final is invoked by fd_quic just before a conn object is
254 : deallocated. Here, we must remove all references to this conn. */
255 :
256 : static void
257 : quic_conn_final( fd_quic_conn_t * conn,
258 0 : void * ctx ) {
259 0 : fd_txsend_tile_t * tile = ctx;
260 0 : quic_conn_deregister( tile, conn );
261 0 : }
262 :
263 : /* quic_conn_close instructs fd_quic to deallocate the given conn and
264 : say goodbye (CONNECTION_CLOSE) to the peer. */
265 :
266 : static void
267 : quic_conn_close( fd_txsend_tile_t * tile,
268 : fd_quic_conn_t * conn,
269 0 : uint reason ) {
270 0 : if( FD_UNLIKELY( !conn ) ) return;
271 : /* Defer a conn close operation */
272 0 : fd_quic_conn_close( conn, reason );
273 0 : quic_conn_deregister( tile, conn );
274 : /* Send out a packet and invoke quic_conn_final */
275 0 : fd_quic_service( tile->quic, fd_log_wallclock() );
276 0 : }
277 :
278 : /* This QUIC servicing is very precarious. Recall a few facts,
279 :
280 : 1) QUIC needs to be serviced periodically to make progress
281 : 2) QUIC servicing may produce outgoing packets that need to be sent
282 : to the network
283 : 3) Elsewhere, the the tile publishes frags to the verify tile to
284 : send our own votes into our leader pipeline
285 :
286 : You could service QUIC in before_credit, as the QUIC tile does, but
287 : this has a problem. If you publish frags in before_credit, you might
288 : overrun the downstream consumer. For net tile, this is OK because it
289 : expects that (as does verify). But the credit counting mechanism
290 : doesn't expect this behavior and will underflow. (That's also not
291 : ideal, in case some plugin wanted to listen reliably on quic->verify
292 : they could not, if it got underflowed). Here though, we want to
293 : avoid dropping outgoing votes to verify, since they might be needed
294 : for liveness of a small cluster.
295 :
296 : We thus take the trade of servicing QUIC in after_credit, which means
297 : it could theoretically get backpressured by verify, however this
298 : isn't realistic in practice, as verify polls round robin and there's
299 : only one vote per slot. */
300 :
301 : static inline void
302 : after_credit( fd_txsend_tile_t * ctx,
303 : fd_stem_context_t * stem,
304 : int * opt_poll_in,
305 0 : int * charge_busy ) {
306 0 : ctx->stem = stem;
307 :
308 0 : if( FD_UNLIKELY( !fd_startup_gate_idle( ctx->startup_gate ) ) ) return;
309 :
310 0 : *charge_busy = fd_quic_service( ctx->quic, fd_log_wallclock() );
311 0 : *opt_poll_in = !*charge_busy; /* refetch credits to prevent above documented situation */
312 :
313 0 : if( FD_UNLIKELY( ctx->leader_schedules<2UL ) ) return;
314 0 : if( FD_UNLIKELY( ctx->voted_slot==ULONG_MAX ) ) return;
315 :
316 0 : fd_pubkey_t const * leaders[ 7UL ];
317 :
318 0 : for( ulong i=0UL; i<7UL; i++ ) {
319 : /* It's possible for leaders[i] to be NULL if target slot is two
320 : epochs ahead of the replay root. This is not possible on mainnet
321 : but can occur on local clusters during warmup epochs. */
322 0 : ulong target_slot = ctx->voted_slot+1UL + i*FD_EPOCH_SLOTS_PER_ROTATION;
323 0 : leaders[ i ] = fd_multi_epoch_leaders_get_leader_for_slot( ctx->mleaders, target_slot );
324 0 : }
325 :
326 : /* Disconnect any QUIC connection to a leader that does not have a
327 : rotation coming up in the next 7 slots. */
328 0 : ulong conn_cnt = ctx->conns_len;
329 0 : for( ulong i=0UL; i<conn_cnt; ) {
330 0 : int keep_conn = 0;
331 0 : for( ulong j=0UL; j<7UL; j++ ) {
332 0 : if( leaders[j] && fd_pubkey_eq( &ctx->conns[ i ].pubkey, leaders[ j ] ) ) {
333 0 : keep_conn = 1;
334 0 : break;
335 0 : }
336 0 : }
337 :
338 0 : if( FD_UNLIKELY( !keep_conn ) ) quic_conn_close( ctx, ctx->conns[ i ].conn, 0 );
339 0 : if( ctx->conns_len==conn_cnt ) i++;
340 0 : conn_cnt = ctx->conns_len;
341 0 : }
342 :
343 : /* Connect to any leader that does not have a connection yet. */
344 0 : for( ulong i=0UL; i<7UL; i++ ) {
345 0 : fd_pubkey_t const * leader = leaders[ i ];
346 0 : if( FD_UNLIKELY( !leader ) ) continue;
347 0 : peer_entry_t * peer = peer_map_ele_query( ctx->peer_map, leader, NULL, ctx->peers );
348 0 : if( FD_UNLIKELY( !peer ) ) continue; /* no contact info */
349 :
350 0 : for( ulong j=0UL; j<2UL; j++ ) {
351 0 : if( FD_UNLIKELY( ctx->conns_len==128UL ) ) break; /* connection limit reached */
352 0 : txsend_conn_t * conn = &peer->quic_conns[ j ];
353 0 : if( FD_LIKELY( conn->quic_conn ) ) continue; /* already connected */
354 0 : if( FD_UNLIKELY( !conn->quic_ip_addr || !conn->quic_port ) ) continue;
355 :
356 : /* Don't try to reconnect more than once every two seconds ...
357 : Basically Agave limits us to 8 connections per minute, so if we
358 : keep trying to reconnect rapidly it's much less effective than
359 : waiting a little bit to ensure we stay under the threshold.
360 :
361 : We should probably make this a bit more sophisticated, with a
362 : simple model that considers past connection attempts, and
363 : future leader slots (e.g. we might still want to burn an
364 : attempt if a leader slot is imminent, even if we recently tried
365 : to connect). For now the dumb logic seems to work well enough. */
366 0 : long now = fd_log_wallclock();
367 0 : if( FD_UNLIKELY( conn->quic_last_connected+2e9L>now ) ) continue;
368 :
369 0 : fd_quic_conn_t * quic_conn =
370 0 : fd_quic_connect( ctx->quic,
371 0 : conn->quic_ip_addr,
372 0 : conn->quic_port,
373 0 : ctx->src_ip_addr,
374 0 : ctx->src_port,
375 0 : now );
376 0 : if( FD_UNLIKELY( !quic_conn ) ) {
377 : /* Should never happen, but handle it gracefully */
378 0 : return;
379 0 : }
380 0 : ctx->conns[ ctx->conns_len ].conn = quic_conn;
381 0 : ctx->conns[ ctx->conns_len ].pubkey = *leader;
382 0 : conn->quic_conn = quic_conn;
383 0 : conn->quic_last_connected = now;
384 0 : ctx->conns_len++;
385 0 : }
386 0 : }
387 0 : }
388 :
389 : static void
390 : send_vote_to_leader( fd_txsend_tile_t * ctx,
391 : fd_pubkey_t const * leader_pubkey,
392 : uchar const * vote_payload,
393 0 : ulong vote_payload_sz ) {
394 0 : peer_entry_t const * peer = peer_map_ele_query_const( ctx->peer_map, leader_pubkey, NULL, ctx->peers );
395 0 : if( FD_UNLIKELY( !peer ) ) return; /* no known contact info */
396 :
397 0 : for( ulong i=0UL; i<2UL; i++ ) {
398 0 : if( FD_UNLIKELY( !peer->udp_ip_addrs[ i ] | !peer->udp_ports[ i ] ) ) continue;
399 :
400 0 : fd_ip4_hdr_t * ip4_hdr = ctx->packet_hdr->ip4;
401 0 : fd_udp_hdr_t * udp_hdr = ctx->packet_hdr->udp;
402 :
403 0 : ip4_hdr->daddr = peer->udp_ip_addrs[ i ];
404 0 : ip4_hdr->net_tot_len = fd_ushort_bswap( (ushort)(vote_payload_sz+sizeof(fd_ip4_hdr_t)+sizeof(fd_udp_hdr_t)) );
405 0 : ip4_hdr->net_id = fd_ushort_bswap( ctx->net_id++ );
406 0 : ip4_hdr->check = 0;
407 0 : ip4_hdr->check = fd_ip4_hdr_check_fast( ip4_hdr );
408 :
409 0 : udp_hdr->net_dport = fd_ushort_bswap( peer->udp_ports[ i ] );
410 0 : udp_hdr->net_len = fd_ushort_bswap( (ushort)( vote_payload_sz+sizeof(fd_udp_hdr_t) ) );
411 0 : send_to_net( ctx, ip4_hdr, udp_hdr, vote_payload, vote_payload_sz, fd_log_wallclock() );
412 0 : }
413 :
414 0 : for( ulong i=0UL; i<2UL; i++ ) {
415 0 : fd_quic_conn_t * conn = peer->quic_conns[ i ].quic_conn;
416 0 : if( FD_UNLIKELY( !conn ) ) continue;
417 :
418 0 : fd_quic_stream_t * stream = fd_quic_conn_new_stream( conn );
419 0 : if( FD_UNLIKELY( !stream ) ) continue;
420 :
421 0 : fd_quic_stream_send( stream, vote_payload, vote_payload_sz, 1 );
422 0 : }
423 0 : }
424 :
425 : /* gossip -> txsend peer synchronization
426 :
427 : The gossip update stream can be replayed to perfectly replicate the
428 : ContactInfo table. The stream is laid out so there are no duplicate
429 : pubkeys.
430 :
431 : However, the txsend tile wants to retain pubkeys past deletion.
432 : Gossip evicts ContactInfos without updates quickly, but txsend should
433 : continue sending to staked leaders even if there is a temporary
434 : gossip outage.
435 :
436 : Therefore, the txsend tile only tombstones entries when the gossip
437 : stream instructs to remove them, instead of deleting. The tombstone
438 : handling is then resolved downstream during ContactInfo updates. */
439 :
440 : static inline void
441 : handle_contact_info_remove( fd_txsend_tile_t * ctx,
442 0 : fd_gossip_update_message_t const * msg ) {
443 0 : FD_TEST( msg->contact_info_remove->idx < FD_CONTACT_INFO_TABLE_SIZE );
444 0 : peer_entry_t * entry = &ctx->peers[ msg->contact_info_remove->idx ];
445 0 : entry->tombstoned = 1;
446 0 : }
447 :
448 : static void
449 : handle_contact_info_update( fd_txsend_tile_t * ctx,
450 0 : fd_gossip_update_message_t const * msg ) {
451 0 : FD_TEST( msg->contact_info->idx < FD_CONTACT_INFO_TABLE_SIZE );
452 :
453 : /* Key updated by the gossip event */
454 0 : fd_pubkey_t key = FD_LOAD( fd_pubkey_t, msg->origin );
455 :
456 : /* Storage entry updated by the gossip event */
457 0 : peer_entry_t * entry = &ctx->peers[ msg->contact_info->idx ];
458 :
459 : /* At this point, entry contains an arbitrary old tombstoned entry or
460 : a previous version of this key. */
461 :
462 0 : if( FD_UNLIKELY( !fd_pubkey_eq( &entry->pubkey, &key ) ) ) {
463 : /* Overwriting an unrelated (tombstoned) entry, free it */
464 0 : quic_conn_close( ctx, entry->quic_conns[ 0 ].quic_conn, 0 );
465 0 : quic_conn_close( ctx, entry->quic_conns[ 1 ].quic_conn, 0 );
466 0 : peer_map_ele_remove( ctx->peer_map, &entry->pubkey, NULL, ctx->peers );
467 0 : memset( entry, 0, sizeof(peer_entry_t) );
468 0 : }
469 :
470 : /* At this point, entry contains a stale version of the same key or is
471 : empty. */
472 :
473 0 : peer_entry_t * stale = peer_map_ele_query( ctx->peer_map, &key, NULL, ctx->peers );
474 0 : if( FD_UNLIKELY( stale && stale!=entry ) ) {
475 : /* The key exists at another entry location, drop that and migrate
476 : it to this slot. */
477 0 : for( ulong i=0UL; i<2UL; i++ ) {
478 0 : entry->quic_conns [ i ] = stale->quic_conns [ i ];
479 0 : entry->udp_ip_addrs[ i ] = stale->udp_ip_addrs[ i ];
480 0 : entry->udp_ports [ i ] = stale->udp_ports [ i ];
481 0 : }
482 0 : peer_map_ele_remove( ctx->peer_map, &stale->pubkey, NULL, ctx->peers );
483 0 : memset( stale, 0, sizeof(peer_entry_t) );
484 0 : fd_memcpy( entry->pubkey.uc, msg->origin, 32UL );
485 0 : FD_TEST( peer_map_ele_insert( ctx->peer_map, entry, ctx->peers ) );
486 0 : } else if( !stale ) {
487 0 : fd_memcpy( entry->pubkey.uc, msg->origin, 32UL );
488 0 : FD_TEST( peer_map_ele_insert( ctx->peer_map, entry, ctx->peers ) );
489 0 : }
490 :
491 0 : entry->tombstoned = 0;
492 :
493 0 : static ulong const quic_socket_idx[ 2UL ] = {
494 0 : FD_GOSSIP_CONTACT_INFO_SOCKET_TPU_VOTE_QUIC,
495 0 : FD_GOSSIP_CONTACT_INFO_SOCKET_TPU_QUIC,
496 0 : };
497 :
498 0 : static ulong const udp_socket_idx[ 2UL ] = {
499 0 : FD_GOSSIP_CONTACT_INFO_SOCKET_TPU_VOTE,
500 0 : FD_GOSSIP_CONTACT_INFO_SOCKET_TPU,
501 0 : };
502 :
503 : /* At this point, entry is not in a map and *entry might still contain
504 : stale endpoint info (from a previous update). Only overwrite it if
505 : update actually contains endpoint info. */
506 :
507 0 : for( ulong i=0UL; i<2UL; i++ ) {
508 0 : if( FD_LIKELY( !msg->contact_info->value->sockets[ quic_socket_idx[ i ] ].is_ipv6 && msg->contact_info->value->sockets[ quic_socket_idx[ i ] ].ip4 ) ) {
509 0 : entry->quic_conns[ i ].quic_ip_addr = msg->contact_info->value->sockets[ quic_socket_idx[ i ] ].ip4;
510 0 : }
511 0 : ushort port = fd_ushort_bswap( msg->contact_info->value->sockets[ quic_socket_idx[ i ] ].port );
512 0 : if( FD_LIKELY( port ) ) {
513 0 : entry->quic_conns[ i ].quic_port = port;
514 0 : }
515 0 : }
516 :
517 0 : for( ulong i=0UL; i<2UL; i++ ) {
518 0 : if( FD_LIKELY( !msg->contact_info->value->sockets[ udp_socket_idx[ i ] ].is_ipv6 && msg->contact_info->value->sockets[ udp_socket_idx[ i ] ].ip4 ) ) {
519 0 : entry->udp_ip_addrs[ i ] = msg->contact_info->value->sockets[ udp_socket_idx[ i ] ].ip4;
520 0 : }
521 0 : if( FD_LIKELY( fd_ushort_bswap( msg->contact_info->value->sockets[ udp_socket_idx[ i ] ].port ) ) ) {
522 0 : entry->udp_ports [ i ] = fd_ushort_bswap( msg->contact_info->value->sockets[ udp_socket_idx[ i ] ].port );
523 0 : }
524 0 : }
525 0 : }
526 :
527 : static void
528 : report_signed_vote( uchar const * payload,
529 : fd_txn_t const * txn,
530 : uchar const * signatures,
531 0 : ulong vote_txn_sz ) {
532 0 : if( FD_LIKELY( !fd_event_tl ) ) return;
533 :
534 0 : if( FD_UNLIKELY( txn->instr_cnt!=1UL ) ) return;
535 0 : if( FD_UNLIKELY( txn->acct_addr_cnt<3UL || txn->acct_addr_cnt>4UL ) ) return;
536 0 : fd_txn_instr_t const * instr = &txn->instr[ 0 ];
537 0 : fd_acct_addr_t const * addrs = fd_txn_get_acct_addrs( txn, payload );
538 0 : if( FD_UNLIKELY( 0!=memcmp( addrs[ instr->program_id ].b, fd_solana_vote_program_id.uc, sizeof(fd_pubkey_t) ) ) ) return;
539 0 : if( FD_UNLIKELY( instr->acct_cnt!=2UL ) ) return;
540 0 : uchar const * instr_addrs = fd_txn_get_instr_accts( instr, payload );
541 0 : uchar const * instr_data = payload + instr->data_off;
542 0 : ulong instr_data_sz = instr->data_sz;
543 0 : if( FD_UNLIKELY( instr_data_sz<sizeof(uint) ) ) return;
544 0 : if( FD_UNLIKELY( FD_LOAD( uint, instr_data )!=FD_VOTE_IX_KIND_TOWER_SYNC ) ) return;
545 0 : instr_data += 4; instr_data_sz -= 4;
546 :
547 0 : fd_compact_tower_sync_serde_t sync[1];
548 0 : if( FD_UNLIKELY( 0!=fd_compact_tower_sync_de( sync, instr_data, instr_data_sz ) ) ) return;
549 0 : if( FD_UNLIKELY( sync->lockouts_cnt==0 || sync->lockouts_cnt>31 ) ) return;
550 :
551 0 : uchar const * rbh = fd_txn_get_recent_blockhash( txn, payload );
552 :
553 0 : fd_event_signed_vote_t ev = {0};
554 0 : FD_TEST( vote_txn_sz<=sizeof(ev.signed_txn) );
555 0 : fd_memcpy( ev.signed_txn, payload, vote_txn_sz );
556 0 : ev.signed_txn_len = vote_txn_sz;
557 0 : fd_memcpy( ev.vote_account, addrs[ instr_addrs[ 0 ] ].b, sizeof(fd_pubkey_t) );
558 0 : fd_memcpy( ev.vote_authority, addrs[ instr_addrs[ 1 ] ].b, sizeof(fd_pubkey_t) );
559 0 : fd_memcpy( ev.fee_payer, addrs[ 0 ].b, sizeof(fd_pubkey_t) );
560 0 : fd_memcpy( ev.signature, signatures, sizeof(fd_ed25519_sig_t) );
561 0 : fd_memcpy( ev.vote_bank_hash, sync->hash.uc, sizeof(fd_hash_t) );
562 0 : fd_memcpy( ev.vote_block_id, sync->block_id.uc, sizeof(fd_hash_t) );
563 0 : fd_memcpy( ev.txn_blockhash, rbh, sizeof(fd_hash_t) );
564 :
565 0 : ulong root_slot = fd_ulong_if( sync->root==ULONG_MAX, 0UL, sync->root );
566 0 : ulong slot = root_slot;
567 0 : for( ulong i=0UL; i<sync->lockouts_cnt; i++ ) {
568 0 : slot += sync->lockouts[ i ].offset;
569 0 : ev.tower[ i ].slot = slot;
570 0 : ev.tower[ i ].confirmation_count = (uchar)sync->lockouts[ i ].confirmation_count;
571 0 : }
572 0 : ev.tower_cnt = sync->lockouts_cnt;
573 0 : ev.vote_slot = slot; /* top of tower */
574 :
575 0 : fd_event_report_signed_vote( &ev );
576 0 : }
577 :
578 : static void
579 : handle_vote_msg( fd_txsend_tile_t * ctx,
580 : fd_stem_context_t * stem,
581 0 : fd_tower_slot_done_t const * slot_done ) {
582 0 : if( FD_UNLIKELY( slot_done->vote_slot==ULONG_MAX ) ) return;
583 0 : if( FD_UNLIKELY( !slot_done->has_vote_txn ) ) return;
584 :
585 0 : ctx->voted_slot = slot_done->vote_slot;
586 :
587 0 : fd_txn_m_t * txnm = fd_chunk_to_laddr( ctx->txsend_out->mem, ctx->txsend_out->chunk );
588 0 : FD_TEST( slot_done->vote_txn_sz<=FD_TXN_MTU );
589 0 : txnm->payload_sz = (ushort)slot_done->vote_txn_sz;
590 0 : txnm->source_ipv4 = ctx->src_ip_addr;
591 0 : txnm->source_tpu = FD_TXN_M_TPU_SOURCE_TXSEND;
592 0 : txnm->block_engine.bundle_id = 0UL;
593 0 : txnm->first_seen_nanos = slot_done->vote_created_nanos;
594 0 : fd_memcpy( fd_txn_m_payload( txnm ), slot_done->vote_txn, slot_done->vote_txn_sz );
595 :
596 0 : uchar txn_mem[ FD_TXN_MAX_SZ ] __attribute__((aligned(alignof(fd_txn_t))));
597 0 : txnm->txn_t_sz = (ushort)fd_txn_parse( slot_done->vote_txn, slot_done->vote_txn_sz, txn_mem, NULL );
598 0 : FD_TEST( txnm->txn_t_sz );
599 :
600 0 : uchar * payload = fd_txn_m_payload( txnm );
601 0 : fd_txn_t const * txn = (fd_txn_t const *)txn_mem;
602 :
603 0 : uchar * signatures = payload + txn->signature_off;
604 0 : uchar const * message = payload + txn->message_off;
605 0 : ulong message_sz = fd_txn_msg_sz( txn, slot_done->vote_txn_sz );
606 0 : fd_keyguard_client_vote_txn_sign( ctx->keyguard_client, signatures, slot_done->authority_idx, message, message_sz );
607 :
608 0 : FD_BASE58_ENCODE_64_BYTES( signatures, vote_sig_b58 );
609 0 : FD_LOG_INFO(( "vote txn for slot %lu created: %s", slot_done->vote_slot, vote_sig_b58 ));
610 :
611 0 : report_signed_vote( payload, txn, signatures, slot_done->vote_txn_sz );
612 :
613 0 : for( ulong i=0UL; i<3UL; i++ ) {
614 0 : ulong target_slot = slot_done->vote_slot+1UL + i*FD_EPOCH_SLOTS_PER_ROTATION;
615 0 : fd_pubkey_t const * leader = fd_multi_epoch_leaders_get_leader_for_slot( ctx->mleaders, target_slot );
616 0 : if( FD_UNLIKELY( !leader ) ) {
617 0 : FD_LOG_WARNING(( "no leader found for slot %lu", target_slot ));
618 0 : continue;
619 0 : }
620 0 : send_vote_to_leader( ctx, leader, payload, slot_done->vote_txn_sz );
621 0 : }
622 :
623 0 : ulong msg_sz = fd_txn_m_realized_footprint( txnm, 0, 0 );
624 0 : ulong tspub_comp = fd_frag_meta_ts_comp( fd_tickcount() );
625 0 : fd_stem_publish( stem, ctx->txsend_out->idx, 1UL, ctx->txsend_out->chunk, msg_sz, 0UL, 0UL, tspub_comp );
626 0 : ctx->txsend_out->chunk = fd_dcache_compact_next( ctx->txsend_out->chunk, msg_sz, ctx->txsend_out->chunk0, ctx->txsend_out->wmark );
627 0 : }
628 :
629 :
630 : static inline int
631 : before_frag( fd_txsend_tile_t * ctx,
632 : ulong in_idx,
633 : ulong seq,
634 0 : ulong sig ) {
635 0 : fd_startup_gate_busy( ctx->startup_gate );
636 :
637 0 : if( FD_LIKELY( ctx->in_kind[ in_idx ]==IN_KIND_TOWER ) ) ctx->tower_in_expect_seq = seq+1UL;
638 0 : if( FD_UNLIKELY( ctx->halt_net_frags && ctx->in_kind[ in_idx ]==IN_KIND_NET ) ) return -1;
639 :
640 0 : if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_GOSSIP ) ) {
641 0 : return sig!=FD_GOSSIP_UPDATE_TAG_CONTACT_INFO && sig!=FD_GOSSIP_UPDATE_TAG_CONTACT_INFO_REMOVE;
642 0 : } else if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_TOWER ) ) {
643 0 : return sig!=FD_TOWER_SIG_SLOT_DONE;
644 0 : }
645 :
646 0 : return 0;
647 0 : }
648 :
649 : static void
650 : during_frag( fd_txsend_tile_t * ctx,
651 : ulong in_idx,
652 : ulong seq,
653 : ulong sig,
654 : ulong chunk,
655 : ulong sz,
656 0 : ulong ctl ) {
657 0 : (void)seq; (void)sig;
658 :
659 0 : ctx->chunk = chunk;
660 :
661 0 : if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_EPOCH ) ) {
662 0 : if( FD_UNLIKELY( chunk<ctx->in[ in_idx ].chunk0 || chunk>ctx->in[ in_idx ].wmark ) )
663 0 : FD_LOG_ERR(( "chunk %lu %lu corrupt, not in range [%lu,%lu,%lu]", chunk, sz, ctx->in[in_idx].chunk0, ctx->in[in_idx].wmark, ctx->in[ in_idx ].mtu ));
664 :
665 0 : fd_epoch_info_msg_t const * msg = fd_chunk_to_laddr_const( ctx->in[ in_idx ].mem, chunk );
666 0 : FD_TEST( msg->staked_vote_cnt<=MAX_STAKE_WEIGHTS ); /* implicit sz verification since sz field on frag_meta too small */
667 0 : FD_TEST( msg->staked_id_cnt<=MAX_STAKE_WEIGHTS );
668 0 : } else {
669 0 : if( FD_UNLIKELY( chunk<ctx->in[ in_idx ].chunk0 || chunk>ctx->in[ in_idx ].wmark || sz>ctx->in[ in_idx ].mtu ) )
670 0 : FD_LOG_ERR(( "chunk %lu %lu corrupt, not in range [%lu,%lu,%lu]", chunk, sz, ctx->in[in_idx].chunk0, ctx->in[in_idx].wmark, ctx->in[ in_idx ].mtu ));
671 0 : }
672 :
673 0 : if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_NET ) ) {
674 0 : void const * src = fd_net_rx_translate_frag( &ctx->net_in_bounds[ in_idx ], chunk, ctl, sz );
675 0 : fd_memcpy( ctx->quic_buf, src, sz );
676 0 : }
677 0 : }
678 :
679 : static void
680 : after_frag( fd_txsend_tile_t * ctx,
681 : ulong in_idx,
682 : ulong seq,
683 : ulong sig,
684 : ulong sz,
685 : ulong tsorig,
686 : ulong tspub,
687 0 : fd_stem_context_t * stem ) {
688 0 : (void)seq; (void)sig; (void)tsorig; (void)tspub;
689 :
690 0 : if( FD_LIKELY( ctx->in_kind[ in_idx ]==IN_KIND_NET ) ) {
691 0 : uchar * ip_packet = ctx->quic_buf+sizeof(fd_eth_hdr_t);
692 0 : ulong ip_packet_sz = sz-sizeof(fd_eth_hdr_t);
693 0 : fd_quic_process_packet( ctx->quic, ip_packet, ip_packet_sz, fd_log_wallclock() );
694 0 : } else if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_GOSSIP ) ) {
695 0 : if( FD_LIKELY( sig==FD_GOSSIP_UPDATE_TAG_CONTACT_INFO ) ) handle_contact_info_update( ctx, fd_chunk_to_laddr_const( ctx->in[ in_idx ].mem, ctx->chunk ) );
696 0 : else handle_contact_info_remove( ctx, fd_chunk_to_laddr_const( ctx->in[ in_idx ].mem, ctx->chunk ) );
697 0 : } else if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_TOWER ) ) {
698 0 : handle_vote_msg( ctx, stem, fd_chunk_to_laddr_const( ctx->in[ in_idx ].mem, ctx->chunk ) );
699 0 : } else if( FD_UNLIKELY( ctx->in_kind[ in_idx ]==IN_KIND_EPOCH ) ) {
700 0 : fd_multi_epoch_leaders_epoch_msg_init( ctx->mleaders, fd_chunk_to_laddr_const( ctx->in[ in_idx ].mem, ctx->chunk ) );
701 0 : fd_multi_epoch_leaders_stake_msg_fini( ctx->mleaders );
702 0 : ctx->leader_schedules++;
703 0 : } else {
704 0 : FD_LOG_ERR(( "unknown in_kind %d on link %lu", ctx->in_kind[ in_idx ], in_idx ));
705 0 : }
706 0 : }
707 :
708 : static void
709 : privileged_init( fd_topo_t const * topo,
710 0 : fd_topo_tile_t const * tile ) {
711 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
712 :
713 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
714 0 : fd_txsend_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_txsend_tile_t), sizeof(fd_txsend_tile_t) );
715 :
716 0 : if( FD_UNLIKELY( !strcmp( tile->txsend.identity_key_path, "" ) ) )
717 0 : FD_LOG_ERR(( "identity_key_path not set" ));
718 :
719 0 : ctx->identity_key[ 0 ] = *(fd_pubkey_t const *)fd_type_pun_const( fd_keyload_load( tile->txsend.identity_key_path, /* pubkey only: */ 1 ) );
720 :
721 0 : FD_TEST( fd_rng_secure( &ctx->seed, sizeof(ctx->seed) ) );
722 0 : }
723 :
724 : static inline fd_txsend_out_t
725 : out1( fd_topo_t const * topo,
726 : fd_topo_tile_t const * tile,
727 0 : char const * name ) {
728 0 : ulong idx = ULONG_MAX;
729 :
730 0 : for( ulong i=0UL; i<tile->out_cnt; i++ ) {
731 0 : fd_topo_link_t const * link = &topo->links[ tile->out_link_id[ i ] ];
732 0 : if( !strcmp( link->name, name ) ) {
733 0 : if( FD_UNLIKELY( idx!=ULONG_MAX ) ) FD_LOG_ERR(( "tile %s:%lu had multiple output links named %s but expected one", tile->name, tile->kind_id, name ));
734 0 : idx = i;
735 0 : }
736 0 : }
737 :
738 0 : if( FD_UNLIKELY( idx==ULONG_MAX ) ) FD_LOG_ERR(( "tile %s:%lu had no output link named %s", tile->name, tile->kind_id, name ));
739 :
740 0 : void * mem = topo->workspaces[ topo->objs[ topo->links[ tile->out_link_id[ idx ] ].dcache_obj_id ].wksp_id ].wksp;
741 0 : ulong chunk0 = fd_dcache_compact_chunk0( mem, topo->links[ tile->out_link_id[ idx ] ].dcache );
742 0 : ulong wmark = fd_dcache_compact_wmark ( mem, topo->links[ tile->out_link_id[ idx ] ].dcache, topo->links[ tile->out_link_id[ idx ] ].mtu );
743 :
744 0 : return (fd_txsend_out_t){ .idx = idx, .mem = mem, .chunk0 = chunk0, .wmark = wmark, .chunk = chunk0 };
745 0 : }
746 :
747 : static void
748 : unprivileged_init( fd_topo_t const * topo,
749 0 : fd_topo_tile_t const * tile ) {
750 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
751 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
752 0 : fd_txsend_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_txsend_tile_t), sizeof(fd_txsend_tile_t) );
753 0 : void * _quic = FD_SCRATCH_ALLOC_APPEND( l, fd_quic_align(), fd_quic_footprint( &quic_limits ) );
754 0 : void * _peer_map = FD_SCRATCH_ALLOC_APPEND( l, peer_map_align(), peer_map_footprint( 2UL*FD_CONTACT_INFO_TABLE_SIZE ) );
755 :
756 0 : ctx->quic = fd_quic_join( fd_quic_new( _quic, &quic_limits ) );
757 0 : FD_TEST( ctx->quic );
758 :
759 0 : ctx->leader_schedules = 0UL;
760 :
761 0 : ctx->mleaders = fd_multi_epoch_leaders_join( fd_multi_epoch_leaders_new( ctx->mleaders_mem ) );
762 0 : FD_TEST( ctx->mleaders );
763 :
764 0 : ctx->peer_map = peer_map_join( peer_map_new( _peer_map, 2UL*FD_CONTACT_INFO_TABLE_SIZE, ctx->seed ) );
765 0 : FD_TEST( ctx->peer_map );
766 :
767 0 : fd_aio_t * quic_tx_aio = fd_aio_join( fd_aio_new( ctx->quic_tx_aio, ctx, quic_tx_aio_send ) );
768 0 : FD_TEST( quic_tx_aio );
769 0 : fd_quic_set_aio_net_tx( ctx->quic, quic_tx_aio );
770 :
771 0 : ctx->quic->config.role = FD_QUIC_ROLE_CLIENT;
772 0 : ctx->quic->config.idle_timeout = 30e9L;
773 0 : ctx->quic->config.ack_delay = 25e6L;
774 0 : ctx->quic->config.keep_alive = 1;
775 0 : ctx->quic->config.sign = quic_tls_cv_sign;
776 0 : ctx->quic->config.sign_ctx = ctx;
777 0 : fd_memcpy( ctx->quic->config.identity_public_key, ctx->identity_key, sizeof(ctx->identity_key) );
778 :
779 0 : ctx->quic->cb.conn_final = quic_conn_final;
780 0 : ctx->quic->cb.quic_ctx = ctx;
781 :
782 0 : FD_TEST( fd_quic_init( ctx->quic ));
783 :
784 0 : for( ulong i=0UL; i<FD_CONTACT_INFO_TABLE_SIZE; i++ ) {
785 0 : ctx->peers[ i ] = (peer_entry_t){0};
786 0 : }
787 :
788 0 : ctx->conns_len = 0UL;
789 0 : ctx->voted_slot = ULONG_MAX;
790 0 : ctx->net_id = 0;
791 :
792 0 : ctx->src_ip_addr = tile->txsend.ip_addr;
793 0 : ctx->src_port = tile->txsend.txsend_src_port;
794 0 : fd_ip4_udp_hdr_init( ctx->packet_hdr, FD_TXN_MTU, ctx->src_ip_addr, ctx->src_port );
795 :
796 0 : fd_startup_gate_init( ctx->startup_gate, topo, tile->in_cnt );
797 :
798 0 : FD_TEST( tile->in_cnt<sizeof(ctx->in_kind)/sizeof(ctx->in_kind[ 0 ]) );
799 0 : for( ulong i=0UL; i<tile->in_cnt; i++ ) {
800 0 : fd_topo_link_t const * link = &topo->links[ tile->in_link_id[ i ] ];
801 0 : fd_topo_wksp_t const * link_wksp = &topo->workspaces[ topo->objs[ link->dcache_obj_id ].wksp_id ];
802 :
803 0 : ctx->in[ i ].mem = link_wksp->wksp;
804 0 : ctx->in[ i ].chunk0 = fd_dcache_compact_chunk0( ctx->in[ i ].mem, link->dcache );
805 0 : ctx->in[ i ].wmark = fd_dcache_compact_wmark ( ctx->in[ i ].mem, link->dcache, link->mtu );
806 0 : ctx->in[ i ].mtu = link->mtu;
807 :
808 0 : if( !strcmp( link->name, "net_txsend" ) ) {
809 0 : fd_net_rx_bounds_init( &ctx->net_in_bounds[ i ], link->dcache );
810 0 : ctx->in_kind[ i ] = IN_KIND_NET;
811 0 : } else if( !strcmp( link->name, "gossip_out" ) ) ctx->in_kind[ i ] = IN_KIND_GOSSIP;
812 0 : else if( !strcmp( link->name, "replay_epoch" ) ) ctx->in_kind[ i ] = IN_KIND_EPOCH;
813 0 : else if( !strcmp( link->name, "tower_out" ) ) ctx->in_kind[ i ] = IN_KIND_TOWER;
814 0 : else if( !strcmp( link->name, "sign_txsend" ) ) ctx->in_kind[ i ] = IN_KIND_SIGN;
815 0 : else FD_LOG_ERR(( "unexpected input link name %s", link->name ));
816 0 : }
817 :
818 0 : *ctx->txsend_out = out1( topo, tile, "txsend_out" );
819 0 : *ctx->net_out = out1( topo, tile, "txsend_net" );
820 :
821 0 : ulong sign_in_idx = fd_topo_find_tile_in_link ( topo, tile, "sign_txsend", tile->kind_id );
822 0 : ulong sign_out_idx = fd_topo_find_tile_out_link( topo, tile, "txsend_sign", tile->kind_id );
823 0 : FD_TEST( sign_in_idx!=ULONG_MAX );
824 0 : fd_topo_link_t const * sign_in = &topo->links[ tile->in_link_id[ sign_in_idx ] ];
825 0 : fd_topo_link_t const * sign_out = &topo->links[ tile->out_link_id[ sign_out_idx ] ];
826 0 : if( FD_UNLIKELY( !fd_keyguard_client_join( fd_keyguard_client_new( ctx->keyguard_client,
827 0 : sign_out->mcache,
828 0 : sign_out->dcache,
829 0 : sign_in->mcache,
830 0 : sign_in->dcache,
831 0 : sign_out->mtu,
832 0 : sign_in->mtu ) ) ) ) {
833 0 : FD_LOG_ERR(( "failed to construct keyguard" ));
834 0 : }
835 :
836 0 : ctx->keyswitch = fd_keyswitch_join( fd_topo_obj_laddr( topo, tile->id_keyswitch_obj_id ) );
837 0 : FD_TEST( ctx->keyswitch );
838 :
839 0 : ctx->av_keyswitch = NULL;
840 0 : if( FD_UNLIKELY( tile->av_keyswitch_obj_id!=ULONG_MAX ) ) {
841 0 : ctx->av_keyswitch = fd_keyswitch_join( fd_topo_obj_laddr( topo, tile->av_keyswitch_obj_id ) );
842 0 : FD_TEST( ctx->av_keyswitch );
843 0 : }
844 :
845 0 : ctx->tower_in_expect_seq = 0UL;
846 0 : ctx->halt_net_frags = 0;
847 :
848 0 : fd_histf_join( fd_histf_new( ctx->quic->metrics.service_duration, FD_MHIST_SECONDS_MIN( TXSEND, SERVICE_DURATION_SECONDS ),
849 0 : FD_MHIST_SECONDS_MAX( TXSEND, SERVICE_DURATION_SECONDS ) ) );
850 0 : fd_histf_join( fd_histf_new( ctx->quic->metrics.receive_duration, FD_MHIST_SECONDS_MIN( TXSEND, RX_DURATION_SECONDS ),
851 0 : FD_MHIST_SECONDS_MAX( TXSEND, RX_DURATION_SECONDS ) ) );
852 :
853 0 : ulong scratch_top = FD_SCRATCH_ALLOC_FINI( l, scratch_align() );
854 0 : if( FD_UNLIKELY( scratch_top > (ulong)scratch + scratch_footprint( tile ) ) )
855 0 : FD_LOG_ERR(( "scratch overflow %lu %lu %lu", scratch_top - (ulong)scratch - scratch_footprint( tile ), scratch_top, (ulong)scratch + scratch_footprint( tile ) ));
856 0 : }
857 :
858 : static ulong
859 : populate_allowed_seccomp( fd_topo_t const * topo FD_PARAM_UNUSED,
860 : fd_topo_tile_t const * tile FD_PARAM_UNUSED,
861 : ulong out_cnt,
862 0 : struct sock_filter * out ) {
863 :
864 0 : populate_sock_filter_policy_fd_txsend_tile( out_cnt, out, (uint)fd_log_private_logfile_fd() );
865 0 : return sock_filter_policy_fd_txsend_tile_instr_cnt;
866 0 : }
867 :
868 : static ulong
869 : populate_allowed_fds( fd_topo_t const * topo FD_PARAM_UNUSED,
870 : fd_topo_tile_t const * tile FD_PARAM_UNUSED,
871 : ulong out_fds_cnt,
872 0 : int * out_fds ) {
873 0 : if( FD_UNLIKELY( out_fds_cnt<2UL ) ) FD_LOG_ERR(( "out_fds_cnt %lu", out_fds_cnt ));
874 :
875 0 : ulong out_cnt = 0;
876 0 : out_fds[ out_cnt++ ] = 2UL; /* stderr */
877 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) )
878 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
879 0 : return out_cnt;
880 0 : }
881 :
882 0 : #define STEM_BURST 1UL
883 0 : #define STEM_LAZY (128L*3000L)
884 :
885 0 : #define STEM_CALLBACK_CONTEXT_TYPE fd_txsend_tile_t
886 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_txsend_tile_t)
887 :
888 0 : #define STEM_CALLBACK_DURING_HOUSEKEEPING during_housekeeping
889 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
890 0 : #define STEM_CALLBACK_AFTER_CREDIT after_credit
891 0 : #define STEM_CALLBACK_BEFORE_FRAG before_frag
892 0 : #define STEM_CALLBACK_DURING_FRAG during_frag
893 0 : #define STEM_CALLBACK_AFTER_FRAG after_frag
894 :
895 : #include "../../disco/stem/fd_stem.c"
896 :
897 : static ulong
898 0 : max_event_sz( fd_topo_tile_t const * tile FD_PARAM_UNUSED ) {
899 0 : return sizeof(fd_event_signed_vote_t);
900 0 : }
901 :
902 : fd_topo_run_tile_t fd_tile_txsend = {
903 : .name = "txsend",
904 : .max_event_sz = max_event_sz,
905 : .populate_allowed_seccomp = populate_allowed_seccomp,
906 : .populate_allowed_fds = populate_allowed_fds,
907 : .scratch_align = scratch_align,
908 : .scratch_footprint = scratch_footprint,
909 : .privileged_init = privileged_init,
910 : .unprivileged_init = unprivileged_init,
911 : .run = stem_run,
912 : };
|