Line data Source code
1 : #define _GNU_SOURCE
2 : #include "utils/fd_ssctrl.h"
3 : #include "utils/fd_ssparse.h"
4 : #include "utils/fd_ssmanifest_parser.h"
5 :
6 : #include "../../disco/topo/fd_topo.h"
7 : #include "../../disco/metrics/fd_metrics.h"
8 : #include "../../flamenco/accdb/fd_accdb_shmem.h"
9 :
10 : #include "generated/fd_snapwr_tile_seccomp.h"
11 :
12 : #include <errno.h>
13 : #include <fcntl.h>
14 : #include <unistd.h>
15 :
16 : #define NAME "snapwr"
17 :
18 0 : #define FD_SNAPWR_WRITE_BUF_SZ (2UL<<20) /* 2MiB */
19 :
20 : struct fd_snapwr_out {
21 : ulong idx;
22 : fd_wksp_t * mem;
23 : ulong chunk0;
24 : ulong wmark;
25 : ulong chunk;
26 : ulong mtu;
27 : };
28 :
29 : typedef struct fd_snapwr_out fd_snapwr_out_t;
30 :
31 : struct fd_snapwr_tile {
32 : int full;
33 : int state;
34 :
35 : ulong lane_cnt;
36 : ulong expected_frame;
37 : ulong pending_control; /* control message expected from snapdc tiles */
38 : uchar control_seen[ FD_TOPO_MAX_TILE_IN_LINKS ];
39 :
40 : ulong partition_sz;
41 :
42 : ulong accounts_off;
43 : ulong flush_off;
44 :
45 : uchar * write_buf;
46 : ulong write_buf_used;
47 :
48 : ulong seed;
49 :
50 : fd_ssparse_t ssparse[1];
51 : fd_ssmanifest_parser_t * manifest_parser;
52 :
53 : struct {
54 : fd_wksp_t * wksp;
55 : ulong chunk0;
56 : ulong wmark;
57 : ulong mtu;
58 : ulong pos;
59 : } in[ FD_TOPO_MAX_TILE_IN_LINKS ];
60 :
61 : fd_snapwr_out_t ct_out;
62 :
63 : struct {
64 : ulong accounts_off;
65 : ulong flush_off;
66 : } recovery;
67 :
68 : struct {
69 : ulong full_bytes_read;
70 : ulong incremental_bytes_read;
71 : ulong bytes_written;
72 : ulong accounts_written;
73 : ulong full_accounts_written;
74 : } metrics;
75 :
76 : fd_snapshot_manifest_t manifest[1];
77 : };
78 :
79 : typedef struct fd_snapwr_tile fd_snapwr_tile_t;
80 :
81 : static inline int
82 0 : should_shutdown( fd_snapwr_tile_t * ctx ) {
83 0 : return ctx->state==FD_SNAPSHOT_STATE_SHUTDOWN;
84 0 : }
85 :
86 : static void
87 0 : metrics_write( fd_snapwr_tile_t * ctx ) {
88 0 : FD_MGAUGE_SET( SNAPWR, FULL_BYTES_READ, ctx->metrics.full_bytes_read );
89 0 : FD_MGAUGE_SET( SNAPWR, INCREMENTAL_BYTES_READ, ctx->metrics.incremental_bytes_read );
90 0 : FD_MGAUGE_SET( SNAPWR, BYTES_WRITTEN, ctx->metrics.bytes_written );
91 0 : FD_MGAUGE_SET( SNAPWR, ACCOUNTS_WRITTEN, ctx->metrics.accounts_written );
92 0 : }
93 :
94 : static ulong
95 0 : scratch_align( void ) {
96 0 : return 512UL;
97 0 : }
98 :
99 : static ulong
100 0 : scratch_footprint( fd_topo_tile_t const * tile ) {
101 0 : (void)tile;
102 0 : ulong l = FD_LAYOUT_INIT;
103 0 : l = FD_LAYOUT_APPEND( l, alignof(fd_snapwr_tile_t), sizeof(fd_snapwr_tile_t) );
104 0 : l = FD_LAYOUT_APPEND( l, fd_ssmanifest_parser_align(), fd_ssmanifest_parser_footprint() );
105 0 : l = FD_LAYOUT_APPEND( l, 1UL, FD_SNAPWR_WRITE_BUF_SZ );
106 0 : return FD_LAYOUT_FINI( l, scratch_align() );
107 0 : }
108 :
109 : static inline void
110 186 : clear_control_barrier( fd_snapwr_tile_t * ctx ) {
111 186 : ctx->pending_control = ULONG_MAX;
112 186 : fd_memset( ctx->control_seen, 0, sizeof(ctx->control_seen) );
113 186 : }
114 :
115 : static void
116 : transition_malformed( fd_snapwr_tile_t * ctx,
117 9 : fd_stem_context_t * stem ) {
118 9 : if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_ERROR ) ) return;
119 9 : ctx->state = FD_SNAPSHOT_STATE_ERROR;
120 9 : fd_stem_publish( stem, ctx->ct_out.idx, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
121 9 : }
122 :
123 : static void
124 9 : buffer_flush( fd_snapwr_tile_t * ctx ) {
125 9 : if( FD_UNLIKELY( !ctx->write_buf_used ) ) return;
126 :
127 0 : ulong sz = ctx->write_buf_used;
128 0 : ulong off = ctx->flush_off;
129 0 : ulong bytes_written = 0UL;
130 0 : while( bytes_written<sz ) {
131 0 : long res = pwrite( FD_ACCDB_FD_RW, ctx->write_buf+bytes_written, sz-bytes_written, (long)(off+bytes_written) );
132 0 : if( FD_UNLIKELY( -1==res ) ) FD_LOG_ERR(( "error writing to disk (%d-%s)", errno, fd_io_strerror( errno ) ));
133 0 : bytes_written += (ulong)res;
134 0 : ctx->metrics.bytes_written += (ulong)res;
135 0 : }
136 0 : ctx->flush_off += sz;
137 0 : ctx->write_buf_used = 0UL;
138 0 : }
139 :
140 : static void
141 : buffer_write( fd_snapwr_tile_t * ctx,
142 : uchar const * data,
143 0 : ulong sz ) {
144 0 : ctx->accounts_off += sz;
145 0 : while( sz ) {
146 0 : ulong avail = FD_SNAPWR_WRITE_BUF_SZ - ctx->write_buf_used;
147 0 : ulong n = fd_ulong_min( sz, avail );
148 0 : fd_memcpy( ctx->write_buf + ctx->write_buf_used, data, n );
149 0 : ctx->write_buf_used += n;
150 0 : data += n;
151 0 : sz -= n;
152 0 : if( FD_UNLIKELY( ctx->write_buf_used==FD_SNAPWR_WRITE_BUF_SZ ) ) buffer_flush( ctx );
153 0 : }
154 0 : }
155 :
156 : static void
157 : buffer_skip( fd_snapwr_tile_t * ctx,
158 0 : ulong sz ) {
159 0 : buffer_flush( ctx );
160 0 : ctx->accounts_off += sz;
161 0 : ctx->flush_off += sz;
162 0 : }
163 :
164 : static void
165 : process_account_header( fd_snapwr_tile_t * ctx,
166 0 : fd_ssparse_advance_result_t * result ) {
167 : /* Ensure header+data does not cross a partition boundary. If it
168 : would, pad with zeros so the account starts at the next one. */
169 0 : ulong account_sz = sizeof(fd_accdb_disk_meta_t) + (ulong)result->account_header.data_len;
170 0 : ulong cur_boundary = ctx->accounts_off / ctx->partition_sz;
171 0 : ulong end_boundary = (ctx->accounts_off + account_sz - 1UL) / ctx->partition_sz;
172 0 : if( FD_UNLIKELY( cur_boundary!=end_boundary ) ) {
173 0 : ulong next = (cur_boundary + 1UL) * ctx->partition_sz;
174 0 : buffer_skip( ctx, next - ctx->accounts_off );
175 0 : }
176 :
177 0 : fd_accdb_disk_meta_t meta;
178 0 : fd_memcpy( meta.pubkey, result->account_header.pubkey, 32UL );
179 0 : meta.size = (uint)result->account_header.data_len;
180 0 : meta.generation = 0U;
181 0 : fd_memcpy( meta.owner, result->account_header.owner, 32UL );
182 0 : buffer_write( ctx, meta.b, sizeof(fd_accdb_disk_meta_t) );
183 0 : ctx->metrics.accounts_written++;
184 0 : }
185 :
186 : static void
187 : process_account_data( fd_snapwr_tile_t * ctx,
188 0 : fd_ssparse_advance_result_t * result ) {
189 0 : buffer_write( ctx, result->account_data.data, result->account_data.data_sz );
190 0 : }
191 :
192 : static int
193 : handle_data_frag( fd_snapwr_tile_t * ctx,
194 : ulong in_idx,
195 : ulong chunk,
196 : ulong sz,
197 60 : fd_stem_context_t * stem ) {
198 60 : if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_FINISHING ) ) {
199 3 : FD_LOG_WARNING(( "received unexpected data frag while in state %s (%lu)",
200 3 : fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
201 3 : transition_malformed( ctx, stem );
202 3 : return 0;
203 3 : }
204 57 : if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_ERROR ) ) {
205 : /* Ignore all data frags after observing an error in the stream until
206 : we receive fail & init control messages to restart processing. */
207 0 : return 0;
208 0 : }
209 57 : if( FD_UNLIKELY( ctx->state!=FD_SNAPSHOT_STATE_PROCESSING ) ) {
210 0 : FD_LOG_ERR(( "received data frag during invalid state %s (%lu)",
211 0 : fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
212 0 : }
213 :
214 57 : FD_TEST( chunk>=ctx->in[ in_idx ].chunk0 && chunk<=ctx->in[ in_idx ].wmark && sz<=ctx->in[ in_idx ].mtu );
215 :
216 69 : for(;;) {
217 69 : if( FD_UNLIKELY( sz-ctx->in[ in_idx ].pos==0UL ) ) break;
218 :
219 12 : uchar const * data = (uchar const *)fd_chunk_to_laddr_const( ctx->in[ in_idx ].wksp, chunk ) + ctx->in[ in_idx ].pos;
220 :
221 12 : fd_ssparse_advance_result_t result[1];
222 12 : int res = fd_ssparse_advance( ctx->ssparse, data, sz-ctx->in[ in_idx ].pos, result );
223 12 : switch( res ) {
224 0 : case FD_SSPARSE_ADVANCE_ERROR:
225 0 : FD_LOG_WARNING(( "error while parsing snapshot stream" ));
226 0 : transition_malformed( ctx, stem );
227 0 : return 0;
228 3 : case FD_SSPARSE_ADVANCE_AGAIN:
229 3 : break;
230 0 : case FD_SSPARSE_ADVANCE_MANIFEST:
231 0 : case FD_SSPARSE_ADVANCE_MANIFEST_DONE: {
232 0 : int res = fd_ssmanifest_parser_consume( ctx->manifest_parser,
233 0 : result->manifest.data,
234 0 : result->manifest.data_sz );
235 0 : if( FD_UNLIKELY( res==FD_SSMANIFEST_PARSER_ADVANCE_ERROR ) ) {
236 0 : FD_LOG_WARNING(( "error while parsing snapshot manifest" ));
237 0 : transition_malformed( ctx, stem );
238 0 : return 0;
239 0 : }
240 0 : break;
241 0 : }
242 0 : case FD_SSPARSE_ADVANCE_STATUS_CACHE:
243 0 : break;
244 0 : case FD_SSPARSE_ADVANCE_ACCOUNT_HEADER:
245 0 : process_account_header( ctx, result );
246 0 : break;
247 0 : case FD_SSPARSE_ADVANCE_ACCOUNT_DATA:
248 0 : process_account_data( ctx, result );
249 0 : break;
250 0 : case FD_SSPARSE_ADVANCE_ACCOUNT_BATCH:
251 0 : FD_TEST( 0 );
252 0 : break;
253 9 : case FD_SSPARSE_ADVANCE_DONE:
254 9 : buffer_flush( ctx );
255 9 : ctx->state = FD_SNAPSHOT_STATE_FINISHING;
256 9 : break;
257 0 : default:
258 0 : FD_LOG_ERR(( "unexpected fd_ssparse_advance result %d", res ));
259 0 : break;
260 12 : }
261 :
262 12 : ctx->in[ in_idx ].pos += result->bytes_consumed;
263 12 : if( FD_LIKELY( ctx->full ) ) ctx->metrics.full_bytes_read += result->bytes_consumed;
264 0 : else ctx->metrics.incremental_bytes_read += result->bytes_consumed;
265 12 : }
266 :
267 57 : int reprocess_frag = ctx->in[ in_idx ].pos<sz;
268 57 : if( FD_LIKELY( !reprocess_frag ) ) ctx->in[ in_idx ].pos = 0UL;
269 57 : return reprocess_frag;
270 57 : }
271 :
272 : static void
273 : handle_control_frag( fd_snapwr_tile_t * ctx,
274 : fd_stem_context_t * stem,
275 96 : ulong sig ) {
276 96 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_META ) ) return;
277 :
278 93 : if( ctx->state==FD_SNAPSHOT_STATE_ERROR && sig!=FD_SNAPSHOT_MSG_CTRL_FAIL ) {
279 : /* Control messages move along the snapshot load pipeline. Since
280 : error conditions can be triggered by any tile in the pipeline,
281 : it is possible to be in error state and still receive otherwise
282 : valid messages. Only a fail message can revert this. */
283 0 : return;
284 93 : };
285 :
286 93 : int forward_msg = 1;
287 :
288 93 : switch( sig ) {
289 0 : case FD_SNAPSHOT_MSG_CTRL_INIT_FULL:
290 6 : case FD_SNAPSHOT_MSG_CTRL_INIT_INCR: {
291 6 : FD_TEST( ctx->state==FD_SNAPSHOT_STATE_IDLE );
292 6 : ctx->state = FD_SNAPSHOT_STATE_PROCESSING;
293 6 : ctx->expected_frame = 0UL;
294 18 : for( ulong i=0UL; i<ctx->lane_cnt; i++ ) {
295 12 : ctx->in[ i ].pos = 0UL;
296 12 : }
297 6 : fd_ssparse_init( ctx->ssparse );
298 6 : fd_ssmanifest_parser_init( ctx->manifest_parser, ctx->manifest );
299 :
300 : /* Rewind metric counters (no-op unless recovering from a fail) */
301 6 : if( sig==FD_SNAPSHOT_MSG_CTRL_INIT_FULL ) {
302 0 : ctx->metrics.full_bytes_read = 0UL;
303 0 : ctx->metrics.incremental_bytes_read = 0UL;
304 0 : ctx->metrics.accounts_written = ctx->metrics.full_accounts_written = 0UL;
305 6 : } else {
306 6 : ctx->metrics.incremental_bytes_read = 0UL;
307 6 : ctx->metrics.accounts_written = ctx->metrics.full_accounts_written;
308 6 : }
309 6 : break;
310 6 : }
311 18 : case FD_SNAPSHOT_MSG_CTRL_FINI: {
312 : /* This is a special case: handle_data_frag must have already
313 : processed FD_SSPARSE_ADVANCE_DONE and moved the state into
314 : FD_SNAPSHOT_STATE_FINISHING. Otherwise, treat this as a
315 : malformed snapshot so that the pipeline can retry. */
316 18 : if( FD_UNLIKELY( ctx->state!=FD_SNAPSHOT_STATE_FINISHING ) ) {
317 3 : FD_LOG_WARNING(( "received FINI while in state %s (%lu), expected FINISHING (possibly truncated tar stream)",
318 3 : fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
319 3 : transition_malformed( ctx, stem );
320 3 : forward_msg = 0;
321 3 : break;
322 3 : }
323 15 : break;
324 18 : }
325 :
326 15 : case FD_SNAPSHOT_MSG_CTRL_NEXT: {
327 6 : FD_TEST( ctx->state==FD_SNAPSHOT_STATE_FINISHING );
328 6 : ctx->state = FD_SNAPSHOT_STATE_IDLE;
329 6 : ctx->recovery.accounts_off = ctx->accounts_off;
330 6 : ctx->recovery.flush_off = ctx->flush_off;
331 6 : ctx->metrics.full_accounts_written = ctx->metrics.accounts_written;
332 6 : break;
333 6 : }
334 :
335 3 : case FD_SNAPSHOT_MSG_CTRL_DONE: {
336 3 : FD_TEST( ctx->state==FD_SNAPSHOT_STATE_FINISHING );
337 3 : ctx->state = FD_SNAPSHOT_STATE_IDLE;
338 3 : break;
339 3 : }
340 :
341 15 : case FD_SNAPSHOT_MSG_CTRL_ERROR: {
342 15 : FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
343 15 : ctx->state = FD_SNAPSHOT_STATE_ERROR;
344 15 : break;
345 15 : }
346 :
347 42 : case FD_SNAPSHOT_MSG_CTRL_FAIL: {
348 42 : FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
349 42 : ctx->write_buf_used = 0UL;
350 42 : ctx->accounts_off = ctx->full ? 0UL : ctx->recovery.accounts_off;
351 42 : ctx->flush_off = ctx->full ? 0UL : ctx->recovery.flush_off;
352 42 : ctx->state = FD_SNAPSHOT_STATE_IDLE;
353 42 : break;
354 42 : }
355 :
356 3 : case FD_SNAPSHOT_MSG_CTRL_SHUTDOWN: {
357 3 : FD_TEST( ctx->state==FD_SNAPSHOT_STATE_IDLE );
358 3 : ctx->state = FD_SNAPSHOT_STATE_SHUTDOWN;
359 3 : break;
360 3 : }
361 :
362 0 : default: {
363 0 : FD_LOG_ERR(( "unexpected control frag %s (%lu) in state %s (%lu)",
364 0 : fd_ssctrl_msg_ctrl_str( sig ), sig,
365 0 : fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
366 0 : break;
367 0 : }
368 93 : }
369 :
370 93 : if( FD_LIKELY( forward_msg ) ) {
371 90 : fd_stem_publish( stem, ctx->ct_out.idx, sig, 0UL, 0UL, 0UL, 0UL, 0UL );
372 90 : }
373 93 : }
374 :
375 : static inline int
376 240 : all_controls_seen( fd_snapwr_tile_t const * ctx ) {
377 240 : int all_seen = 1;
378 978 : for( ulong i=0UL; i<ctx->lane_cnt; i++ ) {
379 738 : all_seen &= !!ctx->control_seen[ i ];
380 738 : }
381 240 : return all_seen;
382 240 : }
383 :
384 : static inline int
385 : before_frag( fd_snapwr_tile_t * ctx,
386 : ulong in_idx,
387 : ulong seq FD_PARAM_UNUSED,
388 162 : ulong sig ) {
389 : /* If we're currently in ERROR state we should only process FAIL
390 : control frags */
391 162 : if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_ERROR ) ) {
392 12 : return sig!=FD_SNAPSHOT_MSG_CTRL_FAIL;
393 12 : }
394 :
395 150 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_CTRL_ERROR ) ) {
396 6 : return 0;
397 6 : }
398 :
399 : /* Once this lane sends the pending control, hold its later frags until
400 : all snapdc lanes send the same control. */
401 144 : if( FD_UNLIKELY( ctx->pending_control!=ULONG_MAX && ctx->control_seen[ in_idx ] ) ) {
402 3 : FD_TEST( sig!=ctx->pending_control );
403 3 : return -1;
404 3 : }
405 :
406 : /* Only accept DATA frags from the expected lane */
407 141 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_DATA && in_idx!=ctx->expected_frame%ctx->lane_cnt ) ) {
408 39 : return -1;
409 39 : }
410 :
411 102 : return 0;
412 141 : }
413 :
414 : static inline int
415 : handle_lane_data_frag( fd_snapwr_tile_t * ctx,
416 : fd_stem_context_t * stem,
417 : ulong in_idx,
418 : ulong chunk,
419 : ulong sz,
420 63 : ulong ctl ) {
421 : /* EOM marks the end of a frame */
422 63 : int eom = !!fd_frag_meta_ctl_eom( ctl );
423 :
424 : /* Tar EOF can precede a zero-byte EOM that closes the frame. */
425 63 : int trailing_eom = ctx->state==FD_SNAPSHOT_STATE_FINISHING && eom && !sz;
426 63 : if( FD_UNLIKELY( !trailing_eom && handle_data_frag( ctx, in_idx, chunk, sz, stem ) ) ) {
427 0 : return 1;
428 0 : }
429 :
430 63 : if( FD_UNLIKELY( eom ) ) {
431 57 : ctx->expected_frame++;
432 57 : }
433 :
434 63 : return 0;
435 63 : }
436 :
437 : static inline void
438 : handle_control_barrier( fd_snapwr_tile_t * ctx,
439 : fd_stem_context_t * stem,
440 : ulong in_idx,
441 255 : ulong sig ) {
442 : /* Error control frags must be immediately handled. */
443 255 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_CTRL_ERROR ) ) {
444 15 : handle_control_frag( ctx, stem, sig );
445 15 : return;
446 15 : }
447 :
448 240 : if( FD_UNLIKELY( sig!=ctx->pending_control ) ) {
449 105 : FD_TEST( ctx->pending_control==ULONG_MAX || sig==FD_SNAPSHOT_MSG_CTRL_FAIL );
450 105 : clear_control_barrier( ctx );
451 :
452 : /* Record the new attempt type on the first INIT copy so FAIL can
453 : roll back correctly if ERROR interrupts a pending INIT message */
454 105 : if( sig==FD_SNAPSHOT_MSG_CTRL_INIT_FULL || sig==FD_SNAPSHOT_MSG_CTRL_INIT_INCR ) {
455 18 : ctx->full = sig==FD_SNAPSHOT_MSG_CTRL_INIT_FULL;
456 18 : }
457 :
458 105 : ctx->pending_control = sig;
459 105 : }
460 :
461 : /* Only process the control frag when all upstream tiles have sent
462 : the same control message. */
463 240 : FD_TEST( !ctx->control_seen[ in_idx ] );
464 240 : ctx->control_seen[ in_idx ] = 1U;
465 240 : if( FD_LIKELY( !all_controls_seen( ctx ) ) ) {
466 159 : return;
467 159 : }
468 :
469 : /* All controls received, process the control frag. */
470 81 : clear_control_barrier( ctx );
471 81 : handle_control_frag( ctx, stem, sig );
472 81 : }
473 :
474 : static inline int
475 : returnable_frag( fd_snapwr_tile_t * ctx,
476 : ulong in_idx,
477 : ulong seq FD_PARAM_UNUSED,
478 : ulong sig,
479 : ulong chunk,
480 : ulong sz,
481 : ulong ctl,
482 : ulong tsorig FD_PARAM_UNUSED,
483 : ulong tspub FD_PARAM_UNUSED,
484 318 : fd_stem_context_t * stem ) {
485 318 : FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
486 :
487 318 : if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_DATA ) ) return handle_lane_data_frag( ctx, stem, in_idx, chunk, sz, ctl );
488 255 : else handle_control_barrier( ctx, stem, in_idx, sig );
489 :
490 255 : return 0;
491 318 : }
492 :
493 : static ulong
494 : populate_allowed_fds( fd_topo_t const * topo FD_PARAM_UNUSED,
495 : fd_topo_tile_t const * tile FD_PARAM_UNUSED,
496 : ulong out_fds_cnt,
497 0 : int * out_fds ) {
498 0 : if( FD_UNLIKELY( out_fds_cnt<3UL ) ) FD_LOG_ERR(( "invalid out_fds_cnt %lu", out_fds_cnt ));
499 :
500 0 : ulong out_cnt = 0;
501 0 : out_fds[ out_cnt++ ] = 2UL; /* stderr */
502 0 : if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) ) {
503 0 : out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
504 0 : }
505 0 : out_fds[ out_cnt++ ] = FD_ACCDB_FD_RW; /* accounts db */
506 :
507 0 : return out_cnt;
508 0 : }
509 :
510 : static ulong
511 : populate_allowed_seccomp( fd_topo_t const * topo,
512 : fd_topo_tile_t const * tile,
513 : ulong out_cnt,
514 0 : struct sock_filter * out ) {
515 0 : (void)topo; (void)tile;
516 :
517 0 : populate_sock_filter_policy_fd_snapwr_tile( out_cnt, out, (uint)fd_log_private_logfile_fd(), (uint)FD_ACCDB_FD_RW );
518 0 : return sock_filter_policy_fd_snapwr_tile_instr_cnt;
519 0 : }
520 :
521 : static inline fd_snapwr_out_t
522 : out1( fd_topo_t const * topo,
523 : fd_topo_tile_t const * tile,
524 0 : char const * name ) {
525 0 : ulong idx = fd_topo_find_tile_out_link( topo, tile, name, 0UL );
526 :
527 0 : if( FD_UNLIKELY( idx==ULONG_MAX ) ) return (fd_snapwr_out_t){ .idx = ULONG_MAX, .mem = NULL, .chunk0 = 0, .wmark = 0, .chunk = 0, .mtu = 0 };
528 :
529 0 : ulong mtu = topo->links[ tile->out_link_id[ idx ] ].mtu;
530 0 : if( FD_UNLIKELY( mtu==0UL ) ) return (fd_snapwr_out_t){ .idx = idx, .mem = NULL, .chunk0 = ULONG_MAX, .wmark = ULONG_MAX, .chunk = ULONG_MAX, .mtu = mtu };
531 :
532 0 : void * mem = topo->workspaces[ topo->objs[ topo->links[ tile->out_link_id[ idx ] ].dcache_obj_id ].wksp_id ].wksp;
533 0 : ulong chunk0 = fd_dcache_compact_chunk0( mem, topo->links[ tile->out_link_id[ idx ] ].dcache );
534 0 : ulong wmark = fd_dcache_compact_wmark ( mem, topo->links[ tile->out_link_id[ idx ] ].dcache, mtu );
535 0 : return (fd_snapwr_out_t){ .idx = idx, .mem = mem, .chunk0 = chunk0, .wmark = wmark, .chunk = chunk0, .mtu = mtu };
536 0 : }
537 :
538 : static void
539 : privileged_init( fd_topo_t const * topo,
540 0 : fd_topo_tile_t const * tile ) {
541 0 : fd_snapwr_tile_t * ctx = fd_topo_obj_laddr( topo, tile->tile_obj_id );
542 0 : FD_TEST( fd_rng_secure( &ctx->seed, 8UL ) );
543 0 : }
544 :
545 : static void
546 : unprivileged_init( fd_topo_t const * topo,
547 0 : fd_topo_tile_t const * tile ) {
548 0 : void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
549 :
550 0 : FD_SCRATCH_ALLOC_INIT( l, scratch );
551 0 : fd_snapwr_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapwr_tile_t), sizeof(fd_snapwr_tile_t) );
552 0 : void * _manifest_parser = FD_SCRATCH_ALLOC_APPEND( l, fd_ssmanifest_parser_align(), fd_ssmanifest_parser_footprint() );
553 0 : void * _write_buf = FD_SCRATCH_ALLOC_APPEND( l, 1UL, FD_SNAPWR_WRITE_BUF_SZ );
554 :
555 0 : ctx->full = 1;
556 0 : ctx->state = FD_SNAPSHOT_STATE_IDLE;
557 0 : ctx->lane_cnt = tile->in_cnt;
558 0 : ctx->expected_frame = 0UL;
559 0 : clear_control_barrier( ctx );
560 :
561 0 : ctx->partition_sz = tile->snapwr.partition_sz;
562 0 : if( FD_UNLIKELY( !ctx->partition_sz ) ) FD_LOG_ERR(( "tile `" NAME "` partition_sz is 0" ));
563 :
564 0 : ctx->accounts_off = 0UL;
565 0 : ctx->flush_off = 0UL;
566 0 : ctx->recovery.accounts_off = 0UL;
567 0 : ctx->recovery.flush_off = 0UL;
568 0 : ctx->write_buf = _write_buf;
569 0 : ctx->write_buf_used = 0UL;
570 :
571 0 : ctx->manifest_parser = fd_ssmanifest_parser_join( fd_ssmanifest_parser_new( _manifest_parser ) );
572 0 : FD_TEST( ctx->manifest_parser );
573 :
574 0 : fd_memset( &ctx->metrics, 0, sizeof(ctx->metrics) );
575 :
576 0 : FD_TEST( tile->in_cnt );
577 :
578 0 : ctx->ct_out = out1( topo, tile, "snapwr_ct" );
579 0 : if( FD_UNLIKELY( ctx->ct_out.idx==ULONG_MAX ) ) FD_LOG_ERR(( "tile `" NAME "` missing required out link `snapwr_ct`" ));
580 :
581 0 : fd_ssparse_init( ctx->ssparse );
582 0 : fd_ssparse_batch_enable( ctx->ssparse, 0 );
583 0 : fd_ssmanifest_parser_init( ctx->manifest_parser, ctx->manifest );
584 :
585 0 : for( ulong i=0UL; i<ctx->lane_cnt; i++ ) {
586 0 : fd_topo_link_t const * in_link = &topo->links[ tile->in_link_id[ i ] ];
587 0 : FD_TEST( 0==strcmp( in_link->name, "snapdc_in" ) );
588 0 : FD_TEST( in_link->kind_id==i );
589 0 : fd_topo_wksp_t const * in_wksp = &topo->workspaces[ topo->objs[ in_link->dcache_obj_id ].wksp_id ];
590 0 : ctx->in[ i ].wksp = in_wksp->wksp;
591 0 : ctx->in[ i ].chunk0 = fd_dcache_compact_chunk0( ctx->in[ i ].wksp, in_link->dcache );
592 0 : ctx->in[ i ].wmark = fd_dcache_compact_wmark( ctx->in[ i ].wksp, in_link->dcache, in_link->mtu );
593 0 : ctx->in[ i ].mtu = in_link->mtu;
594 0 : ctx->in[ i ].pos = 0UL;
595 0 : }
596 0 : }
597 :
598 0 : #define STEM_BURST 1UL
599 :
600 0 : #define STEM_LAZY (128L*3000L)
601 :
602 0 : #define STEM_CALLBACK_CONTEXT_TYPE fd_snapwr_tile_t
603 0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_snapwr_tile_t)
604 :
605 : #define STEM_CALLBACK_SHOULD_SHUTDOWN should_shutdown
606 0 : #define STEM_CALLBACK_METRICS_WRITE metrics_write
607 0 : #define STEM_CALLBACK_BEFORE_FRAG before_frag
608 0 : #define STEM_CALLBACK_RETURNABLE_FRAG returnable_frag
609 :
610 : #include "../../disco/stem/fd_stem.c"
611 :
612 : fd_topo_run_tile_t fd_tile_snapwr = {
613 : .name = NAME,
614 : .populate_allowed_fds = populate_allowed_fds,
615 : .populate_allowed_seccomp = populate_allowed_seccomp,
616 : .scratch_align = scratch_align,
617 : .scratch_footprint = scratch_footprint,
618 : .privileged_init = privileged_init,
619 : .unprivileged_init = unprivileged_init,
620 : .run = stem_run,
621 : };
622 :
623 : #undef NAME
|