Line data Source code
1 : /* The repair command spawns a smaller topology for profiling the repair
2 : tile. This is a standalone application, and it can be run in mainnet,
3 : testnet and/or a private cluster. */
4 :
5 : #include "../../../discof/repair/fd_repair_tile.h"
6 : #include "../../../disco/shred/fd_shred_tile.h"
7 : #include "../../../disco/net/fd_net_tile.h"
8 : #include "../../../disco/tiles.h"
9 : #include "../../../disco/topo/fd_topob.h"
10 : #include "../../../disco/topo/fd_cpu_topo.h"
11 : #include "../../../util/pod/fd_pod_format.h"
12 :
13 : #include "../../firedancer/topology.h"
14 : #include "../../shared/commands/configure/configure.h"
15 : #include "../../shared/commands/run/run.h" /* initialize_workspaces */
16 : #include "../../shared/fd_config.h" /* config_t */
17 : #include "../../shared_dev/commands/dev.h"
18 : #include "../../../disco/topo/fd_topob.h"
19 : #include "../../../util/pod/fd_pod_format.h"
20 : #include "../../../waltz/resolv/fd_io_readline.h"
21 : #include "../../platform/fd_sys_util.h"
22 : #include "../../shared/commands/monitor/helper.h"
23 : #include "../../../disco/metrics/fd_metrics.h"
24 : #include "../../../discof/restore/utils/fd_ssmanifest_parser.h"
25 : #include "../../../discof/genesis/fd_genesi_tile.h"
26 : #include "../../../flamenco/runtime/sysvar/fd_sysvar_epoch_schedule.h"
27 : #include "../../../flamenco/stakes/fd_stake_weight_sort.h"
28 : #include "../../../flamenco/leaders/fd_leaders_base.h"
29 : #include "../../../discof/repair/fd_repair_tile.c"
30 :
31 : #include "gossip.h"
32 : #include "core_subtopo.h"
33 :
34 : #include <unistd.h> /* pause */
35 : #include <fcntl.h>
36 : #include <stdio.h>
37 : #include <stdlib.h>
38 : #include <termios.h>
39 : #include <errno.h>
40 :
41 : extern action_t fd_action_repair;
42 :
43 : struct fd_location_info {
44 : ulong ip4_addr; /* for map key convenience */
45 : char location[ 128 ];
46 : };
47 : typedef struct fd_location_info fd_location_info_t;
48 :
49 : #define MAP_NAME fd_location_table
50 0 : #define MAP_T fd_location_info_t
51 0 : #define MAP_KEY ip4_addr
52 0 : #define MAP_LG_SLOT_CNT 16
53 : #define MAP_MEMOIZE 0
54 : #include "../../../util/tmpl/fd_map.c"
55 :
56 : uchar __attribute__((aligned(alignof(fd_location_info_t)))) location_table_mem[ sizeof(fd_location_info_t) * (1 << 16 ) ];
57 :
58 : static struct termios termios_backup;
59 :
60 : static void
61 0 : restore_terminal( void ) {
62 0 : (void)tcsetattr( STDIN_FILENO, TCSANOW, &termios_backup );
63 0 : }
64 :
65 : extern fd_topo_obj_callbacks_t * CALLBACKS[];
66 :
67 : fd_topo_run_tile_t
68 : fdctl_tile_run( fd_topo_tile_t const * tile );
69 :
70 0 : #define MANIFEST_LOAD_MAX_SZ (2UL * FD_SHMEM_GIGANTIC_PAGE_SZ)
71 :
72 : /* https://github.com/anza-xyz/agave/blob/v3.1.8/runtime/src/snapshot_bank_utils.rs#L632 */
73 : static int
74 0 : repair_verify_epoch_stakes( fd_snapshot_manifest_t const * manifest ) {
75 0 : fd_epoch_schedule_t epoch_schedule = (fd_epoch_schedule_t){
76 0 : .slots_per_epoch = manifest->epoch_schedule_params.slots_per_epoch,
77 0 : .leader_schedule_slot_offset = manifest->epoch_schedule_params.leader_schedule_slot_offset,
78 0 : .warmup = manifest->epoch_schedule_params.warmup,
79 0 : .first_normal_epoch = manifest->epoch_schedule_params.first_normal_epoch,
80 0 : .first_normal_slot = manifest->epoch_schedule_params.first_normal_slot,
81 0 : };
82 :
83 0 : ulong min_required_epoch = fd_slot_to_epoch( &epoch_schedule, manifest->slot, NULL );
84 0 : ulong max_required_epoch = fd_slot_to_leader_schedule_epoch( &epoch_schedule, manifest->slot );
85 :
86 0 : for( ulong i=min_required_epoch; i<=max_required_epoch; i++ ) {
87 0 : int found = 0;
88 0 : for( ulong j=0UL; j<FD_RUNTIME_MANIFEST_EPOCH_STAKES_LEN; j++ ) {
89 0 : if( manifest->epoch_stakes[j].epoch==i ) {
90 0 : found = 1;
91 0 : break;
92 0 : }
93 0 : }
94 0 : if( FD_UNLIKELY( !found ) ) {
95 0 : FD_LOG_WARNING(( "stakes not found for epoch %lu in manifest", i ));
96 0 : return -1;
97 0 : }
98 0 : }
99 0 : return 0;
100 0 : }
101 :
102 : static inline ulong
103 : repair_generate_epoch_info_msg( ulong epoch,
104 : fd_epoch_schedule_t const * epoch_schedule,
105 : fd_snapshot_manifest_epoch_stakes_t const * epoch_stakes,
106 0 : ulong * epoch_info_msg_out ) {
107 0 : fd_epoch_info_msg_t * epoch_info_msg = (fd_epoch_info_msg_t *)fd_type_pun( epoch_info_msg_out );
108 0 : fd_vote_stake_weight_t * stake_weights = fd_epoch_info_msg_stake_weights( epoch_info_msg );
109 :
110 0 : epoch_info_msg->epoch = epoch;
111 0 : epoch_info_msg->start_slot = fd_epoch_slot0( epoch_schedule, epoch );
112 0 : epoch_info_msg->slot_cnt = fd_epoch_slot_cnt( epoch_schedule, epoch );
113 0 : epoch_info_msg->ns_per_slot = FD_SLOT_PARAMS_400MS.ns_per_slot;
114 :
115 0 : fd_memset( &epoch_info_msg->features, 0xFF, sizeof(fd_features_t) );
116 :
117 0 : ulong idx = 0UL;
118 0 : for( ulong i=0UL; i<epoch_stakes->vote_stakes_len; i++ ) {
119 0 : ulong stake = epoch_stakes->vote_stakes[ i ].stake;
120 0 : if( FD_UNLIKELY( !stake ) ) continue;
121 0 : stake_weights[ idx ].stake = stake;
122 0 : memcpy( stake_weights[ idx ].id_key.uc, epoch_stakes->vote_stakes[ i ].identity, sizeof(fd_pubkey_t) );
123 0 : memcpy( stake_weights[ idx ].vote_key.uc, epoch_stakes->vote_stakes[ i ].vote, sizeof(fd_pubkey_t) );
124 0 : idx++;
125 0 : }
126 0 : epoch_info_msg->staked_vote_cnt = idx;
127 0 : sort_vote_weights_by_stake_vote_inplace( stake_weights, idx );
128 :
129 0 : fd_stake_weight_t * id_weights = fd_epoch_info_msg_id_weights( epoch_info_msg );
130 0 : epoch_info_msg->staked_id_cnt = compute_id_weights_from_vote_weights( id_weights, stake_weights, epoch_info_msg->staked_vote_cnt );
131 0 : FD_TEST( idx<=MAX_SHRED_DESTS );
132 :
133 0 : epoch_info_msg->epoch_schedule = *epoch_schedule;
134 0 : return fd_epoch_info_msg_sz( epoch_info_msg->staked_vote_cnt, epoch_info_msg->staked_id_cnt );
135 0 : }
136 :
137 : /* repair_load_manifest loads the snapshot manifest from disk and
138 : pre-populates the snapin_manif and replay_epoch dcache links so
139 : that consumer tiles see the data on their first poll cycle. */
140 : static void
141 : repair_load_manifest( fd_topo_t * topo,
142 0 : char const * manifest_path ) {
143 0 : if( FD_UNLIKELY( !manifest_path || !manifest_path[0] ) ) return;
144 :
145 : /* Parse manifest */
146 :
147 0 : int fd = open( manifest_path, O_RDONLY );
148 0 : if( FD_UNLIKELY( fd<0 ) ) FD_LOG_ERR(( "open(%s) failed (%d-%s)", manifest_path, errno, fd_io_strerror( errno ) ));
149 :
150 0 : fd_snapshot_manifest_t * manifest = aligned_alloc( alignof(fd_snapshot_manifest_t), sizeof(fd_snapshot_manifest_t) );
151 0 : FD_TEST( manifest );
152 0 : for( ulong i=0UL; i<FD_RUNTIME_MANIFEST_EPOCH_STAKES_LEN; i++ ) manifest->epoch_stakes[i].epoch = ULONG_MAX;
153 :
154 0 : uchar * buf = aligned_alloc( 128UL, MANIFEST_LOAD_MAX_SZ );
155 0 : FD_TEST( buf );
156 0 : ulong buf_sz = 0;
157 0 : FD_TEST( !fd_io_read( fd, buf, 0UL, MANIFEST_LOAD_MAX_SZ-1UL, &buf_sz ) );
158 0 : close( fd );
159 :
160 0 : fd_ssmanifest_parser_t * parser = fd_ssmanifest_parser_join( fd_ssmanifest_parser_new(
161 0 : aligned_alloc( fd_ssmanifest_parser_align(), fd_ssmanifest_parser_footprint() ) ) );
162 0 : FD_TEST( parser );
163 0 : fd_ssmanifest_parser_init( parser, manifest );
164 0 : int parser_err = fd_ssmanifest_parser_consume( parser, buf, buf_sz );
165 0 : FD_TEST( parser_err!=FD_SSMANIFEST_PARSER_ADVANCE_ERROR );
166 0 : FD_TEST( fd_ssmanifest_parser_fini( parser )==FD_SSMANIFEST_PARSER_ADVANCE_DONE );
167 0 : free( parser );
168 0 : free( buf );
169 :
170 0 : FD_LOG_NOTICE(( "manifest bank slot %lu", manifest->slot ));
171 0 : FD_TEST( !repair_verify_epoch_stakes( manifest ) );
172 :
173 : /* Update root_slot fseq */
174 :
175 0 : ulong root_slot_obj_id = fd_pod_queryf_ulong( topo->props, ULONG_MAX, "root_slot" );
176 0 : if( FD_LIKELY( root_slot_obj_id!=ULONG_MAX ) ) {
177 0 : ulong * root_fseq = fd_fseq_join( fd_topo_obj_laddr( topo, root_slot_obj_id ) );
178 0 : FD_TEST( root_fseq );
179 0 : fd_fseq_update( root_fseq, manifest->slot );
180 0 : }
181 :
182 : /* Publish manifest to snapin_manif dcache */
183 :
184 0 : ulong snap_link_idx = fd_topo_find_link( topo, "snapin_manif", 0UL );
185 0 : FD_TEST( snap_link_idx!=ULONG_MAX );
186 0 : fd_topo_link_t * snap_link = &topo->links[ snap_link_idx ];
187 0 : fd_wksp_t * snap_mem = topo->workspaces[ topo->objs[ snap_link->dcache_obj_id ].wksp_id ].wksp;
188 0 : ulong snap_chunk0 = fd_dcache_compact_chunk0( snap_mem, snap_link->dcache );
189 0 : ulong snap_wmark = fd_dcache_compact_wmark ( snap_mem, snap_link->dcache, snap_link->mtu );
190 0 : ulong snap_chunk = snap_chunk0;
191 :
192 0 : uchar * snap_dst = fd_chunk_to_laddr( snap_mem, snap_chunk );
193 0 : memcpy( snap_dst, manifest, sizeof(fd_snapshot_manifest_t) );
194 0 : fd_mcache_publish( snap_link->mcache, snap_link->depth, 0UL,
195 0 : fd_ssmsg_sig( FD_SSMSG_MANIFEST_INCREMENTAL ),
196 0 : snap_chunk, sizeof(fd_snapshot_manifest_t), 0UL, 0UL, 0UL );
197 0 : snap_chunk = fd_dcache_compact_next( snap_chunk, sizeof(fd_snapshot_manifest_t), snap_chunk0, snap_wmark );
198 :
199 0 : fd_mcache_publish( snap_link->mcache, snap_link->depth, 1UL,
200 0 : fd_ssmsg_sig( FD_SSMSG_DONE ), 0UL, 0UL, 0UL, 0UL, 0UL );
201 :
202 : /* Publish epoch stake weights to replay_epoch dcache */
203 :
204 0 : ulong epoch_link_idx = fd_topo_find_link( topo, "replay_epoch", 0UL );
205 0 : FD_TEST( epoch_link_idx!=ULONG_MAX );
206 0 : fd_topo_link_t * epoch_link = &topo->links[ epoch_link_idx ];
207 0 : fd_wksp_t * epoch_mem = topo->workspaces[ topo->objs[ epoch_link->dcache_obj_id ].wksp_id ].wksp;
208 0 : ulong epoch_chunk0 = fd_dcache_compact_chunk0( epoch_mem, epoch_link->dcache );
209 0 : ulong epoch_wmark = fd_dcache_compact_wmark ( epoch_mem, epoch_link->dcache, epoch_link->mtu );
210 0 : ulong epoch_chunk = epoch_chunk0;
211 0 : ulong epoch_seq = 0UL;
212 :
213 : /* Construct fd_epoch_schedule_t field-by-field rather than type-punning
214 : from the unpacked manifest struct (fd_epoch_schedule_t is packed). */
215 0 : fd_epoch_schedule_t schedule_local;
216 0 : schedule_local.slots_per_epoch = manifest->epoch_schedule_params.slots_per_epoch;
217 0 : schedule_local.leader_schedule_slot_offset = manifest->epoch_schedule_params.leader_schedule_slot_offset;
218 0 : schedule_local.warmup = manifest->epoch_schedule_params.warmup;
219 0 : schedule_local.first_normal_epoch = manifest->epoch_schedule_params.first_normal_epoch;
220 0 : schedule_local.first_normal_slot = manifest->epoch_schedule_params.first_normal_slot;
221 0 : fd_epoch_schedule_t const * schedule = &schedule_local;
222 0 : ulong epoch = fd_slot_to_epoch( schedule, manifest->slot, NULL );
223 :
224 0 : ulong epoch_stakes_base = epoch > 0UL ? epoch - 1UL : 0UL;
225 0 : ulong leader_schedule_epoch = fd_slot_to_leader_schedule_epoch( schedule, manifest->slot );
226 0 : ulong cur_idx = epoch - epoch_stakes_base;
227 0 : FD_TEST( cur_idx < FD_RUNTIME_MANIFEST_EPOCH_STAKES_LEN );
228 :
229 0 : ulong * epoch_dst = fd_chunk_to_laddr( epoch_mem, epoch_chunk );
230 0 : ulong epoch_sz = repair_generate_epoch_info_msg( epoch, schedule, &manifest->epoch_stakes[cur_idx], epoch_dst );
231 0 : fd_mcache_publish( epoch_link->mcache, epoch_link->depth, epoch_seq,
232 0 : 4UL, epoch_chunk, epoch_sz, 0UL, 0UL, fd_frag_meta_ts_comp( fd_tickcount() ) );
233 0 : epoch_chunk = fd_dcache_compact_next( epoch_chunk, epoch_sz, epoch_chunk0, epoch_wmark );
234 0 : epoch_seq++;
235 0 : FD_LOG_NOTICE(( "sending current epoch stake weights - epoch: %lu", epoch ));
236 :
237 0 : if( leader_schedule_epoch >= epoch + 1UL ) {
238 0 : ulong next_idx = epoch + 1UL - epoch_stakes_base;
239 0 : FD_TEST( next_idx < FD_RUNTIME_MANIFEST_EPOCH_STAKES_LEN );
240 :
241 0 : epoch_dst = fd_chunk_to_laddr( epoch_mem, epoch_chunk );
242 0 : epoch_sz = repair_generate_epoch_info_msg( epoch + 1UL, schedule, &manifest->epoch_stakes[next_idx], epoch_dst );
243 0 : fd_mcache_publish( epoch_link->mcache, epoch_link->depth, epoch_seq,
244 0 : 4UL, epoch_chunk, epoch_sz, 0UL, 0UL, fd_frag_meta_ts_comp( fd_tickcount() ) );
245 0 : epoch_chunk = fd_dcache_compact_next( epoch_chunk, epoch_sz, epoch_chunk0, epoch_wmark );
246 0 : epoch_seq++;
247 0 : FD_LOG_NOTICE(( "sending next epoch stake weights - epoch: %lu", epoch + 1UL ));
248 0 : }
249 0 : (void)epoch_chunk;
250 :
251 0 : free( manifest );
252 0 : }
253 :
254 : /* repair_topo is a subset of "src/app/firedancer/topology.c" at commit
255 : 0d8386f4f305bb15329813cfe4a40c3594249e96, slightly modified to work
256 : as a repair catchup. TODO ideally, one should invoke the firedancer
257 : topology first, and exclude the parts that are not needed, instead of
258 : manually generating new topologies for every command. This would
259 : also guarantee that the catchup is replicating (as close as possible)
260 : the full topology. */
261 : static void
262 0 : repair_topo( config_t * config ) {
263 0 : ulong net_tile_cnt = config->layout.net_tile_count;
264 0 : ulong shred_tile_cnt = config->layout.shred_tile_count;
265 0 : ulong quic_tile_cnt = config->layout.quic_tile_count;
266 0 : ulong sign_tile_cnt = config->firedancer.layout.sign_tile_count;
267 0 : ulong gossvf_tile_cnt = config->firedancer.layout.gossvf_tile_count;
268 :
269 0 : fd_topo_t * topo = { fd_topob_new( &config->topo, config->name ) };
270 0 : topo->max_page_size = fd_cstr_to_shmem_page_sz( config->hugetlbfs.max_page_size );
271 0 : topo->gigantic_page_threshold = config->hugetlbfs.gigantic_page_threshold_mib << 20;
272 :
273 0 : ulong tile_to_cpu[ FD_TILE_MAX ] = {0};
274 0 : ushort parsed_tile_to_cpu[ FD_TILE_MAX ];
275 : /* Unassigned tiles will be floating, unless auto topology is enabled. */
276 0 : for( ulong i=0UL; i<FD_TILE_MAX; i++ ) parsed_tile_to_cpu[ i ] = USHORT_MAX;
277 :
278 0 : int is_auto_affinity = !strcmp( config->layout.affinity, "auto" );
279 0 : int is_bench_auto_affinity = !strcmp( config->development.bench.affinity, "auto" );
280 :
281 0 : if( FD_UNLIKELY( is_auto_affinity != is_bench_auto_affinity ) ) {
282 0 : FD_LOG_ERR(( "The CPU affinity string in the configuration file under [layout.affinity] and [development.bench.affinity] must all be set to 'auto' or all be set to a specific CPU affinity string." ));
283 0 : }
284 :
285 0 : fd_topo_cpus_t cpus[1];
286 0 : fd_topo_cpus_init( cpus );
287 :
288 0 : ulong affinity_tile_cnt = 0UL;
289 0 : if( FD_LIKELY( !is_auto_affinity ) ) affinity_tile_cnt = fd_topob_parse_affinity_cstr( config->layout.affinity, parsed_tile_to_cpu, 0, 1 );
290 :
291 0 : for( ulong i=0UL; i<affinity_tile_cnt; i++ ) {
292 0 : ushort cpu_idx = (ushort)( parsed_tile_to_cpu[ i ] & ~FD_TOPOB_CPU_SHARED );
293 0 : if( FD_UNLIKELY( parsed_tile_to_cpu[ i ]!=USHORT_MAX && cpu_idx>=cpus->cpu_cnt ) )
294 0 : FD_LOG_ERR(( "The CPU affinity string in the configuration file under [layout.affinity] specifies a CPU index of %hu, but the system "
295 0 : "only has %lu CPUs. You should either change the CPU allocations in the affinity string, or increase the number of CPUs "
296 0 : "in the system.",
297 0 : cpu_idx, cpus->cpu_cnt ));
298 0 : tile_to_cpu[ i ] = fd_ulong_if( parsed_tile_to_cpu[ i ]==USHORT_MAX, ULONG_MAX, (ulong)parsed_tile_to_cpu[ i ] );
299 0 : }
300 :
301 0 : fd_core_subtopo( config, tile_to_cpu );
302 0 : fd_gossip_subtopo( config, tile_to_cpu );
303 :
304 : /* topo, name */
305 0 : fd_topob_wksp( topo, "net_shred" );
306 0 : fd_topob_wksp( topo, "net_repair" );
307 0 : fd_topob_wksp( topo, "net_quic" );
308 :
309 0 : fd_topob_wksp( topo, "shred_out" );
310 0 : fd_topob_wksp( topo, "replay_epoch" );
311 :
312 0 : fd_topob_wksp( topo, "poh_shred" );
313 :
314 0 : fd_topob_wksp( topo, "shred_sign" );
315 0 : fd_topob_wksp( topo, "sign_shred" );
316 :
317 0 : fd_topob_wksp( topo, "repair_sign" );
318 0 : fd_topob_wksp( topo, "sign_repair" );
319 0 : fd_topob_wksp( topo, "rnonce" );
320 0 : fd_topob_wksp( topo, "repair_out" );
321 :
322 0 : fd_topob_wksp( topo, "txsend_out" );
323 :
324 0 : fd_topob_wksp( topo, "shred" );
325 0 : fd_topob_wksp( topo, "repair" );
326 0 : fd_topob_wksp( topo, "fec_sets" );
327 0 : fd_topob_wksp( topo, "snapin_manif" );
328 :
329 0 : fd_topob_wksp( topo, "genesi_out" ); /* mock genesi_out for ipecho */
330 :
331 0 : fd_topob_wksp( topo, "tower_out" ); /* mock tower_out for confirmation msgs. Not needed for any topo except eqvoc. */
332 :
333 0 : #define FOR(cnt) for( ulong i=0UL; i<cnt; i++ )
334 :
335 0 : ulong pending_fec_shreds_depth = fd_ulong_min(
336 0 : fd_ulong_pow2_up( fd_ulong_max( config->tiles.shred.max_pending_shred_sets * FD_REEDSOL_DATA_SHREDS_MAX,
337 0 : FD_SHRED_STEM_BURST ) ),
338 0 : USHORT_MAX + 1UL /* dcache max */ );
339 :
340 : /* topo, link_name, wksp_name, depth, mtu, burst */
341 0 : FOR(quic_tile_cnt) fd_topob_link( topo, "quic_net", "net_quic", config->net.ingress_buffer_size, FD_NET_MTU, 1UL );
342 0 : FOR(shred_tile_cnt) fd_topob_link( topo, "shred_net", "net_shred", config->net.ingress_buffer_size, FD_NET_MTU, 1UL );
343 :
344 0 : /**/ fd_topob_link( topo, "replay_epoch", "replay_epoch", 16UL, FD_EPOCH_OUT_MTU, 1UL );
345 :
346 0 : FOR(shred_tile_cnt) fd_topob_link( topo, "shred_sign", "shred_sign", 128UL, 32UL, 1UL );
347 0 : FOR(shred_tile_cnt) fd_topob_link( topo, "sign_shred", "sign_shred", 128UL, 64UL, 1UL );
348 :
349 0 : /**/ fd_topob_link( topo, "repair_net", "net_repair", config->net.ingress_buffer_size, FD_NET_MTU, 1UL );
350 :
351 0 : FOR(shred_tile_cnt) fd_topob_link( topo, "shred_out", "shred_out", pending_fec_shreds_depth, sizeof(fd_shred_message_t), FD_SHRED_STEM_BURST );
352 0 : FOR(sign_tile_cnt-1) fd_topob_link( topo, "repair_sign", "repair_sign", 256UL, FD_REPAIR_MAX_PREIMAGE_SZ, 1UL );
353 0 : FOR(sign_tile_cnt-1) fd_topob_link( topo, "sign_repair", "sign_repair", 128UL, sizeof(fd_ed25519_sig_t), 1UL );
354 :
355 : /**/ fd_topob_link( topo, "repair_out", "repair_out", 128UL, sizeof(fd_repair_fec_complete_t), 1UL );
356 :
357 0 : /**/ fd_topob_link( topo, "poh_shred", "poh_shred", 16384UL, FD_POH_SHRED_MTU, 1UL );
358 :
359 0 : /**/ fd_topob_link( topo, "txsend_out", "txsend_out", 128UL, FD_TXN_MTU, 1UL );
360 :
361 : /**/ fd_topob_link( topo, "snapin_manif", "snapin_manif", 2UL, sizeof(fd_snapshot_manifest_t),1UL );
362 :
363 : /**/ fd_topob_link( topo, "genesi_out", "genesi_out", 1UL, fd_genesi_tile_mtu( config->firedancer.development.genesis.max_file_size_mib<<20 ), 1UL );
364 0 : /**/ fd_topob_link( topo, "tower_out", "tower_out", 1024UL, sizeof(fd_tower_msg_t), 1UL );
365 :
366 0 : FOR(net_tile_cnt) fd_topos_net_rx_link( topo, "net_repair", i, config->net.ingress_buffer_size );
367 0 : FOR(net_tile_cnt) fd_topos_net_rx_link( topo, "net_quic", i, config->net.ingress_buffer_size );
368 0 : FOR(net_tile_cnt) fd_topos_net_rx_link( topo, "net_shred", i, config->net.ingress_buffer_size );
369 :
370 : /* topo, tile_name, tile_wksp, metrics_wksp, cpu_idx, is_agave, uses_id_keyswitch, uses_av_keyswitch */
371 0 : FOR(shred_tile_cnt) fd_topob_tile( topo, "shred", "shred", "metric_in", tile_to_cpu[ topo->tile_cnt ], 0, 1, 0, 0 );
372 0 : /**/ fd_topob_tile( topo, "repair", "repair", "metric_in", tile_to_cpu[ topo->tile_cnt ], 0, 1, 0, 0 );
373 :
374 : /* Setup a shared wksp object for fec sets. */
375 :
376 0 : ulong fec_set_cnt = fd_shred_tile_fec_set_cnt( FD_SHRED_FIREDANCER_FEC_EXPOSURE,
377 0 : config->tiles.shred.max_pending_shred_sets );
378 0 : ulong fec_sets_sz = fec_set_cnt*sizeof(fd_fec_set_t);
379 0 : fd_topo_obj_t * fec_sets_obj = setup_topo_fec_sets( topo, "fec_sets", shred_tile_cnt*fec_sets_sz );
380 0 : for( ulong i=0UL; i<shred_tile_cnt; i++ ) {
381 0 : fd_topo_tile_t * shred_tile = &topo->tiles[ fd_topo_find_tile( topo, "shred", i ) ];
382 0 : fd_topob_tile_uses( topo, shred_tile, fec_sets_obj, FD_SHMEM_JOIN_MODE_READ_WRITE );
383 0 : }
384 0 : FD_TEST( fd_pod_insertf_ulong( topo->props, fec_sets_obj->id, "fec_sets" ) );
385 :
386 : /* There's another special fseq that's used to communicate the shred
387 : version from the Agave boot path to the shred tile. */
388 0 : fd_topo_obj_t * poh_shred_obj = fd_topob_obj( topo, "fseq", "poh_shred" );
389 0 : fd_topo_tile_t * poh_tile = &topo->tiles[ fd_topo_find_tile( topo, "gossip", 0UL ) ];
390 0 : fd_topob_tile_uses( topo, poh_tile, poh_shred_obj, FD_SHMEM_JOIN_MODE_READ_WRITE );
391 :
392 0 : for( ulong i=0UL; i<shred_tile_cnt; i++ ) {
393 0 : fd_topo_tile_t * shred_tile = &topo->tiles[ fd_topo_find_tile( topo, "shred", i ) ];
394 0 : fd_topob_tile_uses( topo, shred_tile, poh_shred_obj, FD_SHMEM_JOIN_MODE_READ_ONLY );
395 0 : }
396 0 : FD_TEST( fd_pod_insertf_ulong( topo->props, poh_shred_obj->id, "poh_shred" ) );
397 :
398 0 : if( FD_LIKELY( !is_auto_affinity ) ) {
399 0 : if( FD_UNLIKELY( affinity_tile_cnt<topo->tile_cnt ) )
400 0 : FD_LOG_ERR(( "The topology you are using has %lu tiles, but the CPU affinity specified in the config tile as [layout.affinity] only provides for %lu cores. "
401 0 : "You should either increase the number of cores dedicated to Firedancer in the affinity string, or decrease the number of cores needed by reducing "
402 0 : "the total tile count. You can reduce the tile count by decreasing individual tile counts in the [layout] section of the configuration file.",
403 0 : topo->tile_cnt, affinity_tile_cnt ));
404 0 : if( FD_UNLIKELY( affinity_tile_cnt>topo->tile_cnt ) )
405 0 : FD_LOG_WARNING(( "The topology you are using has %lu tiles, but the CPU affinity specified in the config tile as [layout.affinity] provides for %lu cores. "
406 0 : "Not all cores in the affinity will be used by Firedancer. You may wish to increase the number of tiles in the system by increasing "
407 0 : "individual tile counts in the [layout] section of the configuration file.",
408 0 : topo->tile_cnt, affinity_tile_cnt ));
409 0 : }
410 :
411 : /* topo, tile_name, tile_kind_id, fseq_wksp, link_name, link_kind_id, reliable, polled */
412 0 : for( ulong j=0UL; j<shred_tile_cnt; j++ )
413 0 : fd_topos_tile_in_net( topo, "metric_in", "shred_net", j, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED ); /* No reliable consumers of networking fragments, may be dropped or overrun */
414 0 : for( ulong j=0UL; j<quic_tile_cnt; j++ )
415 0 : {fd_topos_tile_in_net( topo, "metric_in", "quic_net", j, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED );} /* No reliable consumers of networking fragments, may be dropped or overrun */
416 :
417 0 : /**/ fd_topob_tile_in( topo, "gossip", 0UL, "metric_in", "txsend_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
418 :
419 0 : /**/ fd_topos_tile_in_net( topo, "metric_in", "repair_net", 0UL, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED ); /* No reliable consumers of networking fragments, may be dropped or overrun */
420 :
421 0 : FOR(shred_tile_cnt) for( ulong j=0UL; j<net_tile_cnt; j++ )
422 0 : fd_topob_tile_in( topo, "shred", i, "metric_in", "net_shred", j, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED ); /* No reliable consumers of networking fragments, may be dropped or overrun */
423 0 : FOR(shred_tile_cnt) fd_topob_tile_in( topo, "shred", i, "metric_in", "poh_shred", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
424 0 : FOR(shred_tile_cnt) fd_topob_tile_in( topo, "shred", i, "metric_in", "replay_epoch", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
425 0 : FOR(shred_tile_cnt) fd_topob_tile_in( topo, "shred", i, "metric_in", "gossip_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
426 0 : FOR(shred_tile_cnt) fd_topob_tile_out( topo, "shred", i, "shred_out", i );
427 0 : FOR(shred_tile_cnt) fd_topob_tile_out( topo, "shred", i, "shred_net", i );
428 0 : FOR(shred_tile_cnt) fd_topob_tile_in ( topo, "shred", i, "metric_in", "ipecho_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
429 :
430 : /**/ fd_topob_tile_out( topo, "repair", 0UL, "repair_net", 0UL );
431 :
432 : /* Sign links don't need to be reliable because they are synchronous,
433 : so there's at most one fragment in flight at a time anyway. The
434 : sign links are also not polled by the mux, instead the tiles will
435 : read the sign responses out of band in a dedicated spin loop. */
436 0 : for( ulong i=0UL; i<shred_tile_cnt; i++ ) {
437 0 : /**/ fd_topob_tile_in( topo, "sign", 0UL, "metric_in", "shred_sign", i, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED );
438 0 : /**/ fd_topob_tile_out( topo, "shred", i, "shred_sign", i );
439 0 : /**/ fd_topob_tile_in( topo, "shred", i, "metric_in", "sign_shred", i, FD_TOPOB_UNRELIABLE, FD_TOPOB_UNPOLLED );
440 0 : /**/ fd_topob_tile_out( topo, "sign", 0UL, "sign_shred", i );
441 0 : }
442 0 : FOR(gossvf_tile_cnt) fd_topob_tile_in ( topo, "gossvf", i, "metric_in", "replay_epoch", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
443 :
444 0 : /**/ fd_topob_tile_in ( topo, "gossip", 0UL, "metric_in", "replay_epoch", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
445 :
446 0 : FOR(net_tile_cnt) fd_topob_tile_in( topo, "repair", 0UL, "metric_in", "net_repair", i, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED ); /* No reliable consumers of networking fragments, may be dropped or overrun */
447 0 : /**/ fd_topob_tile_in( topo, "repair", 0UL, "metric_in", "gossip_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
448 0 : fd_topob_tile_in( topo, "repair", 0UL, "metric_in", "snapin_manif", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
449 0 : FOR(shred_tile_cnt) fd_topob_tile_in( topo, "repair", 0UL, "metric_in", "shred_out", i, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
450 0 : FOR(sign_tile_cnt-1) fd_topob_tile_out( topo, "repair", 0UL, "repair_sign", i );
451 0 : FOR(sign_tile_cnt-1) fd_topob_tile_in ( topo, "sign", i+1, "metric_in", "repair_sign", i, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
452 0 : FOR(sign_tile_cnt-1) fd_topob_tile_out( topo, "sign", i+1, "sign_repair", i );
453 0 : FOR(sign_tile_cnt-1) fd_topob_tile_in ( topo, "repair", 0UL, "metric_in", "sign_repair", i, FD_TOPOB_UNRELIABLE, FD_TOPOB_POLLED );
454 :
455 : /**/ fd_topob_tile_out( topo, "repair", 0UL, "repair_out", 0UL );
456 0 : /**/ fd_topob_tile_in ( topo, "gossip", 0UL, "metric_in", "sign_gossip", 0UL, FD_TOPOB_UNRELIABLE, FD_TOPOB_UNPOLLED );
457 0 : /**/ fd_topob_tile_in ( topo, "ipecho", 0UL, "metric_in", "genesi_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
458 0 : /**/ fd_topob_tile_in ( topo, "repair", 0UL, "metric_in", "tower_out", 0UL, FD_TOPOB_RELIABLE, FD_TOPOB_POLLED );
459 :
460 : /* Repair and shred share a secret they use to generate the nonces.
461 : It's not super security sensitive, but for good hygiene, we make it
462 : an object. */
463 0 : if( 1 /* just restrict the scope for these variables in this big function */ ) {
464 0 : fd_topo_obj_t * rnonce_ss_obj = fd_topob_obj( topo, "rnonce_ss", "rnonce" );
465 0 : fd_topo_tile_t * repair_tile = &topo->tiles[ fd_topo_find_tile( topo, "repair", 0UL ) ];
466 0 : fd_topob_tile_uses( topo, repair_tile, rnonce_ss_obj, FD_SHMEM_JOIN_MODE_READ_ONLY );
467 0 : for( ulong i=0UL; i<shred_tile_cnt; i++ ) {
468 0 : fd_topo_tile_t * shred_tile = &topo->tiles[ fd_topo_find_tile( topo, "shred", i ) ];
469 0 : fd_topob_tile_uses( topo, shred_tile, rnonce_ss_obj, FD_SHMEM_JOIN_MODE_READ_ONLY );
470 0 : }
471 0 : FD_TEST( fd_pod_insertf_ulong( topo->props, rnonce_ss_obj->id, "rnonce_ss" ) );
472 0 : }
473 :
474 0 : FD_TEST( fd_link_permit_no_producers( topo, "quic_net" ) == quic_tile_cnt );
475 0 : FD_TEST( fd_link_permit_no_producers( topo, "poh_shred" ) == 1UL );
476 0 : FD_TEST( fd_link_permit_no_producers( topo, "txsend_out" ) == 1UL );
477 0 : FD_TEST( fd_link_permit_no_producers( topo, "genesi_out" ) == 1UL );
478 0 : FD_TEST( fd_link_permit_no_producers( topo, "tower_out" ) == 1UL );
479 0 : FD_TEST( fd_link_permit_no_producers( topo, "replay_epoch" ) == 1UL );
480 0 : FD_TEST( fd_link_permit_no_producers( topo, "snapin_manif" ) == 1UL );
481 0 : FD_TEST( fd_link_permit_no_consumers( topo, "net_quic" ) == net_tile_cnt );
482 0 : FD_TEST( fd_link_permit_no_consumers( topo, "repair_out" ) == 1UL );
483 :
484 0 : config->tiles.txsend.txsend_src_port = 0; /* disable txsend */
485 :
486 0 : FOR(net_tile_cnt) fd_topos_net_tile_finish( topo, i );
487 :
488 0 : for( ulong i=0UL; i<topo->tile_cnt; i++ ) {
489 0 : fd_topo_tile_t * tile = &topo->tiles[ i ];
490 0 : fd_topo_configure_tile( tile, config );
491 0 : }
492 :
493 0 : if( FD_UNLIKELY( is_auto_affinity ) ) fd_topob_auto_layout( topo, 0 );
494 :
495 0 : fd_topob_waker( topo );
496 0 : fd_topob_finish( topo, CALLBACKS );
497 :
498 0 : config->topo = *topo;
499 0 : }
500 :
501 : static char *
502 0 : fmt_count( char buf[ static 64 ], ulong count ) {
503 0 : char tmp[ 64 ];
504 0 : if( FD_LIKELY( count<1000UL ) ) FD_TEST( fd_cstr_printf_check( tmp, 64UL, NULL, "%lu", count ) );
505 0 : else if( FD_LIKELY( count<1000000UL ) ) FD_TEST( fd_cstr_printf_check( tmp, 64UL, NULL, "%.1f K", (double)count/1000.0 ) );
506 0 : else if( FD_LIKELY( count<1000000000UL ) ) FD_TEST( fd_cstr_printf_check( tmp, 64UL, NULL, "%.1f M", (double)count/1000000.0 ) );
507 :
508 0 : FD_TEST( fd_cstr_printf_check( buf, 64UL, NULL, "%12s", tmp ) );
509 0 : return buf;
510 0 : }
511 :
512 : static void
513 : print_histogram_buckets( volatile ulong * metrics,
514 : ulong offset,
515 : int converter,
516 : double histmin,
517 : double histmax,
518 0 : char * title ) {
519 0 : fd_histf_t hist[1];
520 :
521 : /* Create histogram structure only to get bucket edges for display */
522 0 : if( FD_LIKELY( converter == FD_METRICS_CONVERTER_SECONDS ) ) {
523 : /* For SLOT_COMPLETE_TIME: min=0.2, max=2.0 seconds */
524 0 : FD_TEST( fd_histf_new( hist, fd_metrics_convert_seconds_to_ticks( histmin ), fd_metrics_convert_seconds_to_ticks( histmax ) ) );
525 0 : } else if( FD_LIKELY( converter == FD_METRICS_CONVERTER_NONE ) ) {
526 : /* For non-time histograms, we'd need the actual min/max values */
527 0 : FD_TEST( fd_histf_new( hist, (ulong)histmin, (ulong)histmax ) );
528 0 : } else {
529 0 : FD_LOG_ERR(( "unknown converter %i", converter ));
530 0 : }
531 :
532 0 : printf( " +---------------------+--------------------+--------------+\n" );
533 0 : printf( " | %-19s | | Count |\n", title );
534 0 : printf( " +---------------------+--------------------+--------------+\n" );
535 :
536 0 : ulong total_count = 0;
537 0 : for( ulong k = 0; k < FD_HISTF_BUCKET_CNT; k++ ) {
538 0 : ulong bucket_count = metrics[ offset + k ];
539 0 : total_count += bucket_count;
540 0 : }
541 :
542 0 : for( ulong k = 0; k < FD_HISTF_BUCKET_CNT; k++ ) {
543 : /* Get individual bucket count directly from metrics array */
544 0 : ulong bucket_count = metrics[ offset + k ];
545 :
546 0 : char * le_str;
547 0 : char le_buf[ 64 ];
548 0 : if( FD_UNLIKELY( k == FD_HISTF_BUCKET_CNT - 1UL ) ) {
549 0 : le_str = "+Inf";
550 0 : } else {
551 0 : ulong edge = fd_histf_right( hist, k );
552 0 : if( FD_LIKELY( converter == FD_METRICS_CONVERTER_SECONDS ) ) {
553 0 : double edgef = fd_metrics_convert_ticks_to_seconds( edge - 1 );
554 0 : FD_TEST( fd_cstr_printf_check( le_buf, sizeof( le_buf ), NULL, "%.3f", edgef ) );
555 0 : } else {
556 0 : FD_TEST( fd_cstr_printf_check( le_buf, sizeof( le_buf ), NULL, "%.3f", (double)(edge - 1) / 1000000.0 ) );
557 0 : }
558 0 : le_str = le_buf;
559 0 : }
560 :
561 0 : char count_buf[ 64 ];
562 0 : fmt_count( count_buf, bucket_count );
563 :
564 : /* Match visual bar length to the %-18s display column width. */
565 0 : char bar_buf[ 19 ];
566 0 : ulong bar_max = sizeof( bar_buf ) - 1UL;
567 0 : if( bucket_count > 0 && total_count > 0 ) {
568 0 : ulong bar_length = (bucket_count * bar_max) / total_count;
569 0 : if( bar_length == 0 ) bar_length = 1;
570 0 : if( bar_length > bar_max ) bar_length = bar_max;
571 0 : for( ulong i = 0; i < bar_length; i++ ) { bar_buf[ i ] = '|'; }
572 0 : bar_buf[ bar_length ] = '\0';
573 0 : } else {
574 0 : bar_buf[ 0 ] = '\0';
575 0 : }
576 :
577 0 : printf( " | %-19s | %-18s | %s |\n", le_str, bar_buf, count_buf );
578 0 : }
579 :
580 : /* Print sum and total count */
581 0 : char sum_buf[ 64 ];
582 0 : char avg_buf[ 64 ];
583 0 : if( FD_LIKELY( converter == FD_METRICS_CONVERTER_SECONDS ) ) {
584 0 : double sumf = fd_metrics_convert_ticks_to_seconds( metrics[ offset + FD_HISTF_BUCKET_CNT ] );
585 0 : FD_TEST( fd_cstr_printf_check( sum_buf, sizeof( sum_buf ), NULL, "%.6f", sumf ) );
586 0 : double avg = sumf / (double)total_count;
587 0 : FD_TEST( fd_cstr_printf_check( avg_buf, sizeof( avg_buf ), NULL, "%.6f", avg ) );
588 0 : } else {
589 0 : FD_TEST( fd_cstr_printf_check( sum_buf, sizeof( sum_buf ), NULL, "%lu", metrics[ offset + FD_HISTF_BUCKET_CNT ] ));
590 0 : }
591 :
592 0 : printf( " +---------------------+--------------------+---------------+\n" );
593 0 : printf( " | Sum: %-14s | Count: %-11lu | Avg: %-8s |\n", sum_buf, total_count, avg_buf );
594 0 : printf( " +---------------------+--------------------+---------------+\n" );
595 0 : }
596 :
597 : static fd_slot_metrics_t temp_slots[ FD_CATCHUP_METRICS_MAX ];
598 :
599 : static void
600 0 : print_catchup_slots( fd_wksp_t * repair_tile_wksp, ctx_t * repair_ctx, int verbose, int sort_by_slot ) {
601 0 : fd_repair_metrics_t * catchup = repair_ctx->slot_metrics;
602 0 : ulong catchup_gaddr = fd_wksp_gaddr_fast( repair_ctx->wksp, catchup );
603 0 : fd_repair_metrics_t * catchup_table = (fd_repair_metrics_t *)fd_wksp_laddr( repair_tile_wksp, catchup_gaddr );
604 0 : if( FD_LIKELY( sort_by_slot ) ) {
605 0 : fd_repair_metrics_print_sorted( catchup_table, verbose, temp_slots );
606 0 : } else {
607 0 : fd_repair_metrics_print( catchup_table, verbose );
608 0 : }
609 0 : }
610 :
611 : static fd_location_info_t * location_table;
612 : static fd_pubkey_t peers_copy[ FD_REPAIR_PEER_MAX];
613 :
614 : static ulong
615 0 : sort_peers_by_latency( fd_policy_peer_map_t * active_table, fd_policy_peer_dlist_t * peers_dlist, fd_policy_peer_dlist_t * peers_wlist, fd_policy_peer_t * peers_arr ) {
616 0 : ulong i = 0;
617 0 : fd_policy_peer_dlist_iter_t iter = fd_policy_peer_dlist_iter_fwd_init( peers_dlist, peers_arr );
618 0 : while( !fd_policy_peer_dlist_iter_done( iter, peers_dlist, peers_arr ) ) {
619 0 : fd_policy_peer_t * peer = fd_policy_peer_dlist_iter_ele( iter, peers_dlist, peers_arr );
620 0 : if( FD_UNLIKELY( !peer ) ) break;
621 0 : peers_copy[ i++ ] = peer->key;
622 0 : if( FD_UNLIKELY( i >= FD_REPAIR_PEER_MAX ) ) break;
623 0 : iter = fd_policy_peer_dlist_iter_fwd_next( iter, peers_dlist, peers_arr );
624 0 : }
625 0 : ulong fast_cnt = i;
626 0 : iter = fd_policy_peer_dlist_iter_fwd_init( peers_wlist, peers_arr );
627 0 : while( !fd_policy_peer_dlist_iter_done( iter, peers_wlist, peers_arr ) ) {
628 0 : fd_policy_peer_t * peer = fd_policy_peer_dlist_iter_ele( iter, peers_wlist, peers_arr );
629 0 : if( FD_UNLIKELY( !peer ) ) break;
630 0 : peers_copy[ i++ ] = peer->key;
631 0 : if( FD_UNLIKELY( i >= FD_REPAIR_PEER_MAX ) ) break;
632 0 : iter = fd_policy_peer_dlist_iter_fwd_next( iter, peers_wlist, peers_arr );
633 0 : }
634 0 : FD_LOG_NOTICE(( "Fast peers cnt: %lu. Slow peers cnt: %lu.", fast_cnt, i - fast_cnt ));
635 :
636 0 : ulong peer_cnt = i;
637 0 : for( uint i = 0; i < peer_cnt - 1; i++ ) {
638 0 : int swapped = 0;
639 0 : for( uint j = 0; j < peer_cnt - 1 - i; j++ ) {
640 0 : fd_policy_peer_t const * active_j = fd_policy_peer_map_ele_query( active_table, &peers_copy[ j ], NULL, peers_arr );
641 0 : fd_policy_peer_t const * active_j1 = fd_policy_peer_map_ele_query( active_table, &peers_copy[ j + 1 ], NULL, peers_arr );
642 :
643 : /* Skip peers with no responses */
644 0 : double latency_j = 10e9;
645 0 : double latency_j1 = 10e9;
646 0 : if( FD_LIKELY( active_j && active_j->res_cnt > 0 ) ) latency_j = ((double)active_j->total_lat / (double)active_j->res_cnt);
647 0 : if( FD_LIKELY( active_j1 && active_j1->res_cnt > 0 ) ) latency_j1 = ((double)active_j1->total_lat / (double)active_j1->res_cnt);
648 :
649 : /* Swap if j has higher latency than j+1 */
650 0 : if( latency_j > latency_j1 ) {
651 0 : fd_pubkey_t temp = peers_copy[ j ];
652 0 : peers_copy[ j ] = peers_copy[ j + 1 ];
653 0 : peers_copy[ j + 1 ] = temp;
654 0 : swapped = 1;
655 0 : }
656 0 : }
657 0 : if( !swapped ) break;
658 0 : }
659 0 : return peer_cnt;
660 0 : }
661 :
662 : static void
663 0 : print_peer_location_latency( fd_wksp_t * repair_tile_wksp, ctx_t * tile_ctx ) {
664 0 : ulong policy_gaddr = fd_wksp_gaddr_fast( tile_ctx->wksp, tile_ctx->policy );
665 0 : fd_policy_t * policy = fd_wksp_laddr ( repair_tile_wksp, policy_gaddr );
666 0 : ulong peermap_gaddr = fd_wksp_gaddr_fast( tile_ctx->wksp, policy->peers.map );
667 0 : ulong peerarr_gaddr = fd_wksp_gaddr_fast( tile_ctx->wksp, policy->peers.pool );
668 0 : ulong peerlst_gaddr = fd_wksp_gaddr_fast( tile_ctx->wksp, policy->peers.fast );
669 0 : ulong peerwst_gaddr = fd_wksp_gaddr_fast( tile_ctx->wksp, policy->peers.slow );
670 0 : fd_policy_peer_map_t * peers_map = (fd_policy_peer_map_t *) fd_wksp_laddr( repair_tile_wksp, peermap_gaddr );
671 0 : fd_policy_peer_dlist_t * peers_dlist = (fd_policy_peer_dlist_t *)fd_wksp_laddr( repair_tile_wksp, peerlst_gaddr );
672 0 : fd_policy_peer_dlist_t * peers_wlist = (fd_policy_peer_dlist_t *)fd_wksp_laddr( repair_tile_wksp, peerwst_gaddr );
673 0 : fd_policy_peer_t * peers_arr = (fd_policy_peer_t *) fd_wksp_laddr( repair_tile_wksp, peerarr_gaddr );
674 :
675 0 : ulong peer_cnt = sort_peers_by_latency( peers_map, peers_dlist, peers_wlist, peers_arr );
676 0 : printf("\nPeer Location/Latency Information\n");
677 0 : printf( " | %-46s | %-7s | %-8s | %-8s | %-7s | %-7s | %-12s | %s\n", "Pubkey", "Req Cnt", "Req B/s", "Rx B/s", "Rx Rate", "Avg Latency", "Ewma Latency", "Location Info" );
678 0 : for( uint i = 0; i < peer_cnt; i++ ) {
679 0 : fd_policy_peer_t const * active = fd_policy_peer_map_ele_query( peers_map, &peers_copy[ i ], NULL, peers_arr );
680 0 : if( FD_LIKELY( active && active->res_cnt > 0 ) ) {
681 0 : fd_location_info_t * info = fd_location_table_query( location_table, active->ip4, NULL );
682 0 : char * geolocation = info ? info->location : "";
683 0 : double peer_bps = (double)(active->res_cnt * FD_SHRED_MIN_SZ) / ((double)(active->last_resp_ts - active->first_resp_ts) / 1e9);
684 0 : double req_bps = (double)active->req_cnt * 202 / ((double)(active->last_req_ts - active->first_req_ts) / 1e9);
685 0 : FD_BASE58_ENCODE_32_BYTES( active->key.key, key_b58 );
686 0 : printf( "%-5u | %-46s | %-7lu | %-8.2f | %-8.2f | %-7.2f | %10.3fms | %10.3fms | %s\n", i, key_b58, active->req_cnt, req_bps, peer_bps, (double)active->res_cnt / (double)active->req_cnt, ((double)active->total_lat / (double)active->res_cnt) / 1e6, (double)active->ewma_lat / 1e6, geolocation );
687 0 : }
688 0 : }
689 0 : printf("\n");
690 0 : fflush( stdout );
691 0 : }
692 :
693 : static void
694 0 : read_iptable( char * iptable_path, fd_location_info_t * location_table ) {
695 0 : int iptable_fd = open( iptable_path, O_RDONLY );
696 0 : if( FD_UNLIKELY( iptable_fd<0 ) ) return;
697 :
698 : /* read iptable line by line */
699 0 : if( FD_LIKELY( iptable_fd>=0 ) ) {
700 0 : char line[ 256 ];
701 0 : uchar istream_buf[256];
702 0 : fd_io_buffered_istream_t istream[1];
703 0 : fd_io_buffered_istream_init( istream, iptable_fd, istream_buf, sizeof(istream_buf) );
704 0 : for(;;) {
705 0 : int err;
706 0 : if( !fd_io_fgets( line, sizeof(line), istream, &err ) ) break;
707 0 : fd_location_info_t location_info;
708 0 : sscanf( line, "%lu %[^\n]", &location_info.ip4_addr, location_info.location );
709 0 : fd_location_info_t * info = fd_location_table_insert( location_table, location_info.ip4_addr );
710 0 : if( FD_UNLIKELY( info==NULL ) ) break;
711 0 : memcpy( info->location, location_info.location, sizeof(info->location) );
712 0 : }
713 0 : }
714 0 : }
715 :
716 : static void
717 : print_tile_metrics( volatile ulong * shred_metrics,
718 : volatile ulong * repair_metrics,
719 : volatile ulong * repair_metrics_prev, /* for diffing metrics */
720 : volatile ulong ** repair_net_links,
721 : volatile ulong ** net_shred_links,
722 : ulong net_tile_cnt,
723 : ulong * last_sent_cnt,
724 : long last_print_ts,
725 0 : long now ) {
726 0 : char buf2[ 64 ];
727 0 : ulong rcvd = shred_metrics [ MIDX( COUNTER, SHRED, SHRED_REPAIR_RX ) ];
728 0 : ulong sent = repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_WINDOW ) ] +
729 0 : repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_HIGHEST_WINDOW ) ] +
730 0 : repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_ORPHAN ) ];
731 0 : printf(" Requests received: (%lu/%lu) %.1f%% \n", rcvd, sent, (double)rcvd / (double)sent * 100.0 );
732 0 : printf( " +---------------+--------------+\n" );
733 0 : printf( " | Request Type | Count |\n" );
734 0 : printf( " +---------------+--------------+\n" );
735 0 : printf( " | Orphan | %s |\n", fmt_count( buf2, repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_ORPHAN ) ] ) );
736 0 : printf( " | HighestWindow | %s |\n", fmt_count( buf2, repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_HIGHEST_WINDOW ) ] ) );
737 0 : printf( " | Index | %s |\n", fmt_count( buf2, repair_metrics[ MIDX( COUNTER, REPAIR, REQUEST_TX_NEEDED_WINDOW ) ] ) );
738 0 : printf( " +---------------+--------------+\n" );
739 0 : printf( " Send Pkt Rate: %s pps\n", fmt_count( buf2, (ulong)((sent - *last_sent_cnt)*1e9L / (now - last_print_ts) ) ) );
740 0 : *last_sent_cnt = sent;
741 :
742 : /* Sum overrun across all net tiles connected to repair_net */
743 0 : ulong total_overrun = repair_net_links[0][ MIDX( COUNTER, LINK, FRAG_POLLING_OVERRUN ) ]; /* coarse double counting prevention */
744 0 : ulong total_consumed = 0UL;
745 0 : for( ulong i = 0UL; i < net_tile_cnt; i++ ) {
746 0 : volatile ulong * ovar_net_metrics = repair_net_links[i];
747 0 : total_overrun += ovar_net_metrics[ MIDX( COUNTER, LINK, FRAG_READING_OVERRUN ) ];
748 0 : total_consumed += ovar_net_metrics[ MIDX( COUNTER, LINK, FRAG_CONSUMED ) ]; /* consumed is incremented after after_frag is called */
749 0 : }
750 0 : printf( " Outgoing requests overrun: %s\n", fmt_count( buf2, total_overrun ) );
751 0 : printf( " Outgoing requests consumed: %s\n", fmt_count( buf2, total_consumed ) );
752 :
753 0 : total_overrun = net_shred_links[0][ MIDX( COUNTER, LINK, FRAG_READING_OVERRUN ) ];
754 0 : total_consumed = 0UL;
755 0 : for( ulong i = 0UL; i < net_tile_cnt; i++ ) {
756 0 : volatile ulong * ovar_net_metrics = net_shred_links[i];
757 0 : total_overrun += ovar_net_metrics[ MIDX( COUNTER, LINK, FRAG_READING_OVERRUN ) ];
758 0 : total_consumed += ovar_net_metrics[ MIDX( COUNTER, LINK, FRAG_CONSUMED ) ]; /* shred frag filtering happens manually in after_frag, so no need to index every shred_tile. */
759 0 : }
760 :
761 0 : printf( " Incoming shreds overrun: %s\n", fmt_count( buf2, total_overrun ) );
762 0 : printf( " Incoming shreds consumed: %s\n", fmt_count( buf2, total_consumed ) );
763 :
764 0 : print_histogram_buckets( repair_metrics,
765 0 : MIDX( HISTOGRAM, REPAIR, RESPONSE_LATENCY_NANOS ),
766 0 : FD_METRICS_CONVERTER_NONE,
767 0 : FD_METRICS_HISTOGRAM_REPAIR_RESPONSE_LATENCY_NANOS_MIN,
768 0 : FD_METRICS_HISTOGRAM_REPAIR_RESPONSE_LATENCY_NANOS_MAX,
769 0 : "Response Latency" );
770 :
771 0 : printf(" Repair Peers: %lu\n", repair_metrics[ MIDX( COUNTER, REPAIR, PEER_REQUESTED ) ] );
772 0 : printf(" Shreds rejected (no stakes): %lu\n", shred_metrics[ MIDX( COUNTER, SHRED, SHRED_PROCESSED ) ] );
773 : /* Print histogram buckets similar to Prometheus format */
774 0 : print_histogram_buckets( repair_metrics,
775 0 : MIDX( HISTOGRAM, REPAIR, SLOT_COMPLETE_DURATION_SECONDS ),
776 0 : FD_METRICS_CONVERTER_SECONDS,
777 0 : FD_METRICS_HISTOGRAM_REPAIR_SLOT_COMPLETE_DURATION_SECONDS_MIN,
778 0 : FD_METRICS_HISTOGRAM_REPAIR_SLOT_COMPLETE_DURATION_SECONDS_MAX,
779 0 : "Slot Complete Time" );
780 :
781 0 : #define DIFFX(METRIC) repair_metrics[ MIDX( COUNTER, TILE, METRIC ) ] - repair_metrics_prev[ MIDX( COUNTER, TILE, METRIC ) ]
782 0 : ulong hkeep_ticks = DIFFX(REGIME_DURATION_NANOS_CAUGHT_UP_HOUSEKEEPING) + DIFFX(REGIME_DURATION_NANOS_PROCESSING_HOUSEKEEPING) + DIFFX(REGIME_DURATION_NANOS_BACKPRESSURE_HOUSEKEEPING);
783 0 : ulong busy_ticks = DIFFX(REGIME_DURATION_NANOS_PROCESSING_PREFRAG) + DIFFX(REGIME_DURATION_NANOS_PROCESSING_POSTFRAG ) + DIFFX(REGIME_DURATION_NANOS_CAUGHT_UP_PREFRAG);
784 0 : ulong caught_up_ticks = DIFFX(REGIME_DURATION_NANOS_CAUGHT_UP_POSTFRAG);
785 0 : ulong backpressure_ticks = DIFFX(REGIME_DURATION_NANOS_BACKPRESSURE_PREFRAG);
786 0 : ulong total_ticks = hkeep_ticks + busy_ticks + caught_up_ticks + backpressure_ticks;
787 :
788 0 : printf( " Repair Hkeep: %.1f %% Busy: %.1f %% Idle: %.1f %% Backp: %0.1f %%\n",
789 0 : (double)hkeep_ticks/(double)total_ticks*100.0,
790 0 : (double)busy_ticks/(double)total_ticks*100.0,
791 0 : (double)caught_up_ticks/(double)total_ticks*100.0,
792 0 : (double)backpressure_ticks/(double)total_ticks*100.0 );
793 0 : #undef DIFFX
794 0 : fflush( stdout );
795 :
796 0 : printf( " Block failed insert: %lu\n", repair_metrics[ MIDX( COUNTER, REPAIR, BLOCK_INSERT_FAILED ) ] );
797 0 : printf( " Block evicted: %lu\n", repair_metrics[ MIDX( COUNTER, REPAIR, BLOCK_EVICTED ) ] );
798 0 : printf( " slot evicted: %lu\n", repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_LAST_EVICTED ) ] );
799 0 : printf( " slot evicted by: %lu\n", repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_LAST_EVICTION_CAUSE ) ] );
800 0 : printf( " slot failed insert: %lu\n", repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_LAST_INSERT_FAILED ) ] );
801 0 : for( ulong i=0UL; i<FD_METRICS_TOTAL_SZ/sizeof(ulong); i++ ) repair_metrics_prev[ i ] = repair_metrics[ i ];
802 0 : }
803 :
804 : static void
805 : repair_ctx_wksp( args_t * args,
806 : config_t * config,
807 : ctx_t ** repair_ctx,
808 0 : fd_topo_wksp_t ** repair_wksp ) {
809 0 : (void)args;
810 :
811 0 : fd_topo_t * topo = &config->topo;
812 0 : ulong wksp_id = fd_topo_find_wksp( topo, "repair" );
813 0 : if( FD_UNLIKELY( wksp_id==ULONG_MAX ) ) FD_LOG_ERR(( "repair workspace not found" ));
814 :
815 0 : fd_topo_wksp_t * _repair_wksp = &topo->workspaces[ wksp_id ];
816 :
817 0 : ulong tile_id = fd_topo_find_tile( topo, "repair", 0UL );
818 0 : if( FD_UNLIKELY( tile_id==ULONG_MAX ) ) FD_LOG_ERR(( "repair tile not found" ));
819 :
820 0 : fd_topo_join_workspace( topo, _repair_wksp, FD_SHMEM_JOIN_MODE_READ_ONLY, FD_TOPO_CORE_DUMP_LEVEL_DISABLED );
821 :
822 : /* Access the repair tile scratch memory where repair_tile_ctx is stored */
823 0 : fd_topo_tile_t * tile = &topo->tiles[ tile_id ];
824 0 : void * scratch = fd_topo_obj_laddr( &config->topo, tile->tile_obj_id );
825 0 : if( FD_UNLIKELY( !scratch ) ) FD_LOG_ERR(( "Failed to access repair tile scratch memory" ));
826 :
827 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
828 0 : ctx_t * _repair_ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
829 :
830 0 : *repair_ctx = _repair_ctx;
831 0 : *repair_wksp = _repair_wksp;
832 0 : }
833 :
834 : static void
835 : repair_cmd_fn_catchup( args_t * args,
836 0 : config_t * config ) {
837 :
838 0 : memset( &config->topo, 0, sizeof(config->topo) );
839 0 : repair_topo( config );
840 :
841 0 : fd_topo_print_log( 1, &config->topo );
842 :
843 0 : args_t configure_args = {
844 0 : .configure.command = CONFIGURE_CMD_INIT,
845 0 : };
846 0 : for( ulong i=0UL; STAGES[ i ]; i++ ) {
847 0 : configure_args.configure.stages[ i ] = STAGES[ i ];
848 0 : }
849 0 : configure_cmd_fn( &configure_args, config );
850 0 : if( 0==strcmp( config->net.provider, "xdp" ) ) {
851 0 : fd_topo_install_xdp_simple( &config->topo, config->net.bind_address_parsed );
852 0 : }
853 0 : run_firedancer_init( config, 1, 0 );
854 :
855 0 : fd_topo_join_workspaces( &config->topo, FD_SHMEM_JOIN_MODE_READ_WRITE, FD_TOPO_CORE_DUMP_LEVEL_DISABLED );
856 0 : if( 0==strcmp( config->net.provider, "mlx5" ) ) {
857 0 : fd_topo_install_mlx5( &config->topo, NULL );
858 0 : }
859 :
860 0 : fd_topo_fill( &config->topo );
861 :
862 0 : repair_load_manifest( &config->topo, args->repair.manifest_path );
863 :
864 : /* Access repair workspace memory and metrics */
865 :
866 0 : ulong repair_tile_idx = fd_topo_find_tile( &config->topo, "repair", 0UL );
867 0 : ulong shred_tile_idx = fd_topo_find_tile( &config->topo, "shred", 0UL );
868 0 : FD_TEST( repair_tile_idx!=ULONG_MAX );
869 0 : FD_TEST( shred_tile_idx !=ULONG_MAX );
870 0 : fd_topo_tile_t * repair_tile = &config->topo.tiles[ repair_tile_idx ];
871 0 : fd_topo_tile_t * shred_tile = &config->topo.tiles[ shred_tile_idx ];
872 :
873 0 : fd_topo_wksp_t * repair_wksp;
874 0 : ctx_t * repair_ctx;
875 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
876 :
877 0 : volatile ulong * shred_metrics = fd_metrics_tile( shred_tile->metrics );
878 0 : volatile ulong * repair_metrics = fd_metrics_tile( repair_tile->metrics );
879 0 : FD_TEST( repair_metrics );
880 0 : ulong * repair_metrics_prev = aligned_alloc( 8UL, sizeof(ulong) * FD_METRICS_TOTAL_SZ );
881 0 : FD_TEST( repair_metrics_prev );
882 0 : memset( repair_metrics_prev, 0, sizeof(ulong) * FD_METRICS_TOTAL_SZ );
883 :
884 : /* Collect link metrics */
885 :
886 : /* Collect all net tiles and their repair_net link metrics */
887 0 : ulong net_cnt = config->layout.net_tile_count;
888 0 : volatile ulong ** repair_net_links = aligned_alloc( 8UL, net_cnt * sizeof(volatile ulong*) );
889 0 : volatile ulong ** net_shred_links = aligned_alloc( 8UL, net_cnt * sizeof(volatile ulong*) );
890 0 : FD_TEST( repair_net_links );
891 0 : FD_TEST( net_shred_links );
892 :
893 0 : char const * net_tile_name = fd_net_tile_name( config->net.provider );
894 0 : for( ulong i = 0UL; i < net_cnt; i++ ) {
895 0 : ulong tile_idx = fd_topo_find_tile( &config->topo, net_tile_name, i );
896 0 : if( FD_UNLIKELY( tile_idx == ULONG_MAX ) ) FD_LOG_ERR(( "net tile %lu not found", i ));
897 0 : fd_topo_tile_t * tile = &config->topo.tiles[ tile_idx ];
898 :
899 0 : ulong repair_net_in_idx = fd_topo_find_tile_in_link( &config->topo, tile, "repair_net", 0UL );
900 0 : if( FD_UNLIKELY( repair_net_in_idx == ULONG_MAX ) ) FD_LOG_ERR(( "repair_net link not found for net tile %lu", i ));
901 0 : FD_TEST( tile->metrics );
902 0 : repair_net_links[i] = fd_metrics_link_in( tile->metrics, repair_net_in_idx );
903 0 : FD_TEST( repair_net_links[i] );
904 :
905 : /* process all net_shred links */
906 0 : ulong shred_tile_idx = fd_topo_find_tile( &config->topo, "shred", 0 );
907 0 : if( FD_UNLIKELY( shred_tile_idx == ULONG_MAX ) ) FD_LOG_ERR(( "shred tile 0 not found" ));
908 0 : fd_topo_tile_t * shred_tile = &config->topo.tiles[ shred_tile_idx ];
909 :
910 0 : ulong shred_out_in_idx = fd_topo_find_tile_in_link( &config->topo, shred_tile, "net_shred", i );
911 0 : if( FD_UNLIKELY( shred_out_in_idx == ULONG_MAX ) ) FD_LOG_ERR(( "net_shred link not found for shred tile 0" ));
912 0 : FD_TEST( shred_tile->metrics );
913 0 : net_shred_links[i] = fd_metrics_link_in( shred_tile->metrics, shred_out_in_idx );
914 0 : FD_TEST( net_shred_links[i] );
915 0 : }
916 :
917 0 : FD_LOG_NOTICE(( "Repair catchup run" ));
918 :
919 0 : ulong shred_out_link_idx = fd_topo_find_link( &config->topo, "shred_out", 0UL );
920 0 : FD_TEST( shred_out_link_idx!=ULONG_MAX );
921 0 : fd_topo_link_t * shred_out_link = &config->topo.links[ shred_out_link_idx ];
922 0 : fd_frag_meta_t * shred_out_mcache = shred_out_link->mcache;
923 0 : void * shred_out_dcache = config->topo.workspaces[ config->topo.objs[ shred_out_link->dcache_obj_id ].wksp_id ].wksp;
924 :
925 0 : ulong turbine_slot0 = 0;
926 0 : long last_print = fd_log_wallclock();
927 0 : ulong last_sent = 0UL;
928 :
929 0 : fd_topo_run_single_process( &config->topo, 0, config->uid, config->gid, fdctl_tile_run );
930 0 : for(;;) {
931 :
932 0 : if( FD_UNLIKELY( !turbine_slot0 ) ) {
933 0 : fd_frag_meta_t * frag = &shred_out_mcache[0]; /* hack to get first frag */
934 0 : if ( frag->sz > 0 ) {
935 0 : uchar * shred_out_chunk = fd_chunk_to_laddr( shred_out_dcache, frag->chunk );
936 0 : fd_shred_base_t * shred_out_shred = (fd_shred_base_t *)fd_type_pun( shred_out_chunk );
937 0 : turbine_slot0 = shred_out_shred->shred.slot;
938 0 : FD_LOG_NOTICE(("turbine_slot0: %lu", turbine_slot0));
939 0 : }
940 0 : }
941 :
942 : /* print metrics */
943 :
944 0 : long now = fd_log_wallclock();
945 0 : int catchup_finished = 0;
946 0 : if( FD_UNLIKELY( now - last_print > 1e9L ) ) {
947 0 : print_tile_metrics( shred_metrics, repair_metrics, repair_metrics_prev, repair_net_links, net_shred_links, net_cnt, &last_sent, last_print, now );
948 0 : ulong slots_behind = turbine_slot0 > repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_HIGHEST_REPAIRED ) ] ? turbine_slot0 - repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_HIGHEST_REPAIRED ) ] : 0;
949 0 : printf(" Repaired slots: %lu/%lu (slots behind: %lu)\n", repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_HIGHEST_REPAIRED ) ], turbine_slot0, slots_behind );
950 0 : if( turbine_slot0 && !slots_behind ) {
951 0 : catchup_finished = 1;
952 0 : }
953 0 : printf("\n");
954 0 : fflush( stdout );
955 0 : last_print = now;
956 0 : }
957 :
958 0 : if( FD_UNLIKELY( catchup_finished ) ) {
959 : /* repair cmd owned memory */
960 0 : location_table = fd_location_table_join( fd_location_table_new( location_table_mem ) );
961 0 : read_iptable( args->repair.iptable_path, location_table );
962 0 : print_peer_location_latency( repair_wksp->wksp, repair_ctx );
963 0 : print_catchup_slots( repair_wksp->wksp, repair_ctx, 0, 1 );
964 0 : FD_LOG_NOTICE(("Catchup to slot %lu completed successfully", turbine_slot0));
965 0 : fd_sys_util_exit_group( 0 );
966 0 : }
967 0 : }
968 0 : }
969 :
970 : /* Tests equivocation detection & repair path. */
971 : static void
972 : repair_cmd_fn_eqvoc( args_t * args,
973 0 : config_t * config ) {
974 0 : (void)args;
975 0 : memset( &config->topo, 0, sizeof(config->topo) );
976 0 : repair_topo( config );
977 :
978 0 : FD_LOG_NOTICE(( "Repair eqvoc testing init" ));
979 0 : fd_topo_print_log( 1, &config->topo );
980 :
981 0 : args_t configure_args = { .configure.command = CONFIGURE_CMD_INIT, };
982 0 : for( ulong i=0UL; STAGES[ i ]; i++ ) configure_args.configure.stages[ i ] = STAGES[ i ];
983 0 : configure_cmd_fn( &configure_args, config );
984 0 : if( 0==strcmp( config->net.provider, "xdp" ) ) fd_topo_install_xdp_simple( &config->topo, config->net.bind_address_parsed );
985 :
986 0 : run_firedancer_init( config, 1, 0 );
987 0 : fd_topo_join_workspaces( &config->topo, FD_SHMEM_JOIN_MODE_READ_WRITE, FD_TOPO_CORE_DUMP_LEVEL_DISABLED );
988 0 : if( 0==strcmp( config->net.provider, "mlx5" ) ) {
989 0 : fd_topo_install_mlx5( &config->topo, NULL );
990 0 : }
991 0 : fd_topo_fill( &config->topo );
992 :
993 0 : repair_load_manifest( &config->topo, args->repair.manifest_path );
994 :
995 0 : ulong repair_tile_idx = fd_topo_find_tile( &config->topo, "repair", 0UL );
996 0 : fd_topo_tile_t * repair_tile = &config->topo.tiles[ repair_tile_idx ];
997 0 : volatile ulong * repair_metrics = fd_metrics_tile( repair_tile->metrics );
998 :
999 0 : void * scratch = fd_topo_obj_laddr( &config->topo, repair_tile->tile_obj_id );
1000 0 : if( FD_UNLIKELY( !scratch ) ) FD_LOG_ERR(( "Failed to access repair tile scratch memory" ));
1001 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
1002 0 : ctx_t * repair_ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
1003 0 : (void)repair_ctx;
1004 :
1005 : /* read tower_out mcache dcache */
1006 0 : ulong tower_out_link_idx = fd_topo_find_link( &config->topo, "tower_out", 0UL );
1007 0 : FD_TEST( tower_out_link_idx!=ULONG_MAX );
1008 0 : fd_topo_link_t * tower_out_link = &config->topo.links[ tower_out_link_idx ];
1009 0 : fd_frag_meta_t * tower_out_mcache = tower_out_link->mcache;
1010 0 : fd_wksp_t * tower_out_mem = config->topo.workspaces[ config->topo.objs[ tower_out_link->dcache_obj_id ].wksp_id ].wksp;
1011 0 : ulong tower_out_chunk0 = fd_dcache_compact_chunk0( tower_out_mem, tower_out_link->dcache );
1012 0 : ulong tower_out_wmark = fd_dcache_compact_wmark( tower_out_mem, tower_out_link->dcache, tower_out_link->mtu );
1013 0 : ulong tower_out_chunk = tower_out_chunk0;
1014 :
1015 0 : fd_topo_run_single_process( &config->topo, 0, config->uid, config->gid, fdctl_tile_run );
1016 0 : int confirmed = 0;
1017 0 : for(;;) {
1018 : /* publish a confirmation on tower_out */
1019 0 : if( FD_UNLIKELY( !confirmed && repair_metrics[ MIDX( GAUGE, REPAIR, SLOT_HIGHEST_REPAIRED ) ] != 0 ) ) {
1020 0 : fd_tower_slot_confirmed_t * msg = fd_chunk_to_laddr( tower_out_mem, tower_out_chunk );
1021 0 : FD_LOG_NOTICE(( "publishing confirmation for slot %lu", msg->slot ));
1022 0 : fd_mcache_publish( tower_out_mcache, tower_out_link->depth, 0, FD_TOWER_SIG_SLOT_CONFIRMED, tower_out_chunk, sizeof(fd_tower_slot_confirmed_t), 0, 0, 0 );
1023 0 : tower_out_chunk = fd_dcache_compact_next( tower_out_chunk, sizeof(fd_tower_slot_confirmed_t), tower_out_chunk0, tower_out_wmark );
1024 0 : confirmed = 1;
1025 0 : }
1026 0 : sleep( 1 );
1027 0 : }
1028 0 : }
1029 :
1030 : static void
1031 : repair_cmd_fn_metrics( args_t * args,
1032 0 : config_t * config ) {
1033 : //memset( &config->topo, 0, sizeof(config->topo) );
1034 :
1035 0 : fd_topo_join_workspaces( &config->topo, FD_SHMEM_JOIN_MODE_READ_ONLY, FD_TOPO_CORE_DUMP_LEVEL_DISABLED );
1036 0 : fd_topo_fill( &config->topo );
1037 :
1038 0 : ctx_t * repair_ctx;
1039 0 : fd_topo_wksp_t * repair_wksp;
1040 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
1041 :
1042 0 : ulong shred_tile_idx = fd_topo_find_tile( &config->topo, "shred", 0UL );
1043 0 : ulong repair_tile_idx = fd_topo_find_tile( &config->topo, "repair", 0UL );
1044 0 : FD_TEST( shred_tile_idx != ULONG_MAX );
1045 0 : FD_TEST( repair_tile_idx!= ULONG_MAX );
1046 0 : fd_topo_tile_t * shred_tile = &config->topo.tiles[ shred_tile_idx ];
1047 0 : fd_topo_tile_t * repair_tile = &config->topo.tiles[ repair_tile_idx ];
1048 :
1049 0 : volatile ulong * shred_metrics = fd_metrics_tile( shred_tile->metrics );
1050 0 : FD_TEST( shred_metrics );
1051 :
1052 0 : volatile ulong * repair_metrics = fd_metrics_tile( repair_tile->metrics );
1053 0 : FD_TEST( repair_metrics );
1054 0 : ulong * repair_metrics_prev = aligned_alloc( 8UL, sizeof(ulong) * FD_METRICS_TOTAL_SZ );
1055 0 : FD_TEST( repair_metrics_prev );
1056 0 : memset( repair_metrics_prev, 0, sizeof(ulong) * FD_METRICS_TOTAL_SZ );
1057 :
1058 :
1059 0 : ulong net_tile_cnt = config->layout.net_tile_count;
1060 0 : volatile ulong ** repair_net_links = aligned_alloc( 8UL, net_tile_cnt * sizeof(volatile ulong*) );
1061 0 : volatile ulong ** net_shred_links = aligned_alloc( 8UL, net_tile_cnt * sizeof(volatile ulong*) );
1062 0 : FD_TEST( repair_net_links );
1063 0 : FD_TEST( net_shred_links );
1064 :
1065 0 : char const * net_tile_name = fd_net_tile_name( config->net.provider );
1066 0 : for( ulong i = 0UL; i < net_tile_cnt; i++ ) {
1067 : /* process all repair_net links */
1068 0 : ulong tile_idx = fd_topo_find_tile( &config->topo, net_tile_name, i );
1069 0 : if( FD_UNLIKELY( tile_idx == ULONG_MAX ) ) FD_LOG_ERR(( "net tile %lu not found", i ));
1070 0 : fd_topo_tile_t * tile = &config->topo.tiles[ tile_idx ];
1071 :
1072 0 : ulong repair_net_in_idx = fd_topo_find_tile_in_link( &config->topo, tile, "repair_net", 0UL );
1073 0 : if( FD_UNLIKELY( repair_net_in_idx == ULONG_MAX ) ) FD_LOG_ERR(( "repair_net link not found for net tile %lu", i ));
1074 0 : repair_net_links[i] = fd_metrics_link_in( tile->metrics, repair_net_in_idx );
1075 0 : FD_TEST( repair_net_links[i] );
1076 :
1077 : /* process all net_shred links */
1078 0 : tile_idx = fd_topo_find_tile( &config->topo, "shred", 0 );
1079 0 : if( FD_UNLIKELY( tile_idx == ULONG_MAX ) ) FD_LOG_ERR(( "shred tile 0 not found" ));
1080 0 : fd_topo_tile_t * shred_tile = &config->topo.tiles[ tile_idx ];
1081 :
1082 0 : ulong shred_out_in_idx = fd_topo_find_tile_in_link( &config->topo, shred_tile, "net_shred", i );
1083 0 : if( FD_UNLIKELY( shred_out_in_idx == ULONG_MAX ) ) FD_LOG_ERR(( "net_shred link not found for shred tile 0" ));
1084 0 : net_shred_links[i] = fd_metrics_link_in( shred_tile->metrics, shred_out_in_idx );
1085 0 : FD_TEST( net_shred_links[i] );
1086 0 : }
1087 :
1088 0 : long last_print_ts = fd_log_wallclock();
1089 0 : ulong last_sent = 0UL;
1090 0 : for(;;) {
1091 0 : long now = fd_log_wallclock();
1092 0 : if( FD_UNLIKELY( now - last_print_ts > 1e9L ) ) {
1093 0 : print_tile_metrics( shred_metrics, repair_metrics, repair_metrics_prev, repair_net_links, net_shred_links, net_tile_cnt, &last_sent, last_print_ts, now );
1094 0 : last_print_ts = now;
1095 0 : }
1096 0 : }
1097 0 : }
1098 :
1099 : static void
1100 : repair_cmd_fn_forest( args_t * args,
1101 0 : config_t * config ) {
1102 0 : ctx_t * repair_ctx;
1103 0 : fd_topo_wksp_t * repair_wksp;
1104 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
1105 :
1106 0 : ulong forest_gaddr = fd_wksp_gaddr_fast( repair_ctx->wksp, repair_ctx->forest );
1107 0 : fd_forest_t * forest = (fd_forest_t *)fd_wksp_laddr( repair_wksp->wksp, forest_gaddr );
1108 :
1109 0 : for( ;; ) {
1110 0 : fd_forest_print( forest );
1111 0 : sleep( 1 );
1112 0 : }
1113 0 : }
1114 :
1115 : static void
1116 : repair_cmd_fn_inflight( args_t * args,
1117 0 : config_t * config ) {
1118 0 : ctx_t * repair_ctx;
1119 0 : fd_topo_wksp_t * repair_wksp;
1120 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
1121 :
1122 0 : ulong inflights_gaddr = fd_wksp_gaddr_fast( repair_ctx->wksp, repair_ctx->inflights );
1123 0 : fd_inflights_t * inflights = (fd_inflights_t *)fd_wksp_laddr( repair_wksp->wksp, inflights_gaddr );
1124 :
1125 0 : ulong inflight_pool_off = (ulong)inflights->pool - (ulong)repair_ctx->inflights;
1126 0 : fd_inflight_t * inflight_pool = (fd_inflight_t *)fd_wksp_laddr( repair_wksp->wksp, inflights_gaddr + inflight_pool_off );
1127 :
1128 0 : for( ;; ) {
1129 0 : fd_inflights_print( inflights->outstanding_dl, inflight_pool );
1130 0 : printf("popped count: %lu\n", inflights->popped_cnt);
1131 0 : fd_inflights_print( inflights->popped_dl, inflight_pool );
1132 0 : sleep( 1 );
1133 0 : }
1134 0 : }
1135 :
1136 : static void
1137 : repair_cmd_fn_requests( args_t * args,
1138 0 : config_t * config ) {
1139 0 : ctx_t * repair_ctx;
1140 0 : fd_topo_wksp_t * repair_wksp;
1141 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
1142 :
1143 0 : fd_forest_t * forest = fd_forest_join( fd_wksp_laddr( repair_wksp->wksp, fd_wksp_gaddr_fast( repair_ctx->wksp, repair_ctx->forest ) ) );
1144 0 : fd_forest_reqslist_t * dlist = fd_forest_reqslist( forest );
1145 0 : fd_forest_ref_t * pool = fd_forest_reqspool( forest );
1146 :
1147 0 : fd_forest_reqslist_t * orphlist = fd_forest_orphlist( forest );
1148 :
1149 0 : for( ;; ) {
1150 0 : printf("%-15s %-12s %-12s %-12s %-20s %-12s\n",
1151 0 : "Slot", "Buffered Idx", "Complete Idx", "First Shred ts", "Turbine Cnt", "Repair Cnt");
1152 0 : printf("%-15s %-12s %-12s %-12s %-20s %-12s\n",
1153 0 : "---------------", "------------", "------------", "------------",
1154 0 : "--------------------", "------------");
1155 0 : for( fd_forest_reqslist_iter_t iter = fd_forest_reqslist_iter_fwd_init( dlist, pool );
1156 0 : !fd_forest_reqslist_iter_done( iter, dlist, pool );
1157 0 : iter = fd_forest_reqslist_iter_fwd_next( iter, dlist, pool ) ) {
1158 0 : fd_forest_ref_t * req = fd_forest_reqslist_iter_ele( iter, dlist, pool );
1159 0 : fd_forest_blk_t * blk = fd_forest_pool_ele( fd_forest_pool( forest ), req->idx );
1160 :
1161 0 : printf("%-15lu %-12u %-12u %-20ld %-12u %-10u\n",
1162 0 : blk->slot,
1163 0 : blk->buffered_idx,
1164 0 : blk->complete_idx,
1165 0 : blk->first_shred_ts,
1166 0 : blk->turbine_cnt,
1167 0 : blk->repair_cnt);
1168 0 : }
1169 0 : printf("\n");
1170 :
1171 : /* now lets print the orphreqs */
1172 :
1173 0 : printf("Orphan Requests:\n");
1174 0 : printf("%-15s %-12s %-12s %-12s %-20s %-12s %-10s\n",
1175 0 : "Slot", "Consumed Idx", "Buffered Idx", "Complete Idx",
1176 0 : "First Shred Timestamp", "Turbine Cnt", "Repair Cnt");
1177 0 : printf("%-15s %-12s %-12s %-12s %-20s %-12s %-10s\n",
1178 0 : "---------------", "------------", "------------", "------------",
1179 0 : "--------------------", "------------", "----------");
1180 :
1181 0 : for( fd_forest_reqslist_iter_t iter = fd_forest_reqslist_iter_fwd_init( orphlist, pool );
1182 0 : !fd_forest_reqslist_iter_done( iter, orphlist, pool );
1183 0 : iter = fd_forest_reqslist_iter_fwd_next( iter, orphlist, pool ) ) {
1184 0 : fd_forest_ref_t * req = fd_forest_reqslist_iter_ele( iter, orphlist, pool );
1185 0 : fd_forest_blk_t * blk = fd_forest_pool_ele( fd_forest_pool( forest ), req->idx );
1186 0 : printf("%-15lu %-12u %-12u %-20ld %-12u %-10u\n",
1187 0 : blk->slot,
1188 0 : blk->buffered_idx,
1189 0 : blk->complete_idx,
1190 0 : blk->first_shred_ts,
1191 0 : blk->turbine_cnt,
1192 0 : blk->repair_cnt);
1193 0 : }
1194 0 : sleep( 1 );
1195 0 : }
1196 0 : }
1197 :
1198 : static void
1199 : repair_cmd_fn_waterfall( args_t * args,
1200 0 : config_t * config ) {
1201 :
1202 0 : fd_topo_t * topo = &config->topo;
1203 0 : ulong wksp_id = fd_topo_find_wksp( topo, "repair" );
1204 0 : if( FD_UNLIKELY( wksp_id==ULONG_MAX ) ) FD_LOG_ERR(( "repair workspace not found" ));
1205 0 : fd_topo_wksp_t * repair_wksp = &topo->workspaces[ wksp_id ];
1206 0 : fd_topo_join_workspace( topo, repair_wksp, FD_SHMEM_JOIN_MODE_READ_ONLY, FD_TOPO_CORE_DUMP_LEVEL_DISABLED );
1207 :
1208 : /* Access the repair tile scratch memory where repair_tile_ctx is stored */
1209 0 : ulong tile_id = fd_topo_find_tile( topo, "repair", 0UL );
1210 0 : if( FD_UNLIKELY( tile_id==ULONG_MAX ) ) FD_LOG_ERR(( "repair tile not found" ));
1211 0 : fd_topo_tile_t * tile = &topo->tiles[ tile_id ];
1212 0 : void * scratch = fd_topo_obj_laddr( &config->topo, tile->tile_obj_id );
1213 0 : if( FD_UNLIKELY( !scratch ) ) FD_LOG_ERR(( "Failed to access repair tile scratch memory" ));
1214 :
1215 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
1216 0 : ctx_t * repair_ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(ctx_t), sizeof(ctx_t) );
1217 :
1218 : /* catchup cmd owned memory */
1219 0 : location_table = fd_location_table_join( fd_location_table_new( location_table_mem ) );
1220 0 : read_iptable( args->repair.iptable_path, location_table );
1221 :
1222 : // Add terminal setup here - same as monitor.c
1223 0 : atexit( restore_terminal );
1224 0 : if( FD_UNLIKELY( 0!=tcgetattr( STDIN_FILENO, &termios_backup ) ) ) {
1225 0 : FD_LOG_ERR(( "tcgetattr(STDIN_FILENO) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
1226 0 : }
1227 :
1228 : /* Disable character echo and line buffering */
1229 0 : struct termios term = termios_backup;
1230 0 : term.c_lflag &= (tcflag_t)~(ICANON | ECHO);
1231 0 : if( FD_UNLIKELY( 0!=tcsetattr( STDIN_FILENO, TCSANOW, &term ) ) ) {
1232 0 : FD_LOG_WARNING(( "tcsetattr(STDIN_FILENO) failed (%i-%s)", errno, fd_io_strerror( errno ) ));
1233 0 : }
1234 :
1235 0 : int catchup_verbose = 0;
1236 0 : long last_print = 0;
1237 0 : for( ;; ) {
1238 0 : int c = fd_getchar();
1239 0 : if( FD_UNLIKELY( c=='i' ) ) catchup_verbose = !catchup_verbose;
1240 0 : if( FD_UNLIKELY( c=='\x04' ) ) break; /* Ctrl-D */
1241 :
1242 0 : long now = fd_log_wallclock();
1243 0 : if( FD_UNLIKELY( now - last_print > 1e9L ) ) {
1244 0 : last_print = now;
1245 0 : print_catchup_slots( repair_wksp->wksp, repair_ctx, catchup_verbose, args->repair.sort_by_slot );
1246 0 : printf( "catchup slots | Use 'i' to toggle extra slot information" TEXT_NEWLINE );
1247 0 : fflush( stdout );
1248 :
1249 : /* Peer location latency is not that useful post catchup, and also
1250 : requires some concurrent dlist iteration, so only print it when
1251 : in catchup mode. */
1252 0 : }
1253 0 : }
1254 0 : }
1255 :
1256 0 : #define PEERS_DISPLAY_MAX 20
1257 :
1258 : static void
1259 : print_peer_dlist( fd_policy_peer_dlist_t * dlist,
1260 : fd_policy_peer_t * pool,
1261 : fd_policy_peer_dlist_iter_t cursor,
1262 0 : char const * label ) {
1263 0 : ulong cnt = 0;
1264 0 : for( fd_policy_peer_dlist_iter_t it = fd_policy_peer_dlist_iter_fwd_init( dlist, pool );
1265 0 : !fd_policy_peer_dlist_iter_done( it, dlist, pool );
1266 0 : it = fd_policy_peer_dlist_iter_fwd_next( it, dlist, pool ) ) cnt++;
1267 :
1268 0 : printf( "%s (%lu peers)\n", label, cnt );
1269 0 : if( !cnt || fd_policy_peer_dlist_iter_done( cursor, dlist, pool ) ) {
1270 0 : printf( " (empty or iterator not initialized)\n\n" );
1271 0 : return;
1272 0 : }
1273 :
1274 0 : printf( " | %-8s | %-12s | %-12s | %-8s | %-8s\n",
1275 0 : "Idx", "Pubkey", "Ewma Lat", "Avg Lat", "Req/Res" );
1276 0 : printf( "-----+----------+--------------+--------------+----------+---------\n" );
1277 :
1278 0 : fd_policy_peer_dlist_iter_t it = cursor;
1279 0 : for( ulong i = 0; i < PEERS_DISPLAY_MAX && i < cnt; i++ ) {
1280 0 : fd_policy_peer_t * peer = fd_policy_peer_dlist_iter_ele( it, dlist, pool );
1281 :
1282 0 : FD_BASE58_ENCODE_32_BYTES( peer->key.key, b58 );
1283 0 : char pubkey_short[13];
1284 0 : fd_cstr_fini( fd_cstr_append_text( fd_cstr_init( pubkey_short ), b58, 12 ) );
1285 :
1286 0 : double avg_lat_ms = peer->res_cnt ? ((double)peer->total_lat / (double)peer->res_cnt) / 1e6 : 0.0;
1287 0 : double ewma_lat_ms = (double)peer->ewma_lat / 1e6;
1288 :
1289 0 : printf( " %s%c%s | %-8lu | %-12s | %9.3fms | %9.3fms | %lu/%lu\n",
1290 0 : i == 0 ? "\033[1;33m" : "",
1291 0 : i == 0 ? '>' : ' ',
1292 0 : i == 0 ? "\033[0m" : "",
1293 0 : fd_policy_peer_pool_idx( pool, peer ),
1294 0 : pubkey_short,
1295 0 : ewma_lat_ms,
1296 0 : avg_lat_ms,
1297 0 : peer->req_cnt,
1298 0 : peer->res_cnt );
1299 :
1300 0 : it = fd_policy_peer_dlist_iter_fwd_next( it, dlist, pool );
1301 0 : if( fd_policy_peer_dlist_iter_done( it, dlist, pool ) ) {
1302 0 : it = fd_policy_peer_dlist_iter_fwd_init( dlist, pool );
1303 0 : }
1304 0 : }
1305 0 : if( cnt > PEERS_DISPLAY_MAX ) printf( " ... (%lu more)\n", cnt - PEERS_DISPLAY_MAX );
1306 0 : printf( "\n" );
1307 0 : }
1308 :
1309 : static void
1310 : repair_cmd_fn_peers( args_t * args,
1311 0 : config_t * config ) {
1312 0 : ctx_t * repair_ctx;
1313 0 : fd_topo_wksp_t * repair_wksp;
1314 0 : repair_ctx_wksp( args, config, &repair_ctx, &repair_wksp );
1315 :
1316 0 : fd_policy_t * policy = fd_wksp_laddr( repair_wksp->wksp, fd_wksp_gaddr_fast( repair_ctx->wksp, repair_ctx->policy ) );
1317 :
1318 0 : fd_policy_peer_dlist_t * fast_dlist = fd_wksp_laddr( repair_wksp->wksp, fd_wksp_gaddr_fast( repair_ctx->wksp, policy->peers.fast ) );
1319 0 : fd_policy_peer_dlist_t * slow_dlist = fd_wksp_laddr( repair_wksp->wksp, fd_wksp_gaddr_fast( repair_ctx->wksp, policy->peers.slow ) );
1320 0 : fd_policy_peer_t * pool = fd_wksp_laddr( repair_wksp->wksp, fd_wksp_gaddr_fast( repair_ctx->wksp, policy->peers.pool ) );
1321 :
1322 0 : long last_print = 0;
1323 0 : for( ;; ) {
1324 0 : long now = fd_log_wallclock();
1325 0 : if( FD_UNLIKELY( now - last_print > 1e9L ) ) {
1326 0 : last_print = now;
1327 0 : printf( "\033[2J\033[H" );
1328 :
1329 0 : char fast_label[64];
1330 0 : char slow_label[64];
1331 0 : snprintf( fast_label, sizeof(fast_label), "FAST PEERS (ewma < %ldms)", (long)(FD_POLICY_LATENCY_THRESH / 1e6L) );
1332 0 : snprintf( slow_label, sizeof(slow_label), "SLOW PEERS (ewma >= %ldms or no responses)", (long)(FD_POLICY_LATENCY_THRESH / 1e6L) );
1333 0 : print_peer_dlist( fast_dlist, pool, policy->peers.select.fast_iter, fast_label );
1334 0 : print_peer_dlist( slow_dlist, pool, policy->peers.select.slow_iter, slow_label );
1335 :
1336 0 : printf( "select cnt: %u / %u (fast per slow)\n", policy->peers.select.cnt, FD_POLICY_FAST_PER_SLOW );
1337 0 : printf( "pool used: %lu / %lu\n", fd_policy_peer_pool_used( pool ), fd_policy_peer_pool_max( pool ) );
1338 :
1339 0 : fflush( stdout );
1340 0 : }
1341 :
1342 0 : }
1343 0 : }
1344 :
1345 :
1346 :
1347 :
1348 : void
1349 : repair_cmd_args( int * pargc,
1350 : char *** pargv,
1351 0 : args_t * args ) {
1352 :
1353 : /* positional arg */
1354 :
1355 0 : args->repair.pos_arg = (*pargv)[0];
1356 0 : if( FD_UNLIKELY( !args->repair.pos_arg ) ) {
1357 0 : args->repair.help = 1;
1358 0 : return;
1359 0 : }
1360 0 : (*pargc)--;
1361 0 : (*pargv)++;
1362 :
1363 : /* required args */
1364 :
1365 0 : char const * manifest_path = fd_env_strip_cmdline_cstr ( pargc, pargv, "--manifest-path", NULL, NULL );
1366 :
1367 : /* optional args */
1368 :
1369 0 : char const * iptable_path = fd_env_strip_cmdline_cstr ( pargc, pargv, "--iptable", NULL, NULL );
1370 0 : ulong slot = fd_env_strip_cmdline_ulong ( pargc, pargv, "--slot", NULL, ULONG_MAX );
1371 0 : int sort_by_slot = fd_env_strip_cmdline_contains( pargc, pargv, "--sort-by-slot" );
1372 :
1373 0 : if( FD_UNLIKELY( !strcmp( args->repair.pos_arg, "catchup" ) && !manifest_path ) ) {
1374 0 : args->repair.help = 1;
1375 0 : return;
1376 0 : }
1377 :
1378 0 : fd_cstr_fini( fd_cstr_append_cstr_safe( fd_cstr_init( args->repair.manifest_path ), manifest_path, sizeof(args->repair.manifest_path)-1UL ) );
1379 0 : fd_cstr_fini( fd_cstr_append_cstr_safe( fd_cstr_init( args->repair.iptable_path ), iptable_path, sizeof(args->repair.iptable_path )-1UL ) );
1380 0 : args->repair.slot = slot;
1381 0 : args->repair.sort_by_slot = sort_by_slot;
1382 0 : }
1383 :
1384 : static void
1385 : repair_cmd_fn( args_t * args,
1386 0 : config_t * config ) {
1387 :
1388 0 : if( args->repair.help ) {
1389 0 : fd_action_help_print( &fd_action_repair );
1390 0 : return;
1391 0 : }
1392 :
1393 0 : if ( !strcmp( args->repair.pos_arg, "catchup" ) ) repair_cmd_fn_catchup ( args, config );
1394 0 : else if( !strcmp( args->repair.pos_arg, "eqvoc" ) ) repair_cmd_fn_eqvoc ( args, config );
1395 0 : else if( !strcmp( args->repair.pos_arg, "forest" ) ) repair_cmd_fn_forest ( args, config );
1396 0 : else if( !strcmp( args->repair.pos_arg, "inflight" ) ) repair_cmd_fn_inflight ( args, config );
1397 0 : else if( !strcmp( args->repair.pos_arg, "requests" ) ) repair_cmd_fn_requests ( args, config );
1398 0 : else if( !strcmp( args->repair.pos_arg, "waterfall" ) ) repair_cmd_fn_waterfall( args, config );
1399 0 : else if( !strcmp( args->repair.pos_arg, "peers" ) ) repair_cmd_fn_peers ( args, config );
1400 0 : else if( !strcmp( args->repair.pos_arg, "metrics" ) ) repair_cmd_fn_metrics ( args, config );
1401 0 : else fd_action_help_print( &fd_action_repair );
1402 0 : }
1403 :
1404 : static void
1405 0 : repair_args_help( fd_action_help_t * help ) {
1406 0 : fd_action_help_arg( help, "catchup", NULL, "Run a reduced topology that only repairs slots until catchup.\n"
1407 0 : "Requires --manifest-path; accepts --iptable and --sort-by-slot" );
1408 0 : fd_action_help_arg( help, "eqvoc", NULL, "Test equivocation detection and the repair path" );
1409 0 : fd_action_help_arg( help, "forest", NULL, "Print the repair forest. Accepts --slot to drill into a slot" );
1410 0 : fd_action_help_arg( help, "inflight", NULL, "Print the inflight repairs" );
1411 0 : fd_action_help_arg( help, "requests", NULL, "Print the queued repair requests" );
1412 0 : fd_action_help_arg( help, "waterfall", NULL, "Print a waterfall diagram of recent slot completion times and\n"
1413 0 : "response latencies. Accepts --iptable and --sort-by-slot" );
1414 0 : fd_action_help_arg( help, "peers", NULL, "Print the list of slow and fast repair peers" );
1415 0 : fd_action_help_arg( help, "metrics", NULL, "Print repair tile metrics in a digestible format" );
1416 0 : fd_action_help_arg( help, "--manifest-path", "<path>", "Path to manifest file (required by catchup)" );
1417 0 : fd_action_help_arg( help, "--iptable", "<path>", "Path to iptable file (catchup, waterfall)" );
1418 0 : fd_action_help_arg( help, "--slot", "<slot>", "Specific forest slot to drill into (forest)" );
1419 : fd_action_help_arg( help, "--sort-by-slot", NULL, "Sort results by slot (catchup, waterfall)" );
1420 0 : }
1421 :
1422 : action_t fd_action_repair = {
1423 : .name = "repair",
1424 : .args = repair_cmd_args,
1425 : .fn = repair_cmd_fn,
1426 : .perm = dev_cmd_perm,
1427 : .description = "Spawn a reduced topology for inspecting and profiling the repair tile",
1428 : .detail = "Boots a smaller Firedancer topology focused on the repair tile and runs the\n"
1429 : "requested subcommand to drive or inspect repair behavior. Pick one of the\n"
1430 : "subcommands below.",
1431 : .usage = "repair <catchup|eqvoc|forest|inflight|requests|waterfall|peers|metrics> [OPTIONS]",
1432 : .args_help = repair_args_help,
1433 : };
|