GCC Code Coverage Report


Directory: src/
File: src/standalone/thread_pool.c
Date: 2026-09-30 11:11:31
Exec Total Coverage
Lines: 122 122 100.0%
Functions: 9 9 100.0%
Branches: 54 56 96.4%

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