LCOV - code coverage report
Current view: top level - discof/restore - fd_snapdc_tile.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 203 302 67.2 %
Date: 2026-09-17 04:28:31 Functions: 7 28 25.0 %

          Line data    Source code
       1             : #include "utils/fd_ssctrl.h"
       2             : #include "utils/fd_zstd_frame.h"
       3             : 
       4             : #include "../../disco/topo/fd_topo.h"
       5             : #include "../../disco/metrics/fd_metrics.h"
       6             : 
       7             : #include "generated/fd_snapdc_tile_seccomp.h"
       8             : 
       9             : #define ZSTD_STATIC_LINKING_ONLY
      10             : #include <zstd.h>
      11             : 
      12           0 : #define NAME "snapdc"
      13             : 
      14           0 : #define ZSTD_WINDOW_SZ (1UL<<25UL) /* 32MiB */
      15             : 
      16             : /* The snapdc tile is a state machine that decompresses the full and
      17             :    optionally incremental snapshot byte stream that it receives from the
      18             :    snapld tile.  In the event that the snapshot is already uncompressed,
      19             :    this tile simply copies the stream to the next tile in the pipeline. */
      20             : 
      21             : struct fd_snapdc_tile {
      22             :   uint full    : 1;
      23             :   uint is_zstd : 1;
      24             :   uint dirty   : 1;  /* in the middle of a frame? */
      25             :   int state;
      26             : 
      27             :   ulong tile_idx;
      28             :   ulong tile_count;
      29             :   ulong frame_idx;
      30             : 
      31             :   ZSTD_DCtx *     zstd;
      32             :   fd_zstd_frame_t zstd_frame[1];
      33             : 
      34             :   struct {
      35             :     fd_wksp_t * mem;
      36             :     ulong       chunk0;
      37             :     ulong       wmark;
      38             :     ulong       mtu;
      39             :     ulong       frag_pos;
      40             :   } in;
      41             : 
      42             :   struct {
      43             :     fd_wksp_t * mem;
      44             :     ulong       chunk0;
      45             :     ulong       wmark;
      46             :     ulong       chunk;
      47             :     ulong       mtu;
      48             :   } out;
      49             : 
      50             :   struct {
      51             :     struct {
      52             :       ulong compressed_bytes_read;
      53             :       ulong decompressed_bytes_written;
      54             :     } full;
      55             : 
      56             :     struct {
      57             :       ulong compressed_bytes_read;
      58             :       ulong decompressed_bytes_written;
      59             :     } incremental;
      60             :   } metrics;
      61             : };
      62             : typedef struct fd_snapdc_tile fd_snapdc_tile_t;
      63             : 
      64             : FD_FN_PURE static ulong
      65           0 : scratch_align( void ) {
      66           0 :   return fd_ulong_max( alignof(fd_snapdc_tile_t), 32UL );
      67           0 : }
      68             : 
      69             : FD_FN_PURE static ulong
      70           0 : scratch_footprint( fd_topo_tile_t const * tile ) {
      71           0 :   (void)tile;
      72           0 :   ulong l = FD_LAYOUT_INIT;
      73           0 :   l = FD_LAYOUT_APPEND( l, alignof(fd_snapdc_tile_t), sizeof(fd_snapdc_tile_t)                   );
      74           0 :   l = FD_LAYOUT_APPEND( l, 32UL,                      ZSTD_estimateDStreamSize( ZSTD_WINDOW_SZ ) );
      75           0 :   return FD_LAYOUT_FINI( l, scratch_align() );
      76           0 : }
      77             : 
      78             : static inline int
      79           0 : should_shutdown( fd_snapdc_tile_t * ctx ) {
      80           0 :   return ctx->state==FD_SNAPSHOT_STATE_SHUTDOWN;
      81           0 : }
      82             : 
      83             : static void
      84           0 : metrics_write( fd_snapdc_tile_t * ctx ) {
      85           0 :   FD_MGAUGE_SET( SNAPDC, FULL_COMPRESSED_BYTES_READ,              ctx->metrics.full.compressed_bytes_read );
      86           0 :   FD_MGAUGE_SET( SNAPDC, FULL_DECOMPRESSED_BYTES_WRITTEN,         ctx->metrics.full.decompressed_bytes_written );
      87             : 
      88           0 :   FD_MGAUGE_SET( SNAPDC, INCREMENTAL_COMPRESSED_BYTES_READ,       ctx->metrics.incremental.compressed_bytes_read );
      89           0 :   FD_MGAUGE_SET( SNAPDC, INCREMENTAL_DECOMPRESSED_BYTES_WRITTEN,  ctx->metrics.incremental.decompressed_bytes_written );
      90             : 
      91           0 :   FD_MGAUGE_SET( SNAPDC, STATE,                                   (ulong)(ctx->state) );
      92           0 : }
      93             : 
      94             : static void
      95             : transition_malformed( fd_snapdc_tile_t *  ctx,
      96          12 :                       fd_stem_context_t * stem ) {
      97          12 :   if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_ERROR ) ) return;
      98          12 :   ctx->state = FD_SNAPSHOT_STATE_ERROR;
      99          12 :   fd_stem_publish( stem, 0UL, FD_SNAPSHOT_MSG_CTRL_ERROR, 0UL, 0UL, 0UL, 0UL, 0UL );
     100          12 : }
     101             : 
     102             : static inline void
     103             : handle_control_frag( fd_snapdc_tile_t *  ctx,
     104             :                      fd_stem_context_t * stem,
     105             :                      ulong               sig,
     106             :                      ulong               chunk,
     107       24933 :                      ulong               sz ) {
     108       24933 :   if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_LOAD_COMPLETE ) ) return;
     109             : 
     110             :   /* All control messages except META reset the decompression stream */
     111       24927 :   if( FD_UNLIKELY( sig!=FD_SNAPSHOT_MSG_META ) ) {
     112       24921 :     ulong error = ZSTD_DCtx_reset( ctx->zstd, ZSTD_reset_session_only );
     113       24921 :     if( FD_UNLIKELY( ZSTD_isError( error ) ) ) FD_LOG_ERR(( "ZSTD_DCtx_reset failed (%lu-%s)", error, ZSTD_getErrorName( error ) ));
     114       24921 :   }
     115             : 
     116       24927 :   if( ctx->state==FD_SNAPSHOT_STATE_ERROR && sig!=FD_SNAPSHOT_MSG_CTRL_FAIL ) {
     117             :     /* Control messages move along the snapshot load pipeline.  Since
     118             :        error conditions can be triggered by any tile in the pipeline,
     119             :        it is possible to be in error state and still receive otherwise
     120             :        valid messages.  Only a fail message can revert this. */
     121           6 :     return;
     122       24921 :   };
     123             : 
     124       24921 :   if( FD_UNLIKELY( sig==FD_SNAPSHOT_MSG_META ) ) {
     125             :     /* Forward META to snapin so it can update the advertised
     126             :        slot/hash for redirect-based downloads. */
     127           6 :     FD_TEST( sz<=ctx->out.mtu );
     128           6 :     void * dst = fd_chunk_to_laddr( ctx->out.mem, ctx->out.chunk );
     129           6 :     fd_memcpy( dst, fd_chunk_to_laddr_const( ctx->in.mem, chunk ), sz );
     130           6 :     fd_stem_publish( stem, 0UL, sig, ctx->out.chunk, sz, 0UL, 0UL, 0UL );
     131           6 :     ctx->out.chunk = fd_dcache_compact_next( ctx->out.chunk, ctx->out.mtu, ctx->out.chunk0, ctx->out.wmark );
     132           6 :     return;
     133           6 :   }
     134             : 
     135       24915 :   int forward_msg = 1;
     136             : 
     137       24915 :   switch( sig ) {
     138       24861 :     case FD_SNAPSHOT_MSG_CTRL_INIT_FULL:
     139       24870 :     case FD_SNAPSHOT_MSG_CTRL_INIT_INCR: {
     140       24870 :       FD_TEST( ctx->state==FD_SNAPSHOT_STATE_IDLE );
     141       24870 :       ctx->state = FD_SNAPSHOT_STATE_PROCESSING;
     142       24870 :       FD_TEST( sz==sizeof(fd_ssctrl_init_t) );
     143       24870 :       fd_ssctrl_init_t const * msg = fd_chunk_to_laddr_const( ctx->in.mem, chunk );
     144       24870 :       ctx->full = sig==FD_SNAPSHOT_MSG_CTRL_INIT_FULL;
     145       24870 :       ctx->is_zstd = !!msg->zstd;
     146       24870 :       ctx->dirty       = 0;
     147       24870 :       ctx->frame_idx   = 0UL;
     148       24870 :       ctx->in.frag_pos = 0UL;
     149       24870 :       FD_TEST( fd_zstd_frame_new( ctx->zstd_frame ) );
     150       24870 :       if( ctx->full ) {
     151       24861 :         ctx->metrics.full.compressed_bytes_read      = 0UL;
     152       24861 :         ctx->metrics.full.decompressed_bytes_written = 0UL;
     153       24861 :       } else {
     154           9 :         ctx->metrics.incremental.compressed_bytes_read      = 0UL;
     155           9 :         ctx->metrics.incremental.decompressed_bytes_written = 0UL;
     156           9 :       }
     157       24870 :       fd_ssctrl_init_t * msg_out = fd_chunk_to_laddr( ctx->out.mem, ctx->out.chunk );
     158       24870 :       fd_memcpy( msg_out, msg, sz );
     159       24870 :       fd_stem_publish( stem, 0UL, sig, ctx->out.chunk, sz, 0UL, 0UL, 0UL );
     160       24870 :       ctx->out.chunk = fd_dcache_compact_next( ctx->out.chunk, ctx->out.mtu, ctx->out.chunk0, ctx->out.wmark );
     161       24870 :       forward_msg = 0; // we forward the control message in the `fd_ssctrl_init_t` message
     162       24870 :       break;
     163       24870 :     }
     164             : 
     165          12 :     case FD_SNAPSHOT_MSG_CTRL_FINI: {
     166          12 :       FD_TEST( ctx->state==FD_SNAPSHOT_STATE_PROCESSING );
     167          12 :       ctx->state = FD_SNAPSHOT_STATE_FINISHING;
     168          12 :       if( FD_UNLIKELY( ctx->is_zstd && ctx->dirty ) ) {
     169           6 :         FD_LOG_WARNING(( "encountered end-of-file in the middle of a compressed frame for %s snapshot",
     170           6 :                          ctx->full ? "full" : "incremental" ));
     171           6 :         transition_malformed( ctx, stem );
     172           6 :         forward_msg = 0;
     173           6 :         break;
     174           6 :       }
     175           6 :       break;
     176          12 :     }
     177             : 
     178           6 :     case FD_SNAPSHOT_MSG_CTRL_NEXT:
     179           6 :     case FD_SNAPSHOT_MSG_CTRL_DONE: {
     180           6 :       FD_TEST( ctx->state==FD_SNAPSHOT_STATE_FINISHING );
     181           6 :       ctx->state = FD_SNAPSHOT_STATE_IDLE;
     182           6 :       break;
     183           6 :     }
     184             : 
     185           9 :     case FD_SNAPSHOT_MSG_CTRL_ERROR: {
     186           9 :       FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
     187           9 :       ctx->state = FD_SNAPSHOT_STATE_ERROR;
     188           9 :       break;
     189           9 :     }
     190             : 
     191          15 :     case FD_SNAPSHOT_MSG_CTRL_FAIL: {
     192          15 :       FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
     193          15 :       ctx->state = FD_SNAPSHOT_STATE_IDLE;
     194          15 :       break;
     195          15 :     }
     196             : 
     197           3 :     case FD_SNAPSHOT_MSG_CTRL_SHUTDOWN: {
     198           3 :       FD_TEST( ctx->state==FD_SNAPSHOT_STATE_IDLE );
     199           3 :       ctx->state = FD_SNAPSHOT_STATE_SHUTDOWN;
     200           3 :       break;
     201           3 :     }
     202             : 
     203           0 :     default: {
     204           0 :       FD_LOG_ERR(( "unexpected control frag %s (%lu) in state %s (%lu)",
     205           0 :                    fd_ssctrl_msg_ctrl_str( sig ), sig,
     206           0 :                    fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
     207           0 :       break;
     208           0 :     }
     209       24915 :   }
     210             : 
     211             :   /* Forward the control message down the pipeline */
     212       24915 :   if( FD_LIKELY( forward_msg ) ) {
     213          39 :     fd_stem_publish( stem, 0UL, sig, 0UL, 0UL, 0UL, 0UL, 0UL );
     214          39 :   }
     215       24915 : }
     216             : 
     217             : static inline void
     218         327 : finish_frame( fd_snapdc_tile_t * ctx ) {
     219         327 :   ctx->dirty = 0;
     220         327 :   ctx->frame_idx++;
     221         327 :   FD_TEST( fd_zstd_frame_new( ctx->zstd_frame ) );
     222         327 : }
     223             : 
     224             : /* Reads until the end of the current frame (or the end of the frag, if
     225             :    that comes earlier).  Returns 1 if bytes from the next frame remain
     226             :    and stem must reprocess this frag. */
     227             : static inline int
     228             : skip_unowned_frame( fd_snapdc_tile_t *  ctx,
     229             :                     fd_stem_context_t * stem,
     230             :                     uchar const *       data,
     231         369 :                     ulong               sz ) {
     232         369 :   FD_TEST( ctx->frame_idx%ctx->tile_count!=ctx->tile_idx );
     233         369 :   FD_TEST( ctx->dirty || ctx->in.frag_pos<sz );
     234         369 :   ctx->dirty = 1;
     235             : 
     236         369 :   ulong skipped = 0UL;
     237         369 :   int scan_result = fd_zstd_frame_advance( ctx->zstd_frame, data+ctx->in.frag_pos, sz-ctx->in.frag_pos, &skipped );
     238         369 :   switch( scan_result ) {
     239           3 :     case FD_ZSTD_FRAME_ERR: {
     240           3 :       transition_malformed( ctx, stem );
     241           3 :       return 0;
     242           0 :     }
     243         162 :     case FD_ZSTD_FRAME_MORE: {
     244             :       /* Current frame spans until the end of the frag, so we can move
     245             :          onto the next frag. */
     246         162 :       FD_TEST( skipped+ctx->in.frag_pos==sz );
     247         162 :       ctx->in.frag_pos = 0UL;
     248         162 :       return 0;
     249         162 :     }
     250         204 :     case FD_ZSTD_FRAME_END: {
     251             :       /* A frame ends within this frag.  If the frame end coincides with
     252             :          the frag end, wait for the next frag, otherwise reprocess
     253             :          the next frame in this frag. */
     254         204 :       FD_TEST( skipped && skipped<=sz-ctx->in.frag_pos );
     255         204 :       ctx->in.frag_pos += skipped;
     256         204 :       finish_frame( ctx );
     257             : 
     258         204 :       if( FD_UNLIKELY( ctx->in.frag_pos==sz ) ) {
     259         150 :         ctx->in.frag_pos = 0UL;
     260         150 :         return 0;
     261         150 :       }
     262          54 :       return 1;
     263         204 :     }
     264           0 :     default: FD_LOG_ERR(( "unexpected zstd frame scan result %d", scan_result ));
     265         369 :   }
     266         369 : }
     267             : 
     268             : /* Decompresses up to a single frame's data.  Returns 1 if the current
     269             :    frag needs to be reprocessed. */
     270             : static inline int
     271             : process_owned_frame( fd_snapdc_tile_t *  ctx,
     272             :                      fd_stem_context_t * stem,
     273             :                      uchar const *       data,
     274         357 :                      ulong               sz ) {
     275         357 :   FD_TEST( ctx->frame_idx%ctx->tile_count==ctx->tile_idx );
     276         357 :   FD_TEST( ctx->dirty || ctx->in.frag_pos<sz );
     277         357 :   ctx->dirty = 1;
     278             : 
     279         357 :   uchar * out          = fd_chunk_to_laddr( ctx->out.mem, ctx->out.chunk );
     280         357 :   ulong   in_sz        = sz-ctx->in.frag_pos;
     281         357 :   ulong   in_consumed  = 0UL;
     282         357 :   ulong   out_produced = 0UL;
     283         357 :   ulong frame_res = ZSTD_decompressStream_simpleArgs(
     284         357 :       ctx->zstd,
     285         357 :       out,
     286         357 :       ctx->out.mtu,
     287         357 :       &out_produced,
     288         357 :       data+ctx->in.frag_pos,
     289         357 :       in_sz,
     290         357 :       &in_consumed );
     291         357 :   if( FD_UNLIKELY( ZSTD_isError( frame_res ) ) ) {
     292           3 :     FD_LOG_WARNING(( "error while decompressing %s snapshot (%u-%s)",
     293           3 :                      ctx->full ? "full" : "incremental",
     294           3 :                      ZSTD_getErrorCode( frame_res ), ZSTD_getErrorName( frame_res ) ));
     295           3 :     transition_malformed( ctx, stem );
     296           3 :     return 0;
     297           3 :   }
     298             : 
     299         354 :   ctx->in.frag_pos += in_consumed;
     300         354 :   FD_TEST( ctx->in.frag_pos<=sz );
     301             : 
     302         354 :   if( FD_LIKELY( ctx->full ) ) {
     303         348 :     ctx->metrics.full.compressed_bytes_read      += in_consumed;
     304         348 :     ctx->metrics.full.decompressed_bytes_written += out_produced;
     305         348 :   } else {
     306           6 :     ctx->metrics.incremental.compressed_bytes_read      += in_consumed;
     307           6 :     ctx->metrics.incremental.decompressed_bytes_written += out_produced;
     308           6 :   }
     309             : 
     310         354 :   if( FD_UNLIKELY( frame_res && !in_consumed && !out_produced ) ) {
     311           0 :     if( FD_UNLIKELY( ctx->in.frag_pos<sz ) ) {
     312             :       /* No progress with remaining input would retry forever */
     313           0 :       transition_malformed( ctx, stem );
     314           0 :     } else {
     315             :       /* No progress with exhausted input means zstd needs the next frag */
     316           0 :       ctx->in.frag_pos = 0UL;
     317           0 :     }
     318           0 :     return 0;
     319           0 :   }
     320             : 
     321         354 :   if( FD_LIKELY( out_produced || !frame_res ) ) {
     322         258 :     ulong out_ctl = fd_frag_meta_ctl( 0UL, 0, !frame_res, 0 );
     323         258 :     fd_stem_publish( stem, 0UL, FD_SNAPSHOT_MSG_DATA, ctx->out.chunk, out_produced, out_ctl, 0UL, 0UL );
     324         258 :     ctx->out.chunk = fd_dcache_compact_next( ctx->out.chunk, out_produced, ctx->out.chunk0, ctx->out.wmark );
     325         258 :   }
     326             : 
     327         354 :   if( FD_UNLIKELY( !frame_res ) ) finish_frame( ctx );
     328             : 
     329             :   /* frame_res==0 means the frame ended exactly at the output boundary;
     330             :      re-polling then reports "new frame expected" and would mark the
     331             :      stream dirty at a clean EOF. */
     332         354 :   int maybe_more_output = (out_produced==ctx->out.mtu && frame_res!=0UL) || ctx->in.frag_pos<sz;
     333         354 :   if( FD_LIKELY( !maybe_more_output ) ) ctx->in.frag_pos = 0UL;
     334         354 :   return maybe_more_output;
     335         354 : }
     336             : 
     337             : static inline int
     338             : handle_data_frag( fd_snapdc_tile_t *  ctx,
     339             :                   fd_stem_context_t * stem,
     340             :                   ulong               chunk,
     341       27048 :                   ulong               sz ) {
     342       27048 :   if( FD_UNLIKELY( ctx->state==FD_SNAPSHOT_STATE_ERROR ) ) {
     343             :     /* Ignore all data frags after observing an error in the stream until
     344             :        we receive fail & init control messages to restart processing. */
     345           3 :     return 0;
     346           3 :   }
     347       27045 :   if( FD_UNLIKELY( ctx->state!=FD_SNAPSHOT_STATE_PROCESSING ) ) {
     348           0 :     FD_LOG_ERR(( "received unexpected data frag in state %s (%lu)",
     349           0 :                  fd_ssctrl_state_str( (ulong)ctx->state ), (ulong)ctx->state ));
     350           0 :   }
     351             : 
     352       27045 :   FD_TEST( chunk>=ctx->in.chunk0 && chunk<=ctx->in.wmark && sz<=ctx->in.mtu && sz>=ctx->in.frag_pos );
     353       27045 :   uchar const * data = fd_chunk_to_laddr_const( ctx->in.mem, chunk );
     354             : 
     355       27045 :   if( FD_UNLIKELY( !ctx->is_zstd ) ) {
     356       26319 :     if( FD_UNLIKELY( ctx->tile_idx!=0UL ) ) return 0;
     357        1935 :     FD_TEST( ctx->in.frag_pos<sz );
     358        1935 :     uchar const * in  = data+ctx->in.frag_pos;
     359        1935 :     uchar *       out = fd_chunk_to_laddr( ctx->out.mem, ctx->out.chunk );
     360        1935 :     ulong cpy = fd_ulong_min( sz-ctx->in.frag_pos, ctx->out.mtu );
     361        1935 :     fd_memcpy( out, in, cpy );
     362        1935 :     fd_stem_publish( stem, 0UL, FD_SNAPSHOT_MSG_DATA, ctx->out.chunk, cpy, 0UL, 0UL, 0UL );
     363        1935 :     ctx->out.chunk = fd_dcache_compact_next( ctx->out.chunk, cpy, ctx->out.chunk0, ctx->out.wmark );
     364             : 
     365        1935 :     if( FD_LIKELY( ctx->full ) ) {
     366        1920 :       ctx->metrics.full.compressed_bytes_read      += cpy;
     367        1920 :       ctx->metrics.full.decompressed_bytes_written += cpy;
     368        1920 :     } else {
     369          15 :       ctx->metrics.incremental.compressed_bytes_read      += cpy;
     370          15 :       ctx->metrics.incremental.decompressed_bytes_written += cpy;
     371          15 :     }
     372             : 
     373        1935 :     ctx->in.frag_pos += cpy;
     374        1935 :     FD_TEST( ctx->in.frag_pos<=sz );
     375        1935 :     if( FD_UNLIKELY( ctx->in.frag_pos<sz ) ) return 1;
     376         387 :     ctx->in.frag_pos = 0UL;
     377         387 :     return 0;
     378        1935 :   }
     379             : 
     380         726 :   if( ctx->frame_idx%ctx->tile_count!=ctx->tile_idx ) {
     381         369 :     return skip_unowned_frame( ctx, stem, data, sz );
     382         369 :   }
     383             : 
     384         357 :   return process_owned_frame( ctx, stem, data, sz );
     385         726 : }
     386             : 
     387             : static inline int
     388             : returnable_frag( fd_snapdc_tile_t *  ctx,
     389             :                  ulong               in_idx FD_PARAM_UNUSED,
     390             :                  ulong               seq    FD_PARAM_UNUSED,
     391             :                  ulong               sig,
     392             :                  ulong               chunk,
     393             :                  ulong               sz,
     394             :                  ulong               ctl    FD_PARAM_UNUSED,
     395             :                  ulong               tsorig FD_PARAM_UNUSED,
     396             :                  ulong               tspub  FD_PARAM_UNUSED,
     397       51981 :                  fd_stem_context_t * stem ) {
     398       51981 :   FD_TEST( ctx->state!=FD_SNAPSHOT_STATE_SHUTDOWN );
     399             : 
     400       51981 :   if( FD_LIKELY( sig==FD_SNAPSHOT_MSG_DATA ) ) return handle_data_frag( ctx, stem, chunk, sz );
     401       24933 :   else                                                handle_control_frag( ctx, stem, sig, chunk, sz );
     402             : 
     403       24933 :   return 0;
     404       51981 : }
     405             : 
     406             : static ulong
     407             : populate_allowed_fds( fd_topo_t      const * topo FD_PARAM_UNUSED,
     408             :                       fd_topo_tile_t const * tile FD_PARAM_UNUSED,
     409             :                       ulong                  out_fds_cnt,
     410           0 :                       int *                  out_fds ) {
     411           0 :   if( FD_UNLIKELY( out_fds_cnt<2UL ) ) FD_LOG_ERR(( "out_fds_cnt %lu", out_fds_cnt ));
     412             : 
     413           0 :   ulong out_cnt = 0;
     414           0 :   out_fds[ out_cnt++ ] = 2UL; /* stderr */
     415           0 :   if( FD_LIKELY( -1!=fd_log_private_logfile_fd() ) ) {
     416           0 :     out_fds[ out_cnt++ ] = fd_log_private_logfile_fd(); /* logfile */
     417           0 :   }
     418             : 
     419           0 :   return out_cnt;
     420           0 : }
     421             : 
     422             : static ulong
     423             : populate_allowed_seccomp( fd_topo_t const *      topo FD_PARAM_UNUSED,
     424             :                           fd_topo_tile_t const * tile FD_PARAM_UNUSED,
     425             :                           ulong                  out_cnt,
     426           0 :                           struct sock_filter *   out ) {
     427           0 :   populate_sock_filter_policy_fd_snapdc_tile( out_cnt, out, (uint)fd_log_private_logfile_fd() );
     428           0 :   return sock_filter_policy_fd_snapdc_tile_instr_cnt;
     429           0 : }
     430             : 
     431             : static void
     432             : unprivileged_init( fd_topo_t const *      topo,
     433           0 :                    fd_topo_tile_t const * tile ) {
     434           0 :   void * scratch = fd_topo_obj_laddr( topo, tile->tile_obj_id );
     435             : 
     436           0 :   FD_SCRATCH_ALLOC_INIT( l, scratch );
     437           0 :   fd_snapdc_tile_t * ctx = FD_SCRATCH_ALLOC_APPEND( l, alignof(fd_snapdc_tile_t), sizeof(fd_snapdc_tile_t) );
     438           0 :   void * _zstd           = FD_SCRATCH_ALLOC_APPEND( l, 32UL,                      ZSTD_estimateDStreamSize( ZSTD_WINDOW_SZ ) );
     439             : 
     440           0 :   ctx->state      = FD_SNAPSHOT_STATE_IDLE;
     441           0 :   ctx->tile_idx   = tile->kind_id;
     442           0 :   ctx->tile_count = fd_topo_tile_name_cnt( topo, NAME );
     443           0 :   FD_TEST( ctx->tile_count );
     444           0 :   FD_TEST( ctx->tile_idx<ctx->tile_count );
     445             : 
     446           0 :   ctx->zstd = ZSTD_initStaticDStream( _zstd, ZSTD_estimateDStreamSize( ZSTD_WINDOW_SZ ) );
     447           0 :   FD_TEST( ctx->zstd );
     448           0 :   FD_TEST( ctx->zstd==_zstd );
     449             : 
     450           0 :   ctx->dirty       = 0;
     451           0 :   ctx->frame_idx   = 0UL;
     452           0 :   ctx->in.frag_pos = 0UL;
     453           0 :   FD_TEST( fd_zstd_frame_new( ctx->zstd_frame ) );
     454           0 :   fd_memset( &ctx->metrics, 0, sizeof(ctx->metrics) );
     455             : 
     456           0 :   if( FD_UNLIKELY( tile->in_cnt !=1UL ) ) FD_LOG_ERR(( "tile `" NAME "` has %lu ins, expected 1",  tile->in_cnt  ));
     457           0 :   if( FD_UNLIKELY( tile->out_cnt!=1UL ) ) FD_LOG_ERR(( "tile `" NAME "` has %lu outs, expected 1", tile->out_cnt ));
     458             : 
     459           0 :   fd_topo_link_t const * snapin_link = &topo->links[ tile->out_link_id[ 0UL ] ];
     460           0 :   FD_TEST( 0==strcmp( snapin_link->name, "snapdc_in" ) );
     461           0 :   ctx->out.mem    = topo->workspaces[ topo->objs[ snapin_link->dcache_obj_id ].wksp_id ].wksp;
     462           0 :   ctx->out.chunk0 = fd_dcache_compact_chunk0( ctx->out.mem, snapin_link->dcache );
     463           0 :   ctx->out.wmark  = fd_dcache_compact_wmark ( ctx->out.mem, snapin_link->dcache, snapin_link->mtu );
     464           0 :   ctx->out.chunk  = ctx->out.chunk0;
     465           0 :   ctx->out.mtu    = snapin_link->mtu;
     466             : 
     467           0 :   fd_topo_link_t const * in_link = &topo->links[ tile->in_link_id[ 0UL ] ];
     468           0 :   fd_topo_wksp_t const * in_wksp = &topo->workspaces[ topo->objs[ in_link->dcache_obj_id ].wksp_id ];
     469           0 :   ctx->in.mem                    = in_wksp->wksp;
     470           0 :   ctx->in.chunk0                 = fd_dcache_compact_chunk0( ctx->in.mem, in_link->dcache );
     471           0 :   ctx->in.wmark                  = fd_dcache_compact_wmark( ctx->in.mem, in_link->dcache, in_link->mtu );
     472           0 :   ctx->in.mtu                    = in_link->mtu;
     473             : 
     474           0 :   ulong scratch_top = FD_SCRATCH_ALLOC_FINI( l, scratch_align() );
     475           0 :   if( FD_UNLIKELY( scratch_top > (ulong)scratch + scratch_footprint( tile ) ) )
     476           0 :     FD_LOG_ERR(( "scratch overflow %lu %lu %lu",
     477           0 :                  scratch_top - (ulong)scratch - scratch_footprint( tile ),
     478           0 :                  scratch_top,
     479           0 :                  (ulong)scratch + scratch_footprint( tile ) ));
     480           0 : }
     481             : 
     482           0 : #define STEM_BURST 1UL
     483             : 
     484           0 : #define STEM_LAZY  (128L*3000L)
     485             : 
     486           0 : #define STEM_CALLBACK_CONTEXT_TYPE  fd_snapdc_tile_t
     487           0 : #define STEM_CALLBACK_CONTEXT_ALIGN alignof(fd_snapdc_tile_t)
     488             : 
     489             : #define STEM_CALLBACK_SHOULD_SHUTDOWN should_shutdown
     490           0 : #define STEM_CALLBACK_METRICS_WRITE   metrics_write
     491           0 : #define STEM_CALLBACK_RETURNABLE_FRAG returnable_frag
     492             : 
     493             : #include "../../disco/stem/fd_stem.c"
     494             : 
     495             : fd_topo_run_tile_t fd_tile_snapdc = {
     496             :   .name                     = NAME,
     497             :   .populate_allowed_fds     = populate_allowed_fds,
     498             :   .populate_allowed_seccomp = populate_allowed_seccomp,
     499             :   .scratch_align            = scratch_align,
     500             :   .scratch_footprint        = scratch_footprint,
     501             :   .unprivileged_init        = unprivileged_init,
     502             :   .run                      = stem_run,
     503             : };
     504             : 
     505             : #undef NAME

Generated by: LCOV version 1.14