LCOV - code coverage report
Current view: top level - util/tpool - fd_tpool.c (source / functions) Hit Total Coverage
Test: cov.lcov Lines: 231 250 92.4 %
Date: 2026-08-01 05:23:47 Functions: 12 13 92.3 %

          Line data    Source code
       1             : #include "fd_tpool.h"
       2             : 
       3             : #if FD_HAS_THREADS
       4             : #include <pthread.h>
       5             : #endif
       6             : 
       7             : struct fd_tpool_private_worker_cfg {
       8             :   fd_tpool_t * tpool;
       9             :   ulong        tile_idx;
      10             : };
      11             : 
      12             : typedef struct fd_tpool_private_worker_cfg fd_tpool_private_worker_cfg_t;
      13             : 
      14             : static int
      15             : fd_tpool_private_worker_( int     argc,
      16          84 :                           char ** argv ) {
      17          84 :   ulong                           worker_idx = (ulong)(uint)argc;
      18          84 :   fd_tpool_private_worker_cfg_t * cfg        = (fd_tpool_private_worker_cfg_t *)argv;
      19             : 
      20          84 :   fd_tpool_t * tpool    = cfg->tpool;
      21          84 :   ulong        tile_idx = cfg->tile_idx;
      22             : 
      23             :   /* We are BOOT */
      24             : 
      25          84 :   fd_tpool_private_worker_t worker[1];
      26             : 
      27          84 :   memset( worker, 0, sizeof(fd_tpool_private_worker_t) );
      28             : 
      29          84 :   worker->tile_idx = (uint) tile_idx;
      30             : 
      31          84 : # if FD_HAS_THREADS
      32          84 :   int sleeper = !!(tpool->opt & FD_TPOOL_OPT_SLEEP);
      33             : 
      34          84 :   pthread_mutex_t lock[1];
      35          84 :   pthread_cond_t  wake[1];
      36             : 
      37          84 :   if( FD_UNLIKELY( sleeper ) ) {
      38           0 :     if( FD_UNLIKELY( pthread_mutex_init( lock, NULL ) ) ) FD_LOG_ERR(( "pthread_mutex_init failed" ));
      39           0 :     if( FD_UNLIKELY( pthread_cond_init ( wake, NULL ) ) ) FD_LOG_ERR(( "pthread_cond_init failed"  ));
      40           0 :     if( FD_UNLIKELY( pthread_mutex_lock( lock       ) ) ) FD_LOG_ERR(( "pthread_mutex_lock failed" ));
      41           0 :   }
      42             : 
      43          84 :   worker->lock = (ulong)lock;
      44          84 :   worker->wake = (ulong)wake;
      45          84 : # endif
      46             : 
      47          84 :   FD_COMPILER_MFENCE();
      48             : 
      49          84 :   fd_tpool_private_worker( tpool )[ worker_idx ] = worker;
      50             : 
      51          84 :   ulong const * arg  = worker->arg;
      52          84 :   uint          seq1 = worker->seq1;
      53             : 
      54   816531466 :   for(;;) {
      55             : 
      56             :     /* We are IDLE ... see what we should do next */
      57             : 
      58   816531466 : #   if FD_HAS_THREADS
      59   816531466 :     if( FD_UNLIKELY( sleeper ) && FD_UNLIKELY( pthread_cond_wait( wake, lock ) ) )
      60           0 :       FD_LOG_WARNING(( "pthread_cond_wait failed; attempting to continue" ));
      61   816531466 : #   endif
      62             : 
      63   816531466 :     FD_COMPILER_MFENCE();
      64   816531466 :     uint  seq0     = worker->seq0;
      65   816531466 :     FD_COMPILER_MFENCE();
      66   816531466 :     uint  _arg_cnt = worker->arg_cnt;
      67   816531466 :     ulong _task    = worker->task;
      68   816531466 :     FD_COMPILER_MFENCE();
      69             : 
      70   828351461 :     if( FD_UNLIKELY( seq0==seq1 ) ) { /* Got idle */
      71   828185521 :       FD_SPIN_PAUSE();
      72   828185521 :       continue;
      73   828185521 :     }
      74             : 
      75           3 :     if( FD_UNLIKELY( !_task ) ) break; /* Got halt */
      76             : 
      77             :     /* We are EXEC ... do the task and then transition to IDLE */
      78             : 
      79           3 :     if( _arg_cnt==UINT_MAX ) {
      80             : 
      81     7728851 :       fd_tpool_task_t task = (fd_tpool_task_t)_task;
      82             : 
      83     7728851 :       void * task_tpool  = (void *)arg[ 0];
      84     7728851 :       ulong  task_t0     =         arg[ 1]; ulong task_t1     = arg[ 2];
      85     7728851 :       void * task_args   = (void *)arg[ 3];
      86     7728851 :       void * task_reduce = (void *)arg[ 4]; ulong task_stride = arg[ 5];
      87     7728851 :       ulong  task_l0     =         arg[ 6]; ulong task_l1     = arg[ 7];
      88     7728851 :       ulong  task_m0     =         arg[ 8]; ulong task_m1     = arg[ 9];
      89     7728851 :       ulong  task_n0     =         arg[10]; ulong task_n1     = arg[11];
      90             : 
      91     7728851 :       task( task_tpool,task_t0,task_t1, task_args, task_reduce,task_stride, task_l0,task_l1, task_m0,task_m1, task_n0,task_n1 );
      92             : 
      93           5 :     } else {
      94             : 
      95           5 :       fd_tpool_task_v2_t task = (fd_tpool_task_v2_t)_task;
      96             : 
      97           5 :       task( tpool, worker_idx, (ulong)_arg_cnt, arg );
      98             : 
      99           5 :     }
     100             : 
     101           3 :     FD_COMPILER_MFENCE();
     102             : 
     103           3 :     worker->seq1 = seq0;
     104           3 :     seq1 = seq0;
     105           3 :   }
     106             : 
     107             :   /* We are HALT ... clean up and terminate */
     108             : 
     109          84 : # if FD_HAS_THREADS
     110          84 :   if( FD_UNLIKELY( sleeper ) ) {
     111           0 :     if( FD_UNLIKELY( pthread_mutex_unlock ( lock ) ) ) FD_LOG_WARNING(( "pthread_mutex_unlock failed; attempting to continue" ));
     112           0 :     if( FD_UNLIKELY( pthread_cond_destroy ( wake ) ) ) FD_LOG_WARNING(( "pthread_cond_destroy failed; attempting to continue" ));
     113           0 :     if( FD_UNLIKELY( pthread_mutex_destroy( lock ) ) ) FD_LOG_WARNING(( "pthread_mutex_destroy failed; attempting to continue" ));
     114           0 :   }
     115          84 : # endif
     116             : 
     117          84 :   return 0;
     118          84 : }
     119             : 
     120             : #if FD_HAS_THREADS
     121             : void
     122           0 : fd_tpool_private_wake( fd_tpool_private_worker_t * worker ) {
     123           0 :   pthread_mutex_t * lock = (pthread_mutex_t *)worker->lock;
     124           0 :   pthread_cond_t *  wake = (pthread_cond_t  *)worker->wake;
     125           0 :   if( FD_UNLIKELY( pthread_mutex_lock  ( lock ) ) ) FD_LOG_WARNING(( "pthread_mutex_lock failed; attempting to continue" ));
     126           0 :   if( FD_UNLIKELY( pthread_cond_signal ( wake ) ) ) FD_LOG_WARNING(( "pthread_cond_signal failed; attempting to continue" ));
     127           0 :   if( FD_UNLIKELY( pthread_mutex_unlock( lock ) ) ) FD_LOG_WARNING(( "pthread_mutex_unlock failed; attempting to continue" ));
     128           0 : }
     129             : #endif
     130             : 
     131             : ulong
     132        3105 : fd_tpool_align( void ) {
     133        3105 :   return FD_TPOOL_ALIGN;
     134        3105 : }
     135             : 
     136             : ulong
     137     3003099 : fd_tpool_footprint( ulong worker_max ) {
     138     3003099 :   if( FD_UNLIKELY( !((1UL<=worker_max) & (worker_max<=FD_TILE_MAX)) ) ) return 0UL;
     139      474186 :   return fd_ulong_align_up( sizeof(fd_tpool_private_worker_t) +
     140      474186 :                             sizeof(fd_tpool_t) + worker_max*sizeof(fd_tpool_private_worker_t *), FD_TPOOL_ALIGN );
     141     3003099 : }
     142             : 
     143             : fd_tpool_t *
     144             : fd_tpool_init( void * mem,
     145             :                ulong  worker_max,
     146        3105 :                ulong  opt ) {
     147             : 
     148        3105 :   FD_COMPILER_MFENCE();
     149             : 
     150        3105 :   if( FD_UNLIKELY( !mem ) ) {
     151           3 :     FD_LOG_WARNING(( "NULL mem" ));
     152           3 :     return NULL;
     153           3 :   }
     154             : 
     155        3102 :   if( FD_UNLIKELY( !fd_ulong_is_aligned( (ulong)mem, fd_tpool_align() ) ) ) {
     156           3 :     FD_LOG_WARNING(( "bad alignment" ));
     157           3 :     return NULL;
     158           3 :   }
     159             : 
     160        3099 :   ulong footprint = fd_tpool_footprint( worker_max );
     161        3099 :   if( FD_UNLIKELY( !footprint ) ) {
     162           6 :     FD_LOG_WARNING(( "bad worker_max" ));
     163           6 :     return NULL;
     164           6 :   }
     165             : 
     166        3093 :   fd_memset( mem, 0, footprint );
     167             : 
     168        3093 :   fd_tpool_private_worker_t * worker0 = (fd_tpool_private_worker_t *)mem;
     169             : 
     170        3093 :   worker0->seq0 = 1U;
     171        3093 :   worker0->seq1 = 0U;
     172             : 
     173        3093 :   fd_tpool_t * tpool = (fd_tpool_t *)(worker0+1);
     174             : 
     175        3093 :   tpool->opt        = opt;
     176        3093 :   tpool->worker_max = (uint)worker_max;
     177        3093 :   tpool->worker_cnt = 1U;
     178             : 
     179        3093 :   FD_COMPILER_MFENCE();
     180        3093 :   fd_tpool_private_worker( tpool )[0] = worker0;
     181        3093 :   FD_COMPILER_MFENCE();
     182             : 
     183        3093 :   return tpool;
     184        3099 : }
     185             : 
     186             : void *
     187        3096 : fd_tpool_fini( fd_tpool_t * tpool ) {
     188             : 
     189        3096 :   FD_COMPILER_MFENCE();
     190             : 
     191        3096 :   if( FD_UNLIKELY( !tpool ) ) {
     192           3 :     FD_LOG_WARNING(( "NULL tpool" ));
     193           3 :     return NULL;
     194           3 :   }
     195             : 
     196        3156 :   while( fd_tpool_worker_cnt( tpool )>1UL ) {
     197          63 :     if( FD_UNLIKELY( !fd_tpool_worker_pop( tpool ) ) ) {
     198           0 :       FD_LOG_WARNING(( "fd_tpool_worker_pop failed" ));
     199           0 :       return NULL;
     200           0 :     }
     201          63 :   }
     202             : 
     203        3093 :   return (void *)fd_tpool_private_worker0( tpool );
     204        3093 : }
     205             : 
     206             : fd_tpool_t *
     207             : fd_tpool_worker_push( fd_tpool_t * tpool,
     208         120 :                       ulong        tile_idx ) {
     209             : 
     210         120 :   FD_COMPILER_MFENCE();
     211             : 
     212         120 :   if( FD_UNLIKELY( !tpool ) ) {
     213           3 :     FD_LOG_WARNING(( "NULL tpool" ));
     214           3 :     return NULL;
     215           3 :   }
     216             : 
     217         117 :   if( FD_UNLIKELY( !tile_idx ) ) {
     218           3 :     FD_LOG_WARNING(( "cannot push tile_idx 0" ));
     219           3 :     return NULL;
     220           3 :   }
     221             : 
     222         114 :   if( FD_UNLIKELY( tile_idx==fd_tile_idx() ) ) {
     223           3 :     FD_LOG_WARNING(( "cannot push self" ));
     224           3 :     return NULL;
     225           3 :   }
     226             : 
     227         111 :   if( FD_UNLIKELY( tile_idx>=fd_tile_cnt() ) ) {
     228           3 :     FD_LOG_WARNING(( "invalid tile_idx" ));
     229           3 :     return NULL;
     230           3 :   }
     231             : 
     232         108 :   fd_tpool_private_worker_t ** worker     = fd_tpool_private_worker( tpool );
     233         108 :   ulong                        worker_cnt = (ulong)tpool->worker_cnt;
     234             : 
     235         108 :   if( FD_UNLIKELY( worker_cnt>=(ulong)tpool->worker_max ) ) {
     236           3 :     FD_LOG_WARNING(( "too many workers" ));
     237           3 :     return NULL;
     238           3 :   }
     239             : 
     240         507 :   for( ulong worker_idx=0UL; worker_idx<worker_cnt; worker_idx++ )
     241         420 :     if( worker[ worker_idx ]->tile_idx==tile_idx ) {
     242          18 :       FD_LOG_WARNING(( "tile_idx already added to tpool" ));
     243          18 :       return NULL;
     244          18 :     }
     245             : 
     246          87 :   fd_tpool_private_worker_cfg_t cfg[1];
     247             : 
     248          87 :   cfg->tpool    = tpool;
     249          87 :   cfg->tile_idx = tile_idx;
     250             : 
     251          87 :   int     argc = (int)(uint)worker_cnt;
     252          87 :   char ** argv = (char **)fd_type_pun( cfg );
     253             : 
     254          87 :   FD_COMPILER_MFENCE();
     255          87 :   worker[ worker_cnt ] = NULL;
     256          87 :   FD_COMPILER_MFENCE();
     257             : 
     258          87 :   if( FD_UNLIKELY( !fd_tile_exec_new( tile_idx, fd_tpool_private_worker_, argc, argv ) ) ) {
     259           3 :     FD_LOG_WARNING(( "fd_tile_exec_new failed (tile probably already in use)" ));
     260           3 :     return NULL;
     261           3 :   }
     262             : 
     263        2604 :   while( !FD_VOLATILE_CONST( worker[ worker_cnt ] ) ) FD_SPIN_PAUSE();
     264             : 
     265          84 :   tpool->worker_cnt = (uint)(worker_cnt + 1UL);
     266          84 :   return tpool;
     267          87 : }
     268             : 
     269             : fd_tpool_t *
     270          93 : fd_tpool_worker_pop( fd_tpool_t * tpool ) {
     271             : 
     272          93 :   FD_COMPILER_MFENCE();
     273             : 
     274          93 :   if( FD_UNLIKELY( !tpool ) ) {
     275           3 :     FD_LOG_WARNING(( "NULL tpool" ));
     276           3 :     return NULL;
     277           3 :   }
     278             : 
     279          90 :   ulong worker_cnt = (ulong)tpool->worker_cnt;
     280          90 :   if( FD_UNLIKELY( worker_cnt<=1UL ) ) {
     281           3 :     FD_LOG_WARNING(( "no workers to pop" ));
     282           3 :     return NULL;
     283           3 :   }
     284             : 
     285             :   /* Testing for IDLE isn't strictly necessary given requirements to use
     286             :      this and this isn't being done atomically with the actually pop but
     287             :      does help catch obvious user errors. */
     288             : 
     289          87 :   if( FD_UNLIKELY( !fd_tpool_worker_idle( tpool, worker_cnt-1UL ) ) ) {
     290           3 :     FD_LOG_WARNING(( "worker to pop is not idle" ));
     291           3 :     return NULL;
     292           3 :   }
     293             : 
     294             :   /* Send HALT to the worker */
     295             : 
     296          84 :   fd_tpool_private_worker_t * worker = fd_tpool_private_worker( tpool )[ worker_cnt-1UL ];
     297          84 :   uint                        seq0   = worker->seq0 + 1U;
     298          84 :   fd_tile_exec_t *            exec   = fd_tile_exec( worker->tile_idx );
     299             : 
     300          84 :   worker->task = 0UL;
     301          84 :   FD_COMPILER_MFENCE();
     302          84 :   worker->seq0 = seq0;
     303          84 :   FD_COMPILER_MFENCE();
     304          84 :   if( FD_UNLIKELY( tpool->opt & FD_TPOOL_OPT_SLEEP ) ) fd_tpool_private_wake( worker );
     305             : 
     306             :   /* Wait for the worker to shutdown */
     307             : 
     308          84 :   int          ret;
     309          84 :   char const * err = fd_tile_exec_delete( exec, &ret );
     310          84 :   if(      FD_UNLIKELY( err ) ) FD_LOG_WARNING(( "tile err \"%s\" unexpected; attempting to continue", err ));
     311          84 :   else if( FD_UNLIKELY( ret ) ) FD_LOG_WARNING(( "tile ret %i unexpected; attempting to continue", ret ));
     312             : 
     313          84 :   tpool->worker_cnt = (uint)(worker_cnt - 1UL);
     314          84 :   return tpool;
     315          87 : }
     316             : 
     317             : #define FD_TPOOL_EXEC_ALL_IMPL_HDR(style)                                                          \
     318             : void                                                                                               \
     319             : fd_tpool_private_exec_all_##style##_node( void * _node_tpool,                                      \
     320             :                                           ulong  node_t0, ulong node_t1,                           \
     321             :                                           void * args,                                             \
     322             :                                           void * reduce,  ulong stride,                            \
     323             :                                           ulong  l0,      ulong l1,                                \
     324             :                                           ulong  _task,   ulong _tpool,                            \
     325    10432499 :                                           ulong  t0,      ulong t1 ) {                             \
     326    10432499 :   fd_tpool_t *    node_tpool = (fd_tpool_t *   )_node_tpool;                                       \
     327    10432499 :   fd_tpool_task_t task       = (fd_tpool_task_t)_task;                                             \
     328    10432499 :   ulong           wait_cnt   = 0UL;                                                                \
     329    10432499 :   ushort          wait_child[16];   /* Assumes tpool_cnt<=65536 */                                 \
     330    18375122 :   for(;;) {                                                                                        \
     331    18375122 :     ulong node_t_cnt = node_t1 - node_t0;                                                          \
     332    18375122 :     if( node_t_cnt<=1L ) break;                                                                    \
     333    18375122 :     ulong node_ts = node_t0 + fd_tpool_private_split( node_t_cnt );                                \
     334     8475460 :     fd_tpool_exec( node_tpool, node_ts, fd_tpool_private_exec_all_##style##_node,                  \
     335     8475460 :                    node_tpool, node_ts,node_t1, args, reduce,stride, l0,l1, _task,_tpool, t0,t1 ); \
     336     8475460 :     wait_child[ wait_cnt++ ] = (ushort)node_ts;                                                    \
     337     8475460 :     node_t1 = node_ts;                                                                             \
     338     8475460 :   }
     339             : 
     340             : #define FD_TPOOL_EXEC_ALL_IMPL_FTR                                                \
     341    19150980 :   while( wait_cnt ) fd_tpool_wait( node_tpool, (ulong)wait_child[ --wait_cnt ] ); \
     342    10432499 : }
     343             : 
     344      897817 : FD_TPOOL_EXEC_ALL_IMPL_HDR(rrobin)
     345      897817 :   ulong m_stride = t1-t0;
     346      897817 :   ulong m        = l0 + fd_ulong_min( node_t0-t0, ULONG_MAX-l0 ); /* robust against overflow */
     347    57532967 :   while( m<l1 ) {
     348    56635150 :     task( (void *)_tpool,t0,t1, args,reduce,stride, l0,l1, m,m+1UL, node_t0,node_t1 );
     349    56635150 :     m += fd_ulong_min( m_stride, ULONG_MAX-m ); /* robust against overflow */
     350    56635150 :   }
     351      897817 : FD_TPOOL_EXEC_ALL_IMPL_FTR
     352             : 
     353      911093 : FD_TPOOL_EXEC_ALL_IMPL_HDR(block)
     354      911093 :   ulong m0; ulong m1; FD_TPOOL_PARTITION( l0,l1,1UL, node_t0-t0,t1-t0, m0,m1 );
     355    58103238 :   for( ulong m=m0; m<m1; m++ ) task( (void *)_tpool,t0,t1, args,reduce,stride, l0,l1, m,m+1UL, node_t0,node_t1 );
     356      911093 : FD_TPOOL_EXEC_ALL_IMPL_FTR
     357             : 
     358             : #if FD_HAS_ATOMIC
     359      895385 : FD_TPOOL_EXEC_ALL_IMPL_HDR(taskq)
     360      895385 :   ulong * l_next = (ulong *)_tpool;
     361      895385 :   void  * tpool  = (void *)l_next[1];
     362   119326133 :   for(;;) {
     363             : 
     364             :     /* Note that we use an ATOMIC_CAS here instead of an
     365             :        ATOMIC_FETCH_AND_ADD to avoid overflow risks by having threads
     366             :        increment l0 into the tail.  ATOMIC_FETCH_AND_ADD could be used
     367             :        if there is no requirement to the effect that l1+FD_TILE_MAX does
     368             :        not overflow. */
     369             : 
     370   119326133 :     FD_COMPILER_MFENCE();
     371   119326133 :     ulong m0 = *l_next;
     372   119326133 :     FD_COMPILER_MFENCE();
     373             : 
     374   119326133 :     if( FD_UNLIKELY( m0>=l1 ) ) break;
     375   118442103 :     ulong m1 = m0+1UL;
     376   118442103 :     if( FD_UNLIKELY( FD_ATOMIC_CAS( l_next, m0, m1 )!=m0 ) ) {
     377    75628679 :       FD_SPIN_PAUSE();
     378    75628679 :       continue;
     379    75628679 :     }
     380             : 
     381    42813424 :     task( tpool,t0,t1, args,reduce,stride, l0,l1, m0,m1, node_t0,node_t1 );
     382    42813424 :   }
     383      895385 : FD_TPOOL_EXEC_ALL_IMPL_FTR
     384             : #endif
     385             : 
     386      892510 : FD_TPOOL_EXEC_ALL_IMPL_HDR(batch)
     387      892510 :   ulong m0; ulong m1; FD_TPOOL_PARTITION( l0,l1,1UL, node_t0-t0,t1-t0, m0,m1 );
     388      892510 :   task( (void *)_tpool,t0,t1, args,reduce,stride, l0,l1, m0,m1, node_t0,node_t1 );
     389      892510 : FD_TPOOL_EXEC_ALL_IMPL_FTR
     390             : 
     391     6835694 : FD_TPOOL_EXEC_ALL_IMPL_HDR(raw)
     392     6835694 :   task( (void *)_tpool,t0,t1, args,reduce,stride, l0,l1, l0,l1, node_t0,node_t1 );
     393     6835694 : FD_TPOOL_EXEC_ALL_IMPL_FTR
     394             : 
     395             : #undef FD_TPOOL_EXEC_ALL_IMPL_FTR
     396             : #undef FD_TPOOL_EXEC_ALL_IMPL_HDR

Generated by: LCOV version 1.14