| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * thread_pool.c - Reusable thread pool with barrier-based iteration | ||
| 6 | * | ||
| 7 | * All parallelism funnels through iterate: workers wake from a | ||
| 8 | * condvar, enter a barrier-synchronized loop (work → barrier → | ||
| 9 | * leader reduce → barrier → repeat), then return to idle. | ||
| 10 | * parallel_for is iterate with one iteration and no reduce. | ||
| 11 | */ | ||
| 12 | |||
| 13 | #include <stdlib.h> | ||
| 14 | |||
| 15 | #include "standalone/thread_pool.h" | ||
| 16 | |||
| 17 | typedef struct VsWorker | ||
| 18 | { | ||
| 19 | pthread_t thread; | ||
| 20 | uint32_t id; | ||
| 21 | uint32_t start; | ||
| 22 | uint32_t end; | ||
| 23 | } VsWorker; | ||
| 24 | |||
| 25 | struct VsThreadPool | ||
| 26 | { | ||
| 27 | VsWorker *workers; | ||
| 28 | uint32_t nthreads; | ||
| 29 | |||
| 30 | /* Current iterate parameters (set by leader before waking) */ | ||
| 31 | VsParallelForFn work_fn; | ||
| 32 | void *arg; | ||
| 33 | uint32_t max_iterations; | ||
| 34 | volatile bool keep_going; | ||
| 35 | |||
| 36 | /* | ||
| 37 | * SPMD dispatch: when spmd_fn is non-NULL, woken workers run it once (as | ||
| 38 | * participants 1..nthreads) instead of the chunked iterate loop, then | ||
| 39 | * rendezvous at the barrier. The leader runs it as participant 0. Set | ||
| 40 | * before each wake; iterate clears it so workers take the chunked path. | ||
| 41 | */ | ||
| 42 | VsSpmdFn spmd_fn; | ||
| 43 | |||
| 44 | /* Inter-iteration barrier (nthreads + 1 participants) */ | ||
| 45 | pthread_barrier_t barrier; | ||
| 46 | |||
| 47 | /* Idle synchronization */ | ||
| 48 | pthread_mutex_t mutex; | ||
| 49 | pthread_cond_t wake_cv; | ||
| 50 | uint32_t generation; | ||
| 51 | int shutdown; | ||
| 52 | }; | ||
| 53 | |||
| 54 | typedef struct WorkerArg | ||
| 55 | { | ||
| 56 | VsThreadPool *pool; | ||
| 57 | uint32_t id; | ||
| 58 | } WorkerArg; | ||
| 59 | |||
| 60 | static void * | ||
| 61 | 302 | pool_worker_fn(void *raw) | |
| 62 | { | ||
| 63 | 302 | WorkerArg *wa = (WorkerArg *)raw; | |
| 64 | 302 | VsThreadPool *pool = wa->pool; | |
| 65 | 302 | uint32_t id = wa->id; | |
| 66 | 302 | free(wa); | |
| 67 | |||
| 68 | 302 | uint32_t my_gen = 0; | |
| 69 | |||
| 70 | for (;;) | ||
| 71 | { | ||
| 72 | 740 | pthread_mutex_lock(&pool->mutex); | |
| 73 |
4/4✓ Branch 0 taken 724 times.
✓ Branch 1 taken 438 times.
✓ Branch 2 taken 422 times.
✓ Branch 3 taken 302 times.
|
1162 | while (pool->generation == my_gen && !pool->shutdown) |
| 74 | 422 | pthread_cond_wait(&pool->wake_cv, &pool->mutex); | |
| 75 | |||
| 76 |
2/2✓ Branch 0 taken 302 times.
✓ Branch 1 taken 438 times.
|
740 | if (pool->shutdown) |
| 77 | { | ||
| 78 | 302 | pthread_mutex_unlock(&pool->mutex); | |
| 79 | 302 | break; | |
| 80 | } | ||
| 81 | |||
| 82 | 438 | my_gen = pool->generation; | |
| 83 | 438 | pthread_mutex_unlock(&pool->mutex); | |
| 84 | |||
| 85 | /* | ||
| 86 | * SPMD dispatch: run the whole function once as participant id+1 | ||
| 87 | * (the leader is participant 0), then rendezvous. The function | ||
| 88 | * self-synchronizes internally; the pool only bookends the run. | ||
| 89 | */ | ||
| 90 |
2/2✓ Branch 0 taken 234 times.
✓ Branch 1 taken 204 times.
|
438 | if (pool->spmd_fn != NULL) |
| 91 | { | ||
| 92 | 234 | pool->spmd_fn(id + 1, pool->arg); | |
| 93 | 234 | pthread_barrier_wait(&pool->barrier); | |
| 94 | 234 | continue; | |
| 95 | } | ||
| 96 | |||
| 97 | /* Iterate loop with barrier synchronization */ | ||
| 98 |
2/2✓ Branch 0 taken 276 times.
✓ Branch 1 taken 38 times.
|
314 | for (uint32_t iter = 0; iter < pool->max_iterations; iter++) |
| 99 | { | ||
| 100 |
2/2✓ Branch 0 taken 244 times.
✓ Branch 1 taken 32 times.
|
276 | if (pool->workers[id].start < pool->workers[id].end) |
| 101 | 244 | pool->work_fn( | |
| 102 | id, | ||
| 103 | 244 | pool->workers[id].start, | |
| 104 | 244 | pool->workers[id].end, | |
| 105 | pool->arg); | ||
| 106 | |||
| 107 | 276 | pthread_barrier_wait(&pool->barrier); | |
| 108 | |||
| 109 | /* Leader does reduce between these two barriers */ | ||
| 110 | 276 | pthread_barrier_wait(&pool->barrier); | |
| 111 | |||
| 112 |
2/2✓ Branch 0 taken 166 times.
✓ Branch 1 taken 110 times.
|
276 | if (!pool->keep_going) |
| 113 | 166 | break; | |
| 114 | } | ||
| 115 | } | ||
| 116 | |||
| 117 | 302 | return NULL; | |
| 118 | } | ||
| 119 | |||
| 120 | VsThreadPool * | ||
| 121 | 132 | vs_thread_pool_create(uint32_t nthreads) | |
| 122 | { | ||
| 123 | 132 | VsThreadPool *pool = calloc(1, sizeof(VsThreadPool)); | |
| 124 | 132 | pool->nthreads = nthreads; | |
| 125 | 132 | pool->workers = calloc(nthreads, sizeof(VsWorker)); | |
| 126 | |||
| 127 | 132 | pthread_mutex_init(&pool->mutex, NULL); | |
| 128 | 132 | pthread_cond_init(&pool->wake_cv, NULL); | |
| 129 | |||
| 130 |
2/2✓ Branch 0 taken 122 times.
✓ Branch 1 taken 10 times.
|
132 | if (nthreads > 0) |
| 131 | 122 | pthread_barrier_init(&pool->barrier, NULL, nthreads + 1); | |
| 132 | |||
| 133 |
2/2✓ Branch 0 taken 302 times.
✓ Branch 1 taken 132 times.
|
434 | for (uint32_t i = 0; i < nthreads; i++) |
| 134 | { | ||
| 135 | 302 | pool->workers[i].id = i; | |
| 136 | 302 | WorkerArg *wa = malloc(sizeof(WorkerArg)); | |
| 137 | 302 | wa->pool = pool; | |
| 138 | 302 | wa->id = i; | |
| 139 | 302 | pthread_create(&pool->workers[i].thread, NULL, pool_worker_fn, wa); | |
| 140 | } | ||
| 141 | |||
| 142 | 132 | return pool; | |
| 143 | } | ||
| 144 | |||
| 145 | void | ||
| 146 | 50 | vs_thread_pool_iterate( | |
| 147 | VsThreadPool *pool, | ||
| 148 | uint32_t total, | ||
| 149 | VsParallelForFn work_fn, | ||
| 150 | VsReduceFn reduce_fn, | ||
| 151 | void *arg, | ||
| 152 | uint32_t max_iterations) | ||
| 153 | { | ||
| 154 |
3/4✓ Branch 0 taken 46 times.
✓ Branch 1 taken 4 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 46 times.
|
50 | if (total == 0 || max_iterations == 0) |
| 155 | 4 | return; | |
| 156 | |||
| 157 | 46 | uint32_t nt = pool->nthreads; | |
| 158 | |||
| 159 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 42 times.
|
46 | if (nt == 0) |
| 160 | { | ||
| 161 |
1/2✓ Branch 0 taken 8 times.
✗ Branch 1 not taken.
|
8 | for (uint32_t iter = 0; iter < max_iterations; iter++) |
| 162 | { | ||
| 163 | 8 | work_fn(0, 0, total, arg); | |
| 164 |
4/4✓ Branch 0 taken 6 times.
✓ Branch 1 taken 2 times.
✓ Branch 3 taken 4 times.
✓ Branch 4 taken 2 times.
|
8 | if (!reduce_fn || !reduce_fn(arg, iter)) |
| 165 | break; | ||
| 166 | } | ||
| 167 | 4 | return; | |
| 168 | } | ||
| 169 | |||
| 170 | /* Partition across N+1 (workers + leader) */ | ||
| 171 | 42 | uint32_t all = nt + 1; | |
| 172 | 42 | uint32_t per = total / all; | |
| 173 | 42 | uint32_t remainder = total % all; | |
| 174 | 42 | uint32_t off = 0; | |
| 175 | |||
| 176 |
2/2✓ Branch 0 taken 204 times.
✓ Branch 1 taken 42 times.
|
246 | for (uint32_t i = 0; i < nt; i++) |
| 177 | { | ||
| 178 |
2/2✓ Branch 0 taken 8 times.
✓ Branch 1 taken 196 times.
|
204 | uint32_t chunk = per + (i < remainder ? 1 : 0); |
| 179 | 204 | pool->workers[i].start = off; | |
| 180 | 204 | pool->workers[i].end = off + chunk; | |
| 181 | 204 | off += chunk; | |
| 182 | } | ||
| 183 | |||
| 184 | 42 | uint32_t leader_start = off; | |
| 185 | 42 | uint32_t leader_end = total; | |
| 186 | 42 | uint32_t leader_id = nt; | |
| 187 | |||
| 188 | 42 | pool->work_fn = work_fn; | |
| 189 | 42 | pool->arg = arg; | |
| 190 | 42 | pool->max_iterations = max_iterations; | |
| 191 | 42 | pool->keep_going = true; | |
| 192 | 42 | pool->spmd_fn = NULL; | |
| 193 | |||
| 194 | /* Wake workers */ | ||
| 195 | 42 | pthread_mutex_lock(&pool->mutex); | |
| 196 | 42 | pool->generation++; | |
| 197 | 42 | pthread_cond_broadcast(&pool->wake_cv); | |
| 198 | 42 | pthread_mutex_unlock(&pool->mutex); | |
| 199 | |||
| 200 | /* Leader participates in iterate loop */ | ||
| 201 |
2/2✓ Branch 0 taken 60 times.
✓ Branch 1 taken 2 times.
|
62 | for (uint32_t iter = 0; iter < max_iterations; iter++) |
| 202 | { | ||
| 203 |
2/2✓ Branch 0 taken 56 times.
✓ Branch 1 taken 4 times.
|
60 | if (leader_start < leader_end) |
| 204 | 56 | work_fn(leader_id, leader_start, leader_end, arg); | |
| 205 | |||
| 206 | /* Wait for all workers to finish this iteration */ | ||
| 207 | 60 | pthread_barrier_wait(&pool->barrier); | |
| 208 | |||
| 209 | /* Leader does reduction */ | ||
| 210 |
4/4✓ Branch 0 taken 24 times.
✓ Branch 1 taken 36 times.
✓ Branch 3 taken 20 times.
✓ Branch 4 taken 4 times.
|
60 | pool->keep_going = reduce_fn ? reduce_fn(arg, iter) : false; |
| 211 | |||
| 212 | /* Release workers for next iteration (or exit) */ | ||
| 213 | 60 | pthread_barrier_wait(&pool->barrier); | |
| 214 | |||
| 215 |
2/2✓ Branch 0 taken 40 times.
✓ Branch 1 taken 20 times.
|
60 | if (!pool->keep_going) |
| 216 | 40 | break; | |
| 217 | } | ||
| 218 | } | ||
| 219 | |||
| 220 | void | ||
| 221 | 40 | vs_thread_pool_parallel_for( | |
| 222 | VsThreadPool *pool, uint32_t total, VsParallelForFn fn, void *arg) | ||
| 223 | { | ||
| 224 | 40 | vs_thread_pool_iterate(pool, total, fn, NULL, arg, 1); | |
| 225 | 40 | } | |
| 226 | |||
| 227 | void | ||
| 228 | 116 | vs_thread_pool_launch(VsThreadPool *pool, VsSpmdFn fn, void *arg) | |
| 229 | { | ||
| 230 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 114 times.
|
116 | if (pool->nthreads == 0) |
| 231 | 2 | return; /* no workers to wake */ | |
| 232 | |||
| 233 | 114 | pool->spmd_fn = fn; | |
| 234 | 114 | pool->arg = arg; | |
| 235 | |||
| 236 | /* Wake workers — they run fn as participants 1..nthreads, then wait at | ||
| 237 | * the barrier for the leader's join. The leader returns now. */ | ||
| 238 | 114 | pthread_mutex_lock(&pool->mutex); | |
| 239 | 114 | pool->generation++; | |
| 240 | 114 | pthread_cond_broadcast(&pool->wake_cv); | |
| 241 | 114 | pthread_mutex_unlock(&pool->mutex); | |
| 242 | } | ||
| 243 | |||
| 244 | void | ||
| 245 | 116 | vs_thread_pool_join(VsThreadPool *pool) | |
| 246 | { | ||
| 247 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 114 times.
|
116 | if (pool->nthreads == 0) |
| 248 | 2 | return; | |
| 249 | |||
| 250 | /* Rendezvous with the workers, which are waiting at the barrier after | ||
| 251 | * finishing their run. */ | ||
| 252 | 114 | pthread_barrier_wait(&pool->barrier); | |
| 253 | 114 | pool->spmd_fn = NULL; | |
| 254 | } | ||
| 255 | |||
| 256 | void | ||
| 257 | 30 | vs_thread_pool_run_spmd(VsThreadPool *pool, VsSpmdFn fn, void *arg) | |
| 258 | { | ||
| 259 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 28 times.
|
30 | if (pool->nthreads == 0) |
| 260 | { | ||
| 261 | /* Serial: the calling thread is the sole participant. */ | ||
| 262 | 2 | fn(0, arg); | |
| 263 | 2 | return; | |
| 264 | } | ||
| 265 | |||
| 266 | 28 | vs_thread_pool_launch(pool, fn, arg); | |
| 267 | 28 | fn(0, arg); /* leader runs as participant 0 */ | |
| 268 | 28 | vs_thread_pool_join(pool); | |
| 269 | } | ||
| 270 | |||
| 271 | uint32_t | ||
| 272 | 4 | vs_thread_pool_nthreads(const VsThreadPool *pool) | |
| 273 | { | ||
| 274 | 4 | return pool->nthreads; | |
| 275 | } | ||
| 276 | |||
| 277 | void | ||
| 278 | 134 | vs_thread_pool_destroy(VsThreadPool *pool) | |
| 279 | { | ||
| 280 |
2/2✓ Branch 0 taken 2 times.
✓ Branch 1 taken 132 times.
|
134 | if (pool == NULL) |
| 281 | 2 | return; | |
| 282 | |||
| 283 | 132 | pthread_mutex_lock(&pool->mutex); | |
| 284 | 132 | pool->shutdown = 1; | |
| 285 | 132 | pthread_cond_broadcast(&pool->wake_cv); | |
| 286 | 132 | pthread_mutex_unlock(&pool->mutex); | |
| 287 | |||
| 288 |
2/2✓ Branch 0 taken 302 times.
✓ Branch 1 taken 132 times.
|
434 | for (uint32_t i = 0; i < pool->nthreads; i++) |
| 289 | 302 | pthread_join(pool->workers[i].thread, NULL); | |
| 290 | |||
| 291 |
2/2✓ Branch 0 taken 122 times.
✓ Branch 1 taken 10 times.
|
132 | if (pool->nthreads > 0) |
| 292 | 122 | pthread_barrier_destroy(&pool->barrier); | |
| 293 | 132 | pthread_mutex_destroy(&pool->mutex); | |
| 294 | 132 | pthread_cond_destroy(&pool->wake_cv); | |
| 295 | 132 | free(pool->workers); | |
| 296 | 132 | free(pool); | |
| 297 | } | ||
| 298 |