Line data Source code
1 : #define _GNU_SOURCE
2 : #include "fd_snapct_tile.h"
3 : #include "utils/fd_sspeer.h"
4 : #include "utils/fd_ssping.h"
5 : #include "utils/fd_ssctrl.h"
6 : #include "utils/fd_ssarchive.h"
7 : #include "utils/fd_http_resolver.h"
8 : #include "utils/fd_ssmsg.h"
9 : #include "../backup/fd_snap_pool.h"
10 :
11 : #include "../../disco/topo/fd_topo.h"
12 : #include "../../disco/topo/fd_dns_resolve.h"
13 : #include "../../disco/metrics/fd_metrics.h"
14 : #include "../../flamenco/gossip/fd_gossip_message.h"
15 : #include "../../waltz/resolv/fd_netdb.h"
16 : #include "../../waltz/resolv/fd_adns.h"
17 : #include "../../util/fd_hash32.h"
18 :
19 : #include <errno.h>
20 : #include <stdio.h>
21 : #include <fcntl.h>
22 : #include <unistd.h>
23 : #include <sys/stat.h>
24 : #include <netinet/tcp.h>
25 : #include <netinet/in.h>
26 :
27 : #include "generated/fd_snapct_tile_seccomp.h"
28 :
29 : #define NAME "snapct"
30 :
31 : /* FIXME: Do a finishing pass over the default.toml config options / comments */
32 :
33 117 : #define GOSSIP_PEERS_MAX (FD_CONTACT_INFO_TABLE_SIZE)
34 63 : #define SERVER_PEERS_MAX (FD_TOPO_SNAPSHOTS_SERVERS_MAX_RESOLVED)
35 63 : #define TOTAL_PEERS_MAX (GOSSIP_PEERS_MAX + SERVER_PEERS_MAX)
36 :
37 0 : #define IN_KIND_ACK (0)
38 0 : #define IN_KIND_SNAPLD (1)
39 0 : #define IN_KIND_GOSSIP (2)
40 : #define MAX_IN_LINKS (4)
41 :
42 : struct fd_snapct_out_link {
43 : ulong idx;
44 : fd_wksp_t * mem;
45 : ulong chunk0;
46 : ulong wmark;
47 : ulong chunk;
48 : ulong mtu;
49 : };
50 : typedef struct fd_snapct_out_link fd_snapct_out_link_t;
51 :
52 0 : #define FD_SNAPCT_COLLECTING_PEERS_TIMEOUT (90L*1000L*1000L*1000L) /* 1.5 minutes */
53 :
54 : struct gossip_ci_entry {
55 : fd_pubkey_t pubkey;
56 : int allowed;
57 : fd_ip4_port_t rpc_addr;
58 : ulong map_next;
59 : };
60 : typedef struct gossip_ci_entry gossip_ci_entry_t;
61 :
62 : #define MAP_NAME gossip_ci_map
63 21 : #define MAP_KEY pubkey
64 : #define MAP_ELE_T gossip_ci_entry_t
65 : #define MAP_KEY_T fd_pubkey_t
66 24 : #define MAP_NEXT map_next
67 12 : #define MAP_KEY_EQ(k0,k1) fd_pubkey_eq( k0, k1 )
68 75 : #define MAP_KEY_HASH(key,seed) fd_hash32( key->uc, seed )
69 : #include "../../util/tmpl/fd_map_chain.c"
70 :
71 : /* Standalone blacklist keyed on fd_sspeer_key_t (peer identity).
72 : Keying on identity rather than network address is intentional:
73 : a peer's address can change freely, but its identity cannot.
74 : For gossip peers the identity is the pubkey; for URL peers it is
75 : the (hostname, resolved_addr) pair. Blacklisted peers are
76 : permanently banned for the bootstrap lifetime. The blacklist uses
77 : its own dedicated pool so entries are never evicted by pool
78 : pressure, unlike the ssping peer pool. */
79 :
80 : struct fd_sspeer_blacklist_entry {
81 : fd_sspeer_key_t key;
82 : ulong pool_next;
83 : ulong map_next;
84 : };
85 : typedef struct fd_sspeer_blacklist_entry fd_sspeer_blacklist_entry_t;
86 :
87 : #define POOL_NAME blacklist_pool
88 54 : #define POOL_T fd_sspeer_blacklist_entry_t
89 : #define POOL_IDX_T ulong
90 787992 : #define POOL_NEXT pool_next
91 : #include "../../util/tmpl/fd_pool.c"
92 :
93 : #define MAP_NAME blacklist_map
94 18 : #define MAP_KEY key
95 : #define MAP_ELE_T fd_sspeer_blacklist_entry_t
96 : #define MAP_KEY_T fd_sspeer_key_t
97 51 : #define MAP_NEXT map_next
98 39 : #define MAP_KEY_EQ(k0,k1) (fd_sspeer_key_eq(k0,k1))
99 75 : #define MAP_KEY_HASH(key,seed) (fd_sspeer_key_hash(key,seed))
100 : #include "../../util/tmpl/fd_map_chain.c"
101 :
102 : struct fd_snapct_tile {
103 : struct fd_topo_tile_snapct config;
104 : int gossip_enabled;
105 : int download_enabled;
106 :
107 : fd_netdb_fds_t netdb_fds[1];
108 :
109 : fd_adns_t * adns;
110 : struct {
111 : char hostname[ FD_FQDN_BUF_MAX ];
112 : ushort port; /* net order */
113 : int is_https;
114 : int resolved;
115 : long retry_nanos; /* re-queue due time, 0 while in flight */
116 : } dns_servers[ FD_TOPO_SNAPSHOTS_SERVERS_MAX ];
117 : struct {
118 : char hostname[ FD_FQDN_BUF_MAX ];
119 : ushort port;
120 : int resolved;
121 : long retry_nanos;
122 : } dns_entrypoints[ FD_TOPO_GOSSIP_ENTRYPOINTS_MAX ];
123 :
124 : ulong resolved_servers_cnt;
125 : struct {
126 : fd_ip4_port_t addr;
127 : char hostname[ FD_FQDN_BUF_MAX ];
128 : int is_https;
129 : } resolved_servers[ FD_TOPO_SNAPSHOTS_SERVERS_MAX_RESOLVED ];
130 :
131 : ulong resolved_entrypoints_cnt;
132 : fd_ip4_port_t resolved_entrypoints[ FD_TOPO_GOSSIP_ENTRYPOINTS_MAX ];
133 :
134 : fd_ssping_t * ssping;
135 : fd_http_resolver_t * ssresolver;
136 : fd_sspeer_selector_t * selector;
137 : ulong ssping_seed;
138 : ulong gossip_ci_seed;
139 : ulong selector_seed;
140 : ulong blacklist_seed;
141 :
142 : fd_sspeer_blacklist_entry_t * blacklist_pool;
143 : blacklist_map_t * blacklist_map;
144 :
145 : int state;
146 : int malformed;
147 : int load_complete;
148 : long deadline_nanos;
149 : int flush_ack;
150 : int flush_ack_cnt;
151 : fd_sspeer_t peer;
152 :
153 : struct {
154 : int dir_fd;
155 : int full_snapshot_fd;
156 : int incremental_snapshot_fd;
157 : char full_snapshot_name[ FD_SNAP_NAME_MAX ];
158 : char incremental_snapshot_name[ FD_SNAP_NAME_MAX ];
159 : } local_out;
160 :
161 : char http_full_snapshot_name[ PATH_MAX ];
162 : char http_incr_snapshot_name[ PATH_MAX ];
163 :
164 : void const * gossip_in_mem;
165 : void const * snapld_in_mem;
166 : uchar in_kind[ MAX_IN_LINKS ];
167 :
168 : struct {
169 : ulong full_slot;
170 : ulong slot;
171 : int pending;
172 : } predicted_incremental;
173 :
174 : struct {
175 : ulong full_snapshot_slot;
176 : char full_snapshot_path[ PATH_MAX ];
177 : ulong full_snapshot_size;
178 : int full_snapshot_zstd;
179 :
180 : uchar full_snapshot_hash[ FD_HASH_FOOTPRINT ];
181 : uchar incremental_snapshot_hash[ FD_HASH_FOOTPRINT ];
182 :
183 : ulong incremental_snapshot_slot;
184 : char incremental_snapshot_path[ PATH_MAX ];
185 : ulong incremental_snapshot_size;
186 : int incremental_snapshot_zstd;
187 : } local_in;
188 :
189 : struct {
190 : struct {
191 : ulong bytes_read;
192 : ulong bytes_written;
193 : ulong bytes_total;
194 : uint num_retries;
195 : } full;
196 :
197 : struct {
198 : ulong bytes_read;
199 : ulong bytes_written;
200 : ulong bytes_total;
201 : uint num_retries;
202 : } incremental;
203 : } metrics;
204 :
205 : struct {
206 : gossip_ci_entry_t * ci_table; /* flat array of all gossip entries, allowed or not */
207 : gossip_ci_map_t * ci_map; /* map from pubkey to only allowed gossip entries */
208 : ulong allowed_cnt; /* number of allowed entries in ci_map */
209 : int saturated;
210 : } gossip;
211 :
212 : long snapshot_start_timestamp_ns;
213 :
214 : fd_snapct_out_link_t out_ld;
215 : fd_snapct_out_link_t out_gui;
216 : fd_snapct_out_link_t out_rp;
217 : };
218 : typedef struct fd_snapct_tile fd_snapct_tile_t;
219 :
220 : static int
221 0 : gossip_enabled( fd_topo_tile_t const * tile ) {
222 0 : return tile->snapct.sources.gossip.allow_any || tile->snapct.sources.gossip.allow_list_cnt>0UL;
223 0 : }
224 :
225 : static int
226 0 : download_enabled( fd_topo_tile_t const * tile ) {
227 0 : return gossip_enabled( tile ) || tile->snapct.sources.servers_cnt>0UL;
228 0 : }
229 :
230 0 : #define ADNS_REQS_MAX (FD_TOPO_SNAPSHOTS_SERVERS_MAX+FD_TOPO_GOSSIP_ENTRYPOINTS_MAX)
231 :
232 : static ulong
233 126 : scratch_align( void ) {
234 126 : return fd_ulong_max( alignof(fd_snapct_tile_t),
235 126 : fd_ulong_max( fd_ssping_align(),
236 126 : fd_ulong_max( alignof(gossip_ci_entry_t),
237 126 : fd_ulong_max( gossip_ci_map_align(),
238 126 : fd_ulong_max( fd_http_resolver_align(),
239 126 : fd_ulong_max( fd_sspeer_selector_align(),
240 126 : fd_ulong_max( blacklist_pool_align(),
241 126 : fd_ulong_max( blacklist_map_align(),
242 126 : fd_adns_align() ) ) ) ) ) ) ) );
243 126 : }
244 :
245 : static ulong
246 42 : scratch_footprint( fd_topo_tile_t const * tile FD_PARAM_UNUSED ) {
247 42 : ulong l = FD_LAYOUT_INIT;
248 42 : l = FD_LAYOUT_APPEND( l, alignof(fd_snapct_tile_t), sizeof(fd_snapct_tile_t) );
249 42 : l = FD_LAYOUT_APPEND( l, fd_ssping_align(), fd_ssping_footprint( TOTAL_PEERS_MAX ) );
250 42 : l = FD_LAYOUT_APPEND( l, alignof(gossip_ci_entry_t), sizeof(gossip_ci_entry_t) * GOSSIP_PEERS_MAX );
251 42 : l = FD_LAYOUT_APPEND( l, gossip_ci_map_align(), gossip_ci_map_footprint( gossip_ci_map_chain_cnt_est( GOSSIP_PEERS_MAX ) ) );
252 42 : l = FD_LAYOUT_APPEND( l, fd_http_resolver_align(), fd_http_resolver_footprint( SERVER_PEERS_MAX ) );
253 42 : l = FD_LAYOUT_APPEND( l, fd_sspeer_selector_align(), fd_sspeer_selector_footprint( TOTAL_PEERS_MAX ) );
254 42 : l = FD_LAYOUT_APPEND( l, blacklist_pool_align(), blacklist_pool_footprint( TOTAL_PEERS_MAX ) );
255 42 : l = FD_LAYOUT_APPEND( l, blacklist_map_align(), blacklist_map_footprint( blacklist_map_chain_cnt_est( TOTAL_PEERS_MAX ) ) );
256 42 : l = FD_LAYOUT_APPEND( l, fd_adns_align(), fd_adns_footprint( ADNS_REQS_MAX ) );
257 42 : l = FD_LAYOUT_APPEND( l, fd_alloc_align(), fd_alloc_footprint() );
258 42 : return FD_LAYOUT_FINI( l, scratch_align() );
259 42 : }
260 :
261 : static inline int
262 0 : should_shutdown( fd_snapct_tile_t * ctx ) {
263 0 : return ctx->state==FD_SNAPCT_STATE_SHUTDOWN;
264 0 : }
265 :
266 : static void
267 0 : metrics_write( fd_snapct_tile_t * ctx ) {
268 0 : FD_MGAUGE_SET( SNAPCT, FULL_BYTES_READ, ctx->metrics.full.bytes_read );
269 0 : FD_MGAUGE_SET( SNAPCT, FULL_BYTES_WRITTEN, ctx->metrics.full.bytes_written );
270 0 : FD_MGAUGE_SET( SNAPCT, FULL_SIZE_BYTES, ctx->metrics.full.bytes_total );
271 0 : FD_MGAUGE_SET( SNAPCT, FULL_RETRY, ctx->metrics.full.num_retries );
272 :
273 0 : FD_MGAUGE_SET( SNAPCT, INCREMENTAL_BYTES_READ, ctx->metrics.incremental.bytes_read );
274 0 : FD_MGAUGE_SET( SNAPCT, INCREMENTAL_BYTES_WRITTEN, ctx->metrics.incremental.bytes_written );
275 0 : FD_MGAUGE_SET( SNAPCT, INCREMENTAL_SIZE_BYTES, ctx->metrics.incremental.bytes_total );
276 0 : FD_MGAUGE_SET( SNAPCT, INCREMENTAL_RETRY, ctx->metrics.incremental.num_retries );
277 :
278 0 : FD_MGAUGE_SET( SNAPCT, PREDICTED_SLOT, ctx->predicted_incremental.slot );
279 :
280 :
281 0 : FD_MGAUGE_SET( SNAPCT, STATE, (ulong)ctx->state );
282 0 : }
283 :
284 : static void
285 : snapshot_path_gui_publish( fd_snapct_tile_t * ctx,
286 : fd_stem_context_t * stem,
287 : char const * path,
288 0 : int is_full ) {
289 : /* The messages below cannot be obtained directly from metrics. */
290 0 : fd_snapct_update_t * out = fd_chunk_to_laddr( ctx->out_gui.mem, ctx->out_gui.chunk );
291 0 : FD_TEST( fd_cstr_printf_check( out->read_path, PATH_MAX, NULL, "%s", path ) );
292 0 : out->is_download = 0;
293 0 : out->type = fd_int_if( is_full, FD_SNAPCT_SNAPSHOT_TYPE_FULL, FD_SNAPCT_SNAPSHOT_TYPE_INCREMENTAL );
294 0 : fd_stem_publish( stem, ctx->out_gui.idx, 0UL, ctx->out_gui.chunk, sizeof(fd_snapct_update_t) , 0UL, 0UL, 0UL );
295 0 : ctx->out_gui.chunk = fd_dcache_compact_next( ctx->out_gui.chunk, sizeof(fd_snapct_update_t), ctx->out_gui.chunk0, ctx->out_gui.wmark );
296 0 : }
297 :
298 : static void
299 0 : predict_incremental( fd_snapct_tile_t * ctx ) {
300 0 : if( FD_UNLIKELY( !ctx->config.incremental_snapshots ) ) return;
301 0 : if( FD_UNLIKELY( ctx->predicted_incremental.full_slot==FD_SSPEER_SLOT_UNKNOWN ) ) return;
302 :
303 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 1, ctx->predicted_incremental.full_slot );
304 :
305 0 : if( FD_LIKELY( best.addr.l ) ) {
306 0 : if( FD_UNLIKELY( ctx->predicted_incremental.slot!=best.incr_slot ) ) {
307 0 : ctx->predicted_incremental.slot = best.incr_slot;
308 0 : ctx->predicted_incremental.pending = 1;
309 0 : }
310 0 : }
311 0 : }
312 :
313 : static void
314 : on_resolve( void * _ctx,
315 : fd_sspeer_key_t const * key,
316 : fd_ip4_port_t addr,
317 : ulong full_slot,
318 : ulong incr_slot,
319 : uchar full_hash[ FD_HASH_FOOTPRINT ],
320 0 : uchar incr_hash[ FD_HASH_FOOTPRINT ] ) {
321 0 : fd_snapct_tile_t * ctx = (fd_snapct_tile_t *)_ctx;
322 :
323 0 : if( FD_UNLIKELY( full_slot!=FD_SSPEER_SLOT_UNKNOWN && full_slot>=FD_SSPEER_PLAUSIBLE_MAX_SLOT ) ) return;
324 0 : if( FD_UNLIKELY( incr_slot!=FD_SSPEER_SLOT_UNKNOWN && incr_slot>=FD_SSPEER_PLAUSIBLE_MAX_SLOT ) ) return;
325 :
326 : /* Do not update peers that have been permanently blacklisted. */
327 0 : if( FD_UNLIKELY( key && blacklist_map_ele_query( ctx->blacklist_map, key, NULL, ctx->blacklist_pool ) ) ) return;
328 : /* Do not re-add peers whose addr is temporarily banned by ssping. */
329 0 : if( FD_UNLIKELY( fd_ssping_is_invalidated( ctx->ssping, addr ) ) ) return;
330 :
331 : /* add() handles both new and existing peers, so peers removed from
332 : the selector during a previous blacklist, timeout, or failed
333 : re-resolve are re-added with the freshly resolved data. */
334 0 : ulong score = fd_sspeer_selector_add( ctx->selector, key, addr, FD_SSPEER_LATENCY_UNKNOWN,
335 0 : full_slot, incr_slot, full_hash, incr_hash );
336 0 : if( FD_UNLIKELY( score==FD_SSPEER_SCORE_INVALID ) ) {
337 0 : if( FD_UNLIKELY( key==NULL ) ) {
338 0 : FD_LOG_DEBUG(( "selector add on resolve failed for NULL peer key" ));
339 0 : } else {
340 0 : if( FD_UNLIKELY( !key->is_url ) ) {
341 0 : FD_BASE58_ENCODE_32_BYTES( key->pubkey->key, pubkey_b58 );
342 0 : FD_LOG_DEBUG(( "selector add on resolve failed for peer with pubkey %s", pubkey_b58 ));
343 0 : } else {
344 0 : FD_LOG_DEBUG(( "selector add on resolve failed for peer %s with addr " FD_IP4_ADDR_FMT ":%hu", key->url.hostname,
345 0 : FD_IP4_ADDR_FMT_ARGS( key->url.resolved_addr.addr ), fd_ushort_bswap( key->url.resolved_addr.port ) ));
346 0 : }
347 0 : }
348 0 : }
349 0 : fd_sspeer_selector_process_cluster_slot( ctx->selector );
350 0 : predict_incremental( ctx );
351 0 : }
352 :
353 : static void
354 : on_ping( void * _ctx,
355 : fd_ip4_port_t addr,
356 0 : ulong latency ) {
357 0 : fd_snapct_tile_t * ctx = (fd_snapct_tile_t *)_ctx;
358 :
359 0 : ulong cnt = fd_sspeer_selector_update_on_ping( ctx->selector, addr, latency );
360 0 : if( FD_UNLIKELY( !cnt ) ) {
361 : /* The update may fail in normal operation, e.g. after a peer has
362 : been removed from the selector. The log level is set to a
363 : minimum accordingly. */
364 0 : FD_LOG_DEBUG(( "selector update on ping did not find address " FD_IP4_ADDR_FMT ":%hu",
365 0 : FD_IP4_ADDR_FMT_ARGS( addr.addr ), fd_ushort_bswap( addr.port ) ));
366 0 : }
367 0 : predict_incremental( ctx );
368 0 : }
369 :
370 : static void
371 : on_snapshot_hash( fd_snapct_tile_t * ctx,
372 : fd_sspeer_key_t const * key,
373 : fd_ip4_port_t addr,
374 0 : fd_gossip_update_message_t const * msg ) {
375 0 : ulong full_slot = msg->snapshot_hashes->full_slot;
376 0 : ulong incr_slot = FD_SSPEER_SLOT_UNKNOWN;
377 0 : uchar const * incr_hash = NULL;
378 :
379 0 : for( ulong i=0UL; i<msg->snapshot_hashes->incremental_len; i++ ) {
380 0 : if( FD_LIKELY( !incr_hash || msg->snapshot_hashes->incremental[ i ].slot>incr_slot ) ) {
381 0 : incr_slot = msg->snapshot_hashes->incremental[ i ].slot;
382 0 : incr_hash = msg->snapshot_hashes->incremental[ i ].hash;
383 0 : }
384 0 : }
385 :
386 0 : if( FD_UNLIKELY( full_slot>=FD_SSPEER_PLAUSIBLE_MAX_SLOT ) ) return;
387 0 : if( FD_UNLIKELY( incr_slot!=FD_SSPEER_SLOT_UNKNOWN && incr_slot>=FD_SSPEER_PLAUSIBLE_MAX_SLOT ) ) return;
388 :
389 0 : if( FD_UNLIKELY( !addr.l ) ) {
390 : /* A peer that does not advertise an rpc_addr cannot be added to
391 : the selector: if previously added, remove it. The remove
392 : operation becomes a no-op if the peer is not found. */
393 0 : fd_sspeer_selector_remove( ctx->selector, key );
394 0 : return;
395 0 : }
396 : /* Do not re-add peers that have been permanently blacklisted. */
397 0 : if( FD_UNLIKELY( blacklist_map_ele_query( ctx->blacklist_map, key, NULL, ctx->blacklist_pool ) ) ) return;
398 : /* Do not re-add peers whose addr is temporarily banned by ssping. */
399 0 : if( FD_UNLIKELY( fd_ssping_is_invalidated( ctx->ssping, addr ) ) ) return;
400 : /* The add may fail due to capacity/pool exhaustion. The cluster
401 : slot is recomputed from tracked peers only, so a failed add will
402 : not influence the cluster slot. */
403 0 : fd_sspeer_selector_add( ctx->selector, key, addr, FD_SSPEER_LATENCY_UNKNOWN,
404 0 : full_slot, incr_slot,
405 0 : msg->snapshot_hashes->full_hash, incr_hash );
406 0 : fd_sspeer_selector_process_cluster_slot( ctx->selector );
407 0 : predict_incremental( ctx );
408 0 : }
409 :
410 : static void
411 : send_expected_slot( fd_snapct_tile_t * ctx,
412 : fd_stem_context_t * stem,
413 0 : ulong slot ) {
414 0 : uint tsorig; uint tspub;
415 0 : fd_ssmsg_slot_to_frag( slot, &tsorig, &tspub );
416 0 : fd_stem_publish( stem, ctx->out_rp.idx, FD_SSMSG_EXPECTED_SLOT, 0UL, 0UL, 0UL, tsorig, tspub );
417 0 : }
418 :
419 : /* snapshot_pool_select picks a snapshot file to overwrite with newly
420 : downloaded data. pool[return value] gives the inode to recycle. */
421 :
422 : static uint
423 : snapshot_pool_select( fd_backup_inode_t const * pool,
424 : uint slot0,
425 0 : uint slot1 ) {
426 0 : uint selected = UINT_MAX;
427 0 : ulong oldest = ULONG_MAX;
428 0 : for( uint i=slot0; i<slot1; i++ ) {
429 0 : if( FD_UNLIKELY( pool[ i ].full_slot==ULONG_MAX ) ) return i;
430 0 : ulong slot = pool[ i ].incr_slot==ULONG_MAX ? pool[ i ].full_slot : pool[ i ].incr_slot;
431 0 : if( slot<oldest ) {
432 0 : selected = i;
433 0 : oldest = slot;
434 0 : }
435 0 : }
436 0 : return selected;
437 0 : }
438 :
439 : /* snapshot_output_prepare recycles a 'partial' snapshot file or
440 : overwrites an older snapshot. */
441 :
442 : static void
443 : snapshot_output_prepare( fd_snapct_tile_t * ctx,
444 0 : int full ) {
445 0 : int fd = full ? ctx->local_out.full_snapshot_fd : ctx->local_out.incremental_snapshot_fd;
446 0 : char * name = full ? ctx->local_out.full_snapshot_name : ctx->local_out.incremental_snapshot_name;
447 :
448 0 : FD_TEST( ctx->local_out.dir_fd!=-1 && fd!=-1 );
449 0 : FD_TEST( fd>=FD_SNAP_FD( 0U ) && fd<FD_SNAP_FD( FD_SNAP_MAX ) );
450 :
451 0 : char partial_name[ FD_SNAP_NAME_MAX ];
452 0 : fd_snap_pool_partial_name( partial_name, (uint)(fd-FD_SNAP_FD( 0U )) );
453 0 : if( strcmp( name, partial_name ) ) {
454 0 : char const * local_path = full ? ctx->local_in.full_snapshot_path : ctx->local_in.incremental_snapshot_path;
455 0 : char const * local_name = strrchr( local_path, '/' );
456 0 : local_name = local_name ? local_name+1 : local_path;
457 0 : if( !strcmp( name, local_name ) ) {
458 0 : if( full ) ctx->local_in.full_snapshot_slot = ULONG_MAX;
459 0 : else ctx->local_in.incremental_snapshot_slot = ULONG_MAX;
460 0 : }
461 :
462 0 : if( FD_UNLIKELY( -1==renameat( ctx->local_out.dir_fd, name, ctx->local_out.dir_fd, partial_name ) ) )
463 0 : FD_LOG_ERR(( "renameat(%s, %s) failed (%i-%s)", name, partial_name, errno, fd_io_strerror( errno ) ));
464 0 : fd_cstr_ncpy( name, partial_name, FD_SNAP_NAME_MAX );
465 0 : }
466 :
467 0 : if( FD_UNLIKELY( -1==ftruncate( fd, 0UL ) ) )
468 0 : FD_LOG_ERR(( "ftruncate(%s) failed (%i-%s)", name, errno, fd_io_strerror( errno ) ));
469 0 : if( FD_UNLIKELY( -1==lseek( fd, 0L, SEEK_SET ) ) )
470 0 : FD_LOG_ERR(( "lseek(%s) failed (%i-%s)", name, errno, fd_io_strerror( errno ) ));
471 0 : }
472 :
473 : static void
474 0 : rename_full_snapshot( fd_snapct_tile_t * ctx ) {
475 0 : FD_TEST( -1!=ctx->local_out.dir_fd );
476 :
477 0 : if( FD_LIKELY( -1!=ctx->local_out.full_snapshot_fd && ctx->http_full_snapshot_name[ 0 ]!='\0' ) ) {
478 0 : int err = renameat2( ctx->local_out.dir_fd, ctx->local_out.full_snapshot_name,
479 0 : ctx->local_out.dir_fd, ctx->http_full_snapshot_name, RENAME_NOREPLACE );
480 0 : if( FD_UNLIKELY( err && errno!=EEXIST ) )
481 0 : FD_LOG_ERR(( "renameat2() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
482 0 : if( FD_LIKELY( !err ) )
483 0 : fd_cstr_ncpy( ctx->local_out.full_snapshot_name, ctx->http_full_snapshot_name, FD_SNAP_NAME_MAX );
484 0 : }
485 0 : }
486 :
487 : static void
488 0 : rename_incr_snapshot( fd_snapct_tile_t * ctx ) {
489 0 : FD_TEST( -1!=ctx->local_out.dir_fd );
490 :
491 0 : if( FD_LIKELY( -1!=ctx->local_out.incremental_snapshot_fd && ctx->http_incr_snapshot_name[ 0 ]!='\0' ) ) {
492 0 : int err = renameat2( ctx->local_out.dir_fd, ctx->local_out.incremental_snapshot_name,
493 0 : ctx->local_out.dir_fd, ctx->http_incr_snapshot_name, RENAME_NOREPLACE );
494 0 : if( FD_UNLIKELY( err && errno!=EEXIST ) )
495 0 : FD_LOG_ERR(( "renameat2() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
496 0 : if( FD_LIKELY( !err ) )
497 0 : fd_cstr_ncpy( ctx->local_out.incremental_snapshot_name, ctx->http_incr_snapshot_name, FD_SNAP_NAME_MAX );
498 0 : }
499 0 : }
500 :
501 : static ulong
502 : rlimit_file_cnt( fd_topo_t const * topo FD_PARAM_UNUSED,
503 0 : fd_topo_tile_t const * tile ) {
504 0 : ulong cnt = 1UL + /* stderr */
505 0 : 1UL + /* logfile */
506 0 : 1UL; /* boot control pipe */
507 0 : if( download_enabled( tile ) ) {
508 0 : cnt += FD_SSPING_FD_CNT + /* ssping sockets */
509 0 : 2UL + /* dirfd + full snapshot pool fd */
510 0 : tile->snapct.sources.servers_cnt; /* http resolver peer full sockets */
511 0 : if( tile->snapct.incremental_snapshots ) {
512 0 : cnt += 1UL + /* incremental snapshot pool fd */
513 0 : tile->snapct.sources.servers_cnt; /* http resolver peer incr sockets */
514 0 : }
515 0 : }
516 0 : return cnt;
517 0 : }
518 :
519 : static ulong
520 : populate_allowed_seccomp( fd_topo_t const * topo,
521 : fd_topo_tile_t const * tile,
522 : ulong out_cnt,
523 0 : struct sock_filter * out ) {
524 :
525 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
526 :
527 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
528 0 : fd_snapct_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapct_tile_t), sizeof(fd_snapct_tile_t) );
529 :
530 0 : int min_ping_fd = INT_MAX;
531 0 : int max_ping_fd = 0;
532 0 : if( download_enabled( tile ) ) {
533 0 : min_ping_fd = FD_SSPING_FD_MIN;
534 0 : max_ping_fd = FD_SSPING_FD_MIN + (int)FD_SSPING_FD_CNT - 1;
535 0 : }
536 :
537 0 : populate_sock_filter_policy_fd_snapct_tile( out_cnt, out, (uint)fd_log_private_logfile_fd(), (uint)ctx->local_out.dir_fd, (uint)ctx->local_out.full_snapshot_fd, (uint)ctx->local_out.incremental_snapshot_fd, (uint)min_ping_fd, (uint)max_ping_fd, (uint)ctx->netdb_fds->etc_hosts, (uint)ctx->netdb_fds->etc_resolv_conf );
538 0 : return sock_filter_policy_fd_snapct_tile_instr_cnt;
539 0 : }
540 :
541 : static ulong
542 : populate_allowed_fds( fd_topo_t const * topo,
543 : fd_topo_tile_t const * tile,
544 : ulong out_fds_cnt,
545 0 : int * out_fds ) {
546 0 : if( FD_UNLIKELY( out_fds_cnt<7UL ) ) FD_LOG_ERR(( "out_fds_cnt %lu is too small", out_fds_cnt ));
547 :
548 0 : ulong out_cnt = 0;
549 0 : out_fds[ out_cnt++ ] = 2UL; /* stderr */
550 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) ) {
551 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
552 0 : }
553 :
554 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
555 :
556 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
557 0 : fd_snapct_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapct_tile_t), sizeof(fd_snapct_tile_t) );
558 0 : if( FD_LIKELY( -1!=ctx->local_out.dir_fd ) ) out_fds[ out_cnt++ ] = ctx->local_out.dir_fd;
559 0 : if( FD_LIKELY( -1!=ctx->local_out.full_snapshot_fd ) ) out_fds[ out_cnt++ ] = ctx->local_out.full_snapshot_fd;
560 0 : if( FD_LIKELY( -1!=ctx->local_out.incremental_snapshot_fd ) ) out_fds[ out_cnt++ ] = ctx->local_out.incremental_snapshot_fd;
561 0 : if( FD_LIKELY( -1!=ctx->netdb_fds->etc_hosts ) ) out_fds[ out_cnt++ ] = ctx->netdb_fds->etc_hosts;
562 0 : if( FD_LIKELY( -1!=ctx->netdb_fds->etc_resolv_conf ) ) out_fds[ out_cnt++ ] = ctx->netdb_fds->etc_resolv_conf;
563 0 : if( FD_LIKELY( download_enabled( tile ) ) ) {
564 0 : if( FD_UNLIKELY( out_cnt+FD_SSPING_FD_CNT > out_fds_cnt ) ) {
565 0 : FD_LOG_ERR(( "out_fds_cnt %lu must be at least %lu", out_fds_cnt, out_cnt + FD_SSPING_FD_CNT ));
566 0 : }
567 0 : for( int i=FD_SSPING_FD_MIN; i<FD_SSPING_FD_MIN+(int)FD_SSPING_FD_CNT; i++ ) {
568 0 : out_fds[ out_cnt++ ] = i;
569 0 : }
570 0 : }
571 :
572 0 : return out_cnt;
573 0 : }
574 :
575 : static void
576 : init_load( fd_snapct_tile_t * ctx,
577 : fd_stem_context_t * stem,
578 : int full,
579 0 : int file ) {
580 0 : ctx->snapshot_start_timestamp_ns = fd_log_wallclock();
581 0 : fd_ssctrl_init_t * out = fd_chunk_to_laddr( ctx->out_ld.mem, ctx->out_ld.chunk );
582 0 : out->file = file;
583 0 : out->zstd = !file || (full ? ctx->local_in.full_snapshot_zstd : ctx->local_in.incremental_snapshot_zstd);
584 0 : if( file ) {
585 0 : out->slot = full ? ctx->local_in.full_snapshot_slot : ctx->local_in.incremental_snapshot_slot;
586 0 : if( full ) fd_memcpy( out->snapshot_hash, ctx->local_in.full_snapshot_hash, FD_HASH_FOOTPRINT );
587 0 : else fd_memcpy( out->snapshot_hash, ctx->local_in.incremental_snapshot_hash, FD_HASH_FOOTPRINT );
588 0 : } else {
589 0 : out->slot = full ? ctx->predicted_incremental.full_slot : ctx->predicted_incremental.slot;
590 0 : if( full ) fd_memcpy( out->snapshot_hash, ctx->peer.full_hash, FD_HASH_FOOTPRINT );
591 0 : else fd_memcpy( out->snapshot_hash, ctx->peer.incr_hash, FD_HASH_FOOTPRINT );
592 0 : }
593 :
594 0 : out->is_redirect = !file; /* always use redirect for HTTP downloads */
595 :
596 0 : if( file ) out->file_sz = full ? ctx->local_in.full_snapshot_size : ctx->local_in.incremental_snapshot_size;
597 0 : else out->file_sz = 0UL;
598 :
599 0 : if( !file ) {
600 0 : out->addr = ctx->peer.addr;
601 0 : if( full ) {
602 0 : FD_TEST( fd_cstr_printf_check( out->path, PATH_MAX, &out->path_len, "/snapshot.tar.bz2" ) );
603 0 : FD_TEST( fd_cstr_printf_check( ctx->http_full_snapshot_name, PATH_MAX, NULL, "snapshot.tar.bz2" ) );
604 0 : } else {
605 0 : FD_TEST( fd_cstr_printf_check( out->path, PATH_MAX, &out->path_len, "/incremental-snapshot.tar.bz2" ) );
606 0 : FD_TEST( fd_cstr_printf_check( ctx->http_incr_snapshot_name, PATH_MAX, NULL, "incremental-snapshot.tar.bz2" ) );
607 0 : }
608 :
609 0 : out->is_https = 0; /* if not found in the config list, it's not https */
610 0 : out->hostname[0] = '\0'; /* .. and it doesn't have a hostname either. */
611 0 : for( ulong i=0UL; i<ctx->resolved_servers_cnt; i++ ) {
612 0 : if( FD_UNLIKELY( ctx->peer.addr.l==ctx->resolved_servers[ i ].addr.l ) ) {
613 0 : fd_cstr_ncpy( out->hostname, ctx->resolved_servers[ i ].hostname, sizeof(out->hostname) );
614 0 : out->is_https = ctx->resolved_servers[ i ].is_https;
615 0 : break;
616 0 : }
617 0 : }
618 0 : }
619 0 : fd_stem_publish( stem, ctx->out_ld.idx, full ? FD_SNAPSHOT_MSG_CTRL_INIT_FULL : FD_SNAPSHOT_MSG_CTRL_INIT_INCR, ctx->out_ld.chunk, sizeof(fd_ssctrl_init_t), 0UL, 0UL, 0UL );
620 0 : ctx->out_ld.chunk = fd_dcache_compact_next( ctx->out_ld.chunk, sizeof(fd_ssctrl_init_t), ctx->out_ld.chunk0, ctx->out_ld.wmark );
621 0 : ctx->flush_ack = 0;
622 0 : ctx->load_complete = 0;
623 :
624 0 : if( !file ) snapshot_output_prepare( ctx, full );
625 :
626 : /* Clear stale http_*_snapshot_name for file loads before rename
627 : functions run. GUI publish is deferred to the META handler. */
628 0 : if( file ) {
629 0 : if( full ) fd_cstr_fini( ctx->http_full_snapshot_name );
630 0 : else fd_cstr_fini( ctx->http_incr_snapshot_name );
631 0 : }
632 0 : }
633 :
634 : static void
635 : log_download( fd_snapct_tile_t * ctx,
636 : int full,
637 : fd_ip4_port_t addr,
638 0 : ulong slot ) {
639 0 : for( gossip_ci_map_iter_t iter = gossip_ci_map_iter_init( ctx->gossip.ci_map, ctx->gossip.ci_table );
640 0 : !gossip_ci_map_iter_done( iter, ctx->gossip.ci_map, ctx->gossip.ci_table );
641 0 : iter = gossip_ci_map_iter_next( iter, ctx->gossip.ci_map, ctx->gossip.ci_table ) ) {
642 0 : gossip_ci_entry_t const * ci_entry = gossip_ci_map_iter_ele_const( iter, ctx->gossip.ci_map, ctx->gossip.ci_table );
643 0 : if( ci_entry->rpc_addr.l==addr.l ) {
644 0 : FD_TEST( ci_entry->allowed );
645 :
646 0 : int is_entrypoint = 0;
647 0 : for( ulong i=0UL; i<ctx->resolved_entrypoints_cnt; i++ ) {
648 0 : if( FD_UNLIKELY( ctx->resolved_entrypoints[ i ].addr==addr.addr ) ) { is_entrypoint = 1; break; }
649 0 : }
650 0 : char const * kind = ctx->config.sources.gossip.allow_any ? "untrusted gossip peer" : "trusted gossip peer";
651 0 : if( FD_UNLIKELY( is_entrypoint ) ) kind = "entrypoint gossip peer";
652 :
653 0 : FD_BASE58_ENCODE_32_BYTES( ci_entry->pubkey.uc, pubkey_b58 );
654 0 : FD_LOG_NOTICE(( "downloading %s snapshot from %s %s%s%s",
655 0 : full ? "full" : "incremental", kind, fd_log_style_bold(), pubkey_b58, fd_log_style_normal() ));
656 0 : FD_LOG_INFO(( "downloading %s snapshot at slot %lu from %s %s at " FD_IP4_ADDR_FMT ":%hu",
657 0 : full ? "full" : "incremental", slot, kind, pubkey_b58,
658 0 : FD_IP4_ADDR_FMT_ARGS( addr.addr ), fd_ushort_bswap( addr.port ) ));
659 0 : return;
660 0 : }
661 0 : }
662 :
663 0 : for( ulong i=0UL; i<ctx->resolved_servers_cnt; i++ ) {
664 0 : if( addr.l==ctx->resolved_servers[ i ].addr.l ) {
665 0 : char const * scheme = ctx->resolved_servers[ i ].is_https ? "https" : "http";
666 0 : if( ctx->resolved_servers[ i ].hostname[ 0 ] ) {
667 0 : FD_LOG_NOTICE(( "downloading %s snapshot from %s%s://%s:%hu%s",
668 0 : full ? "full" : "incremental",
669 0 : fd_log_style_bold(), scheme, ctx->resolved_servers[ i ].hostname, fd_ushort_bswap( addr.port ), fd_log_style_normal() ));
670 0 : FD_LOG_INFO(( "downloading %s snapshot at slot %lu from configured server with index %lu at %s://%s:%hu",
671 0 : full ? "full" : "incremental", slot, i,
672 0 : scheme, ctx->resolved_servers[ i ].hostname, fd_ushort_bswap( addr.port ) ));
673 0 : } else {
674 0 : FD_LOG_NOTICE(( "downloading %s snapshot from %s%s://" FD_IP4_ADDR_FMT ":%hu%s",
675 0 : full ? "full" : "incremental",
676 0 : fd_log_style_bold(), scheme, FD_IP4_ADDR_FMT_ARGS( addr.addr ), fd_ushort_bswap( addr.port ), fd_log_style_normal() ));
677 0 : FD_LOG_INFO(( "downloading %s snapshot at slot %lu from configured server with index %lu at %s://" FD_IP4_ADDR_FMT ":%hu",
678 0 : full ? "full" : "incremental", slot, i,
679 0 : scheme, FD_IP4_ADDR_FMT_ARGS( addr.addr ), fd_ushort_bswap( addr.port ) ));
680 0 : }
681 0 : return;
682 0 : }
683 0 : }
684 :
685 0 : FD_TEST( 0 ); /* should not be possible */
686 0 : }
687 :
688 : static void
689 : log_completion( fd_snapct_tile_t * ctx,
690 0 : int full ) {
691 0 : double elapsed = (double)(fd_log_wallclock() - ctx->snapshot_start_timestamp_ns) / 1e9;
692 0 : FD_LOG_INFO(( "%s snapshot load completed in %.3f seconds", full ? "full" : "incremental", elapsed ));
693 0 : }
694 :
695 : /* Blacklist the current peer: invalidate in ssping, remove from the
696 : selector, and add to the dedicated blacklist map. Skips the insert
697 : if the peer is already blacklisted (dedup). Warns if the pool is
698 : full (the ssping ban still provides temporary protection). */
699 : static void
700 24 : blacklist_peer( fd_snapct_tile_t * ctx ) {
701 24 : fd_ssping_invalidate( ctx->ssping, ctx->peer.addr, fd_log_wallclock() );
702 24 : fd_sspeer_selector_remove_by_addr( ctx->selector, ctx->peer.addr );
703 24 : fd_sspeer_selector_process_cluster_slot( ctx->selector );
704 24 : if( FD_UNLIKELY( blacklist_map_ele_query( ctx->blacklist_map, &ctx->peer.key, NULL, ctx->blacklist_pool ) ) ) return;
705 21 : if( FD_LIKELY( blacklist_pool_free( ctx->blacklist_pool ) ) ) {
706 18 : fd_sspeer_blacklist_entry_t * bl = blacklist_pool_ele_acquire( ctx->blacklist_pool );
707 18 : bl->key = ctx->peer.key;
708 18 : blacklist_map_ele_insert( ctx->blacklist_map, bl, ctx->blacklist_pool );
709 18 : if( FD_UNLIKELY( !ctx->peer.key.is_url ) ) {
710 18 : FD_BASE58_ENCODE_32_BYTES( ctx->peer.key.pubkey->uc, pubkey_b58 );
711 18 : FD_LOG_WARNING(( "permanently blacklisted peer identity %s", pubkey_b58 ));
712 18 : } else {
713 0 : FD_LOG_WARNING(( "permanently blacklisted peer identity %s",
714 0 : ctx->peer.key.url.hostname[0] ? ctx->peer.key.url.hostname : "(none)" ));
715 0 : }
716 18 : } else {
717 3 : FD_LOG_WARNING(( "blacklist pool full, peer banned via ssping only" ));
718 3 : }
719 21 : }
720 :
721 0 : #define DNS_RETRY_NANOS (15L*1000L*1000L*1000L)
722 0 : #define DNS_REQ_ID_ENTRYPOINT (0x100UL) /* req_id: server idx, or this bit + entrypoint idx */
723 :
724 : static void
725 : dns_queue( fd_snapct_tile_t * ctx,
726 0 : long now ) {
727 0 : for( ulong i=0UL; i<ctx->config.sources.servers_cnt; i++ ) {
728 0 : if( FD_LIKELY( ctx->dns_servers[ i ].resolved || !ctx->dns_servers[ i ].retry_nanos || ctx->dns_servers[ i ].retry_nanos>now ) ) continue;
729 0 : if( FD_UNLIKELY( fd_adns_resolve( ctx->adns, ctx->dns_servers[ i ].hostname, i ) ) ) break;
730 0 : ctx->dns_servers[ i ].retry_nanos = 0L; /* in flight */
731 0 : }
732 0 : for( ulong i=0UL; i<ctx->config.entrypoints_cnt; i++ ) {
733 0 : if( FD_LIKELY( ctx->dns_entrypoints[ i ].resolved || !ctx->dns_entrypoints[ i ].retry_nanos || ctx->dns_entrypoints[ i ].retry_nanos>now ) ) continue;
734 0 : if( FD_UNLIKELY( fd_adns_resolve( ctx->adns, ctx->dns_entrypoints[ i ].hostname, DNS_REQ_ID_ENTRYPOINT|i ) ) ) break;
735 0 : ctx->dns_entrypoints[ i ].retry_nanos = 0L;
736 0 : }
737 0 : }
738 :
739 : static void
740 : dns_advance( fd_snapct_tile_t * ctx,
741 0 : long now ) {
742 0 : dns_queue( ctx, now );
743 :
744 0 : fd_adns_result_t res[ 1 ];
745 0 : while( fd_adns_advance( ctx->adns, now, res ) ) {
746 0 : if( FD_UNLIKELY( res->req_id & DNS_REQ_ID_ENTRYPOINT ) ) {
747 0 : ulong i = res->req_id & ~DNS_REQ_ID_ENTRYPOINT;
748 0 : if( FD_UNLIKELY( res->err ) ) {
749 0 : FD_LOG_WARNING(( "could not resolve [gossip.entrypoints] entry \"%s\" (%s), retrying", ctx->config.entrypoints[ i ], fd_gai_strerror( res->err ) ));
750 0 : ctx->dns_entrypoints[ i ].retry_nanos = now+DNS_RETRY_NANOS;
751 0 : continue;
752 0 : }
753 0 : ctx->dns_entrypoints[ i ].resolved = 1;
754 0 : ctx->resolved_entrypoints[ ctx->resolved_entrypoints_cnt++ ] = (fd_ip4_port_t){ .addr=res->addrs[ 0 ], .port=ctx->dns_entrypoints[ i ].port };
755 0 : continue;
756 0 : }
757 :
758 0 : ulong i = res->req_id;
759 0 : if( FD_UNLIKELY( res->err ) ) {
760 0 : FD_LOG_WARNING(( "could not resolve [snapshots.sources.servers] entry \"%s\" (%s), retrying", ctx->config.sources.servers[ i ], fd_gai_strerror( res->err ) ));
761 0 : ctx->dns_servers[ i ].retry_nanos = now+DNS_RETRY_NANOS;
762 0 : continue;
763 0 : }
764 0 : ctx->dns_servers[ i ].resolved = 1;
765 :
766 0 : ulong addr_cnt = fd_ulong_min( res->addr_cnt, FD_TOPO_MAX_RESOLVED_ADDRS );
767 0 : for( ulong j=0UL; j<addr_cnt; j++ ) {
768 0 : ulong k = ctx->resolved_servers_cnt++;
769 0 : ctx->resolved_servers[ k ].addr = (fd_ip4_port_t){ .addr=res->addrs[ j ], .port=ctx->dns_servers[ i ].port };
770 0 : ctx->resolved_servers[ k ].is_https = ctx->dns_servers[ i ].is_https;
771 0 : fd_cstr_ncpy( ctx->resolved_servers[ k ].hostname, ctx->dns_servers[ i ].hostname, sizeof(ctx->resolved_servers[ k ].hostname) );
772 :
773 : /* The peer needs to be added to resolver and to selector.
774 : Only if this succeeds, add the peer to ssping list. */
775 0 : if( FD_LIKELY( !fd_http_resolver_add( ctx->ssresolver,
776 0 : ctx->resolved_servers[ k ].addr,
777 0 : ctx->resolved_servers[ k ].hostname,
778 0 : ctx->resolved_servers[ k ].is_https,
779 0 : ctx->selector ) ) ) {
780 0 : fd_ssping_add( ctx->ssping, ctx->resolved_servers[ k ].addr );
781 0 : }
782 0 : }
783 0 : }
784 0 : }
785 :
786 : static void
787 : after_credit( fd_snapct_tile_t * ctx,
788 : fd_stem_context_t * stem,
789 : int * opt_poll_in FD_PARAM_UNUSED,
790 0 : int * charge_busy FD_PARAM_UNUSED ) {
791 0 : long now = fd_log_wallclock();
792 :
793 0 : if( FD_LIKELY( ctx->adns ) ) dns_advance( ctx, now );
794 0 : if( FD_LIKELY( ctx->ssping ) ) fd_ssping_advance( ctx->ssping, now, ctx->selector );
795 0 : if( FD_LIKELY( ctx->ssresolver ) ) fd_http_resolver_advance( ctx->ssresolver, now, ctx->selector );
796 :
797 : /* Advances above may remove peers, making cluster_slot dirty.
798 : Recompute so best() calls below use up-to-date scores.
799 : No-op when the dirty flag is not set internally. */
800 0 : fd_sspeer_selector_process_cluster_slot( ctx->selector );
801 :
802 : /* send an expected slot message as the predicted incremental
803 : could have changed as a result of the pinger, resolver, or from
804 : processing gossip frags in gossip_frag. */
805 0 : if( FD_LIKELY( ctx->predicted_incremental.pending ) ) {
806 0 : send_expected_slot( ctx, stem, ctx->predicted_incremental.slot );
807 0 : ctx->predicted_incremental.pending = 0;
808 0 : }
809 :
810 : /* Note: All state transitions should occur within this switch
811 : statement to make it easier to reason about the state management. */
812 :
813 0 : switch ( ctx->state ) {
814 :
815 : /* ============================================================== */
816 0 : case FD_SNAPCT_STATE_INIT: {
817 0 : if( FD_UNLIKELY( !ctx->download_enabled ) ) {
818 0 : ulong local_slot = ctx->config.incremental_snapshots ? ctx->local_in.incremental_snapshot_slot : ctx->local_in.full_snapshot_slot;
819 0 : send_expected_slot( ctx, stem, local_slot );
820 0 : FD_LOG_NOTICE(( "reading full snapshot from file %s%s%s", fd_log_style_dim(), ctx->local_in.full_snapshot_path, fd_log_style_normal() ));
821 0 : FD_LOG_INFO(( "reading full snapshot at slot %lu from local file `%s`", ctx->local_in.full_snapshot_slot, ctx->local_in.full_snapshot_path ));
822 0 : ctx->predicted_incremental.full_slot = ctx->local_in.full_snapshot_slot;
823 0 : ctx->state = FD_SNAPCT_STATE_READING_FULL_FILE;
824 0 : init_load( ctx, stem, 1, 1 );
825 0 : break;
826 0 : }
827 0 : ctx->deadline_nanos = now+ctx->config.wait_for_peers_timeout_nanos;
828 0 : ctx->state = FD_SNAPCT_STATE_WAITING_FOR_PEERS;
829 0 : break;
830 0 : }
831 :
832 : /* ============================================================== */
833 0 : case FD_SNAPCT_STATE_WAITING_FOR_PEERS: {
834 0 : if( FD_UNLIKELY( now>ctx->deadline_nanos ) ) FD_LOG_ERR(( "timed out waiting for peers." ));
835 :
836 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 0, FD_SSPEER_SLOT_UNKNOWN );
837 0 : if( FD_LIKELY( best.addr.l ) ) {
838 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS;
839 0 : ctx->deadline_nanos = now+FD_SNAPCT_COLLECTING_PEERS_TIMEOUT;
840 0 : }
841 0 : break;
842 0 : }
843 :
844 : /* ============================================================== */
845 0 : case FD_SNAPCT_STATE_WAITING_FOR_PEERS_INCREMENTAL: {
846 0 : if( FD_UNLIKELY( now>ctx->deadline_nanos ) ) FD_LOG_ERR(( "timed out waiting for incremental snapshot peers." ));
847 :
848 0 : FD_TEST( ctx->predicted_incremental.full_slot!=FD_SSPEER_SLOT_UNKNOWN );
849 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 1, ctx->predicted_incremental.full_slot );
850 0 : if( FD_LIKELY( best.addr.l ) ) {
851 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL;
852 0 : ctx->deadline_nanos = now;
853 0 : }
854 0 : break;
855 0 : }
856 :
857 : /* ============================================================== */
858 0 : case FD_SNAPCT_STATE_COLLECTING_PEERS: {
859 0 : if( FD_UNLIKELY( !ctx->gossip.saturated && now<ctx->deadline_nanos ) ) break;
860 :
861 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 0, FD_SSPEER_SLOT_UNKNOWN );
862 0 : if( FD_UNLIKELY( !best.addr.l ) ) {
863 0 : if( !ctx->gossip_enabled ) {
864 0 : FD_LOG_ERR(( "no peers are available and discovery of new peers via gossip is disabled. aborting." ));
865 0 : }
866 0 : ctx->deadline_nanos = now + ctx->config.wait_for_peers_timeout_nanos;
867 0 : ctx->state = FD_SNAPCT_STATE_WAITING_FOR_PEERS;
868 0 : break;
869 0 : }
870 :
871 0 : fd_sscluster_slot_t cluster = fd_sspeer_selector_cluster_slot( ctx->selector );
872 0 : if( FD_UNLIKELY( cluster.incremental==FD_SSPEER_SLOT_UNKNOWN && ctx->config.incremental_snapshots ) ) {
873 : /* We must have a cluster full slot to be in this state. */
874 0 : FD_TEST( cluster.full!=FD_SSPEER_SLOT_UNKNOWN );
875 : /* fall back to full snapshot only if the highest cluster slot
876 : is a full snapshot only */
877 0 : FD_LOG_WARNING(( "incremental snapshots were enabled via [snapshots.incremental_snapshots], but no incremental snapshot is available in the cluster. "
878 0 : "falling back to full snapshots only." ));
879 0 : ctx->config.incremental_snapshots = 0;
880 0 : }
881 :
882 0 : ulong cluster_slot = ctx->config.incremental_snapshots ? cluster.incremental : cluster.full;
883 :
884 : /* Determine the best effective slot achievable using the local
885 : full snapshot. When incrementals are disabled, the effective
886 : slot is the full snapshot slot itself. When enabled, it is the
887 : best incremental we can pair with the local full (either from
888 : a local file or downloaded from a peer). */
889 :
890 0 : ulong local_effective_slot = ULONG_MAX;
891 0 : if( FD_LIKELY( ctx->local_in.full_snapshot_slot!=ULONG_MAX ) ) {
892 0 : if( FD_LIKELY( ctx->config.incremental_snapshots ) ) {
893 0 : ulong local_incr = ctx->local_in.incremental_snapshot_slot;
894 0 : if( local_incr!=ULONG_MAX && local_incr>=fd_ulong_sat_sub( cluster_slot, ctx->config.sources.max_local_incremental_age ) ) {
895 0 : local_effective_slot = local_incr;
896 0 : } else {
897 0 : fd_sspeer_t best_incr = fd_sspeer_selector_best( ctx->selector, 1, ctx->local_in.full_snapshot_slot );
898 0 : if( FD_LIKELY( best_incr.addr.l ) ) {
899 0 : ctx->predicted_incremental.slot = best_incr.incr_slot;
900 0 : ctx->local_in.incremental_snapshot_slot = ULONG_MAX; /* don't use the local incremental */
901 0 : local_effective_slot = best_incr.incr_slot;
902 0 : }
903 0 : }
904 0 : } else {
905 0 : local_effective_slot = ctx->local_in.full_snapshot_slot;
906 0 : }
907 0 : }
908 :
909 0 : int can_use_local_full = local_effective_slot!=ULONG_MAX &&
910 0 : local_effective_slot>=fd_ulong_sat_sub( cluster_slot, ctx->config.sources.max_local_full_effective_age );
911 0 : if( FD_LIKELY( can_use_local_full ) ) {
912 0 : send_expected_slot( ctx, stem, local_effective_slot );
913 :
914 0 : FD_LOG_NOTICE(( "reading full snapshot from file %s%s%s", fd_log_style_dim(), ctx->local_in.full_snapshot_path, fd_log_style_normal() ));
915 0 : FD_LOG_INFO(( "reading full snapshot at slot %lu with cluster slot %lu from local file `%s`",
916 0 : ctx->local_in.full_snapshot_slot, cluster_slot, ctx->local_in.full_snapshot_path ));
917 0 : ctx->predicted_incremental.full_slot = ctx->local_in.full_snapshot_slot;
918 0 : ctx->state = FD_SNAPCT_STATE_READING_FULL_FILE;
919 0 : init_load( ctx, stem, 1, 1 );
920 0 : } else {
921 0 : if( FD_LIKELY( ctx->local_in.full_snapshot_slot!=ULONG_MAX ) ) {
922 0 : if( local_effective_slot==ULONG_MAX ) {
923 0 : if( ctx->local_in.incremental_snapshot_slot!=ULONG_MAX ) {
924 0 : FD_LOG_INFO(( "local full snapshot at slot %lu cannot be used because local incremental snapshot at slot %lu "
925 0 : "is too old and no downloadable incremental could be found (cluster slot %lu), downloading instead",
926 0 : ctx->local_in.full_snapshot_slot, ctx->local_in.incremental_snapshot_slot, cluster_slot ));
927 0 : } else {
928 0 : FD_LOG_NOTICE(( "local full snapshot at slot %lu cannot be used because no matching incremental snapshot "
929 0 : "could be found (cluster slot %lu), downloading instead",
930 0 : ctx->local_in.full_snapshot_slot, cluster_slot ));
931 0 : }
932 0 : } else {
933 0 : FD_LOG_NOTICE(( "local full snapshot at slot %lu (effective slot %lu) is too old for cluster slot %lu max age %u, downloading instead",
934 0 : ctx->local_in.full_snapshot_slot, local_effective_slot, cluster_slot, ctx->config.sources.max_local_full_effective_age ));
935 0 : }
936 0 : } else {
937 0 : FD_LOG_INFO(( "no local snapshot available, downloading from peer" ));
938 0 : }
939 :
940 0 : if( FD_UNLIKELY( !ctx->config.incremental_snapshots ) ) {
941 0 : send_expected_slot( ctx, stem, best.full_slot );
942 0 : } else {
943 0 : fd_sspeer_t best_incremental = fd_sspeer_selector_best( ctx->selector, 1, best.full_slot );
944 0 : if( FD_LIKELY( best_incremental.addr.l ) ) {
945 0 : ctx->predicted_incremental.slot = best_incremental.incr_slot;
946 0 : send_expected_slot( ctx, stem, best_incremental.incr_slot );
947 0 : }
948 0 : }
949 :
950 0 : ctx->peer = best;
951 0 : ctx->state = FD_SNAPCT_STATE_READING_FULL_HTTP;
952 0 : ctx->predicted_incremental.full_slot = best.full_slot;
953 0 : init_load( ctx, stem, 1, 0 );
954 0 : log_download( ctx, 1, best.addr, best.full_slot );
955 0 : }
956 0 : break;
957 0 : }
958 :
959 : /* ============================================================== */
960 0 : case FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL: {
961 0 : if( FD_UNLIKELY( now<ctx->deadline_nanos ) ) break;
962 :
963 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 1, ctx->predicted_incremental.full_slot );
964 0 : if( FD_UNLIKELY( !best.addr.l ) ) {
965 0 : if( !ctx->gossip_enabled ) {
966 0 : FD_LOG_ERR(( "no incremental snapshot peers are available and discovery of new peers via gossip is disabled. aborting." ));
967 0 : }
968 0 : ctx->deadline_nanos = now + ctx->config.wait_for_peers_timeout_nanos;
969 0 : ctx->state = FD_SNAPCT_STATE_WAITING_FOR_PEERS_INCREMENTAL;
970 0 : break;
971 0 : }
972 :
973 : /* decide whether to use the local incremental snapshot if one
974 : exists and is not too old, otherwise download a new incremental
975 : snapshot. */
976 0 : ulong cluster_slot = fd_sspeer_selector_cluster_slot( ctx->selector ).incremental;
977 0 : ulong local_slot = ctx->local_in.incremental_snapshot_slot;
978 0 : int local_too_old = local_slot<fd_ulong_sat_sub( cluster_slot, ctx->config.sources.max_local_incremental_age );
979 0 : if( FD_LIKELY( local_slot!=ULONG_MAX && !local_too_old ) ) {
980 0 : ctx->predicted_incremental.slot = local_slot;
981 0 : send_expected_slot( ctx, stem, local_slot );
982 :
983 0 : FD_LOG_NOTICE(( "reading incremental snapshot from file %s%s%s", fd_log_style_dim(), ctx->local_in.incremental_snapshot_path, fd_log_style_normal() ));
984 0 : FD_LOG_INFO(( "reading incremental snapshot at slot %lu from local file `%s`", ctx->local_in.incremental_snapshot_slot, ctx->local_in.incremental_snapshot_path ));
985 0 : ctx->state = FD_SNAPCT_STATE_READING_INCREMENTAL_FILE;
986 0 : init_load( ctx, stem, 0, 1 );
987 0 : } else {
988 0 : ctx->predicted_incremental.slot = best.incr_slot;
989 0 : send_expected_slot( ctx, stem, best.incr_slot );
990 :
991 0 : ctx->peer = best;
992 0 : ctx->state = FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP;
993 0 : init_load( ctx, stem, 0, 0 );
994 0 : log_download( ctx, 0, best.addr, best.incr_slot );
995 0 : }
996 0 : break;
997 0 : }
998 :
999 : /* ============================================================== */
1000 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI:
1001 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1002 0 : ctx->malformed = 0;
1003 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1004 0 : ctx->flush_ack = 0;
1005 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET;
1006 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from local file `%s`",
1007 0 : ctx->local_in.incremental_snapshot_slot, ctx->local_in.incremental_snapshot_path ));
1008 0 : break;
1009 0 : }
1010 :
1011 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1012 :
1013 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE;
1014 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_DONE, 0UL, 0UL, 0UL, 0UL, 0UL );
1015 0 : ctx->flush_ack = 0;
1016 0 : break;
1017 :
1018 : /* ============================================================== */
1019 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE:
1020 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1021 0 : ctx->malformed = 0;
1022 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1023 0 : ctx->flush_ack = 0;
1024 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET;
1025 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from local file `%s`",
1026 0 : ctx->local_in.incremental_snapshot_slot, ctx->local_in.incremental_snapshot_path ));
1027 0 : break;
1028 0 : }
1029 :
1030 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1031 :
1032 0 : log_completion( ctx, 0/*incr*/ );
1033 0 : ctx->state = FD_SNAPCT_STATE_SHUTDOWN;
1034 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_SHUTDOWN, 0UL, 0UL, 0UL, 0UL, 0UL );
1035 0 : break;
1036 :
1037 : /* ============================================================== */
1038 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI:
1039 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1040 0 : ctx->malformed = 0;
1041 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1042 0 : ctx->flush_ack = 0;
1043 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET;
1044 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1045 0 : "blacklisting peer due to download failure.",
1046 0 : ctx->predicted_incremental.slot,
1047 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_incr_snapshot_name ));
1048 0 : blacklist_peer( ctx );
1049 0 : break;
1050 0 : }
1051 :
1052 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1053 :
1054 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE;
1055 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_DONE, 0UL, 0UL, 0UL, 0UL, 0UL );
1056 0 : ctx->flush_ack = 0;
1057 0 : break;
1058 :
1059 : /* ============================================================== */
1060 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE:
1061 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1062 0 : ctx->malformed = 0;
1063 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1064 0 : ctx->flush_ack = 0;
1065 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET;
1066 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1067 0 : "blacklisting peer due to download failure.",
1068 0 : ctx->predicted_incremental.slot,
1069 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_incr_snapshot_name ));
1070 0 : blacklist_peer( ctx );
1071 0 : break;
1072 0 : }
1073 :
1074 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1075 :
1076 0 : log_completion( ctx, 0/*incr*/ );
1077 0 : ctx->state = FD_SNAPCT_STATE_SHUTDOWN;
1078 0 : rename_incr_snapshot( ctx );
1079 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_SHUTDOWN, 0UL, 0UL, 0UL, 0UL, 0UL );
1080 0 : break;
1081 :
1082 : /* ============================================================== */
1083 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI:
1084 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1085 0 : ctx->malformed = 0;
1086 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1087 0 : ctx->flush_ack = 0;
1088 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET;
1089 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from local file `%s`",
1090 0 : ctx->local_in.full_snapshot_slot, ctx->local_in.full_snapshot_path ));
1091 0 : break;
1092 0 : }
1093 :
1094 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1095 :
1096 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE;
1097 0 : ulong sig = ctx->config.incremental_snapshots &&
1098 0 : (ctx->local_in.incremental_snapshot_slot!=ULONG_MAX || ctx->download_enabled) ? FD_SNAPSHOT_MSG_CTRL_NEXT : FD_SNAPSHOT_MSG_CTRL_DONE;
1099 0 : if( sig==FD_SNAPSHOT_MSG_CTRL_DONE && ctx->config.incremental_snapshots ) {
1100 : /* set incremental snapshots to 0 if there is no local
1101 : incremental snapshot and download is not enabled. */
1102 0 : FD_LOG_INFO(( "incremental snapshots were enabled via [snapshots.incremental_snapshots] "
1103 0 : "but no incremental snapshot exists on disk and no snapshot peers are configured. "
1104 0 : "skipping incremental snapshot load." ));
1105 0 : ctx->config.incremental_snapshots = 0;
1106 0 : }
1107 0 : fd_stem_publish( stem, ctx->out_ld.idx, sig, 0UL, 0UL, 0UL, 0UL, 0UL );
1108 0 : ctx->flush_ack = 0;
1109 0 : break;
1110 :
1111 : /* ============================================================== */
1112 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE:
1113 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1114 0 : ctx->malformed = 0;
1115 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1116 0 : ctx->flush_ack = 0;
1117 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET;
1118 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from local file `%s`",
1119 0 : ctx->local_in.full_snapshot_slot, ctx->local_in.full_snapshot_path ));
1120 0 : break;
1121 0 : }
1122 :
1123 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1124 :
1125 0 : log_completion( ctx, 1/*full*/ );
1126 0 : if( FD_LIKELY( !ctx->config.incremental_snapshots ) ) {
1127 0 : ctx->state = FD_SNAPCT_STATE_SHUTDOWN;
1128 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_SHUTDOWN, 0UL, 0UL, 0UL, 0UL, 0UL );
1129 0 : break;
1130 0 : }
1131 :
1132 0 : if( FD_LIKELY( ctx->download_enabled ) ) {
1133 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL;
1134 0 : ctx->deadline_nanos = 0L;
1135 0 : } else {
1136 0 : FD_LOG_NOTICE(( "reading incremental snapshot at slot %lu from local file `%s`", ctx->local_in.incremental_snapshot_slot, ctx->local_in.incremental_snapshot_path ));
1137 0 : ctx->state = FD_SNAPCT_STATE_READING_INCREMENTAL_FILE;
1138 0 : init_load( ctx, stem, 0, 1 );
1139 0 : }
1140 0 : break;
1141 :
1142 : /* ============================================================== */
1143 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI:
1144 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1145 0 : ctx->malformed = 0;
1146 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1147 0 : ctx->flush_ack = 0;
1148 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET;
1149 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1150 0 : "blacklisting peer due to download failure.",
1151 0 : ctx->predicted_incremental.full_slot,
1152 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_full_snapshot_name ));
1153 0 : blacklist_peer( ctx );
1154 0 : break;
1155 0 : }
1156 :
1157 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1158 :
1159 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE;
1160 0 : fd_stem_publish( stem, ctx->out_ld.idx, ctx->config.incremental_snapshots ? FD_SNAPSHOT_MSG_CTRL_NEXT : FD_SNAPSHOT_MSG_CTRL_DONE, 0UL, 0UL, 0UL, 0UL, 0UL );
1161 0 : ctx->flush_ack = 0;
1162 0 : break;
1163 :
1164 : /* ============================================================== */
1165 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE:
1166 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1167 0 : ctx->malformed = 0;
1168 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1169 0 : ctx->flush_ack = 0;
1170 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET;
1171 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1172 0 : "blacklisting peer due to download failure.",
1173 0 : ctx->predicted_incremental.full_slot,
1174 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_full_snapshot_name ));
1175 0 : blacklist_peer( ctx );
1176 0 : break;
1177 0 : }
1178 :
1179 0 : if( ctx->flush_ack < ctx->flush_ack_cnt ) break;
1180 :
1181 0 : rename_full_snapshot( ctx );
1182 :
1183 0 : log_completion( ctx, 1/*full*/ );
1184 0 : if( FD_LIKELY( !ctx->config.incremental_snapshots ) ) {
1185 0 : ctx->state = FD_SNAPCT_STATE_SHUTDOWN;
1186 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_SHUTDOWN, 0UL, 0UL, 0UL, 0UL, 0UL );
1187 0 : break;
1188 0 : }
1189 :
1190 : /* Get the best incremental peer to download from */
1191 0 : fd_sspeer_t best = fd_sspeer_selector_best( ctx->selector, 1, ctx->predicted_incremental.full_slot );
1192 0 : if( FD_UNLIKELY( !best.addr.l ) ) {
1193 0 : ctx->deadline_nanos = now;
1194 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL;
1195 0 : break;
1196 0 : }
1197 :
1198 0 : ctx->predicted_incremental.slot = best.incr_slot;
1199 0 : send_expected_slot( ctx, stem, best.incr_slot );
1200 :
1201 0 : ctx->peer = best;
1202 0 : ctx->state = FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP;
1203 0 : init_load( ctx, stem, 0, 0 );
1204 0 : log_download( ctx, 0, best.addr, best.incr_slot );
1205 0 : break;
1206 :
1207 : /* ============================================================== */
1208 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET:
1209 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET:
1210 0 : if( FD_UNLIKELY( ctx->flush_ack<ctx->flush_ack_cnt ) ) break;
1211 :
1212 0 : if( ctx->metrics.full.num_retries==ctx->config.max_retry_abort ) {
1213 0 : FD_LOG_ERR(( "hit retry limit of %u for full snapshot, aborting", ctx->config.max_retry_abort ));
1214 0 : }
1215 :
1216 0 : ctx->metrics.full.num_retries++;
1217 0 : FD_LOG_NOTICE(( "retrying full snapshot download (attempt %u/%u)",
1218 0 : ctx->metrics.full.num_retries, ctx->config.max_retry_abort ));
1219 :
1220 0 : ctx->metrics.full.bytes_read = 0UL;
1221 0 : ctx->metrics.full.bytes_written = 0UL;
1222 0 : ctx->metrics.full.bytes_total = 0UL;
1223 :
1224 0 : ctx->metrics.incremental.bytes_read = 0UL;
1225 0 : ctx->metrics.incremental.bytes_written = 0UL;
1226 0 : ctx->metrics.incremental.bytes_total = 0UL;
1227 :
1228 0 : if( !ctx->download_enabled ) {
1229 : /* if we are unable to download new snapshots and unable to load
1230 : our local snapshot, we must shutdown the validator. */
1231 0 : FD_LOG_ERR(( "unable to load local snapshot %s and no snapshot peers were configured. aborting.", ctx->local_in.full_snapshot_path ));
1232 0 : } else {
1233 0 : if( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ) ctx->local_in.full_snapshot_slot = ULONG_MAX;
1234 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS;
1235 0 : ctx->deadline_nanos = 0L;
1236 0 : }
1237 0 : break;
1238 :
1239 : /* ============================================================== */
1240 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET:
1241 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET:
1242 0 : if( FD_UNLIKELY( ctx->flush_ack<ctx->flush_ack_cnt ) ) break;
1243 :
1244 0 : if( ctx->metrics.incremental.num_retries==ctx->config.max_retry_abort ) {
1245 0 : FD_LOG_ERR(("hit retry limit of %u for incremental snapshot. aborting", ctx->config.max_retry_abort ));
1246 0 : }
1247 :
1248 0 : ctx->metrics.incremental.num_retries++;
1249 0 : FD_LOG_NOTICE(( "retrying incremental snapshot download (attempt %u/%u)",
1250 0 : ctx->metrics.incremental.num_retries, ctx->config.max_retry_abort ));
1251 :
1252 0 : ctx->metrics.incremental.bytes_read = 0UL;
1253 0 : ctx->metrics.incremental.bytes_written = 0UL;
1254 0 : ctx->metrics.incremental.bytes_total = 0UL;
1255 :
1256 0 : if( !ctx->download_enabled ) {
1257 : /* if we are unable to download new snapshots and unable to load
1258 : our local snapshot, we must shutdown the validator. */
1259 0 : FD_LOG_ERR(( "unable to load local snapshot %s and no snapshot peers were configured. aborting.", ctx->local_in.incremental_snapshot_path ));
1260 0 : } else {
1261 0 : if( ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ctx->local_in.incremental_snapshot_slot = ULONG_MAX;
1262 0 : ctx->state = FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL;
1263 0 : ctx->deadline_nanos = 0L;
1264 0 : }
1265 0 : break;
1266 :
1267 : /* ============================================================== */
1268 0 : case FD_SNAPCT_STATE_READING_FULL_FILE:
1269 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1270 0 : ctx->malformed = 0;
1271 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1272 0 : ctx->flush_ack = 0;
1273 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET;
1274 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from local file `%s`",
1275 0 : ctx->local_in.full_snapshot_slot, ctx->local_in.full_snapshot_path ));
1276 0 : break;
1277 0 : }
1278 0 : if( FD_UNLIKELY( ctx->flush_ack < ctx->flush_ack_cnt ) ) break;
1279 0 : if( FD_UNLIKELY( ctx->load_complete ) ) {
1280 0 : ctx->load_complete = 0;
1281 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FINI, 0UL, 0UL, 0UL, 0UL, 0UL );
1282 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI;
1283 0 : ctx->flush_ack = 0;
1284 0 : }
1285 0 : break;
1286 :
1287 : /* ============================================================== */
1288 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE:
1289 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1290 0 : ctx->malformed = 0;
1291 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1292 0 : ctx->flush_ack = 0;
1293 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET;
1294 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from local file `%s`",
1295 0 : ctx->local_in.incremental_snapshot_slot, ctx->local_in.incremental_snapshot_path ));
1296 0 : break;
1297 0 : }
1298 0 : if( FD_UNLIKELY( ctx->flush_ack < ctx->flush_ack_cnt ) ) break;
1299 0 : if( FD_UNLIKELY( ctx->load_complete ) ) {
1300 0 : ctx->load_complete = 0;
1301 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FINI, 0UL, 0UL, 0UL, 0UL, 0UL );
1302 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI;
1303 0 : ctx->flush_ack = 0;
1304 0 : }
1305 0 : break;
1306 :
1307 : /* ============================================================== */
1308 0 : case FD_SNAPCT_STATE_READING_FULL_HTTP:
1309 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1310 0 : ctx->malformed = 0;
1311 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1312 0 : ctx->flush_ack = 0;
1313 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET;
1314 0 : FD_LOG_WARNING(( "failed to load full snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1315 0 : "blacklisting peer due to download failure",
1316 0 : ctx->predicted_incremental.full_slot,
1317 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_full_snapshot_name ));
1318 0 : blacklist_peer( ctx );
1319 0 : break;
1320 0 : }
1321 0 : if( FD_UNLIKELY( ctx->flush_ack < ctx->flush_ack_cnt ) ) break;
1322 0 : if( FD_UNLIKELY( ctx->load_complete ) ) {
1323 0 : ctx->load_complete = 0;
1324 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FINI, 0UL, 0UL, 0UL, 0UL, 0UL );
1325 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI;
1326 0 : ctx->flush_ack = 0;
1327 0 : }
1328 0 : break;
1329 :
1330 : /* ============================================================== */
1331 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP:
1332 0 : if( FD_UNLIKELY( ctx->malformed ) ) {
1333 0 : ctx->malformed = 0;
1334 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FAIL, 0UL, 0UL, 0UL, 0UL, 0UL );
1335 0 : ctx->flush_ack = 0;
1336 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET;
1337 0 : FD_LOG_WARNING(( "failed to load incremental snapshot at slot %lu from http://" FD_IP4_ADDR_FMT ":%hu/%s. "
1338 0 : "blacklisting peer due to download failure",
1339 0 : ctx->predicted_incremental.slot,
1340 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ), ctx->http_incr_snapshot_name ));
1341 0 : blacklist_peer( ctx );
1342 0 : break;
1343 0 : }
1344 0 : if( FD_UNLIKELY( ctx->flush_ack < ctx->flush_ack_cnt ) ) break;
1345 0 : if( FD_UNLIKELY( ctx->load_complete ) ) {
1346 0 : ctx->load_complete = 0;
1347 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_FINI, 0UL, 0UL, 0UL, 0UL, 0UL );
1348 0 : ctx->state = FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI;
1349 0 : ctx->flush_ack = 0;
1350 0 : }
1351 0 : break;
1352 :
1353 : /* ============================================================== */
1354 0 : case FD_SNAPCT_STATE_SHUTDOWN:
1355 : /* Transitioning to the shutdown state indicates snapshot load is
1356 : completed without errors. Otherwise, snapct would have aborted
1357 : earlier. */
1358 0 : break;
1359 :
1360 : /* ============================================================== */
1361 0 : default: FD_LOG_ERR(( "unexpected state %s", fd_snapct_state_str( ctx->state ) ));
1362 0 : }
1363 0 : }
1364 :
1365 : static void
1366 : gossip_frag( fd_snapct_tile_t * ctx,
1367 : ulong sig,
1368 : ulong sz FD_PARAM_UNUSED,
1369 45 : ulong chunk ) {
1370 45 : FD_TEST( ctx->gossip_enabled );
1371 :
1372 45 : if( FD_UNLIKELY( sig==FD_GOSSIP_UPDATE_TAG_PEER_SATURATED ) ) {
1373 0 : ctx->gossip.saturated = 1;
1374 0 : return;
1375 0 : }
1376 :
1377 45 : if( !( sig==FD_GOSSIP_UPDATE_TAG_CONTACT_INFO ||
1378 45 : sig==FD_GOSSIP_UPDATE_TAG_CONTACT_INFO_REMOVE ||
1379 45 : sig==FD_GOSSIP_UPDATE_TAG_SNAPSHOT_HASHES ) ) return;
1380 :
1381 45 : fd_gossip_update_message_t const * msg = fd_chunk_to_laddr_const( ctx->gossip_in_mem, chunk );
1382 45 : switch( msg->tag ) {
1383 39 : case FD_GOSSIP_UPDATE_TAG_CONTACT_INFO: {
1384 39 : FD_TEST( msg->contact_info->idx<GOSSIP_PEERS_MAX );
1385 39 : fd_pubkey_t const * pubkey = (fd_pubkey_t const *)msg->origin;
1386 39 : gossip_ci_entry_t * entry = ctx->gossip.ci_table + msg->contact_info->idx;
1387 39 : if( FD_UNLIKELY( !fd_pubkey_eq( &entry->pubkey, pubkey ) ) ) {
1388 : /* Initialize the new gossip entry, which may or may not be allowed */
1389 30 : FD_TEST( fd_pubkey_check_zero( &entry->pubkey ) );
1390 30 : entry->pubkey = *pubkey;
1391 30 : entry->rpc_addr.l = 0UL;
1392 30 : if( ctx->config.sources.gossip.allow_any ) {
1393 21 : entry->allowed = 1;
1394 21 : for( ulong i=0UL; i<ctx->config.sources.gossip.block_list_cnt; i++ ) {
1395 3 : if( fd_pubkey_eq( pubkey, &ctx->config.sources.gossip.block_list[ i ] ) ) {
1396 3 : entry->allowed = 0;
1397 3 : break;
1398 3 : }
1399 3 : }
1400 21 : } else {
1401 9 : entry->allowed = 0;
1402 9 : for( ulong i=0UL; i<ctx->config.sources.gossip.allow_list_cnt; i++ ) {
1403 3 : if( fd_pubkey_eq( pubkey, &ctx->config.sources.gossip.allow_list[ i ] ) ) {
1404 3 : entry->allowed = 1;
1405 3 : break;
1406 3 : }
1407 3 : }
1408 9 : }
1409 30 : FD_TEST( ULONG_MAX==gossip_ci_map_idx_query_const( ctx->gossip.ci_map, pubkey, ULONG_MAX, ctx->gossip.ci_table ) );
1410 30 : if( entry->allowed ) {
1411 21 : gossip_ci_map_idx_insert( ctx->gossip.ci_map, msg->contact_info->idx, ctx->gossip.ci_table );
1412 21 : ctx->gossip.allowed_cnt++;
1413 : /* Allow-list shortcut: if an explicit allow list is
1414 : configured and all expected peers have arrived, declare
1415 : saturation immediately without waiting for the gossip
1416 : tile's general saturation signal. */
1417 21 : if( FD_UNLIKELY( !ctx->config.sources.gossip.allow_any &&
1418 21 : ctx->config.sources.gossip.allow_list_cnt>0UL &&
1419 21 : ctx->gossip.allowed_cnt==ctx->config.sources.gossip.allow_list_cnt ) ) {
1420 3 : FD_LOG_NOTICE(( "all %lu allowed gossip peers discovered", ctx->config.sources.gossip.allow_list_cnt ));
1421 3 : ctx->gossip.saturated = 1;
1422 3 : }
1423 21 : }
1424 30 : }
1425 39 : if( !entry->allowed ) break;
1426 : /* Maybe update the RPC address of a new or existing allowed gossip peer */
1427 30 : fd_ip4_port_t cur_addr = entry->rpc_addr;
1428 30 : fd_ip4_port_t new_addr;
1429 30 : new_addr.addr = msg->contact_info->value->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_RPC ].is_ipv6 ? 0 : msg->contact_info->value->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_RPC ].ip4;
1430 30 : new_addr.port = msg->contact_info->value->sockets[ FD_GOSSIP_CONTACT_INFO_SOCKET_RPC ].port;
1431 : /* Sanitize non-public/multicast/broadcast addresses to zero so
1432 : they are never stored in entry->rpc_addr or handed to
1433 : downstream code paths (e.g. on_snapshot_hash). Normalizing
1434 : before comparing against cur_addr avoids repeated selector
1435 : removals and log warnings on every gossip refresh. */
1436 30 : uint unfiltered_addr = new_addr.addr;
1437 30 : if( FD_UNLIKELY( !fd_ip4_addr_is_public( new_addr.addr ) ||
1438 30 : fd_ip4_addr_is_mcast( new_addr.addr ) ||
1439 30 : fd_ip4_addr_is_bcast( new_addr.addr ) ) ) {
1440 18 : new_addr.l = 0;
1441 18 : }
1442 :
1443 30 : if( FD_UNLIKELY( new_addr.l!=cur_addr.l ) ) {
1444 15 : fd_sspeer_key_t entry_key = {0};
1445 15 : *entry_key.pubkey = entry->pubkey;
1446 15 : entry_key.is_url = 0;
1447 15 : if( FD_LIKELY( !!new_addr.l ) ) {
1448 : /* Do not re-add peers that have been permanently blacklisted
1449 : or whose new address is temporarily banned by ssping.
1450 : Bookkeeping (ssping cleanup, rpc_addr update) still runs
1451 : so the CI table and ssping pool stay consistent. */
1452 12 : if( FD_LIKELY( !blacklist_map_ele_query( ctx->blacklist_map, &entry_key, NULL, ctx->blacklist_pool ) &&
1453 12 : !fd_ssping_is_invalidated( ctx->ssping, new_addr ) ) ) {
1454 9 : if( FD_LIKELY( FD_SSPEER_SCORE_INVALID!=fd_sspeer_selector_add( ctx->selector, &entry_key, new_addr,
1455 9 : FD_SSPEER_LATENCY_UNKNOWN, FD_SSPEER_SLOT_UNKNOWN,
1456 9 : FD_SSPEER_SLOT_UNKNOWN, NULL, NULL ) ) ) {
1457 9 : fd_ssping_add( ctx->ssping, new_addr );
1458 9 : }
1459 9 : } else {
1460 : /* Evict stale old address from selector to stay in sync
1461 : with the CI table update below. Increment ssping
1462 : refcnt for new_addr to balance the fd_ssping_remove
1463 : on the next address change or peer removal. */
1464 3 : fd_sspeer_selector_remove( ctx->selector, &entry_key );
1465 3 : fd_ssping_add( ctx->ssping, new_addr );
1466 3 : }
1467 12 : } else {
1468 3 : if( FD_UNLIKELY( !!unfiltered_addr ) ) {
1469 3 : FD_LOG_WARNING(( "filtered non-public RPC address " FD_IP4_ADDR_FMT " from gossip peer", FD_IP4_ADDR_FMT_ARGS( unfiltered_addr ) ));
1470 3 : }
1471 3 : fd_sspeer_selector_remove( ctx->selector, &entry_key );
1472 3 : }
1473 15 : if( FD_LIKELY( !!cur_addr.l ) ) {
1474 6 : fd_ssping_remove( ctx->ssping, cur_addr );
1475 6 : }
1476 15 : entry->rpc_addr = new_addr;
1477 15 : if( !ctx->config.sources.gossip.allow_any ) {
1478 0 : FD_BASE58_ENCODE_32_BYTES( pubkey->uc, pubkey_b58 );
1479 0 : if( FD_LIKELY( !!new_addr.l ) ) {
1480 0 : FD_LOG_NOTICE(( "allowed gossip peer added with public key `%s` and RPC address `" FD_IP4_ADDR_FMT ":%hu`",
1481 0 : pubkey_b58, FD_IP4_ADDR_FMT_ARGS( new_addr.addr ), fd_ushort_bswap( new_addr.port ) ));
1482 0 : } else {
1483 0 : FD_LOG_WARNING(( "allowed gossip peer with public key `%s` does not advertise an RPC address", pubkey_b58 ));
1484 0 : }
1485 0 : }
1486 15 : }
1487 30 : break;
1488 39 : }
1489 6 : case FD_GOSSIP_UPDATE_TAG_CONTACT_INFO_REMOVE: {
1490 6 : FD_TEST( msg->contact_info_remove->idx<GOSSIP_PEERS_MAX );
1491 6 : gossip_ci_entry_t * entry = ctx->gossip.ci_table + msg->contact_info_remove->idx;
1492 6 : ulong rem_idx = gossip_ci_map_idx_remove( ctx->gossip.ci_map, &entry->pubkey, ULONG_MAX, ctx->gossip.ci_table );
1493 6 : if( rem_idx!=ULONG_MAX ) {
1494 3 : FD_TEST( entry->allowed && rem_idx==msg->contact_info_remove->idx );
1495 3 : ctx->gossip.allowed_cnt--;
1496 3 : fd_ip4_port_t addr = entry->rpc_addr;
1497 3 : if( FD_LIKELY( !!addr.l ) ) {
1498 0 : fd_ssping_remove( ctx->ssping, addr );
1499 0 : fd_sspeer_key_t entry_key = {0};
1500 0 : *entry_key.pubkey = entry->pubkey;
1501 0 : entry_key.is_url = 0;
1502 0 : fd_sspeer_selector_remove( ctx->selector, &entry_key );
1503 0 : }
1504 3 : if( !ctx->config.sources.gossip.allow_any ) {
1505 0 : FD_BASE58_ENCODE_32_BYTES( entry->pubkey.uc, pubkey_b58 );
1506 0 : FD_LOG_WARNING(( "allowed gossip peer removed with public key `%s` and RPC address `" FD_IP4_ADDR_FMT ":%hu`",
1507 0 : pubkey_b58, FD_IP4_ADDR_FMT_ARGS( addr.addr ), fd_ushort_bswap( addr.port ) ));
1508 0 : }
1509 3 : }
1510 6 : fd_memset( entry, 0, sizeof(*entry) );
1511 6 : break;
1512 6 : }
1513 0 : case FD_GOSSIP_UPDATE_TAG_SNAPSHOT_HASHES: {
1514 0 : ulong idx = gossip_ci_map_idx_query_const( ctx->gossip.ci_map, (fd_pubkey_t const *)msg->origin, ULONG_MAX, ctx->gossip.ci_table );
1515 0 : if( FD_LIKELY( idx!=ULONG_MAX ) ) {
1516 0 : gossip_ci_entry_t * entry = ctx->gossip.ci_table + idx;
1517 0 : FD_TEST( entry->allowed );
1518 0 : fd_sspeer_key_t entry_key = {0};
1519 0 : *entry_key.pubkey = entry->pubkey;
1520 0 : entry_key.is_url = 0;
1521 0 : on_snapshot_hash( ctx, &entry_key, entry->rpc_addr, msg );
1522 0 : }
1523 0 : break;
1524 0 : }
1525 0 : default:
1526 0 : FD_LOG_ERR(( "snapct: unexpected gossip tag %u", (uint)msg->tag ));
1527 0 : break;
1528 45 : }
1529 45 : }
1530 :
1531 : /* Validate and handle a pipeline control ack for INIT_{FULL,INCR},
1532 : NEXT, DONE, and FINI. Returns 0 on success, and -1 if the sig was
1533 : unrecognized or the state was unexpected. */
1534 : static int
1535 : process_ctrl_ack( fd_snapct_tile_t * ctx,
1536 0 : ulong sig ) {
1537 0 : switch( sig ) {
1538 0 : case FD_SNAPSHOT_MSG_CTRL_INIT_FULL:
1539 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_READING_FULL_HTTP ||
1540 0 : ctx->state==FD_SNAPCT_STATE_READING_FULL_FILE ) ) {
1541 0 : ctx->flush_ack++;
1542 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1543 0 : } else if( FD_UNLIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1544 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ) ) {
1545 : /* Safe to ignore -- stale ack from before RESET. */
1546 0 : } else return -1;
1547 0 : return 0;
1548 :
1549 0 : case FD_SNAPSHOT_MSG_CTRL_INIT_INCR:
1550 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP ||
1551 0 : ctx->state==FD_SNAPCT_STATE_READING_INCREMENTAL_FILE ) ) {
1552 0 : ctx->flush_ack++;
1553 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1554 0 : } else if( FD_UNLIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET ||
1555 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ) {
1556 : /* Safe to ignore -- stale ack from before RESET. */
1557 0 : } else return -1;
1558 0 : return 0;
1559 :
1560 0 : case FD_SNAPSHOT_MSG_CTRL_NEXT:
1561 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE ||
1562 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE ) ) {
1563 0 : ctx->flush_ack++;
1564 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1565 0 : } else if( FD_UNLIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1566 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ) ) {
1567 : /* Safe to ignore -- stale ack from before RESET. */
1568 0 : } else return -1;
1569 0 : return 0;
1570 :
1571 0 : case FD_SNAPSHOT_MSG_CTRL_DONE:
1572 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE ||
1573 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE ||
1574 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE ||
1575 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE ) ) {
1576 0 : ctx->flush_ack++;
1577 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1578 0 : } else if( FD_UNLIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1579 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ||
1580 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET ||
1581 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ) {
1582 : /* Safe to ignore -- stale ack from before RESET. */
1583 0 : } else return -1;
1584 0 : return 0;
1585 :
1586 0 : case FD_SNAPSHOT_MSG_CTRL_FINI:
1587 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI ||
1588 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI ||
1589 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI ||
1590 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI ) ) {
1591 0 : ctx->flush_ack++;
1592 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1593 0 : } else if( FD_UNLIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1594 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ||
1595 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET ||
1596 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ) {
1597 : /* Safe to ignore -- stale ack from before RESET. */
1598 0 : } else return -1;
1599 0 : return 0;
1600 :
1601 0 : default:
1602 0 : return -1;
1603 0 : }
1604 0 : }
1605 :
1606 : static void
1607 : snapld_frag( fd_snapct_tile_t * ctx,
1608 : ulong sig,
1609 : ulong sz,
1610 : ulong chunk,
1611 15 : fd_stem_context_t * stem ) {
1612 15 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_META ) ) {
1613 : /* Before snapld starts sending down data fragments, it first sends
1614 : a metadata message containing the total size of the snapshot.
1615 : Both file and HTTP paths send this message. */
1616 0 : int full, file;
1617 0 : switch( ctx->state ) {
1618 0 : case FD_SNAPCT_STATE_READING_FULL_FILE: full = 1; file = 1; break;
1619 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE: full = 0; file = 1; break;
1620 0 : case FD_SNAPCT_STATE_READING_FULL_HTTP: full = 1; file = 0; break;
1621 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP: full = 0; file = 0; break;
1622 :
1623 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET:
1624 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET:
1625 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET:
1626 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET:
1627 0 : return; /* Ignore */
1628 0 : default: FD_LOG_ERR(( "invalid meta frag in state %s", fd_snapct_state_str( ctx->state ) ));
1629 0 : }
1630 :
1631 0 : FD_TEST( sz==sizeof(fd_ssctrl_meta_t) );
1632 0 : fd_ssctrl_meta_t const * meta = fd_chunk_to_laddr_const( ctx->snapld_in_mem, chunk );
1633 :
1634 0 : if( FD_UNLIKELY( meta->total_sz==0UL ) ) {
1635 0 : if( FD_UNLIKELY( !ctx->malformed ) ) {
1636 0 : FD_LOG_WARNING(( "received zero-length metadata for %s snapshot, marking malformed", full ? "full" : "incremental" ));
1637 0 : ctx->malformed = 1;
1638 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
1639 0 : }
1640 0 : return;
1641 0 : }
1642 :
1643 0 : if( full ) ctx->metrics.full.bytes_total = meta->total_sz;
1644 0 : else ctx->metrics.incremental.bytes_total = meta->total_sz;
1645 :
1646 : /* Publish snapshot path to GUI. For file loads, use the local
1647 : path directly. For HTTP downloads, construct the full URL from
1648 : the resolved snapshot name. */
1649 0 : if( file ) {
1650 0 : if( FD_LIKELY( !!ctx->out_gui.mem ) ) {
1651 0 : snapshot_path_gui_publish( ctx, stem, full ? ctx->local_in.full_snapshot_path : ctx->local_in.incremental_snapshot_path, full );
1652 0 : }
1653 0 : } else {
1654 0 : if( full ) fd_cstr_ncpy( ctx->http_full_snapshot_name, meta->resolved_name, PATH_MAX );
1655 0 : else fd_cstr_ncpy( ctx->http_incr_snapshot_name, meta->resolved_name, PATH_MAX );
1656 :
1657 0 : if( FD_LIKELY( !!ctx->out_gui.mem ) ) {
1658 0 : char snapshot_path[ PATH_MAX+30UL ]; /* 30 is fd_cstr_nlen( "https://255.255.255.255:65536/", ULONG_MAX ) */
1659 0 : FD_TEST( fd_cstr_printf_check( snapshot_path, sizeof(snapshot_path), NULL, "http://" FD_IP4_ADDR_FMT ":%hu/%s",
1660 0 : FD_IP4_ADDR_FMT_ARGS( ctx->peer.addr.addr ), fd_ushort_bswap( ctx->peer.addr.port ),
1661 0 : full ? ctx->http_full_snapshot_name : ctx->http_incr_snapshot_name ) );
1662 0 : snapshot_path_gui_publish( ctx, stem, snapshot_path, full );
1663 0 : }
1664 0 : }
1665 :
1666 0 : return;
1667 0 : }
1668 15 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_CTRL_FAIL ) ) {
1669 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1670 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ||
1671 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET ||
1672 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ) {
1673 0 : ctx->flush_ack++;
1674 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1675 0 : } else {
1676 0 : FD_LOG_ERR(( "unexpected control frag %lu (%s) in state %d (%s)", sig, fd_ssctrl_msg_ctrl_str( sig ), ctx->state, fd_snapct_state_str( ctx->state ) ));
1677 0 : }
1678 0 : return;
1679 0 : }
1680 15 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_CTRL_ERROR ) ) {
1681 : /* CTRL_ERROR message directly from snapld can be snapld-generated
1682 : or snapct-generated-and-forwarded. */
1683 0 : switch( ctx->state ) {
1684 0 : case FD_SNAPCT_STATE_READING_FULL_FILE:
1685 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI:
1686 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE:
1687 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE:
1688 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI:
1689 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE:
1690 0 : case FD_SNAPCT_STATE_READING_FULL_HTTP:
1691 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI:
1692 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE:
1693 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP:
1694 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI:
1695 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE:
1696 0 : FD_LOG_WARNING(( "received error from snapld in state %d (%s)",
1697 0 : ctx->state, fd_snapct_state_str( ctx->state ) ));
1698 0 : ctx->malformed = 1;
1699 0 : break;
1700 0 : default:
1701 0 : break;
1702 0 : }
1703 0 : return;
1704 0 : }
1705 15 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_LOAD_COMPLETE ) ) {
1706 15 : int full = 0;
1707 15 : switch( ctx->state ) {
1708 3 : case FD_SNAPCT_STATE_READING_FULL_FILE:
1709 6 : case FD_SNAPCT_STATE_READING_FULL_HTTP: full = 1; break;
1710 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE:
1711 3 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP: full = 0; break;
1712 3 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET:
1713 3 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET:
1714 3 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET:
1715 6 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET:
1716 6 : return; /* Ignore during reset. */
1717 0 : default:
1718 0 : FD_LOG_ERR(( "invalid load_complete in state %s", fd_snapct_state_str( ctx->state ) ));
1719 0 : return;
1720 15 : }
1721 9 : if( FD_UNLIKELY( ctx->malformed ) ) return;
1722 : /* Validate that all expected bytes were received. */
1723 6 : ulong bytes_read = full ? ctx->metrics.full.bytes_read : ctx->metrics.incremental.bytes_read;
1724 6 : ulong bytes_total = full ? ctx->metrics.full.bytes_total : ctx->metrics.incremental.bytes_total;
1725 6 : if( FD_UNLIKELY( !bytes_total || bytes_read!=bytes_total ) ) {
1726 0 : ctx->malformed = 1;
1727 0 : FD_LOG_WARNING(( "load_complete but bytes_read %lu != bytes_total %lu for %s snapshot",
1728 0 : bytes_read, bytes_total, full ? "full" : "incremental" ));
1729 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
1730 0 : return;
1731 0 : }
1732 6 : ctx->load_complete = 1;
1733 6 : return;
1734 6 : }
1735 0 : if( FD_UNLIKELY( sig!=FD_SNAPSHOT_MSG_DATA ) ) {
1736 0 : if( process_ctrl_ack( ctx, sig ) ) {
1737 0 : FD_LOG_ERR(( "unexpected control frag %lu (%s) in state %d (%s)", sig, fd_ssctrl_msg_ctrl_str( sig ), ctx->state, fd_snapct_state_str( ctx->state ) ));
1738 0 : }
1739 0 : return;
1740 0 : }
1741 :
1742 0 : int full, file;
1743 0 : switch( ctx->state ) {
1744 : /* Expected cases, fall through below */
1745 0 : case FD_SNAPCT_STATE_READING_FULL_FILE: full = 1; file = 1; break;
1746 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE: full = 0; file = 1; break;
1747 0 : case FD_SNAPCT_STATE_READING_FULL_HTTP: full = 1; file = 0; break;
1748 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP: full = 0; file = 0; break;
1749 :
1750 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET:
1751 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET:
1752 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET:
1753 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET:
1754 : /* We are waiting for a reset to fully propagate through the
1755 : pipeline, just throw away any trailing data frags. */
1756 0 : return;
1757 :
1758 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI:
1759 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI:
1760 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI:
1761 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI:
1762 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE:
1763 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE:
1764 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE:
1765 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE:
1766 : /* Based on previously received data frags, we expected that the
1767 : current full / incremental snapshot was finished, but then we
1768 : received additional data frags. Unsafe to continue so throw
1769 : away the whole snapshot. Do not publish MSG_CTRL_ERROR here:
1770 : data forwarding is already complete, and the state machine
1771 : will publish MSG_CTRL_FAIL on the next tick. */
1772 0 : if( !ctx->malformed ) {
1773 0 : ctx->malformed = 1;
1774 0 : FD_LOG_WARNING(( "complete snapshot loaded but read %lu extra bytes", sz ));
1775 0 : }
1776 0 : return;
1777 :
1778 0 : case FD_SNAPCT_STATE_WAITING_FOR_PEERS:
1779 0 : case FD_SNAPCT_STATE_WAITING_FOR_PEERS_INCREMENTAL:
1780 0 : case FD_SNAPCT_STATE_COLLECTING_PEERS:
1781 0 : case FD_SNAPCT_STATE_COLLECTING_PEERS_INCREMENTAL:
1782 0 : case FD_SNAPCT_STATE_SHUTDOWN:
1783 0 : default:
1784 0 : FD_LOG_ERR(( "invalid data frag in state %s", fd_snapct_state_str( ctx->state ) ));
1785 0 : return;
1786 0 : }
1787 :
1788 0 : if( FD_UNLIKELY( full && ctx->metrics.full.bytes_total==0UL ) ) {
1789 0 : if( !ctx->malformed ) {
1790 0 : ctx->malformed = 1;
1791 0 : FD_LOG_WARNING(( "received data frag for full snapshot with zero bytes_total" ));
1792 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
1793 0 : }
1794 0 : return;
1795 0 : }
1796 0 : if( FD_UNLIKELY( !full && ctx->metrics.incremental.bytes_total==0UL ) ) {
1797 0 : if( !ctx->malformed ) {
1798 0 : ctx->malformed = 1;
1799 0 : FD_LOG_WARNING(( "received data frag for incremental snapshot with zero bytes_total" ));
1800 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
1801 0 : }
1802 0 : return;
1803 0 : }
1804 :
1805 0 : if( full ) ctx->metrics.full.bytes_read += sz;
1806 0 : else ctx->metrics.incremental.bytes_read += sz;
1807 :
1808 0 : if( !file && -1!=ctx->local_out.dir_fd ) {
1809 0 : uchar const * data = fd_chunk_to_laddr_const( ctx->snapld_in_mem, chunk );
1810 0 : int fd = full ? ctx->local_out.full_snapshot_fd : ctx->local_out.incremental_snapshot_fd;
1811 0 : ulong written_sz = 0;
1812 0 : while( written_sz<sz ) {
1813 0 : long result = write( fd, data+written_sz, sz-written_sz );
1814 0 : if( FD_UNLIKELY( -1==result && errno==EINTR ) ) continue;
1815 0 : if( FD_UNLIKELY( -1==result && errno==ENOSPC ) ) {
1816 0 : FD_LOG_ERR(( "Out of disk space when writing out snapshot data to `%s`", ctx->config.snapshots_path ));
1817 0 : } else if( FD_UNLIKELY( 0L>result ) ) {
1818 0 : FD_LOG_ERR(( "write() failed (%i-%s)", errno, fd_io_strerror( errno ) ));
1819 0 : } else if( FD_UNLIKELY( 0L==result ) ) {
1820 0 : FD_LOG_ERR(( "write() returned 0 for non-zero count" ));
1821 0 : }
1822 :
1823 0 : written_sz += (ulong)result;
1824 0 : }
1825 0 : if( full ) ctx->metrics.full.bytes_written += sz;
1826 0 : else ctx->metrics.incremental.bytes_written += sz;
1827 0 : }
1828 :
1829 0 : if( FD_UNLIKELY( ( full && ctx->metrics.full.bytes_read > ctx->metrics.full.bytes_total ) ||
1830 0 : (!full && ctx->metrics.incremental.bytes_read > ctx->metrics.incremental.bytes_total ) ) ) {
1831 0 : if( !ctx->malformed ) {
1832 0 : ctx->malformed = 1;
1833 0 : FD_LOG_WARNING(( "expected %s snapshot size of %lu bytes but read %lu bytes",
1834 0 : full ? "full" : "incremental",
1835 0 : full ? ctx->metrics.full.bytes_total : ctx->metrics.incremental.bytes_total,
1836 0 : full ? ctx->metrics.full.bytes_read : ctx->metrics.incremental.bytes_read ));
1837 0 : fd_stem_publish( stem, ctx->out_ld.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
1838 0 : }
1839 0 : return;
1840 0 : }
1841 0 : }
1842 :
1843 : static void
1844 : ctrl_ack_frag( fd_snapct_tile_t * ctx,
1845 0 : ulong sig ) {
1846 0 : switch( sig ) {
1847 0 : case FD_SNAPSHOT_MSG_CTRL_FAIL:
1848 0 : if( FD_LIKELY( ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_RESET ||
1849 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_FULL_FILE_RESET ||
1850 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_RESET ||
1851 0 : ctx->state==FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_RESET ) ) {
1852 0 : ctx->flush_ack++;
1853 0 : FD_TEST( ctx->flush_ack <= ctx->flush_ack_cnt );
1854 0 : } else {
1855 0 : FD_LOG_ERR(( "unexpected control frag %lu (%s) in state %d (%s)", sig, fd_ssctrl_msg_ctrl_str( sig ), ctx->state, fd_snapct_state_str( ctx->state ) ));
1856 0 : }
1857 0 : return;
1858 :
1859 0 : case FD_SNAPSHOT_MSG_CTRL_SHUTDOWN:
1860 0 : return;
1861 :
1862 0 : case FD_SNAPSHOT_MSG_CTRL_ERROR:
1863 0 : switch( ctx->state ) {
1864 0 : case FD_SNAPCT_STATE_READING_FULL_FILE:
1865 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_FINI:
1866 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_FILE_DONE:
1867 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_FILE:
1868 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_FINI:
1869 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_FILE_DONE:
1870 0 : case FD_SNAPCT_STATE_READING_FULL_HTTP:
1871 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_FINI:
1872 0 : case FD_SNAPCT_STATE_FLUSHING_FULL_HTTP_DONE:
1873 0 : case FD_SNAPCT_STATE_READING_INCREMENTAL_HTTP:
1874 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_FINI:
1875 0 : case FD_SNAPCT_STATE_FLUSHING_INCREMENTAL_HTTP_DONE:
1876 : /* Do not publish MSG_CTRL_ERROR: the error originated
1877 : downstream, so re-publishing would be redundant.
1878 : MSG_CTRL_FAIL follows on the next state machine tick. */
1879 0 : FD_LOG_WARNING(( "received error from downstream tile while in state %s",
1880 0 : fd_snapct_state_str( ctx->state ) ));
1881 0 : ctx->malformed = 1;
1882 0 : break;
1883 0 : default:
1884 0 : break;
1885 0 : }
1886 0 : return;
1887 :
1888 0 : default:
1889 0 : break;
1890 0 : }
1891 0 : if( process_ctrl_ack( ctx, sig ) ) {
1892 0 : FD_LOG_ERR(( "unexpected control frag %lu (%s) in state %d (%s)", sig, fd_ssctrl_msg_ctrl_str( sig ), ctx->state, fd_snapct_state_str( ctx->state ) ));
1893 0 : }
1894 0 : }
1895 :
1896 : static int
1897 : returnable_frag( fd_snapct_tile_t * ctx,
1898 : ulong in_idx,
1899 : ulong seq FD_PARAM_UNUSED,
1900 : ulong sig,
1901 : ulong chunk,
1902 : ulong sz,
1903 : ulong ctl FD_PARAM_UNUSED,
1904 : ulong tsorig FD_PARAM_UNUSED,
1905 : ulong tspub FD_PARAM_UNUSED,
1906 0 : fd_stem_context_t * stem ) {
1907 0 : if( FD_LIKELY( ctx->in_kind[ in_idx ]==IN_KIND_GOSSIP ) ) {
1908 0 : gossip_frag( ctx, sig, sz, chunk );
1909 0 : } else if( ctx->in_kind[ in_idx ]==IN_KIND_SNAPLD ) {
1910 0 : snapld_frag( ctx, sig, sz, chunk, stem );
1911 0 : } else if( ctx->in_kind[ in_idx ]==IN_KIND_ACK ) {
1912 0 : ctrl_ack_frag( ctx, sig );
1913 0 : } else FD_LOG_ERR(( "invalid in_kind %lu %u", in_idx, (uint)ctx->in_kind[ in_idx ] ));
1914 0 : return 0;
1915 0 : }
1916 :
1917 : static void
1918 : privileged_init( fd_topo_t const * topo,
1919 0 : fd_topo_tile_t const * tile ) {
1920 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
1921 :
1922 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
1923 0 : fd_snapct_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapct_tile_t), sizeof(fd_snapct_tile_t) );
1924 0 : void * _ssping = FD_SCRATCH_ALLOC_APPEND( l, fd_ssping_align(), fd_ssping_footprint( TOTAL_PEERS_MAX ) );
1925 0 : FD_SCRATCH_ALLOC_APPEND( l, alignof(gossip_ci_entry_t), sizeof(gossip_ci_entry_t)*GOSSIP_PEERS_MAX );
1926 0 : FD_SCRATCH_ALLOC_APPEND( l, gossip_ci_map_align(), gossip_ci_map_footprint( gossip_ci_map_chain_cnt_est( GOSSIP_PEERS_MAX ) ) );
1927 0 : void * _ssresolver = FD_SCRATCH_ALLOC_APPEND( l, fd_http_resolver_align(), fd_http_resolver_footprint( SERVER_PEERS_MAX ) );
1928 0 : FD_SCRATCH_ALLOC_APPEND( l, fd_sspeer_selector_align(), fd_sspeer_selector_footprint( TOTAL_PEERS_MAX ) );
1929 0 : FD_SCRATCH_ALLOC_APPEND( l, blacklist_pool_align(), blacklist_pool_footprint( TOTAL_PEERS_MAX ) );
1930 0 : FD_SCRATCH_ALLOC_APPEND( l, blacklist_map_align(), blacklist_map_footprint( blacklist_map_chain_cnt_est( TOTAL_PEERS_MAX ) ) );
1931 0 : FD_SCRATCH_ALLOC_APPEND( l, fd_adns_align(), fd_adns_footprint( ADNS_REQS_MAX ) );
1932 :
1933 0 : FD_TEST( fd_rng_secure( &ctx->ssping_seed, sizeof(ctx->ssping_seed) ) );
1934 0 : FD_TEST( fd_rng_secure( &ctx->gossip_ci_seed, sizeof(ctx->gossip_ci_seed) ) );
1935 0 : FD_TEST( fd_rng_secure( &ctx->selector_seed, sizeof(ctx->selector_seed) ) );
1936 0 : FD_TEST( fd_rng_secure( &ctx->blacklist_seed, sizeof(ctx->blacklist_seed) ) );
1937 :
1938 : /* The resolver only needs a trust store if some server is https. */
1939 0 : int any_https = 0;
1940 0 : for( ulong i=0UL; i<tile->snapct.sources.servers_cnt; i++ ) {
1941 0 : char hostname[ FD_FQDN_BUF_MAX ];
1942 0 : ushort port;
1943 0 : int is_https;
1944 0 : fd_dns_peer_parse( tile->snapct.sources.servers[ i ], "snapshots.sources.servers", hostname, &port, &is_https );
1945 0 : any_https |= is_https;
1946 0 : }
1947 :
1948 0 : ctx->ssping = NULL;
1949 0 : if( FD_LIKELY( download_enabled( tile ) ) ) ctx->ssping = fd_ssping_join( fd_ssping_new( _ssping, TOTAL_PEERS_MAX, ctx->ssping_seed, on_ping, ctx ) );
1950 0 : if( FD_LIKELY( tile->snapct.sources.servers_cnt ) ) ctx->ssresolver = fd_http_resolver_join( fd_http_resolver_new( _ssresolver, SERVER_PEERS_MAX, tile->snapct.incremental_snapshots, any_https, on_resolve, ctx ) );
1951 0 : else ctx->ssresolver = NULL;
1952 :
1953 0 : ctx->netdb_fds->etc_hosts = -1;
1954 0 : ctx->netdb_fds->etc_resolv_conf = -1;
1955 0 : if( FD_LIKELY( download_enabled( tile ) ) ) FD_TEST( fd_netdb_open_fds( ctx->netdb_fds ) );
1956 0 : ctx->resolved_entrypoints_cnt = 0UL;
1957 0 : ctx->resolved_servers_cnt = 0UL;
1958 :
1959 0 : fd_snap_pool_layout_t layout = fd_snap_pool_layout(
1960 0 : tile->snapct.max_full_snapshots_to_keep,
1961 0 : tile->snapct.max_incremental_snapshots_to_keep,
1962 0 : tile->snapct.incremental_snapshots,
1963 0 : download_enabled( tile ) );
1964 0 : uint snap_full_max = (uint)layout.full_max;
1965 0 : uint snap_incr_max = (uint)layout.incr_max;
1966 0 : uint retained_snap_max = (uint)layout.retained_max;
1967 0 : uint scratch_full_cnt = (uint)layout.scratch_full;
1968 0 : uint scratch_incr_cnt = (uint)layout.scratch_incr;
1969 0 : uint snap_max = (uint)layout.max;
1970 :
1971 0 : ulong full_slot = ULONG_MAX;
1972 0 : ulong incremental_slot = ULONG_MAX;
1973 0 : int full_is_zstd = 0;
1974 0 : int incremental_is_zstd = 0;
1975 0 : char full_path[ PATH_MAX ] = {0};
1976 0 : char incremental_path[ PATH_MAX ] = {0};
1977 0 : uchar full_snapshot_hash[ FD_HASH_FOOTPRINT ] = {0};
1978 0 : uchar incremental_snapshot_hash[ FD_HASH_FOOTPRINT ] = {0};
1979 0 : if( FD_UNLIKELY( -1==fd_ssarchive_latest_pair( tile->snapct.snapshots_path,
1980 0 : tile->snapct.incremental_snapshots,
1981 0 : &full_slot,
1982 0 : &incremental_slot,
1983 0 : full_path,
1984 0 : incremental_path,
1985 0 : &full_is_zstd,
1986 0 : &incremental_is_zstd,
1987 0 : full_snapshot_hash,
1988 0 : incremental_snapshot_hash ) ) ) {
1989 0 : if( FD_UNLIKELY( !download_enabled( tile ) ) ) {
1990 0 : FD_LOG_ERR(( "No snapshots found in `%s` and no download sources are enabled. "
1991 0 : "Please enable downloading via [snapshots.sources] and restart.", tile->snapct.snapshots_path ));
1992 0 : }
1993 0 : ctx->local_in.full_snapshot_slot = ULONG_MAX;
1994 0 : ctx->local_in.incremental_snapshot_slot = ULONG_MAX;
1995 0 : ctx->local_in.full_snapshot_size = 0UL;
1996 0 : ctx->local_in.incremental_snapshot_size = 0UL;
1997 0 : ctx->local_in.full_snapshot_zstd = 0;
1998 0 : ctx->local_in.incremental_snapshot_zstd = 0;
1999 0 : fd_cstr_fini( ctx->local_in.full_snapshot_path );
2000 0 : fd_cstr_fini( ctx->local_in.incremental_snapshot_path );
2001 0 : fd_memset( ctx->local_in.full_snapshot_hash, 0, FD_HASH_FOOTPRINT );
2002 0 : fd_memset( ctx->local_in.incremental_snapshot_hash, 0, FD_HASH_FOOTPRINT );
2003 0 : } else {
2004 0 : FD_TEST( full_slot!=ULONG_MAX );
2005 :
2006 0 : ctx->local_in.full_snapshot_slot = full_slot;
2007 0 : ctx->local_in.incremental_snapshot_slot = incremental_slot;
2008 0 : ctx->local_in.full_snapshot_zstd = full_is_zstd;
2009 0 : ctx->local_in.incremental_snapshot_zstd = incremental_is_zstd;
2010 :
2011 0 : fd_cstr_ncpy( ctx->local_in.full_snapshot_path, full_path, PATH_MAX );
2012 0 : fd_memcpy( ctx->local_in.full_snapshot_hash, full_snapshot_hash, FD_HASH_FOOTPRINT );
2013 0 : struct stat full_stat;
2014 0 : if( FD_UNLIKELY( -1==stat( ctx->local_in.full_snapshot_path, &full_stat ) ) ) FD_LOG_ERR(( "stat() failed `%s` (%i-%s)", full_path, errno, fd_io_strerror( errno ) ));
2015 0 : if( FD_UNLIKELY( !S_ISREG( full_stat.st_mode ) ) ) FD_LOG_ERR(( "full snapshot path `%s` is not a regular file", full_path ));
2016 0 : ctx->local_in.full_snapshot_size = (ulong)full_stat.st_size;
2017 :
2018 0 : if( FD_LIKELY( incremental_slot!=ULONG_MAX ) ) {
2019 0 : fd_cstr_ncpy( ctx->local_in.incremental_snapshot_path, incremental_path, PATH_MAX );
2020 0 : fd_memcpy( ctx->local_in.incremental_snapshot_hash, incremental_snapshot_hash, FD_HASH_FOOTPRINT );
2021 0 : struct stat incremental_stat;
2022 0 : if( FD_UNLIKELY( -1==stat( ctx->local_in.incremental_snapshot_path, &incremental_stat ) ) ) FD_LOG_ERR(( "stat() failed `%s` (%i-%s)", incremental_path, errno, fd_io_strerror( errno ) ));
2023 0 : if( FD_UNLIKELY( !S_ISREG( incremental_stat.st_mode ) ) ) FD_LOG_ERR(( "incremental snapshot path `%s` is not a regular file", incremental_path ));
2024 0 : ctx->local_in.incremental_snapshot_size = (ulong)incremental_stat.st_size;
2025 0 : } else {
2026 0 : ctx->local_in.incremental_snapshot_size = 0UL;
2027 0 : fd_cstr_fini( ctx->local_in.incremental_snapshot_path );
2028 0 : }
2029 0 : }
2030 :
2031 0 : ctx->local_out.dir_fd = -1;
2032 0 : ctx->local_out.full_snapshot_fd = -1;
2033 0 : ctx->local_out.incremental_snapshot_fd = -1;
2034 0 : fd_cstr_fini( ctx->local_out.full_snapshot_name );
2035 0 : fd_cstr_fini( ctx->local_out.incremental_snapshot_name );
2036 0 : if( FD_LIKELY( download_enabled( tile ) ) ) {
2037 0 : ctx->local_out.dir_fd = open( tile->snapct.snapshots_path, O_RDONLY|O_DIRECTORY|O_CLOEXEC );
2038 0 : if( FD_UNLIKELY( -1==ctx->local_out.dir_fd ) )
2039 0 : FD_LOG_ERR(( "open(%s) failed (%i-%s)", tile->snapct.snapshots_path, errno, fd_io_strerror( errno ) ));
2040 :
2041 0 : FD_TEST( snap_max && snap_max<=FD_SNAP_MAX );
2042 0 : fd_backup_inode_t pool[ FD_SNAP_MAX ] = {0};
2043 0 : if( FD_UNLIKELY( snap_max<FD_SNAP_MAX && -1!=fcntl( FD_SNAP_FD( snap_max ), F_GETFD ) ) ) {
2044 0 : FD_LOG_ERR(( "inherited snapshot pool is larger than the expected %u slots", snap_max ));
2045 0 : }
2046 0 : fd_snap_pool_recover( ctx->local_out.dir_fd, tile->snapct.snapshots_path, pool, snap_max );
2047 0 : if( FD_LIKELY( snap_full_max ) ) {
2048 0 : uint full_idx = snapshot_pool_select( pool, 0U, snap_full_max );
2049 0 : FD_TEST( full_idx!=UINT_MAX );
2050 0 : ctx->local_out.full_snapshot_fd = FD_SNAP_FD( full_idx );
2051 0 : fd_cstr_ncpy( ctx->local_out.full_snapshot_name, pool[ full_idx ].name, FD_SNAP_NAME_MAX );
2052 0 : } else {
2053 0 : uint full_idx = retained_snap_max; /* scratch full slot */
2054 0 : FD_TEST( scratch_full_cnt && full_idx<snap_max );
2055 0 : ctx->local_out.full_snapshot_fd = FD_SNAP_FD( full_idx );
2056 0 : fd_cstr_ncpy( ctx->local_out.full_snapshot_name, pool[ full_idx ].name, FD_SNAP_NAME_MAX );
2057 0 : }
2058 :
2059 0 : if( FD_LIKELY( tile->snapct.incremental_snapshots && snap_incr_max ) ) {
2060 0 : uint incr_idx = snapshot_pool_select( pool, snap_full_max, retained_snap_max );
2061 0 : FD_TEST( incr_idx!=UINT_MAX );
2062 0 : ctx->local_out.incremental_snapshot_fd = FD_SNAP_FD( incr_idx );
2063 0 : fd_cstr_ncpy( ctx->local_out.incremental_snapshot_name, pool[ incr_idx ].name, FD_SNAP_NAME_MAX );
2064 0 : } else if( FD_LIKELY( tile->snapct.incremental_snapshots ) ) {
2065 0 : uint incr_idx = retained_snap_max + scratch_full_cnt; /* scratch incremental slot */
2066 0 : FD_TEST( scratch_incr_cnt && incr_idx<snap_max );
2067 0 : ctx->local_out.incremental_snapshot_fd = FD_SNAP_FD( incr_idx );
2068 0 : fd_cstr_ncpy( ctx->local_out.incremental_snapshot_name, pool[ incr_idx ].name, FD_SNAP_NAME_MAX );
2069 0 : }
2070 :
2071 0 : FD_TEST( ctx->local_out.full_snapshot_fd!=ctx->local_out.incremental_snapshot_fd );
2072 0 : }
2073 :
2074 0 : if( FD_LIKELY( fd_sandbox_gettid()==fd_sandbox_getpid() ) ) {
2075 0 : for( uint i=0U; i<snap_max; i++ ) {
2076 0 : int fd = FD_SNAP_FD( i );
2077 0 : if( fd==ctx->local_out.full_snapshot_fd || fd==ctx->local_out.incremental_snapshot_fd ) continue;
2078 0 : if( FD_UNLIKELY( close( fd ) ) )
2079 0 : FD_LOG_ERR(( "close(snapshot pool fd %d) failed (%i-%s)", fd, errno, fd_io_strerror( errno ) ));
2080 0 : }
2081 0 : }
2082 0 : }
2083 :
2084 : static inline fd_snapct_out_link_t
2085 : out1( fd_topo_t const * topo,
2086 : fd_topo_tile_t const * tile,
2087 0 : char const * name ) {
2088 0 : ulong idx = ULONG_MAX;
2089 :
2090 0 : for( ulong i=0UL; i<tile->out_cnt; i++ ) {
2091 0 : fd_topo_link_t const * link = &topo->links[ tile->out_link_id[ i ] ];
2092 0 : if( !strcmp( link->name, name ) ) {
2093 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 ));
2094 0 : idx = i;
2095 0 : }
2096 0 : }
2097 :
2098 0 : if( FD_UNLIKELY( idx==ULONG_MAX ) ) return (fd_snapct_out_link_t){ .idx = ULONG_MAX, .mem = NULL, .chunk0 = 0, .wmark = 0, .chunk = 0, .mtu = 0 };
2099 :
2100 0 : ulong mtu = topo->links[ tile->out_link_id[ idx ] ].mtu;
2101 0 : if( FD_UNLIKELY( mtu==0UL ) ) return (fd_snapct_out_link_t){ .idx = idx, .mem = NULL, .chunk0 = ULONG_MAX, .wmark = ULONG_MAX, .chunk = ULONG_MAX, .mtu = mtu };
2102 :
2103 0 : void * mem = topo->workspaces[ topo->objs[ topo->links[ tile->out_link_id[ idx ] ].dcache_obj_id ].wksp_id ].wksp;
2104 0 : ulong chunk0 = fd_dcache_compact_chunk0( mem, topo->links[ tile->out_link_id[ idx ] ].dcache );
2105 0 : ulong wmark = fd_dcache_compact_wmark ( mem, topo->links[ tile->out_link_id[ idx ] ].dcache, mtu );
2106 0 : return (fd_snapct_out_link_t){ .idx = idx, .mem = mem, .chunk0 = chunk0, .wmark = wmark, .chunk = chunk0, .mtu = mtu };
2107 0 : }
2108 :
2109 : static void
2110 : unprivileged_init( fd_topo_t const * topo,
2111 0 : fd_topo_tile_t const * tile ) {
2112 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
2113 :
2114 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
2115 0 : fd_snapct_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapct_tile_t), sizeof(fd_snapct_tile_t) );
2116 0 : FD_SCRATCH_ALLOC_APPEND( l, fd_ssping_align(), fd_ssping_footprint( TOTAL_PEERS_MAX ) );
2117 0 : void * _ci_table = FD_SCRATCH_ALLOC_APPEND( l, alignof(gossip_ci_entry_t), sizeof(gossip_ci_entry_t) * GOSSIP_PEERS_MAX );
2118 0 : void * _ci_map = FD_SCRATCH_ALLOC_APPEND( l, gossip_ci_map_align(), gossip_ci_map_footprint( gossip_ci_map_chain_cnt_est( GOSSIP_PEERS_MAX ) ) );
2119 0 : FD_SCRATCH_ALLOC_APPEND( l, fd_http_resolver_align(), fd_http_resolver_footprint( SERVER_PEERS_MAX ) );
2120 0 : void * _selector = FD_SCRATCH_ALLOC_APPEND( l, fd_sspeer_selector_align(), fd_sspeer_selector_footprint( TOTAL_PEERS_MAX ) );
2121 0 : void * _bl_pool = FD_SCRATCH_ALLOC_APPEND( l, blacklist_pool_align(), blacklist_pool_footprint( TOTAL_PEERS_MAX ) );
2122 0 : void * _bl_map = FD_SCRATCH_ALLOC_APPEND( l, blacklist_map_align(), blacklist_map_footprint( blacklist_map_chain_cnt_est( TOTAL_PEERS_MAX ) ) );
2123 0 : void * _adns = FD_SCRATCH_ALLOC_APPEND( l, fd_adns_align(), fd_adns_footprint( ADNS_REQS_MAX ) );
2124 :
2125 0 : ctx->config = tile->snapct;
2126 0 : ctx->gossip_enabled = gossip_enabled( tile );
2127 0 : ctx->download_enabled = download_enabled( tile );
2128 :
2129 0 : ctx->adns = NULL;
2130 0 : if( FD_LIKELY( ctx->download_enabled ) ) {
2131 0 : ctx->adns = fd_adns_join( fd_adns_new( _adns, ADNS_REQS_MAX ) );
2132 0 : FD_TEST( ctx->adns );
2133 0 : for( ulong i=0UL; i<ctx->config.sources.servers_cnt; i++ ) {
2134 0 : fd_dns_peer_parse( ctx->config.sources.servers[ i ], "snapshots.sources.servers", ctx->dns_servers[ i ].hostname, &ctx->dns_servers[ i ].port, &ctx->dns_servers[ i ].is_https );
2135 0 : ctx->dns_servers[ i ].resolved = 0;
2136 0 : FD_TEST( !fd_adns_resolve( ctx->adns, ctx->dns_servers[ i ].hostname, i ) );
2137 0 : ctx->dns_servers[ i ].retry_nanos = 0L; /* in flight */
2138 0 : }
2139 0 : for( ulong i=0UL; i<ctx->config.entrypoints_cnt; i++ ) {
2140 0 : fd_dns_peer_parse( ctx->config.entrypoints[ i ], "gossip.entrypoints", ctx->dns_entrypoints[ i ].hostname, &ctx->dns_entrypoints[ i ].port, NULL );
2141 0 : ctx->dns_entrypoints[ i ].resolved = 0;
2142 0 : FD_TEST( !fd_adns_resolve( ctx->adns, ctx->dns_entrypoints[ i ].hostname, DNS_REQ_ID_ENTRYPOINT|i ) );
2143 0 : ctx->dns_entrypoints[ i ].retry_nanos = 0L;
2144 0 : }
2145 0 : }
2146 :
2147 0 : ctx->selector = fd_sspeer_selector_join( fd_sspeer_selector_new( _selector, TOTAL_PEERS_MAX, ctx->selector_seed ) );
2148 0 : ctx->blacklist_pool = blacklist_pool_join( blacklist_pool_new( _bl_pool, TOTAL_PEERS_MAX ) );
2149 0 : ctx->blacklist_map = blacklist_map_join( blacklist_map_new( _bl_map, blacklist_map_chain_cnt_est( TOTAL_PEERS_MAX ), ctx->blacklist_seed ) );
2150 :
2151 0 : if( FD_UNLIKELY( !ctx->config.incremental_snapshots ) ) {
2152 0 : FD_LOG_WARNING(( "incremental snapshots disabled via [snapshots.incremental_snapshots]." ));
2153 0 : }
2154 :
2155 0 : ctx->state = FD_SNAPCT_STATE_INIT;
2156 0 : ctx->malformed = 0;
2157 0 : ctx->load_complete = 0;
2158 0 : FD_CHECK_ERR( ctx->config.wait_for_peers_timeout_nanos>0L, "snapct wait_for_peers_timeout_nanos must be positive" );
2159 0 : ctx->deadline_nanos = fd_log_wallclock() + ctx->config.wait_for_peers_timeout_nanos;
2160 0 : ctx->flush_ack = 0;
2161 0 : ctx->flush_ack_cnt = 0;
2162 0 : ctx->peer.addr.l = 0UL;
2163 :
2164 0 : fd_memset( ctx->http_full_snapshot_name, 0, PATH_MAX );
2165 0 : fd_memset( ctx->http_incr_snapshot_name, 0, PATH_MAX );
2166 :
2167 0 : ctx->gossip_in_mem = NULL;
2168 0 : int has_snapld_dc = 0, ack_cnt = 0;
2169 0 : FD_TEST( tile->in_cnt<=MAX_IN_LINKS );
2170 0 : for( ulong i=0UL; i<(tile->in_cnt); i++ ) {
2171 0 : fd_topo_link_t const * in_link = &topo->links[ tile->in_link_id[ i ] ];
2172 0 : if( 0==strcmp( in_link->name, "gossip_out" ) ) {
2173 0 : ctx->in_kind[ i ] = IN_KIND_GOSSIP;
2174 0 : ctx->gossip_in_mem = topo->workspaces[ topo->objs[ in_link->dcache_obj_id ].wksp_id ].wksp;
2175 0 : } else if( 0==strcmp( in_link->name, "snapld_dc" ) ) {
2176 0 : ctx->in_kind[ i ] = IN_KIND_SNAPLD;
2177 0 : ctx->snapld_in_mem = topo->workspaces[ topo->objs[ in_link->dcache_obj_id ].wksp_id ].wksp;
2178 0 : FD_TEST( !has_snapld_dc );
2179 0 : has_snapld_dc = 1;
2180 0 : } else if( 0==strcmp( in_link->name, "snapin_ct" ) || 0==strcmp( in_link->name, "snapwr_ct" ) ){
2181 0 : ctx->in_kind[ i ] = IN_KIND_ACK;
2182 0 : ack_cnt++;
2183 0 : }
2184 0 : }
2185 0 : FD_TEST( has_snapld_dc && ack_cnt>0 );
2186 0 : ctx->flush_ack_cnt = ack_cnt + 1; /* +1 for snapld (acks via snapld_dc) */
2187 0 : FD_TEST( ctx->gossip_enabled==(ctx->gossip_in_mem!=NULL) );
2188 :
2189 0 : ctx->predicted_incremental.full_slot = FD_SSPEER_SLOT_UNKNOWN;
2190 0 : ctx->predicted_incremental.slot = FD_SSPEER_SLOT_UNKNOWN;
2191 0 : ctx->predicted_incremental.pending = 0;
2192 :
2193 0 : fd_memset( &ctx->metrics, 0, sizeof(ctx->metrics) );
2194 :
2195 0 : fd_memset( _ci_table, 0, sizeof(gossip_ci_entry_t) * GOSSIP_PEERS_MAX );
2196 0 : ctx->gossip.ci_table = _ci_table;
2197 0 : ctx->gossip.ci_map = gossip_ci_map_join( gossip_ci_map_new( _ci_map, gossip_ci_map_chain_cnt_est( GOSSIP_PEERS_MAX ), ctx->gossip_ci_seed ) );
2198 0 : ctx->gossip.allowed_cnt = 0UL;
2199 0 : ctx->gossip.saturated = !ctx->gossip_enabled;
2200 :
2201 0 : if( FD_UNLIKELY( tile->out_cnt<2UL || tile->out_cnt>3UL ) ) FD_LOG_ERR(( "tile `" NAME "` has %lu outs, expected 2-3", tile->out_cnt ));
2202 0 : ctx->out_ld = out1( topo, tile, "snapct_ld" );
2203 0 : ctx->out_gui = out1( topo, tile, "snapct_gui" );
2204 0 : ctx->out_rp = out1( topo, tile, "snapct_repr" );
2205 0 : }
2206 :
2207 : /* after_credit can result in as many as 5 stem publishes in some code
2208 : paths, and returnable_frag can result in 1. */
2209 0 : #define STEM_BURST 6UL
2210 :
2211 0 : #define STEM_LAZY 1000L
2212 :
2213 0 : #define STEM_CALLBACK_CONTEXT_TYPE fd_snapct_tile_t
2214 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_snapct_tile_t)
2215 :
2216 : #define STEM_CALLBACK_SHOULD_SHUTDOWN should_shutdown
2217 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
2218 0 : #define STEM_CALLBACK_AFTER_CREDIT after_credit
2219 0 : #define STEM_CALLBACK_RETURNABLE_FRAG returnable_frag
2220 :
2221 : #include "../../disco/stem/fd_stem.c"
2222 :
2223 : #ifndef FD_TILE_TEST
2224 : fd_topo_run_tile_t fd_tile_snapct = {
2225 : .name = NAME,
2226 : .rlimit_file_cnt_fn = rlimit_file_cnt,
2227 : .populate_allowed_seccomp = populate_allowed_seccomp,
2228 : .populate_allowed_fds = populate_allowed_fds,
2229 : .scratch_align = scratch_align,
2230 : .scratch_footprint = scratch_footprint,
2231 : .privileged_init = privileged_init,
2232 : .unprivileged_init = unprivileged_init,
2233 : .run = stem_run,
2234 : .keep_host_networking = 1,
2235 : .allow_connect = 1,
2236 : .allow_renameat = 1,
2237 : };
2238 : #endif
2239 :
2240 : #undef NAME
|