GCC Code Coverage Report


Directory: src/
File: src/index/parallel_build_leader.c
Date: 2026-09-30 11:11:31
Exec Total Coverage
Lines: 314 327 96.0%
Functions: 5 5 100.0%
Branches: 89 121 73.6%

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_leader.c - PG parallel index build, leader side
6 *
7 * The leader sets up the DSM, launches workers, participates in the
8 * sampling + k-means phases as worker_id 0, streams the centroid tree to
9 * pages, then (phase 3) merges the workers' cluster-keyed sorted runs and
10 * builds each posting list. Head blocks are formula-derived (cluster c's head
11 * is first_posting + c), so no per-partition head array is materialized.
12 * Returns false if parallelism could not start, so the caller falls back to a
13 * serial build.
14 */
15
16 #ifdef VS_STANDALONE
17 #include "standalone/instr_time.h"
18 #include "standalone/parallel_ctx.h" /* ParallelContext + lifecycle */
19 #include "standalone/pg_compat.h"
20 #else
21 #include <postgres.h>
22
23 #include <access/parallel.h>
24 #include <access/table.h>
25 #include <access/tableam.h>
26 #include <access/xloginsert.h>
27 #include <catalog/index.h>
28 #include <commands/progress.h>
29 #include <common/pg_prng.h>
30 #include <executor/instrument.h>
31 #include <miscadmin.h>
32 #include <pgstat.h>
33 #include <portability/instr_time.h>
34 #include <tcop/tcopprot.h>
35 #include <utils/backend_progress.h>
36 #include <utils/memutils.h>
37 #include <utils/rel.h>
38 #include <utils/sampling.h>
39 #include <utils/wait_event.h>
40 #endif
41
42 #include <inttypes.h>
43 #include <math.h>
44
45 #include "algo/distance.h"
46 #include "algo/hkmeans.h"
47 #include "algo/kmeans.h"
48 #include "algo/kmeans_internal.h"
49 #include "algo/vecops.h"
50 #include "core/log.h"
51 #include "core/memory.h"
52 #include "index/build_progress.h"
53 #include "index/centroid_build.h"
54 #include "index/centroid_page.h"
55 #include "index/index_build.h"
56 #include "index/parallel_build.h"
57 #include "index/posting_build.h"
58 #include "index/posting_page.h"
59 #include "quant/fastscan.h"
60 #include "types/vec16.h"
61 #include "types/vec32.h"
62
63 #ifndef VS_STANDALONE
64 #include "build.h"
65 #include "meta.h"
66 #include "storage.h"
67 #include "support_pg.h"
68 #endif
69
70 /*
71 * Build every cluster's posting list from a populated (not yet performsorted)
72 * cluster-keyed sorter. Shared by the serial build and the parallel leader:
73 * performsort, then read entries grouped by cluster and write each list with a
74 * single resident page builder (head at first_posting + c, continuations
75 * appended at the relation's end and chained). Ends the sorter.
76 *
77 * Why iterate clusters 0..nlist rather than just draining the sorter until it
78 * is empty: EVERY cluster needs a posting-list head, including clusters that
79 * received no vectors. The centroid tree's leaf entries reference
80 * first_posting
81 * + c for every c in [0, nlist) (see prism_write_centroid_tree), so an empty
82 * cluster still needs a valid (empty) head block to point at. A loop over the
83 * cluster index emits a head for every cluster uniformly — an empty cluster's
84 * inner while simply does not run and finish() returns an empty head. A
85 * drain-until-empty loop would only produce heads for clusters present in the
86 * stream and would have to separately backfill empty heads for gap clusters
87 * and for all trailing clusters past the last one seen.
88 *
89 * This REQUIRES (and assumes) the sorter returns entries in ascending cluster
90 * order, matching the 0..nlist iteration order: the sorter is keyed on the
91 * cluster id (PG tuplesort on the int4 key; standalone qsort by cluster), so
92 * all entries for a cluster are contiguous and clusters appear in increasing
93 * order. Thus while walking c upward, any remaining entry has cur_cluster >= c
94 * (it equals c for a non-empty cluster, or is greater when c is empty). The
95 * asserts below make that contract explicit: an out-of-order key would pair
96 * entries with the wrong cluster, and an out-of-range id (>= nlist) would
97 * never match any c and be silently dropped (leaving the sorter non-empty at
98 * the end).
99 */
100 void
101 276 prism_posting_build_lists(
102 PrismSorter *sorter,
103 VsStorage *storage,
104 uint32_t nlist,
105 Dimension dim,
106 bool fastscan,
107 const RaBitQParams *rq_params,
108 BlockNumber first_posting)
109 {
110 276 prism_pbuild_sort_performsort(sorter);
111
112 /*
113 * Each list's head was pre-written with its pt_centroid (P^T * centroid,
114 * the RaBitQ encode reference) during the centroid build: adopt it and
115 * append the sorted entries in place. No second head construction, no
116 * in-RAM float centroids — the tree can be freed before this runs — and
117 * a cluster with no entries keeps its on-disk head untouched.
118 */
119 276 uint32_t cur_cluster = 0;
120 276 const void *entry = NULL;
121 276 bool have = prism_pbuild_sort_getnext(sorter, &cur_cluster, &entry);
122
3/3
✓ Branch 0 taken 1114 times.
✓ Branch 1 taken 9633 times.
✓ Branch 2 taken 194 times.
10941 for (uint32_t c = 0; c < nlist; c++)
123 {
124 /* Head block of cluster c is the formula first_posting + c; the head
125 * region [first_posting, first_posting + nlist) was pre-extended and
126 * its pt_centroid written during the centroid streaming. Continuation
127 * pages are appended at the relation's end and chained (no per-cluster
128 * reserve), so no O(nlist) reserve arrays are needed. Adopt the
129 * pre-written head and append in place: no second head construction,
130 * and a cluster with no entries keeps its on-disk head untouched. */
131 10665 BlockNumber head_blk = first_posting + c;
132 9551 PrismPostingBuilder hb;
133 10665 prism_posting_builder_adopt_head(
134 &hb, storage, rq_params, dim, c, head_blk, fastscan);
135
136 /* Ascending, grouped order: a pending entry is never for a cluster we
137 * already finished. If it were < c we would have skipped its head. */
138
4/4
✓ Branch 0 taken 1149 times.
✓ Branch 1 taken 9516 times.
✓ Branch 2 taken 9516 times.
✓ Branch 3 taken 1114 times.
10665 Assert(!have || cur_cluster >= c);
139
140
4/4
✓ Branch 0 taken 574869 times.
✓ Branch 1 taken 302 times.
✓ Branch 2 taken 564506 times.
✓ Branch 3 taken 10363 times.
575171 while (have && cur_cluster == c)
141 {
142 564506 prism_posting_entry_add(&hb, entry, dim);
143 564506 have = prism_pbuild_sort_getnext(sorter, &cur_cluster, &entry);
144 }
145
146 10665 (void)prism_posting_builder_finish(&hb);
147 10665 prism_posting_builder_cleanup(&hb);
148 }
149
150 /* Every entry must have landed in some cluster's list. A leftover entry
151 * means an id >= nlist (out of range) that matched no c and would
152 * otherwise be silently dropped. */
153
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 276 times.
276 Assert(!have);
154 276 prism_pbuild_sort_end(sorter);
155 276 }
156
157 /* ----------------------------------------------------------------
158 * Batched streaming tree build — leader-side callbacks
159 * ---------------------------------------------------------------- */
160
161 /* PLAN-pass batch callback: record each subtree's leaf count + centroid-page
162 * count (no writes) so the leader can size the reserve + block layout, and
163 * keep the blob in the spillable store for the streaming pass to read back
164 * in the same (child) order. */
165 typedef struct PlanCbArg
166 {
167 uint32_t *nleaves_arr; /* [km_k] */
168 uint32_t *pages_arr; /* [km_k] */
169 uint32_t max_ent;
170 uint32_t subtree_nlevels; /* actual depth of the (uniform) subtrees */
171 Dimension dim;
172 double *leaf_sum; /* [dim] leaf-centroid sum across all subtrees, for
173 * the leaf_mean the leader uses as global_mean */
174 PrismBlobStore *store; /* subtree blobs, in child order */
175 } PlanCbArg;
176
177 static void
178 102 plan_batch_cb(
179 void *arg,
180 const uint32_t *children,
181 uint32_t bs,
182 char *base,
183 uint64_t slot_size)
184 {
185 102 PlanCbArg *a = (PlanCbArg *)arg;
186
2/2
✓ Branch 0 taken 237 times.
✓ Branch 1 taken 102 times.
339 for (uint32_t s = 0; s < bs; s++)
187 {
188 111 const HKMeansResult *sub = (const HKMeansResult *)
189 237 prism_dsm_child_subtree(base, s, slot_size);
190 237 uint32_t child = children[s];
191 237 BlockNumber *nfb = vs_alloc((size_t)sub->nnodes * sizeof(BlockNumber));
192 /* Subtrees are built to a uniform depth, so any non-empty one gives
193 * the streamed tree's subtree depth (full depth = this + 1 root
194 * level). */
195
1/2
✓ Branch 0 taken 237 times.
✗ Branch 1 not taken.
237 if (sub->nleaves > 0)
196 237 a->subtree_nlevels = sub->nlevels;
197 237 a->nleaves_arr[child] = sub->nleaves;
198 474 a->pages_arr[child] = (uint32_t)
199 237 prism_compute_centroid_layout(sub, a->max_ent, 0, nfb);
200 237 vs_free(nfb);
201
202
1/2
✓ Branch 0 taken 237 times.
✗ Branch 1 not taken.
237 if (a->leaf_sum != NULL)
203 {
204 237 const float *lc = hk_leaf_centroids(sub);
205
2/2
✓ Branch 0 taken 1783 times.
✓ Branch 1 taken 237 times.
2020 for (uint32_t l = 0; l < sub->nleaves; l++)
206
2/2
✓ Branch 0 taken 45726 times.
✓ Branch 1 taken 1783 times.
47509 for (Dimension d = 0; d < a->dim; d++)
207 45726 a->leaf_sum[d] += lc[(size_t)l * a->dim + d];
208 }
209
210 237 prism_pbuild_blobstore_put(a->store, sub, sub->total_size);
211 }
212 102 }
213
214 /* What both routing-tree shapes hand back to do_parallel_build: the layout
215 * facts the leader publishes to the workers before phase 3. */
216 typedef struct TreeLayout
217 {
218 BlockNumber first_posting; /* leaf c's head = first_posting + c */
219 BlockNumber root_blk; /* tree root page */
220 uint8_t nlevels; /* written depth */
221 uint32_t nlist; /* actual leaf count */
222 } TreeLayout;
223
224 /*
225 * Assemble and write the routing tree for the hierarchical shape
226 * (nlevels >= 2): each participant builds the subtree of every root child
227 * it owns; the leader records each batch's layout counts (PLAN), computes
228 * the block layout and the leaf-centroid mean, pre-extends the relation,
229 * replays the spilled subtree blobs to pages, and writes the root page
230 * above them. Fills global_mean and the layout the caller publishes.
231 */
232 static bool
233 50 build_routing_tree_batched(
234 PrismBuildShared *shared,
235 VsStorage *storage,
236 PrismBuildProgress *prog,
237 PrismDsmSamples *dsm_samples,
238 PrismDsmRootAssign *dsm_ra,
239 float *cents,
240 uint32_t km_k,
241 uint32_t nlist,
242 uint32_t fan_out,
243 Barrier *barrier,
244 RaBitQParams *rq_params,
245 uint32_t max_ent,
246 BlockNumber first_centroid,
247 float *global_mean,
248 PrismExactCentroidCollector *collector,
249 TreeLayout *out)
250 {
251 50 Dimension dim = shared->dim;
252 50 int nparticipants = shared->nparticipants;
253 50 PrismCentroidFormat fmt = shared->centroid_format;
254
255 /* Size the ring slot from the actual per-child sample counts: a subtree
256 * can never hold more leaves than the samples routed to its root child,
257 * which caps the slot far below the analytic worst case (nlist here is
258 * the worst-case sizing bound, up to fan_out x the requested count). */
259 50 uint32_t max_cc = 1;
260 {
261 50 uint32_t *child_count = vs_alloc0((size_t)km_k * sizeof(uint32_t));
262 50 prism_pbuild_count_children(
263 dsm_samples, dsm_ra, nparticipants, km_k, child_count);
264
3/3
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 145 times.
✓ Branch 2 taken 20 times.
291 for (uint32_t c = 0; c < km_k; c++)
265
2/2
✓ Branch 0 taken 76 times.
✓ Branch 1 taken 50 times.
241 if (child_count[c] > max_cc)
266 76 max_cc = child_count[c];
267 50 vs_free(child_count);
268 }
269 50 uint32_t nlist_c = (nlist + fan_out - 1) / fan_out;
270 20 uint64_t slot_size =
271 50 vs_hkmeans_max_blob_size_capped(nlist_c, fan_out, dim, max_cc);
272
273 /* The slot must fit one allocation (the leader reads each spilled blob
274 * back through a palloc'd buffer), which also keeps the blob format's
275 * 32-bit interior offsets valid. */
276
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 50 times.
50 if (slot_size > (uint64_t)MaxAllocSize)
277 ✗ vs_error(
278 "prism: subtree slot %llu MB exceeds the allocation limit "
279 "(nlist %u, fan_out %u); increase fan_out or decrease nlist",
280 (unsigned long long)(slot_size >> 20),
281 nlist,
282 fan_out);
283 /* The ring is the build's only region outside the sample budget; hold
284 * it to the same standard. work_mem_kb == 0 = unbudgeted back-end. */
285
2/2
✓ Branch 0 taken 20 times.
✓ Branch 1 taken 30 times.
50 if (shared->work_mem_kb > 0 &&
286 20 (uint64_t)nparticipants * slot_size >
287
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 20 times.
20 (uint64_t)shared->work_mem_kb * 1024)
288 ✗ vs_error(
289 "prism: subtree ring %llu MB exceeds maintenance_work_mem "
290 "(nlist %u, fan_out %u); increase maintenance_work_mem or "
291 "fan_out",
292 (unsigned long long)(((uint64_t)nparticipants * slot_size) >>
293 20),
294 nlist,
295 fan_out);
296
297 /* Ring barrier: create + publish (handle, slot size) before arriving;
298 * the workers attach after. */
299 50 void *ring_seg = NULL;
300 50 char *subtrees_base = prism_pbuild_subtree_ring_create(
301 shared, nparticipants, slot_size, &ring_seg);
302
1/2
✗ Branch 1 not taken.
✓ Branch 2 taken 20 times.
50 vs_debug(
303 "prism: subtree ring %d x %llu KB (largest child %u samples)",
304 nparticipants,
305 (unsigned long long)(slot_size >> 10),
306 max_cc);
307 50 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
308
309 50 uint32_t *nleaves_arr = vs_alloc0((size_t)km_k * sizeof(uint32_t));
310 50 uint32_t *pages_arr = vs_alloc0((size_t)km_k * sizeof(uint32_t));
311
312 /* PLAN pass: build subtrees, discover leaf + page counts (and the
313 * leaf-centroid sum for global_mean). */
314 50 prism_build_report_phase(prog, PRISM_BUILD_PHASE_SUBTREES);
315 147 PlanCbArg planarg = {
316 .nleaves_arr = nleaves_arr,
317 .pages_arr = pages_arr,
318 .max_ent = max_ent,
319 .subtree_nlevels = 1,
320 .dim = dim,
321 49 .leaf_sum = vs_alloc0((size_t)dim * sizeof(double)),
322 49 .store = prism_pbuild_blobstore_begin(),
323 };
324 /* The batch schedule is largest-first; the blob store receives the
325 * subtrees in that order, so the replay below needs the same order
326 * to place each blob at its child's reserved block range. */
327 49 uint32_t *child_order = vs_alloc((size_t)km_k * sizeof(uint32_t));
328 49 prism_pbuild_stream_subtrees(
329 0,
330 nparticipants,
331 dsm_samples,
332 dsm_ra,
333 cents,
334 km_k,
335 nlist,
336 fan_out,
337 dim,
338 shared->metric,
339 shared->km_max_iterations,
340 subtrees_base,
341 slot_size,
342 barrier,
343 plan_batch_cb,
344 &planarg,
345 child_order);
346
347 /* Every batch is consumed into the blob store; the ring is dead before
348 * the write pass starts. */
349 49 prism_pbuild_subtree_ring_release(ring_seg);
350
351 /* Root page(s) occupy the reserved block(s) at first_centroid;
352 * subtrees follow, so meta.first_centroid stays 1 (root written last,
353 * in place). */
354 49 uint32_t root_pages = (km_k + max_ent - 1) / max_ent;
355 49 uint32_t *leaf_off = vs_alloc((size_t)km_k * sizeof(uint32_t));
356 49 uint32_t *block_off = vs_alloc((size_t)km_k * sizeof(uint32_t));
357 49 uint32_t lo = 0;
358 49 BlockNumber bo = 0;
359
3/3
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 141 times.
✓ Branch 2 taken 19 times.
286 for (uint32_t c = 0; c < km_k; c++)
360 {
361 237 leaf_off[c] = lo;
362 237 block_off[c] = (uint32_t)bo;
363 237 lo += nleaves_arr[c];
364 237 bo += pages_arr[c];
365 }
366 49 uint32_t actual_nlist = lo;
367 49 BlockNumber subtree_base = first_centroid + root_pages;
368 49 BlockNumber first_posting = subtree_base + bo;
369
370 /* The centroid-page region [first_centroid, first_posting) is now
371 * sized; the streaming pass below collects every internal node's
372 * exact centroids into it (the root included). */
373
2/2
✓ Branch 0 taken 45 times.
✓ Branch 1 taken 4 times.
49 if (collector != NULL)
374
1/2
✓ Branch 0 taken 19 times.
✗ Branch 1 not taken.
45 prism_exact_centroid_collector_init(
375 collector,
376 dim,
377 fmt,
378 first_centroid,
379 45 (uint32_t)(first_posting - first_centroid),
380
1/2
✓ Branch 0 taken 19 times.
✗ Branch 1 not taken.
45 prism_exact_centroid_budget(shared->work_mem_kb),
381 prism_exact_centroid_expected_slots(actual_nlist, fan_out));
382
383 /* Leaf-centroid mean from the PLAN pass -> the encoder centering. */
384
2/2
✓ Branch 0 taken 1068 times.
✓ Branch 1 taken 49 times.
1117 for (Dimension d = 0; d < dim; d++)
385 1068 global_mean[d] = actual_nlist > 0 ? (float)(planarg.leaf_sum[d] /
386 1068 (double)actual_nlist)
387
1/2
✓ Branch 0 taken 1068 times.
✗ Branch 1 not taken.
1068 : 0.0f;
388
2/2
✓ Branch 0 taken 3 times.
✓ Branch 1 taken 46 times.
49 if (shared->metric == DISTANCE_COSINE)
389 3 vs_l2_normalize(global_mean, dim);
390 49 vs_free(planarg.leaf_sum);
391 49 planarg.leaf_sum = NULL;
392
393 /* Head blocks are formula-derived: leaf c's head is first_posting + c,
394 * a contiguous head region of actual_nlist pages. Pre-extend the
395 * relation to cover the centroid pages + the head region so the write
396 * pass can write both at reserved blocks; continuation pages are
397 * appended past it during prism_posting_build_lists. No O(nlist) reserve
398 * arrays. */
399 49 prism_build_reserve_layout(storage, first_posting + actual_nlist);
400
401 /* Streaming pass: read each subtree blob back from the store (child
402 * order matches the append order) and stream its centroid + head
403 * pages to the reserved block range, then the root page. Leader-only:
404 * the PLAN pass already produced every subtree, so the workers have
405 * nothing to contribute here and run no barriers for this phase. */
406 49 prism_build_report_phase(prog, PRISM_BUILD_PHASE_CENTROID);
407 49 BlockNumber *subtree_root_blk = vs_alloc(
408 30 (size_t)km_k * sizeof(BlockNumber));
409 19 PrismHeadWriteCtx head;
410 68 prism_head_write_ctx_init(
411 49 &head, storage, rq_params, dim, shared->fastscan, first_posting);
412 49 HKMeansResult *blob = vs_alloc(slot_size);
413 49 prism_pbuild_blobstore_rewind(planarg.store);
414
3/3
✓ Branch 0 taken 126 times.
✓ Branch 1 taken 141 times.
✓ Branch 2 taken 19 times.
286 for (uint32_t i = 0; i < km_k; i++)
415 {
416 /* Blobs arrive in the largest-first batch order; c is the child
417 * whose reserved block range this blob belongs to. */
418 237 uint32_t c = child_order[i];
419 237 (void)prism_pbuild_blobstore_get(planarg.store, blob, slot_size);
420 237 BlockNumber base_blk = subtree_base + block_off[c];
421 237 subtree_root_blk[c] = prism_routing_subtree_write(
422 storage,
423 blob,
424 dim,
425 shared->metric,
426 fan_out,
427 1, /* subtrees hang off the level-0 root */
428 fmt,
429 rq_params,
430 global_mean,
431 first_posting,
432 237 leaf_off[c],
433 base_blk,
434 prism_write_leaf_head,
435 &head,
436 collector);
437 }
438 49 vs_free(blob);
439 49 vs_free(child_order);
440 49 prism_pbuild_blobstore_end(planarg.store);
441 49 planarg.store = NULL;
442
443 /* Root centroid page at the reserved first_centroid (children =
444 * subtree roots). Written last, but in place, so first_centroid
445 * stays 1. */
446 49 BlockNumber root_blk = first_centroid;
447 49 prism_centroid_write_node(
448 storage,
449 dim,
450 cents,
451 km_k,
452 fmt,
453 0,
454 0,
455 30 (uint16_t)fan_out,
456 rq_params,
457 global_mean,
458 subtree_root_blk,
459 NULL,
460 root_blk,
461 collector);
462
463 49 prism_head_write_ctx_cleanup(&head);
464 49 vs_free(subtree_root_blk);
465 49 vs_free(leaf_off);
466 49 vs_free(block_off);
467 49 vs_free(nleaves_arr);
468 49 vs_free(pages_arr);
469
470 49 out->first_posting = first_posting;
471 49 out->root_blk = root_blk;
472 49 out->nlevels = (uint8_t)(planarg.subtree_nlevels + 1);
473 49 out->nlist = actual_nlist;
474 49 return true;
475 }
476
477 /*
478 * Assemble and write the routing tree for the flat shape (nlevels == 1):
479 * the root k-means centroids already ARE the leaves, so stream the
480 * one-level tree directly (root = leaf-parent page at first_centroid).
481 * Fills global_mean and the layout the caller publishes; false when the
482 * flat carrier cannot be built.
483 */
484 static bool
485 76 build_routing_tree_flat(
486 PrismBuildShared *shared,
487 VsStorage *storage,
488 PrismBuildProgress *prog,
489 float *cents,
490 uint32_t km_k,
491 uint32_t fan_out,
492 RaBitQParams *rq_params,
493 uint32_t max_ent,
494 BlockNumber first_centroid,
495 float *global_mean,
496 TreeLayout *out)
497 {
498 76 Dimension dim = shared->dim;
499 76 PrismCentroidFormat fmt = shared->centroid_format;
500
501 /* Flat (nlevels == 1): cents already holds every leaf centroid; stream
502 * the one-level tree directly (root = leaf-parent at first_centroid).
503 */
504
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 72 times.
76 if (shared->metric == DISTANCE_COSINE)
505
2/2
✓ Branch 0 taken 41 times.
✓ Branch 1 taken 4 times.
45 for (uint32_t c = 0; c < km_k; c++)
506 41 vs_l2_normalize(cents + (size_t)c * dim, dim);
507
508 76 HKMeansResult *flat = vs_hkmeans_build_flat(cents, km_k, fan_out, dim);
509
2/2
✓ Branch 0 taken 24 times.
✓ Branch 1 taken 52 times.
76 if (flat == NULL)
510 ✗ return false;
511
512 76 uint32_t actual_nlist = flat->nleaves;
513 76 BlockNumber *nfb = vs_alloc((size_t)flat->nnodes * sizeof(BlockNumber));
514 24 uint32_t centroid_pages = (uint32_t)
515 76 prism_compute_centroid_layout(flat, max_ent, 0, nfb);
516 76 vs_free(nfb);
517 76 BlockNumber first_posting = first_centroid + centroid_pages;
518
519 /* Leaf-centroid mean -> the encoder centering (flat tree in hand). */
520 76 vec32_mean(hk_leaf_centroids(flat), actual_nlist, dim, global_mean);
521
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 72 times.
76 if (shared->metric == DISTANCE_COSINE)
522 4 vs_l2_normalize(global_mean, dim);
523
524 /* Head region: actual_nlist pages at first_posting (leaf c -> head
525 * first_posting + c). Pre-extend to cover centroid + head region;
526 * continuations append past it. No O(nlist) reserve. */
527 76 prism_build_reserve_layout(storage, first_posting + actual_nlist);
528
529 76 prism_build_report_phase(prog, PRISM_BUILD_PHASE_CENTROID);
530 24 PrismHeadWriteCtx head;
531 100 prism_head_write_ctx_init(
532 76 &head, storage, rq_params, dim, shared->fastscan, first_posting);
533 76 BlockNumber root_blk = prism_routing_subtree_write(
534 storage,
535 flat,
536 dim,
537 shared->metric,
538 fan_out,
539 0, /* the flat tree IS the root level */
540 fmt,
541 rq_params,
542 global_mean,
543 first_posting,
544 0,
545 first_centroid,
546 prism_write_leaf_head,
547 &head,
548 /* the flat tree's only level is the leaf level — nothing
549 * internal to collect */
550 NULL);
551 76 prism_head_write_ctx_cleanup(&head);
552
553 76 out->first_posting = first_posting;
554 76 out->root_blk = root_blk;
555 76 out->nlevels = (uint8_t)flat->nlevels;
556 76 out->nlist = actual_nlist;
557 76 vs_free(flat);
558 76 return true;
559 }
560
561 bool
562 127 do_parallel_build(
563 Relation heap,
564 Relation index,
565 struct IndexInfo *index_info,
566 const PrismBuildConfig *config,
567 VsStorage *storage,
568 struct PrismBuildProgress *prog,
569 uint32_t *out_nlist,
570 uint8_t *out_tree_nlevels,
571 double *out_heap_tuples,
572 double *out_indtuples,
573 double *out_soar_dupes,
574 float **out_global_mean,
575 BlockNumber *out_first_posting)
576 {
577 127 int nworkers = index_info->ii_ParallelWorkers;
578
579 127 *out_nlist = 0;
580 127 *out_tree_nlevels = 0;
581
1/2
✓ Branch 0 taken 127 times.
✗ Branch 1 not taken.
127 if (out_first_posting)
582 127 *out_first_posting = InvalidBlockNumber;
583
1/2
✓ Branch 0 taken 127 times.
✗ Branch 1 not taken.
127 if (out_global_mean)
584 127 *out_global_mean = NULL;
585
586 45 PrismPBuildLeader lead;
587
1/2
✗ Branch 1 not taken.
✓ Branch 2 taken 127 times.
127 if (!prism_pbuild_setup_shared(&lead, heap, index, config, nworkers))
588 ✗ return false;
589
590 127 ParallelContext *pcxt = lead.pcxt;
591 127 PrismBuildShared *shared = lead.shared;
592 127 Barrier *barrier = lead.barrier;
593 127 PrismDsmSamples *dsm_samples = lead.dsm_samples;
594 127 char *centroids_base = lead.centroids_base;
595 127 float *cents = lead.cents;
596 127 char *km_workers_base = lead.km_workers_base;
597 127 PrismDsmRootAssign *dsm_ra = lead.dsm_ra;
598 127 WalUsage *walusage = lead.walusage;
599 127 BufferUsage *bufferusage = lead.bufferusage;
600 127 int nparticipants = lead.nparticipants;
601 127 uint32_t km_k = lead.km_k;
602 127 Dimension dim = lead.dim;
603 127 uint32_t nlist = lead.nlist;
604 127 uint64_t rabitq_seed = lead.rabitq_seed;
605 127 uint32_t fan_out = lead.fan_out;
606
607 /*
608 * Introspection: name the dataset-scaling allocations up front (always
609 * logged, regardless of the GUC) so an OOM in any of them is
610 * pre-explained, and record the committed DSM size so the per-phase memory
611 * lines can add it. nlist here is the worst-case upper bound (the tree is
612 * not built yet).
613 */
614 127 prism_build_report_dsm_bytes(prog, (uint64_t)lead.dsm_total);
615 127 prism_build_report_planned_alloc(
616 prog,
617 127 (uint64_t)prism_dsm_samples_size(
618 nparticipants, lead.max_per_worker, dim),
619 /* The tree is streamed to pages from a bounded ring of subtree
620 * slots (part of dsm_total); it is never held whole in memory. */
621 0,
622 127 (uint64_t)lead.dsm_total);
623
624 45 instr_time t_launch_start;
625 127 INSTR_TIME_SET_CURRENT(t_launch_start);
626
627 /*
628 * The leader participates in every phase, so it joins the dynamic barrier
629 * as a party (workers attach the same way as they start). Attaching before
630 * launching workers guarantees the leader is counted before any worker can
631 * arrive, so the barrier never advances a phase without it.
632 */
633 127 BarrierAttach(barrier);
634
635 /* ---- Launch workers + wait until they've all attached ---- */
636
2/2
✓ Branch 1 taken 1 times.
✓ Branch 2 taken 126 times.
127 if (!prism_pbuild_launch(pcxt, barrier, shared))
637 {
638 /* Teardown already ran; hand the sample segment back too so the
639 * serial fallback starts from a clean budget. */
640 1 prism_pbuild_samples_release(dsm_samples, lead.sample_seg);
641 1 return false;
642 }
643
644 /* The launch may have narrowed the participant count to the party that
645 * actually attached; partition the phases below over that count. */
646 126 nparticipants = shared->nparticipants;
647
648 /* ==== Leader participates in all phases as worker_id=0 ==== */
649
650 /* ---- Phase 1: sampling (leader runs the worker body as participant 0)
651 * ---- */
652 126 prism_build_report_phase(prog, PRISM_BUILD_PHASE_SAMPLE);
653 126 prism_pbuild_exec_sampling(
654 0, heap, index, index_info, shared, dsm_samples, barrier);
655
656 44 instr_time t_sample_end;
657 126 INSTR_TIME_SET_CURRENT(t_sample_end);
658 126 INSTR_TIME_SUBTRACT(t_sample_end, t_launch_start);
659
1/2
✗ Branch 1 not taken.
✓ Branch 2 taken 44 times.
126 vs_debug(
660 "prism: phase 1 (sampling) %.1fms",
661 INSTR_TIME_GET_MILLISEC(t_sample_end));
662
663 44 instr_time t_km_start;
664 126 INSTR_TIME_SET_CURRENT(t_km_start);
665
666 /* ---- Phase 2: root k-means. The leader runs the worker body as
667 * participant 0; inside it additionally seeds the initial centroids and
668 * reduces the per-iteration accumulators (gated on participant 0). ---- */
669 126 prism_build_report_phase(prog, PRISM_BUILD_PHASE_KMEANS);
670 126 uint32_t km_iters = prism_pbuild_exec_kmeans(
671 0, shared, dsm_samples, centroids_base, km_workers_base, barrier);
672
673 {
674 44 instr_time t_km_elapsed;
675 126 INSTR_TIME_SET_CURRENT(t_km_elapsed);
676 126 INSTR_TIME_SUBTRACT(t_km_elapsed, t_km_start);
677
1/2
✗ Branch 1 not taken.
✓ Branch 2 taken 44 times.
126 vs_debug(
678 "prism: root kmeans %.1fms (%u iters, k=%u)",
679 INSTR_TIME_GET_MILLISEC(t_km_elapsed),
680 km_iters,
681 km_k);
682 }
683
684 /*
685 * Compute nlevels to decide flat vs hierarchical assembly. For
686 * nlevels == 1 (flat) the root k-means already produced every leaf, so
687 * the tree is built from the centroids directly. For nlevels >= 2 each
688 * participant builds the full subtree — to whatever depth nlist/fan_out
689 * needs — for the root children it owns, and the leader streams each to
690 * centroid pages a batch at a time (no in-RAM whole-tree assembly).
691 */
692 126 uint32_t nlevels = vs_hkmeans_nlevels(nlist, fan_out);
693
694 /* ---- Phase 2b: root assignment (leader as participant 0). ---- */
695 126 prism_pbuild_exec_root_assign(
696 0, shared, dsm_samples, dsm_ra, centroids_base, barrier);
697
698 /* ---- Batched streaming tree build --------------------------------------
699 * Workers build per-root-child subtrees into a bounded ring of slots; the
700 * leader records each batch's layout counts, keeps the blobs in a
701 * spillable store, and streams them to centroid pages. Peak subtree DSM is
702 * nparticipants slots, independent of nlist.
703 * ------------------------------ */
704 126 prism_build_report_phase(prog, PRISM_BUILD_PHASE_SETUP);
705 126 RaBitQParams *rq_params = vs_rabitq_create(dim, rabitq_seed);
706 126 PrismCentroidFormat fmt = shared->centroid_format;
707 126 uint32_t max_ent = prism_centroid_max_entries_fmt(dim, fmt);
708
709 /*
710 * global_mean = mean of the LEAF centroids (the encoder centering) -- an
711 * unweighted per-cluster mean, not the per-vector sample mean. The
712 * quantization quality of every centroid and posting code depends on this
713 * anchor. Pages are written only after the PLAN pass has built every
714 * subtree, so
715 * the leaf-centroid sum is accumulated there (plan_batch_cb) and the mean
716 * is ready before any page is encoded. Filled per branch below.
717 */
718 126 const size_t vec_nbytes = (size_t)dim * sizeof(float);
719 126 float *global_mean = vs_alloc(vec_nbytes);
720
721 126 BlockNumber first_centroid = PRISM_FIRST_CENTROID_BLKNO;
722 126 BlockNumber first_posting = 0;
723 126 BlockNumber root_blk = InvalidBlockNumber;
724
725 /* Exact internal-centroid collection for the phase-2.5/3 build
726 * descent (see PrismExactCentroidCollector in index_build.h). */
727 126 PrismExactCentroidCollector exact_centroids = {0};
728
2/2
✓ Branch 0 taken 20 times.
✓ Branch 1 taken 24 times.
126 PrismExactCentroidCollector *collector =
729 106 prism_exact_centroid_enabled(nlevels, fmt) ? &exact_centroids
730
2/2
✓ Branch 0 taken 26 times.
✓ Branch 1 taken 56 times.
82 : NULL;
731
732 44 TreeLayout layout;
733 44 bool tree_ok;
734
2/2
✓ Branch 0 taken 50 times.
✓ Branch 1 taken 76 times.
126 if (nlevels >= 2)
735 50 tree_ok = build_routing_tree_batched(
736 shared,
737 storage,
738 prog,
739 dsm_samples,
740 dsm_ra,
741 cents,
742 km_k,
743 nlist,
744 fan_out,
745 barrier,
746 rq_params,
747 max_ent,
748 first_centroid,
749 global_mean,
750 collector,
751 &layout);
752 else
753 76 tree_ok = build_routing_tree_flat(
754 shared,
755 storage,
756 prog,
757 cents,
758 km_k,
759 fan_out,
760 rq_params,
761 max_ent,
762 first_centroid,
763 global_mean,
764 &layout);
765
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 125 times.
125 if (!tree_ok)
766 {
767 ✗ if (collector != NULL)
768 ✗ prism_exact_centroid_collector_cleanup(collector);
769 ✗ vs_free(global_mean);
770 ✗ prism_pbuild_samples_release(dsm_samples, lead.sample_seg);
771 ✗ WaitForParallelWorkersToFinish(pcxt);
772 ✗ prism_pbuild_teardown(pcxt);
773 ✗ return false;
774 }
775 125 nlist = layout.nlist;
776 125 first_posting = layout.first_posting;
777 125 root_blk = layout.root_blk;
778 125 const uint8_t out_nlevels = layout.nlevels;
779
780 /* Publish the exact internal-node centroids the tree write collected,
781 * for the workers' phase-2.5/3 build descent (exact-centroid seam; must
782 * precede the tree-ready barrier). Without a collector the collection
783 * is an empty header and the workers' scoring hook stays inert. */
784 125 void *exact_seg = NULL;
785 125 char *exact_cents = prism_pbuild_exact_centroids_create(
786 shared,
787 prism_exact_centroid_collection_size(collector),
788 &exact_seg);
789 125 prism_exact_centroid_collection_write(collector, exact_cents);
790
2/2
✓ Branch 0 taken 45 times.
✓ Branch 1 taken 80 times.
125 if (collector != NULL)
791 45 prism_exact_centroid_collector_cleanup(collector);
792
793 /* Publish the routing state the workers read in phase 3: nlist, the tree
794 * root block + depth, the posting-head base (leaf c's head = first_posting
795 * + c), the global mean, and the page store. */
796 125 shared->nlist = nlist;
797 125 shared->first_centroid = root_blk;
798 125 shared->nlevels = out_nlevels;
799 125 shared->first_posting = first_posting;
800
1/2
✓ Branch 0 taken 125 times.
✗ Branch 1 not taken.
125 if (out_first_posting)
801 125 *out_first_posting = first_posting;
802 {
803 43 float *dsm_gmean =
804 125 shm_toc_lookup(pcxt->toc, PRISM_DSM_KEY_GLOBAL_MEAN, false);
805 125 memcpy(dsm_gmean, global_mean, vec_nbytes);
806 }
807 125 prism_pbuild_publish_storage(shared, storage);
808
809
1/2
✓ Branch 0 taken 125 times.
✗ Branch 1 not taken.
125 if (out_global_mean)
810 125 *out_global_mean = global_mean; /* caller owns it (metadata write) */
811 else
812 ✗ vs_free(global_mean);
813
814 /* No in-RAM tree exists; the caller's metadata write needs only the
815 * shape. */
816 125 *out_nlist = nlist;
817 125 *out_tree_nlevels = out_nlevels;
818
819 /* The refine decision, now that the actual leaf count is known: refine
820 * only when the sample was bounded below the table AND the leaves are
821 * sample-thin -- a leaf's encode reference is a sample mean whose error
822 * shrinks with its sample count, so past the threshold the full-table
823 * scan recomputes what the sample already got right. Published before
824 * the ready barrier; workers gate the refine phase (and its barriers)
825 * on it after that barrier. */
826 {
827 125 uint64_t collected = 0;
828 125 uint64_t seen = 0;
829
2/2
✓ Branch 0 taken 315 times.
✓ Branch 1 taken 125 times.
440 for (int t = 0; t < nparticipants; t++)
830 {
831 315 collected += prism_dsm_sample_counts(dsm_samples)[t];
832 315 seen += prism_dsm_sample_seen(dsm_samples)[t];
833 }
834 /* kept < seen means the sampling scans skipped real rows, so the
835 * sample is a strict subset of the table (kept == seen means the
836 * sample IS the table -- nothing to refine from). */
837
1/2
✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
3 shared->refine = shared->refine_threshold > 0 && collected < seen &&
838
4/6
✓ Branch 0 taken 3 times.
✓ Branch 1 taken 122 times.
✓ Branch 2 taken 3 times.
✗ Branch 3 not taken.
✓ Branch 4 taken 3 times.
✗ Branch 5 not taken.
128 shared->refine_tile_cap > 0 && nlist > 0 &&
839
2/2
✓ Branch 0 taken 1 times.
✓ Branch 1 taken 2 times.
3 collected / nlist <
840 ✗ (uint64_t)shared->refine_threshold;
841
2/5
✗ Branch 0 not taken.
✓ Branch 1 taken 82 times.
✓ Branch 2 taken 43 times.
✗ Branch 3 not taken.
✗ Branch 4 not taken.
125 vs_debug(
842 "prism: refine gate: kept=%" PRIu64 " seen=%" PRIu64
843 " nlist=%u threshold=%u -> %s",
844 collected,
845 seen,
846 nlist,
847 shared->refine_threshold,
848 shared->refine ? "refine" : "skip");
849 }
850
851 /* The samples are dead (their last readers were the subtree builders);
852 * lay the refine accumulator over them before the ready barrier so the
853 * workers -- who pass that barrier ahead of the refine phase -- see an
854 * initialized header. */
855 125 PrismDsmRefineAccum *refine_accum = NULL;
856
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 123 times.
125 if (shared->refine)
857 {
858 2 refine_accum = prism_pbuild_refine_overlay(dsm_samples);
859 2 refine_accum->nleaves = shared->refine_tile_cap;
860 2 refine_accum->dim = dim;
861 }
862
863 /* Initialize the shared cluster sorter for the launched-worker count
864 * BEFORE the ready barrier, so it is ready when workers attach in phase 3.
865 * The leader merges only (it does not sort a share). */
866 43 void *sortshared =
867 125 shm_toc_lookup(pcxt->toc, PRISM_DSM_KEY_SORTSHARED, false);
868 125 prism_pbuild_sort_shared_init(
869 125 sortshared, pcxt->nworkers_launched, pcxt->seg);
870
871 /* Barrier: centroid/head pages written + routing state published + sorter
872 * ready; workers build their page-backed router next. */
873 125 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
874
875 /* ---- Phase 2.5: page-backed full-table refine (only when subsampled)
876 * ---- The workers route + accumulate; the leader (participant 0) clears
877 * the tiled accumulator, resets the scan per tile, and rewrites each
878 * leaf's head-page pt_centroid to the full-table mean. Gated on
879 * shared->refine, matching the workers, so the internal barriers stay
880 * in lockstep. */
881
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 123 times.
125 if (shared->refine)
882 {
883 2 prism_build_report_phase(prog, PRISM_BUILD_PHASE_REFINE);
884 2 PrismDsmRefineAccum *accum = refine_accum;
885 2 PrismHeadWriteCtx rhead;
886 4 prism_head_write_ctx_init(
887 &rhead,
888 storage,
889 rq_params,
890 dim,
891 2 shared->fastscan,
892 first_posting);
893 2 prism_pbuild_exec_refine_paged(
894 0,
895 heap,
896 index,
897 index_info,
898 shared,
899 NULL,
900 first_posting,
901 accum,
902 barrier,
903 prism_write_leaf_head,
904 &rhead);
905 2 prism_head_write_ctx_cleanup(&rhead);
906 }
907
908 /* Phase 3: workers scan + route page-backed + encode + sort; the leader
909 * merges. Reported exactly once -- the seam fires the "prism-build-load"
910 * test hook, and a 'wait' attached there must pause the build a single
911 * time -- and before the scan-reset barrier below releases the workers,
912 * so progress reflects the whole (multi-hour at scale) scan. */
913 125 prism_build_report_phase(prog, PRISM_BUILD_PHASE_SCAN_PARALLEL);
914
915 /* The samples (and the refine overlay riding in them) are dead; hand the
916 * segment back before the posting sort claims its own memory budget. */
917 125 prism_pbuild_samples_release(dsm_samples, lead.sample_seg);
918 125 dsm_samples = NULL;
919
920 /* Re-init the scan for the posting phase (the refine passes above consumed
921 * it). Guarded by the barrier below so no worker scans before the reset.
922 */
923 125 prism_pbuild_rescan(heap, shared);
924
925 /* Barrier: scan reset for the posting phase; workers start the page-backed
926 * posting scan+sort. */
927 125 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
928
929 43 instr_time t_scan_start;
930 125 INSTR_TIME_SET_CURRENT(t_scan_start);
931
932 /* Barrier: wait until every worker has finished sorting its run, then
933 * merge and build. The leader does not scan; the workers cover the heap.
934 */
935 125 BarrierArriveAndWait(barrier, WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
936 125 BarrierDetach(barrier);
937
938 125 WaitForParallelWorkersToFinish(pcxt);
939
940 43 instr_time t_scan_end;
941 125 INSTR_TIME_SET_CURRENT(t_scan_end);
942 125 INSTR_TIME_SUBTRACT(t_scan_end, t_scan_start);
943
944
2/2
✓ Branch 0 taken 190 times.
✓ Branch 1 taken 125 times.
315 for (int i = 0; i < pcxt->nworkers_launched; i++)
945 84 InstrAccumParallelQuery(&bufferusage[i], &walusage[i]);
946
947 125 *out_heap_tuples = shared->reltuples;
948 125 *out_indtuples = shared->indtuples;
949 125 *out_soar_dupes = shared->soar_dupes;
950
951 /* The workers cover the heap cooperatively while the leader blocks on the
952 * scan barrier above, so there is no leader-side loop to advance the % mid
953 * scan; publish the final scanned count now that the workers have
954 * finished.
955 */
956 125 prism_build_report_progress(prog, shared->indtuples);
957
958 43 instr_time t_merge_start;
959 125 INSTR_TIME_SET_CURRENT(t_merge_start);
960
961 /*
962 * Merge the workers' sorted runs and build each cluster's posting list
963 * with a single resident page builder, in cluster order (same loop as the
964 * serial path). Entries arrive grouped by cluster, so there is no
965 * partial-page fold and no chain to splice.
966 */
967 125 prism_build_report_phase(prog, PRISM_BUILD_PHASE_POSTING);
968 293 PrismSorter *sorter = prism_pbuild_sort_begin(
969 sortshared,
970 125 pcxt->seg,
971 0,
972 pcxt->nworkers_launched,
973 true,
974 125 (uint32_t)prism_posting_entry_size(dim),
975 shared->work_mem_kb);
976 125 prism_posting_build_lists(
977 sorter,
978 storage,
979 nlist,
980 dim,
981 125 shared->fastscan,
982 rq_params,
983 first_posting);
984
985 125 uint32_t total_pages = RelationGetNumberOfBlocks(index) - first_posting;
986
987 43 instr_time t_merge_end;
988 125 INSTR_TIME_SET_CURRENT(t_merge_end);
989 125 INSTR_TIME_SUBTRACT(t_merge_end, t_merge_start);
990
991
1/2
✗ Branch 1 not taken.
✓ Branch 2 taken 43 times.
125 vs_debug(
992 "prism: parallel streaming build with %d workers, "
993 "%u clusters, %u pages, "
994 "scan+drain %.1fms, finalize %.1fms",
995 pcxt->nworkers_launched,
996 nlist,
997 total_pages,
998 INSTR_TIME_GET_MILLISEC(t_scan_end),
999 INSTR_TIME_GET_MILLISEC(t_merge_end));
1000
1001 /* Workers detached from the exact-centroid collection when routing ended;
1002 * the leader's release is the last and frees it. */
1003 125 prism_pbuild_exact_centroids_release(exact_seg);
1004
1005 125 prism_pbuild_teardown(pcxt);
1006
1007 125 return true;
1008 }
1009