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
|