GCC Code Coverage Report


Directory: src/
File: src/index/parallel_build_worker.c
Date: 2026-09-30 11:11:31
Exec Total Coverage
Lines: 325 327 99.4%
Functions: 12 12 100.0%
Branches: 116 126 92.1%

Line Branch Exec Source
1 /*
2 * Copyright (c) 2026 Tiger Data, Inc.
3 * Licensed under the PostgreSQL License. See LICENSE for details.
4 *
5 * parallel_build_worker.c - Parallel index build, per-participant bodies
6 *
7 * The leader runs these as participant 0 alongside the workers:
8 * Sampling: cooperative heap scan, each participant fills its own
9 * sample slot.
10 * K-means: leader seeds the initial centroids, then iterative
11 * assignment + accumulation over the per-participant samples,
12 * barrier-synchronized with a leader reduce.
13 * Subtrees: batched largest-first schedule; each participant clusters
14 * its batch's child into a ring slot the leader consumes.
15 * Refine (when subsampled): route the full table page-backed and
16 * accumulate per-leaf means into the shared tiled accumulator.
17 * Posting: cooperative heap scan, page-backed routing + RaBitQ encode
18 * into the cluster-keyed shared sort.
19 */
20
21 #ifdef VS_STANDALONE
22 #include "standalone/parallel_ctx.h" /* ParallelWorkerNumber */
23 #include "standalone/pg_compat.h"
24 #else
25 #include <postgres.h>
26
27 #include <access/parallel.h>
28 #include <catalog/index.h>
29 #include <miscadmin.h>
30 #include <storage/barrier.h>
31 #include <storage/bufmgr.h>
32 #include <storage/shm_toc.h>
33 #include <utils/rel.h>
34 #include <utils/wait_event.h>
35 #endif
36
37 #include <inttypes.h>
38 #include <math.h>
39
40 #include "algo/hkmeans.h"
41 #include "algo/kmeans_internal.h"
42 #include "algo/vecops.h"
43 #include "core/log.h"
44 #include "core/memory.h"
45 #include "index/build_progress.h"
46 #include "index/index_build.h"
47 #include "index/parallel_build.h"
48 #include "index/posting_build.h"
49 #include "quant/rabitq.h"
50
51 /* ----------------------------------------------------------------
52 * Phase 1: Sampling callback
53 *
54 * Collects vectors into a per-worker sample slot (stride-subsampled,
55 * normalized for cosine). Initial centroids are seeded later by the
56 * leader from the pooled samples.
57 * ---------------------------------------------------------------- */
58
59 void
60 279842 prism_sample_cb(void *state, ItemPointerData tid, const float *vec)
61 {
62 279842 SampleCbState *sc = (SampleCbState *)state;
63
64 109480 (void)tid;
65
66 279842 sc->seen++;
67 /* The sampling pass walks the whole heap; surface that in the progress
68 * view in coarse batches (the leader publishes exact counts at phase
69 * boundaries). */
70
2/2
✓ Branch 0 taken 205 times.
✓ Branch 1 taken 279637 times.
279842 if ((sc->seen & 1023) == 0)
71 205 prism_build_progress_incr_tuples(1024);
72
73 /* Stride-based subsampling */
74
2/2
✓ Branch 0 taken 22941 times.
✓ Branch 1 taken 256901 times.
279842 if (sc->stride_counter > 0)
75 {
76 22941 sc->stride_counter--;
77 22941 return;
78 }
79 256901 sc->stride_counter = sc->stride - 1;
80
81
2/2
✓ Branch 0 taken 88222 times.
✓ Branch 1 taken 168679 times.
256901 if (sc->count >= sc->max_samples)
82 21488 return;
83
84 215608 Dimension dim = sc->dim;
85 215608 float *dest = sc->samples + (size_t)sc->count * dim;
86
87 /* Normalize for cosine. Zero-norm vectors are dropped from the
88 * sample: they carry no direction, so under cosine they sit at the
89 * origin -- equidistant from every real point -- and k-means drags
90 * centroids toward them, warping the encode references (and so the
91 * distance estimates) for every vector in the affected subtree.
92 * They stay in the posting lists like any other row; they just
93 * don't get a vote on the clustering. The check runs on the source
94 * so a dropped vector costs no copy, and the scale writes straight
95 * into the sample slot, normalizing and copying in one pass. */
96
2/2
✓ Branch 0 taken 5159 times.
✓ Branch 1 taken 210449 times.
215608 if (sc->metric == DISTANCE_COSINE)
97 {
98 5159 float norm_sq = vs_l2_norm_squared(vec, dim);
99
100
2/2
✓ Branch 0 taken 3358 times.
✓ Branch 1 taken 1801 times.
5159 if (norm_sq == 0.0f)
101 80 return;
102 5078 vec32_scale(vec, 1.0f / sqrtf(norm_sq), dest, dim);
103 }
104 else
105 210449 memcpy(dest, vec, dim * sizeof(float));
106
107 215527 sc->count++;
108 }
109
110 /* ----------------------------------------------------------------
111 * Phase 2: K-means assignment — thin wrappers around shared kernel
112 * ---------------------------------------------------------------- */
113
114 void
115 3333 prism_km_assign_and_accumulate(
116 const float *samples,
117 uint32_t nsamples,
118 const float *centroids,
119 const float *norms_c,
120 uint32_t nlist,
121 Dimension dim,
122 DistanceMetric metric,
123 float *out_sums,
124 uint32_t *out_cnts,
125 float *out_cost)
126 {
127 3333 memset(out_sums, 0, (size_t)nlist * dim * sizeof(float));
128 3333 memset(out_cnts, 0, nlist * sizeof(uint32_t));
129 3333 *out_cost = 0.0f;
130 3333 kmeans_assign_accumulate(
131 samples,
132 NULL,
133 0,
134 nsamples,
135 centroids,
136 norms_c,
137 nlist,
138 dim,
139 metric,
140 NULL,
141 0,
142 out_sums,
143 out_cnts,
144 out_cost);
145 3333 }
146
147 /*
148 * Build ONE root-child's subtree into `slot` (a DSM ring slot indexed by
149 * participant, not by child, so the batched streaming build keeps only
150 * nparticipants subtrees resident). The subtree is clustered from this child's
151 * root-assigned samples, indexed in place (no contiguous per-child copy).
152 */
153 void
154 239 prism_build_child_subtree(
155 uint32_t child,
156 uint32_t child_count,
157 int nparticipants,
158 PrismDsmSamples *dsm_samples,
159 PrismDsmRootAssign *dsm_ra,
160 const float *root_cents,
161 uint32_t nlist,
162 uint32_t fan_out,
163 Dimension dim,
164 DistanceMetric metric,
165 uint32_t km_max_iterations,
166 char *slot,
167 uint64_t slot_size)
168 {
169 239 uint32_t nlist_c = (nlist + fan_out - 1) / fan_out;
170 /* The caller's histogram already counted this child's samples (one pass
171 * over the assignments instead of one per child). */
172 239 uint32_t cc = child_count;
173
174 239 KMeansOptions opts = VS_KMEANS_OPTIONS_DEFAULT;
175 239 opts.max_iterations = km_max_iterations;
176 239 opts.algorithm = KMEANS_ALGO_LLOYD;
177 239 opts.initial_centroids = NULL;
178
179 239 const size_t vec_nbytes = (size_t)dim * sizeof(float);
180
181 239 HKMeansResult *sub = NULL;
182
2/2
✓ Branch 0 taken 3 times.
✓ Branch 1 taken 236 times.
239 if (cc == 0)
183 {
184 3 float *seed = vs_alloc(vec_nbytes);
185 3 memcpy(seed, root_cents + (size_t)child * dim, vec_nbytes);
186 3 sub = vs_hkmeans_f32(
187 seed, 1, NULL, dim, nlist_c, fan_out, metric, &opts);
188 3 vs_free(seed);
189 }
190 else
191 {
192 236 const float *vbase = prism_dsm_worker_samples(dsm_samples, 0);
193 236 uint32_t mpw = dsm_samples->max_per_worker;
194 236 uint32_t *idx = vs_alloc((size_t)cc * sizeof(uint32_t));
195 236 uint32_t g = 0;
196
3/3
✓ Branch 0 taken 364 times.
✓ Branch 1 taken 455 times.
✓ Branch 2 taken 110 times.
929 for (int t = 0; t < nparticipants; t++)
197 {
198 693 const uint32_t *ra = prism_dsm_root_assignments(dsm_ra, t);
199 693 uint32_t n = prism_dsm_sample_counts(dsm_samples)[t];
200
2/2
✓ Branch 0 taken 1229412 times.
✓ Branch 1 taken 693 times.
1230105 for (uint32_t i = 0; i < n; i++)
201
2/2
✓ Branch 0 taken 140526 times.
✓ Branch 1 taken 1088886 times.
1229412 if (ra[i] == child)
202 140526 idx[g++] = (uint32_t)t * mpw + i;
203 }
204 236 sub = vs_hkmeans_f32(
205 vbase, cc, idx, dim, nlist_c, fan_out, metric, &opts);
206 236 vs_free(idx);
207 }
208
209
1/2
✓ Branch 0 taken 239 times.
✗ Branch 1 not taken.
239 if (sub != NULL)
210 {
211
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 239 times.
239 if ((uint64_t)sub->total_size > slot_size)
212 ✗ vs_error(
213 "prism: subtree blob %u exceeds slot "
214 "(%u > %" PRIu64 ")",
215 child,
216 sub->total_size,
217 slot_size);
218 239 memcpy(slot, sub, sub->total_size);
219 239 vs_free(sub);
220 }
221 239 }
222
223 void
224 191 prism_pbuild_count_children(
225 PrismDsmSamples *dsm_samples,
226 PrismDsmRootAssign *dsm_ra,
227 int nparticipants,
228 uint32_t km_k,
229 uint32_t *out_counts)
230 {
231 191 memset(out_counts, 0, (size_t)km_k * sizeof(uint32_t));
232
2/2
✓ Branch 0 taken 575 times.
✓ Branch 1 taken 191 times.
766 for (int t = 0; t < nparticipants; t++)
233 {
234 575 const uint32_t *ra = prism_dsm_root_assignments(dsm_ra, t);
235 575 uint32_t n = prism_dsm_sample_counts(dsm_samples)[t];
236
2/2
✓ Branch 0 taken 589562 times.
✓ Branch 1 taken 575 times.
590137 for (uint32_t i = 0; i < n; i++)
237 589562 out_counts[ra[i]]++;
238 }
239 191 }
240
241 void
242 141 prism_pbuild_stream_subtrees(
243 int participant_id,
244 int nparticipants,
245 PrismDsmSamples *dsm_samples,
246 PrismDsmRootAssign *dsm_ra,
247 const float *root_cents,
248 uint32_t km_k,
249 uint32_t nlist,
250 uint32_t fan_out,
251 Dimension dim,
252 DistanceMetric metric,
253 uint32_t km_max_iterations,
254 char *subtrees_base,
255 uint64_t slot_size,
256 Barrier *barrier,
257 PrismBatchCb batch_cb,
258 void *cb_arg,
259 uint32_t *out_child_order)
260 {
261 141 uint32_t np = (uint32_t)nparticipants;
262 141 uint32_t nbatches = (km_k + np - 1) / np;
263 57 char *slot =
264 141 prism_dsm_child_subtree(subtrees_base, participant_id, slot_size);
265
266 /* One histogram pass over the root assignments feeds every child's
267 * sample count (prism_build_child_subtree needs it, and counting per
268 * child would re-scan the assignments km_k times). */
269 141 uint32_t *child_count = vs_alloc0((size_t)km_k * sizeof(uint32_t));
270 141 prism_pbuild_count_children(
271 dsm_samples, dsm_ra, nparticipants, km_k, child_count);
272
273 /* Schedule children largest-first (LPT): every batch waits for its
274 * slowest subtree at two barriers, so the skewed children must land in
275 * the full batches, leaving the tail batch the small ones. Counts are
276 * identical for every participant (same shared assignments) and ties
277 * break on the child id, so all participants and the leader's blob
278 * replay derive the same order with no coordination. */
279 141 uint32_t *order = vs_alloc((size_t)km_k * sizeof(uint32_t));
280
3/3
✓ Branch 0 taken 364 times.
✓ Branch 1 taken 424 times.
✓ Branch 2 taken 57 times.
845 for (uint32_t i = 0; i < km_k; i++)
281 704 order[i] = i;
282
2/2
✓ Branch 0 taken 563 times.
✓ Branch 1 taken 141 times.
704 for (uint32_t i = 1; i < km_k; i++)
283 {
284 563 uint32_t id = order[i];
285 563 uint32_t j = i;
286
6/6
✓ Branch 0 taken 1348 times.
✓ Branch 1 taken 209 times.
✓ Branch 2 taken 994 times.
✓ Branch 3 taken 354 times.
✓ Branch 4 taken 64 times.
✓ Branch 5 taken 134 times.
1557 while (j > 0 && (child_count[order[j - 1]] < child_count[id] ||
287
2/2
✓ Branch 0 taken 9 times.
✓ Branch 1 taken 211 times.
220 (child_count[order[j - 1]] == child_count[id] &&
288
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 9 times.
9 order[j - 1] > id)))
289 {
290 994 order[j] = order[j - 1];
291 994 j--;
292 }
293 563 order[j] = id;
294 }
295
2/2
✓ Branch 0 taken 49 times.
✓ Branch 1 taken 92 times.
141 if (out_child_order != NULL)
296 49 memcpy(out_child_order, order, (size_t)km_k * sizeof(uint32_t));
297
298
2/2
✓ Branch 0 taken 285 times.
✓ Branch 1 taken 139 times.
424 for (uint32_t b = 0; b < nbatches; b++)
299 {
300 285 uint32_t my_idx = b * np + (uint32_t)participant_id;
301
2/2
✓ Branch 0 taken 239 times.
✓ Branch 1 taken 46 times.
285 if (my_idx < km_k)
302 {
303 239 uint32_t my_child = order[my_idx];
304 239 prism_build_child_subtree(
305 my_child,
306 239 child_count[my_child],
307 nparticipants,
308 dsm_samples,
309 dsm_ra,
310 root_cents,
311 nlist,
312 fan_out,
313 dim,
314 metric,
315 km_max_iterations,
316 slot,
317 slot_size);
318 }
319
320 /* All participants have built this batch's subtrees into their slots.
321 */
322 285 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
323
324 /* Leader consumes the batch before the slots are reused. */
325
3/4
✓ Branch 0 taken 102 times.
✓ Branch 1 taken 181 times.
✓ Branch 2 taken 56 times.
✗ Branch 3 not taken.
283 if (participant_id == 0 && batch_cb != NULL)
326 {
327 102 uint32_t base = b * np;
328 102 uint32_t bs = km_k - base;
329
2/2
✓ Branch 0 taken 26 times.
✓ Branch 1 taken 30 times.
102 if (bs > np)
330 26 bs = np;
331 102 batch_cb(cb_arg, order + base, bs, subtrees_base, slot_size);
332 }
333
334 /* Leader done with the batch; slots free for the next batch. */
335 283 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
336 }
337
338 139 vs_free(order);
339 139 vs_free(child_count);
340 139 }
341
342 /* ----------------------------------------------------------------
343 * Phases 1, 2, 2b: per-participant execution, shared by leader and workers
344 *
345 * The leader runs as participant 0 and calls these exactly as a worker does,
346 * rather than duplicating the bodies. The phase barriers live inside each
347 * function, so the leader and workers stay in lockstep by construction; the
348 * leader-only coordination (seeding the initial centroids, reducing the
349 * per-iteration accumulators) is gated on participant_id == 0 and runs at the
350 * barrier rendezvous while the other participants wait.
351 * ---------------------------------------------------------------- */
352
353 void
354 318 prism_pbuild_exec_sampling(
355 int participant_id,
356 Relation heap,
357 Relation index,
358 struct IndexInfo *index_info,
359 PrismBuildShared *shared,
360 PrismDsmSamples *dsm_samples,
361 Barrier *barrier)
362 {
363 318 Dimension dim = shared->dim;
364
365 /* Stride to subsample ~max_samples_per_worker from this participant's
366 * share of the heap. */
367 318 double est_rows = RelationGetNumberOfBlocks(heap) *
368 318 (BLCKSZ / (double)(dim * sizeof(float) + 32));
369 318 double est_per = est_rows / shared->nparticipants;
370 318 uint32_t stride = 1;
371
2/2
✓ Branch 0 taken 99 times.
✓ Branch 1 taken 219 times.
318 if (est_per > shared->max_samples_per_worker)
372 99 stride = (uint32_t)(est_per / shared->max_samples_per_worker);
373
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 188 times.
273 if (stride < 1)
374 130 stride = 1;
375
376 318 SampleCbState sc = {
377 318 .samples = prism_dsm_worker_samples(dsm_samples, participant_id),
378 .count = 0,
379 .seen = 0,
380 188 .max_samples = shared->max_samples_per_worker,
381 .stride = stride,
382 .stride_counter = 0,
383 .dim = dim,
384 318 .metric = shared->metric,
385 };
386
387 /* progress: only the leader (participant 0) drives the PG progress view.
388 * The scan count is discarded: the posting pass is the one full scan that
389 * feeds shared->reltuples, so the extra passes must not double-count. */
390 318 (void)prism_build_scan(
391 heap,
392 index,
393 index_info,
394 shared,
395 true,
396 participant_id == 0,
397 prism_sample_cb,
398 &sc);
399
400 318 prism_dsm_sample_counts(dsm_samples)[participant_id] = sc.count;
401 318 prism_dsm_sample_seen(dsm_samples)[participant_id] = sc.seen;
402
403 /* Flush the sub-batch remainder so the participant's contribution to
404 * the progress view is exact (and so every participant makes at least
405 * one flush per scan, whatever share the work-stealing gave it). */
406 318 prism_build_progress_incr_tuples((int64_t)(sc.seen & 1023));
407
408 /* Barrier: all participants done sampling. */
409 318 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
410 318 }
411
412 uint32_t
413 318 prism_pbuild_exec_kmeans(
414 int participant_id,
415 PrismBuildShared *shared,
416 PrismDsmSamples *dsm_samples,
417 char *centroids_base,
418 char *km_workers_base,
419 Barrier *barrier)
420 {
421 318 Dimension dim = shared->dim;
422 318 uint32_t km_k = shared->km_k;
423 318 int nparticipants = shared->nparticipants;
424 318 float *cents = prism_dsm_centroids(centroids_base);
425
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
318 float *norms_c = prism_dsm_norms_c(centroids_base, km_k, dim);
426 /* One seed vector, and the whole km_k-centroid array. */
427 318 const size_t vec_nbytes = (size_t)dim * sizeof(float);
428 318 const size_t cents_nbytes = (size_t)km_k * vec_nbytes;
429
430 /*
431 * Leader-only seed: pick km_k initial centroids from the pooled samples,
432 * spread evenly across the concatenation of every participant's slot.
433 * Centralizing it (rather than slicing by participant) keeps init
434 * independent of how many workers launched and robust to a sample-starved
435 * one. The result varies run to run (work-stealing scan order), so the
436 * parallel build is not bit-reproducible.
437 */
438
2/2
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
318 if (participant_id == 0)
439 {
440 82 uint32_t total_ns = 0;
441
2/2
✓ Branch 0 taken 318 times.
✓ Branch 1 taken 126 times.
444 for (int t = 0; t < nparticipants; t++)
442 318 total_ns += prism_dsm_sample_counts(dsm_samples)[t];
443
444
3/4
✓ Branch 0 taken 124 times.
✓ Branch 1 taken 2 times.
✓ Branch 2 taken 82 times.
✗ Branch 3 not taken.
126 uint32_t step = (km_k > 0 && total_ns >= km_k) ? total_ns / km_k : 1;
445
2/2
✓ Branch 0 taken 742 times.
✓ Branch 1 taken 126 times.
868 for (uint32_t i = 0; i < km_k; i++)
446 {
447
2/2
✓ Branch 0 taken 741 times.
✓ Branch 1 taken 1 times.
742 uint32_t gidx = (total_ns > 0) ? (i * step) % total_ns : 0;
448
449 742 int t = 0;
450 742 uint32_t base = 0;
451
2/2
✓ Branch 0 taken 1466 times.
✓ Branch 1 taken 1 times.
1467 while (t < nparticipants &&
452
3/3
✓ Branch 0 taken 509 times.
✓ Branch 1 taken 777 times.
✓ Branch 2 taken 180 times.
1466 base + prism_dsm_sample_counts(dsm_samples)[t] <= gidx)
453 {
454 725 base += prism_dsm_sample_counts(dsm_samples)[t];
455 725 t++;
456 }
457
2/2
✓ Branch 0 taken 741 times.
✓ Branch 1 taken 1 times.
742 if (t < nparticipants)
458 741 memcpy(cents + (size_t)i * dim,
459 741 prism_dsm_worker_samples(dsm_samples, t) +
460 741 (size_t)(gidx - base) * dim,
461 vec_nbytes);
462 }
463
464
2/2
✓ Branch 0 taken 117 times.
✓ Branch 1 taken 9 times.
126 if (shared->metric == DISTANCE_L2)
465
2/2
✓ Branch 0 taken 685 times.
✓ Branch 1 taken 117 times.
802 for (uint32_t j = 0; j < km_k; j++)
466 685 norms_c[j] = vs_l2_norm_squared(cents + (size_t)j * dim, dim);
467 }
468
469 /* Barrier: initial centroids + norms ready; guards the reads below against
470 * the leader's seed write. */
471 318 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
472
473 318 float *my_samples = prism_dsm_worker_samples(dsm_samples, participant_id);
474
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
318 uint32_t my_n = prism_dsm_sample_counts(dsm_samples)[participant_id];
475
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
318 float *my_sums = prism_dsm_km_worker_sums(
476 km_workers_base, km_k, dim, participant_id);
477
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
318 uint32_t *my_cnts = prism_dsm_km_worker_cnts(
478 km_workers_base, km_k, dim, participant_id);
479
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
318 float *my_cost = prism_dsm_km_worker_cost(
480 km_workers_base, km_k, dim, participant_id);
481
482 /* Reduce scratch is leader-only. */
483
2/2
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
318 float *old_cents = participant_id == 0 ? vs_alloc(cents_nbytes) : NULL;
484 318 Size km_sz = prism_dsm_km_workers_size(nparticipants, km_k, dim);
485
486 318 uint32_t iters = 0;
487
2/2
✓ Branch 0 taken 3333 times.
✓ Branch 1 taken 107 times.
3440 for (uint32_t iter = 0; iter < shared->km_max_iterations; iter++)
488 {
489 3333 iters++;
490
491 /* Every participant assigns + accumulates over its own samples. */
492 3333 prism_km_assign_and_accumulate(
493 my_samples,
494 my_n,
495 cents,
496 norms_c,
497 km_k,
498 dim,
499 shared->metric,
500 my_sums,
501 my_cnts,
502 my_cost);
503
504 /* Barrier: all accumulators written; the leader reduces. */
505 3333 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
506
507
2/2
✓ Branch 0 taken 1174 times.
✓ Branch 1 taken 2159 times.
3333 if (participant_id == 0)
508 {
509 1174 memcpy(old_cents, cents, cents_nbytes);
510
511 1174 const float **all_sums = vs_alloc(nparticipants * sizeof(float *));
512 1174 const uint32_t **all_cnts = vs_alloc(
513 nparticipants * sizeof(uint32_t *));
514 1174 float *all_costs = vs_alloc(nparticipants * sizeof(float));
515
3/3
✓ Branch 0 taken 1840 times.
✓ Branch 1 taken 2173 times.
✓ Branch 2 taken 494 times.
4507 for (int t = 0; t < nparticipants; t++)
516 {
517 3333 all_sums[t] = prism_dsm_km_worker_sums(
518 km_workers_base, km_k, dim, t);
519 3333 all_cnts[t] = prism_dsm_km_worker_cnts(
520 km_workers_base, km_k, dim, t);
521 3333 all_costs[t] = *prism_dsm_km_worker_cost(
522 km_workers_base, km_k, dim, t);
523 }
524
525 494 float total_cost;
526 1174 float shift_sq = kmeans_merge_centroids(
527 cents,
528 norms_c,
529 old_cents,
530 all_sums,
531 all_cnts,
532 all_costs,
533 nparticipants,
534 km_k,
535 dim,
536 shared->metric,
537 &total_cost);
538
539 1174 float tol_sq = shared->km_tolerance * shared->km_tolerance;
540 1174 shared->km_converged = (shift_sq < tol_sq);
541 1174 memset(km_workers_base, 0, km_sz);
542
543 1174 vs_free(all_sums);
544 1174 vs_free(all_cnts);
545 1174 vs_free(all_costs);
546 }
547
548 /* Barrier: updated centroids + convergence flag visible to all. */
549 3333 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
550
551
2/2
✓ Branch 0 taken 1506 times.
✓ Branch 1 taken 1827 times.
3333 if (shared->km_converged)
552 112 break;
553 }
554
555
2/2
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 192 times.
318 if (old_cents != NULL)
556 126 vs_free(old_cents);
557
558 318 return iters;
559 }
560
561 void
562 318 prism_pbuild_exec_root_assign(
563 int participant_id,
564 PrismBuildShared *shared,
565 PrismDsmSamples *dsm_samples,
566 PrismDsmRootAssign *dsm_ra,
567 char *centroids_base,
568 Barrier *barrier)
569 {
570 318 Dimension dim = shared->dim;
571 318 uint32_t km_k = shared->km_k;
572 318 float *cents = prism_dsm_centroids(centroids_base);
573 318 float *norms_c = prism_dsm_norms_c(centroids_base, km_k, dim);
574
575 506 kmeans_assign(
576 318 prism_dsm_worker_samples(dsm_samples, participant_id),
577 0,
578 318 prism_dsm_sample_counts(dsm_samples)[participant_id],
579 cents,
580 norms_c,
581 km_k,
582 dim,
583 shared->metric,
584 prism_dsm_root_assignments(dsm_ra, participant_id));
585
586 /* Barrier: all root assignments written; subtree gather can read them. */
587 318 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
588 318 }
589
590 /* ----------------------------------------------------------------
591 * Phase 2.5: parallel full-table leaf-encode-reference refinement
592 *
593 * Page-backed mirror of the serial serial_refine_heads: route every row
594 * exactly as the query/insert do (prism_query_route k=1 over the centroid
595 * pages, then head -> leaf), accumulate per-leaf means into the tiled DSM
596 * accumulator, and rewrite each leaf's head-page pt_centroid to the full-table
597 * mean. No in-RAM tree.
598 * ---------------------------------------------------------------- */
599
600 typedef struct RefineCbState
601 {
602 PrismBuildShared *shared;
603 PrismQueryState *qs; /* page-backed router (workers) */
604 BlockNumber first_posting; /* head -> leaf: leaf = head - first_posting */
605 double *sums; /* shared accumulator, indexed leaf - tile_lo */
606 uint64_t *counts; /* shared accumulator */
607 Dimension dim;
608 bool cosine;
609 uint32_t tile_lo; /* accumulate only leaves in [tile_lo, tile_hi) */
610 uint32_t tile_hi;
611 float *scratch; /* per-participant normalized copy (cosine) */
612 } RefineCbState;
613
614 static void
615 12600 prism_refine_cb(void *state, ItemPointerData tid, const float *vec)
616 {
617 12600 RefineCbState *rs = (RefineCbState *)state;
618 12600 Dimension dim = rs->dim;
619
620 12600 (void)tid;
621
622 12600 uint32_t idx;
623 25200 const float *v = prism_refine_route_row(
624 12600 rs->qs,
625 rs->first_posting,
626 vec,
627 dim,
628 12600 rs->cosine,
629 rs->scratch,
630 rs->tile_lo,
631 rs->tile_hi,
632 &idx);
633
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 12600 times.
12600 if (v == NULL)
634 ✗ return;
635
636 12600 uint32_t stripe = idx % PRISM_REFINE_LOCK_STRIPES;
637
638 12600 prism_pbuild_accum_lock(rs->shared, stripe);
639 12600 double *sum = rs->sums + (size_t)idx * dim;
640
2/2
✓ Branch 0 taken 1228800 times.
✓ Branch 1 taken 12600 times.
1241400 for (Dimension j = 0; j < dim; j++)
641 1228800 sum[j] += v[j];
642 12600 rs->counts[idx]++;
643 12600 prism_pbuild_accum_unlock(rs->shared, stripe);
644 }
645
646 void
647 6 prism_pbuild_exec_refine_paged(
648 int participant_id,
649 Relation heap,
650 Relation index,
651 struct IndexInfo *index_info,
652 PrismBuildShared *shared,
653 struct PrismQueryState *qs,
654 BlockNumber first_posting,
655 PrismDsmRefineAccum *accum,
656 Barrier *barrier,
657 PrismLeafWriteFn write_head,
658 void *write_head_ctx)
659 {
660 6 Dimension dim = shared->dim;
661 6 uint32_t nleaves = shared->nlist; /* actual leaf count (published) */
662 6 double *sums = prism_dsm_refine_sums(accum);
663 6 uint64_t *counts = prism_dsm_refine_counts(accum);
664
665 /* The accumulator holds at most accum->nleaves leaves (the bounded tile
666 * capacity), so leaves are processed in tiles, re-scanning the heap per
667 * tile. nleaves <= capacity is a single tile (the common case). Leader and
668 * workers derive `tile` identically from the DSM capacity, so the barrier
669 * sequence below stays in lockstep. */
670 6 uint32_t tile = accum->nleaves;
671
672 12 RefineCbState rs = {
673 .shared = shared,
674 .qs = qs,
675 .first_posting = first_posting,
676 .sums = sums,
677 .counts = counts,
678 .dim = dim,
679 6 .cosine = (shared->metric == DISTANCE_COSINE),
680 6 .scratch = vs_alloc((size_t)dim * sizeof(float)),
681 };
682
683 /* Single pass: routing reads only the centroid pages, which refine
684 * never rewrites, so the assignment is a fixed point and one pass
685 * reaches it. */
686 {
687
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 6 times.
12 for (uint32_t lo = 0; lo < nleaves; lo += tile)
688 {
689 6 uint32_t hi = (lo + tile < nleaves) ? lo + tile : nleaves;
690 6 rs.tile_lo = lo;
691 6 rs.tile_hi = hi;
692
693 /* Leader zeroes the (tile-sized) accumulator and resets the scan.
694 */
695
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 4 times.
6 if (participant_id == 0)
696 {
697 2 memset(sums, 0, (size_t)(hi - lo) * dim * sizeof(double));
698 2 memset(counts, 0, (size_t)(hi - lo) * sizeof(uint64_t));
699 2 prism_pbuild_rescan(heap, shared);
700 }
701 /* Barrier: accumulator cleared + scan reset before anyone scans.
702 */
703 6 BarrierArriveAndWait(
704 barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
705
706 /* Workers cooperatively scan the heap and accumulate page-backed;
707 * the leader does not route (it has no qs), it only
708 * clears/divides, mirroring the phase-3 division of labor. */
709
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 2 times.
6 if (participant_id != 0)
710 4 (void)prism_build_scan(
711 heap,
712 index,
713 index_info,
714 shared,
715 true,
716 false,
717 prism_refine_cb,
718 &rs);
719
720 /* Barrier: every row accumulated before the leader divides. */
721 6 BarrierArriveAndWait(
722 barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
723
724 /* Leader rewrites this tile's leaf head pages = per-leaf means. */
725
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 4 times.
6 if (participant_id == 0)
726 2 prism_refine_write_means(
727 sums,
728 counts,
729 lo,
730 hi,
731 dim,
732 rs.scratch,
733 write_head,
734 write_head_ctx);
735 /* Barrier: refined heads written before the next tile/pass. */
736 6 BarrierArriveAndWait(
737 barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
738 }
739 }
740
741 6 vs_free(rs.scratch);
742 6 }
743
744 /* ----------------------------------------------------------------
745 * Phase 3: posting scan -> cluster-keyed sort (sort-seam path)
746 *
747 * Each worker routes every vector page-backed (the same prism_query_route the
748 * query and insert paths use), RaBitQ-encodes against the target list's head
749 * pt_centroid, and feeds the compact entry into the shared cluster-keyed
750 * sorter (primary + optional SOAR / boundary secondary), via the shared
751 * PrismBuildRouteCtx helper. The leader merges and builds the pages. Memory is
752 * bounded by maintenance_work_mem inside the sorter.
753 * ---------------------------------------------------------------- */
754 static void
755 278842 route_scan_cb(void *state, ItemPointerData tid, const float *vec)
756 {
757 278842 PrismBuildRouteCtx *ctx = (PrismBuildRouteCtx *)state;
758
759 278842 prism_build_route_emit(ctx, vec, tid);
760
2/2
✓ Branch 0 taken 204 times.
✓ Branch 1 taken 278638 times.
278842 if (((uint64_t)ctx->indtuples & 1023) == 0)
761 204 prism_build_progress_incr_tuples(1024);
762 278842 }
763
764 /* ----------------------------------------------------------------
765 * Worker entry point — multi-phase build
766 * ---------------------------------------------------------------- */
767
768 void
769 192 prism_parallel_build_main(dsm_segment *seg, shm_toc *toc)
770 {
771 86 PrismPBuildWorker w;
772 192 prism_pbuild_worker_attach(toc, &w);
773
774 192 PrismBuildShared *shared = w.shared;
775 192 Barrier *barrier = w.barrier;
776 192 Relation heapRel = w.heapRel;
777 192 Relation indexRel = w.indexRel;
778 192 int worker_id = w.worker_id;
779 192 Dimension dim = w.dim;
780
781 /* ---- Phases 1, 2, 2b: the shared per-participant bodies (the leader runs
782 * the very same code as participant 0). Each call includes its phase
783 * barrier(s). ---- */
784 192 void *sample_seg = NULL;
785 86 PrismDsmSamples *dsm_samples =
786 192 prism_pbuild_samples_attach(toc, shared, &sample_seg);
787 192 char *centroids_base = shm_toc_lookup(toc, PRISM_DSM_KEY_CENTROIDS, false);
788 192 float *cents = prism_dsm_centroids(centroids_base);
789 86 char *km_workers_base =
790 192 shm_toc_lookup(toc, PRISM_DSM_KEY_KM_WORKERS, false);
791 86 PrismDsmRootAssign *dsm_ra =
792 192 shm_toc_lookup(toc, PRISM_DSM_KEY_ROOT_ASSIGN, false);
793 192 IndexInfo *indexInfo = BuildIndexInfo(indexRel);
794 #ifndef VS_STANDALONE
795 /* The leader marked the build concurrent and scans with an MVCC snapshot;
796 * the worker's freshly built IndexInfo defaults to non-concurrent, so it
797 * must be aligned or heapam's snapshot/OldestXmin assert trips in the scan
798 * (a valid OldestXmin paired with an MVCC snapshot). Matches nbtsort. */
799 86 indexInfo->ii_Concurrent = shared->concurrent;
800 #endif
801
802 192 prism_pbuild_exec_sampling(
803 worker_id,
804 heapRel,
805 indexRel,
806 indexInfo,
807 shared,
808 dsm_samples,
809 barrier);
810 192 prism_pbuild_exec_kmeans(
811 worker_id,
812 shared,
813 dsm_samples,
814 centroids_base,
815 km_workers_base,
816 barrier);
817 192 prism_pbuild_exec_root_assign(
818 worker_id, shared, dsm_samples, dsm_ra, centroids_base, barrier);
819
820 /* ---- Phase 2c: batched streaming subtree build (page-backed) ----
821 *
822 * For a hierarchical tree (>= 2 levels) the workers build per-root-child
823 * subtrees into a bounded ring of slots (subtree DSM is nparticipants
824 * slots, independent of the partition count); the leader consumes each
825 * batch, recording layout counts and keeping the blob in a spillable
826 * store, and later streams every subtree's pages by itself from that
827 * store. Workers pass no callback (leader-only consumes each batch). The
828 * barrier sequence is inside prism_pbuild_stream_subtrees, identical for
829 * leader and workers. A flat (1-level) build has no subtrees — the leader
830 * writes the single level directly, and neither side runs the subtree
831 * barriers. */
832
2/2
✓ Branch 1 taken 92 times.
✓ Branch 2 taken 100 times.
192 if (vs_hkmeans_nlevels(shared->nlist, shared->fan_out) >= 2)
833 {
834 /* Ring barrier: the leader creates the subtree ring (sized from the
835 * per-child sample counts) and publishes its handle + slot size
836 * before arriving; attach only after. */
837 92 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
838 92 void *ring_seg = NULL;
839 38 char *subtrees_base =
840 92 prism_pbuild_subtree_ring_attach(shared, &ring_seg);
841 92 prism_pbuild_stream_subtrees(
842 worker_id,
843 shared->nparticipants,
844 dsm_samples,
845 dsm_ra,
846 cents,
847 shared->km_k,
848 shared->nlist,
849 shared->fan_out,
850 dim,
851 shared->metric,
852 shared->km_max_iterations,
853 subtrees_base,
854 shared->subtree_slot_size,
855 barrier,
856 NULL,
857 NULL,
858 NULL);
859 90 prism_pbuild_subtree_ring_release(ring_seg);
860 }
861
862 /* Barrier: leader finished streaming the centroid tree + published the
863 * routing state + initialized the sorter; workers build their page-backed
864 * router next (from the just-published shared state). */
865 190 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
866
867 /* ---- Phase 3 setup: page-backed router shared by refine + posting scan
868 * --- */
869 190 void *sortshared = shm_toc_lookup(toc, PRISM_DSM_KEY_SORTSHARED, false);
870 190 BlockNumber first_posting = shared->first_posting;
871 84 const float *global_mean =
872 190 shm_toc_lookup(toc, PRISM_DSM_KEY_GLOBAL_MEAN, false);
873
874 /* Exact internal-node centroids, published by the leader before the
875 * tree-ready barrier above: the build descent scores the internal
876 * tree levels exactly (an empty collection leaves the hook inert). */
877 190 void *exact_seg = NULL;
878 84 PrismExactInternalCentroids exact_centroids;
879 190 prism_exact_centroid_collection_view(
880 190 prism_pbuild_exact_centroids_attach(shared, &exact_seg),
881 &exact_centroids);
882
883 190 uint32_t entry_size = (uint32_t)prism_posting_entry_size(dim);
884 190 RaBitQParams *rq_params = vs_rabitq_create(dim, shared->rabitq_seed);
885
886 /* Per-worker storage over the index for page-backed head/centroid reads
887 * (PG opens one on the worker's indexRel; standalone shares the leader's).
888 */
889 190 VsStorage *storage = prism_pbuild_worker_storage(&w);
890
891 /* Routing base — the same PrismIndexBase the query/insert build, so the
892 * worker routes each row identically. nlevels + first_centroid (the
893 * streamed tree's root block) come from the shared state the leader
894 * published; the scales + global mean + fastscan bits also from shared. */
895 84 PrismIndexBase base;
896 274 prism_build_router_base_init(
897 &base,
898 rq_params,
899 storage,
900 dim,
901 190 shared->nlevels,
902 shared->first_centroid,
903 shared->metric,
904 shared->centroid_format,
905 shared->fastscan_bits,
906 190 shared->centroid_error_scale,
907 190 shared->centroid_beam_scale,
908 shared->fan_out,
909 shared->nlist,
910 shared->rabitq_seed,
911 global_mean,
912 190 vs_alloc((size_t)dim * sizeof(float)));
913 /* Build-only accuracy hook: exact scoring of the internal tree levels
914 * (the query and insert paths never set this). */
915 190 base.exact_internal = &exact_centroids;
916
917 84 PrismQueryState qs;
918 190 prism_query_state_init(&qs, &base, 1, PRISM_SECONDARY_TOPK);
919
920 /* ---- Phase 2.5: page-backed full-table refine (only when subsampled)
921 * ---- Workers route + accumulate; the leader clears/divides and rewrites
922 * heads. Gated on shared->refine (identical on both sides) so the
923 * barrier sequence stays in lockstep. */
924
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 186 times.
190 if (shared->refine)
925 {
926 /* The accumulator overlays the sample region (dead since the subtree
927 * phase); the leader initialized its header before the tree-ready
928 * barrier above. */
929 4 PrismDsmRefineAccum *accum = prism_pbuild_refine_overlay(dsm_samples);
930 4 prism_pbuild_exec_refine_paged(
931 worker_id,
932 heapRel,
933 indexRel,
934 indexInfo,
935 shared,
936 &qs,
937 first_posting,
938 accum,
939 barrier,
940 NULL,
941 NULL);
942 }
943
944 /* The samples (and the refine overlay riding in them) are dead; hand the
945 * segment back before the posting sort claims its own memory budget. */
946 190 prism_pbuild_samples_release(dsm_samples, sample_seg);
947 190 dsm_samples = NULL;
948
949 /* ---- Phase 3: posting scan -> cluster-keyed sort (page-backed) ---- */
950
951 /* worker_id is 1..N for launched workers; the sorter's 0-based worker
952 * index is worker_id - 1. Only the workers sort concurrently -- the
953 * leader's merge-only sorter opens after every worker sort finishes and
954 * uses the full budget -- so the budget splits across the workers
955 * alone; splitting across all participants would strand one share for
956 * the whole sort phase. */
957 190 int nsorters = shared->nparticipants > 1 ? shared->nparticipants - 1 : 1;
958 190 int worker_wm = shared->work_mem_kb / nsorters;
959
1/2
✓ Branch 0 taken 106 times.
✗ Branch 1 not taken.
190 if (worker_wm < 64)
960 106 worker_wm = 64;
961 190 PrismSorter *sorter = prism_pbuild_sort_begin(
962 sortshared, seg, worker_id - 1, 0, false, entry_size, worker_wm);
963
964 84 PrismBuildRouteCtx route;
965 190 prism_build_route_ctx_init(
966 &route,
967 &qs,
968 sorter,
969 rq_params,
970 storage,
971 first_posting,
972 dim,
973 shared->soar_lambda,
974 shared->boundary_epsilon);
975
976 /* Barrier: the leader reset the scan for the posting phase (after refine
977 * consumed it); workers may now scan. */
978 190 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
979
980 190 double heap_tuples = prism_build_scan(
981 heapRel,
982 indexRel,
983 indexInfo,
984 shared,
985 true,
986 false,
987 route_scan_cb,
988 &route);
989
990 /* Sub-batch progress remainder, as in the sampling scan. */
991 190 prism_build_progress_incr_tuples(
992 190 (int64_t)((uint64_t)route.indtuples & 1023));
993
994 190 prism_pbuild_sort_performsort(sorter);
995
996 /* Barrier: every worker has finished sorting; the leader merges next. */
997 190 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
998 190 BarrierDetach(barrier);
999
1000 190 prism_pbuild_worker_add_counts(
1001 shared, route.indtuples, route.soar_dupes, heap_tuples);
1002
1003 190 prism_pbuild_sort_end(sorter);
1004 190 prism_build_route_ctx_cleanup(&route);
1005 190 prism_query_state_cleanup(&qs);
1006 190 vs_free(base.pt_global_mean);
1007 190 prism_pbuild_exact_centroids_release(exact_seg);
1008 190 prism_pbuild_worker_storage_release(storage);
1009
1010 190 prism_pbuild_worker_detach(toc, &w);
1011 190 }
1012