GCC Code Coverage Report


Directory: src/
File: src/index/parallel_build.h
Date: 2026-09-30 11:11:31
Exec Total Coverage
Lines: 58 61 95.1%
Functions: 17 20 85.0%
Branches: 17 19 89.5%

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.h - Parallel index build (worker + leader shared decls)
6 *
7 * Phased parallel build, barrier-synchronized:
8 * Sampling: workers cooperatively scan the heap into a bounded,
9 * budget-sized sample region (a dedicated segment, released once
10 * clustering is done).
11 * Clustering: root k-means over the sample (leader reduces between
12 * iterations), then per-root-child subtrees in batched ring slots;
13 * the leader records each batch's layout and spills the subtree
14 * blobs, then streams every centroid + head page to disk by itself.
15 * Refine (when subsampled): workers route the full table page-backed
16 * and accumulate per-leaf means; the leader rewrites the heads.
17 * Posting: workers route + encode into a cluster-keyed shared sort;
18 * the leader merges runs and writes the posting lists.
19 */
20
21 #ifndef PRISM_PARALLEL_BUILD_H
22 #define PRISM_PARALLEL_BUILD_H
23
24 /*
25 * Back-end primitives the parallel build runs on. The standalone shims live in
26 * src/standalone/; a PostgreSQL build uses the real PG headers. This guard is
27 * the single place that switch is made — the standalone headers carry no PG
28 * branch of their own.
29 */
30 #ifdef VS_STANDALONE
31 #include "standalone/barrier.h"
32 #include "standalone/pg_compat.h" /* Size, BlockNumber, ItemPointerData, Relation */
33 #include "standalone/shm_toc.h"
34 #else
35 #include <postgres.h>
36
37 #include <storage/barrier.h>
38 #include <storage/block.h>
39 #include <storage/dsm.h>
40 #include <storage/itemptr.h>
41 #include <storage/shm_toc.h>
42 #include <utils/rel.h>
43 #endif
44
45 #include "algo/hkmeans.h"
46 #include "core/memory.h"
47 #include "core/types.h"
48 #include "index/centroid_page.h" /* PrismCentroidFormat */
49 #include "index/posting_build.h"
50 #include "index/storage.h" /* VsStorage */
51 #include "quant/rabitq.h"
52
53 /* Forward decl so prism_build_scan's prototype can reference it without
54 * pulling in the executor headers; callers that pass one already have the full
55 * type.
56 */
57 struct IndexInfo;
58
59 /*
60 * Per-vector build scan callback (back-end-neutral): the PG scan adapter
61 * unwraps each heap tuple's Datum to the vector's data pointer, and the
62 * standalone work-stealing scan points straight into its vector array. The tid
63 * is the real heap TID under PG and a reversibly-synthesized one in
64 * standalone.
65 */
66 typedef void (*PrismBuildScanCb)(
67 void *state, ItemPointerData tid, const float *vec);
68
69 /* ----------------------------------------------------------------
70 * shm_toc region keys (the toc is a DSM segment in PG, a heap arena in
71 * standalone). A few regions are PG-only (WAL/buffer usage, query text) but
72 * their keys live here with the rest for one contiguous numbering.
73 * ---------------------------------------------------------------- */
74
75 #define PRISM_DSM_KEY_SHARED UINT64CONST(0xB000000000000001)
76 #define PRISM_DSM_KEY_WAL_USAGE UINT64CONST(0xB000000000000006)
77 #define PRISM_DSM_KEY_BUFFER_USAGE UINT64CONST(0xB000000000000007)
78 #define PRISM_DSM_KEY_QUERY_TEXT UINT64CONST(0xB000000000000008)
79 #define PRISM_DSM_KEY_BARRIER UINT64CONST(0xB000000000000009)
80 #define PRISM_DSM_KEY_SAMPLES UINT64CONST(0xB00000000000000A)
81 #define PRISM_DSM_KEY_CENTROIDS UINT64CONST(0xB00000000000000B)
82 #define PRISM_DSM_KEY_KM_WORKERS UINT64CONST(0xB00000000000000C)
83 #define PRISM_DSM_KEY_ROOT_ASSIGN UINT64CONST(0xB00000000000000E)
84 #define PRISM_DSM_KEY_SORTSHARED UINT64CONST(0xB000000000000012)
85 /* Page-backed phase-3 routing: the global mean, published by the leader before
86 * the tree-ready barrier so workers route exactly as the query/insert paths
87 * do. The posting-head base (leaf c's head = first_posting + c) is a scalar in
88 * PrismBuildShared, not a shared array. */
89 #define PRISM_DSM_KEY_GLOBAL_MEAN UINT64CONST(0xB000000000000014)
90
91 /* ----------------------------------------------------------------
92 * PrismBuildShared — back-end-neutral shared build state
93 *
94 * Immutable config set by the leader before the workers start, plus counters
95 * the workers update concurrently during the posting scan under a lock the
96 * back-end owns (prism_pbuild_worker_add_counts). Each back-end embeds this as
97 * the first member of its own struct carrying the back-end-specific state —
98 * for PG, the relation OIDs, query id, the spinlock, and a trailing
99 * ParallelTableScanDesc (see PrismBuildSharedPg in parallel_backend.c).
100 * ---------------------------------------------------------------- */
101
102 typedef struct PrismBuildShared
103 {
104 /* Immutable — set by the leader before launch */
105 Dimension dim;
106 DistanceMetric metric;
107 uint32_t nlist;
108 uint32_t fan_out;
109 double soar_lambda;
110 double boundary_epsilon;
111 bool fastscan;
112 PrismCentroidFormat centroid_format;
113 uint64_t rabitq_seed;
114 /* Planned as 1 + planned workers (it sizes the per-participant DSM
115 * regions), then narrowed by prism_pbuild_launch to 1 + the workers that
116 * actually started when the launch falls short. Work partitioned by
117 * participant must use this count; region sizing keeps the planned
118 * value and leaves the tail slots unused. */
119 int nparticipants;
120 /* CREATE INDEX CONCURRENTLY: the leader scans with an MVCC snapshot and
121 * participants take weak relation locks. Workers must mark their rebuilt
122 * IndexInfo concurrent too, or heapam's snapshot/OldestXmin check trips.
123 */
124 bool concurrent;
125 /* Posting sort work budget (KB); PG sets it from maintenance_work_mem,
126 * standalone leaves it 0 (its in-memory sorter ignores it). The shared
127 * phase-3 code reads this instead of the PG-only GUC. */
128 int work_mem_kb;
129
130 /* K-means config */
131 uint32_t max_samples_per_worker;
132 uint32_t km_max_iterations;
133 float km_tolerance;
134 uint32_t km_k; /* root k-means k (= fan_out) */
135
136 /* Refine gate input, set at setup: the sample-per-leaf threshold below
137 * which a subsampled build refines the leaf encode references on the
138 * full table (0 = refinement off; standalone leaves it 0). Whether the
139 * build IS subsampled is decided empirically from the sampling pass:
140 * each participant counts the live rows its scan saw next to the rows
141 * it kept, and kept < seen means the sample excludes real rows. */
142 uint32_t refine_threshold;
143
144 /* The refine decision. The LEADER makes it after clustering -- the
145 * actual leaf count is only known then, and dividing the collected
146 * samples by the worst-case leaf bound would understate samples-per-leaf
147 * and refine too eagerly -- and publishes it before the tree-ready
148 * barrier. Workers read it after that barrier, so leader and workers
149 * gate the refine phase (and its barriers) identically. */
150 bool refine;
151
152 /* Tile capacity (leaves) of the refine accumulator that overlays the
153 * sample region once sampling is done; sized at setup
154 * (the back-end bounds it by its memory budget AND the sample region's
155 * size, since the overlay lives inside that region). */
156 uint32_t refine_tile_cap;
157
158 /* Counters — updated concurrently under the back-end's lock */
159 double reltuples;
160 double indtuples;
161 double soar_dupes;
162
163 /* K-means convergence — set by leader between barriers */
164 bool km_converged;
165
166 /* Per-child subtree blob slot size (bytes) in the child-subtrees region;
167 * set by the leader before launch so workers can index their slot. */
168 /* Written by the leader after root assignment (see the subtree-ring
169 * seam); workers read it after the ring barrier. */
170 uint64_t subtree_slot_size;
171
172 /* Page-backed routing knobs (mirror the prism.centroid_* GUCs), so
173 * phase-3 workers build a PrismIndexBase that routes identically to
174 * the query path. fastscan_bits is the FASTSCAN centroid bit width
175 * (base.fastscan when the centroid format is FASTSCAN; 0 otherwise). */
176 float centroid_error_scale;
177 float centroid_beam_scale;
178 int fastscan_bits;
179
180 /* Published by the leader after the streaming tree write: the
181 * centroid-tree root block (workers' phase-3
182 * PrismIndexBase.first_centroid) and the tree depth (base.nlevels). The
183 * tree itself lives only on pages.
184 */
185 BlockNumber first_centroid;
186 uint8_t nlevels;
187
188 /* Published by the leader before phase 3: the first posting-head block.
189 * Cluster c's head is first_posting + c (formula), so workers map a routed
190 * head block back to its leaf by subtraction — no O(nlist) head array. */
191 BlockNumber first_posting;
192 } PrismBuildShared;
193
194 /* ----------------------------------------------------------------
195 * K-means sampling: per-worker sample slots in DSM
196 *
197 * Layout: [sample_counts[nparticipants]] [rows_seen[nparticipants]] then
198 * [samples[worker_id][max_per_worker * dim]] packed.
199 * ---------------------------------------------------------------- */
200
201 typedef struct PrismDsmSamples
202 {
203 uint32_t nparticipants;
204 uint32_t max_per_worker;
205 Dimension dim;
206 /* counts[nparticipants] (rows kept), then MAXALIGN'd
207 * seen[nparticipants] (live rows the sampling scan visited; kept < seen
208 * means the sample is a strict subset of the table -- 64-bit, since a
209 * participant's scan share is unbounded and a wrapped counter would
210 * silently flip the refine gate), then the per-worker sample blocks
211 * (samples[worker][max_per_worker * dim]) packed right after. */
212 uint32_t counts[];
213 } PrismDsmSamples;
214
215 static inline uint32_t *
216 4537 prism_dsm_sample_counts(PrismDsmSamples *s)
217 {
218
4/4
✓ Branch 0 taken 509 times.
✓ Branch 1 taken 561 times.
✓ Branch 2 taken 44 times.
✓ Branch 3 taken 86 times.
4407 return s->counts;
219 }
220
221 static inline uint64_t *
222 2564 prism_dsm_sample_seen(PrismDsmSamples *s)
223 {
224 1503 return (uint64_t *)MAXALIGN(s->counts + s->nparticipants);
225 }
226
227 static inline float *
228 1931 prism_dsm_worker_samples(PrismDsmSamples *s, int worker_id)
229 {
230 /* The sample blocks begin after counts[] and seen[]. */
231 1931 float *samples = (float *)(prism_dsm_sample_seen(s) + s->nparticipants);
232 1931 return samples + (size_t)worker_id * s->max_per_worker * s->dim;
233 }
234
235 static inline Size
236 336 prism_dsm_samples_size(
237 int nparticipants, uint32_t max_per_worker, Dimension dim)
238 {
239 336 Size sz = offsetof(PrismDsmSamples, counts);
240 336 sz += (Size)nparticipants * sizeof(uint32_t);
241 336 sz = MAXALIGN(sz);
242 336 sz += (Size)nparticipants * sizeof(uint64_t);
243 336 sz += (Size)nparticipants * max_per_worker * dim * sizeof(float);
244
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 45 times.
336 return sz;
245 }
246
247 /* ----------------------------------------------------------------
248 * Root assignments: per-worker uint32_t[max_per_worker] in DSM
249 *
250 * After root k-means converges, each worker stores its root
251 * assignments here. Used to filter samples during child k-means.
252 * ---------------------------------------------------------------- */
253
254 typedef struct PrismDsmRootAssign
255 {
256 uint32_t nparticipants;
257 uint32_t max_per_worker;
258 /* assignments[worker][max_per_worker] packed contiguously. */
259 uint32_t assignments[];
260 } PrismDsmRootAssign;
261
262 static inline uint32_t *
263 1586 prism_dsm_root_assignments(PrismDsmRootAssign *ra, int worker_id)
264 {
265 1586 return ra->assignments + (size_t)worker_id * ra->max_per_worker;
266 }
267
268 static inline Size
269 254 prism_dsm_root_assign_size(int nparticipants, uint32_t max_per_worker)
270 {
271 254 Size sz = offsetof(PrismDsmRootAssign, assignments);
272 254 sz += (Size)nparticipants * max_per_worker * sizeof(uint32_t);
273 254 return sz;
274 }
275
276 /* ----------------------------------------------------------------
277 * K-means shared centroids in DSM
278 *
279 * centroids[nlist * dim] + norms_c[nlist]
280 * Written by leader between barriers, read by workers during
281 * assignment.
282 * ---------------------------------------------------------------- */
283
284 static inline Size
285 209 prism_dsm_centroids_size(uint32_t nlist, Dimension dim)
286 {
287 209 return (Size)nlist * dim * sizeof(float) + (Size)nlist * sizeof(float);
288 }
289
290 static inline float *
291 940 prism_dsm_centroids(char *base)
292 {
293 940 return (float *)base;
294 }
295
296 /* ----------------------------------------------------------------
297 * Child subtrees: per-root-child HKMeansResult blobs in DSM
298 *
299 * After the root k-means splits the samples into fan_out groups, each group's
300 * subtree is built independently and written, as a contiguous HKMeansResult,
301 * into a fixed-size slot. The batched streaming build keeps only a bounded
302 * ring of nparticipants slots resident (each participant owns slot
303 * participant_id; the leader streams each batch to pages before the next batch
304 * reuses the ring), so the region is O(nparticipants * slot_size), independent
305 * of the partition count. The leader sizes the slot only after root
306 * assignment -- a subtree can never hold more leaves than its child's
307 * samples, which caps the slot far below the analytic worst case -- and
308 * creates the ring as its own segment then (the subtree-ring seam below);
309 * the slot size travels in PrismBuildShared.subtree_slot_size.
310 * ---------------------------------------------------------------- */
311
312 static inline Size
313 50 prism_dsm_child_subtrees_size(uint32_t nslots, uint64_t slot_size)
314 {
315 50 return (Size)nslots * (Size)slot_size;
316 }
317
318 static inline char *
319 378 prism_dsm_child_subtree(char *base, uint32_t slot, uint64_t slot_size)
320 {
321 378 return base + (Size)slot * (Size)slot_size;
322 }
323
324 /*
325 * Worst-case leaf count for a tree targeting `nlist` leaves at this fan_out:
326 * fan_out^nlevels (clamped to >= nlist). The exact leaf count isn't known
327 * until k-means runs, so the parallel build uses this to size its DSM regions
328 * up front, then narrows to the real tree->nleaves afterward.
329 */
330 static inline uint32_t
331 127 prism_max_nlist(uint32_t nlist, uint32_t fan_out)
332 {
333 127 uint32_t nlevels = vs_hkmeans_nlevels(nlist, fan_out);
334 127 uint32_t max_nlist = 1;
335
3/3
✓ Branch 0 taken 122 times.
✓ Branch 1 taken 154 times.
✓ Branch 2 taken 45 times.
321 for (uint32_t l = 0; l < nlevels; l++)
336 194 max_nlist *= fan_out;
337 127 return max_nlist < nlist ? nlist : max_nlist;
338 }
339
340 static inline float *
341 636 prism_dsm_norms_c(char *base, uint32_t nlist, Dimension dim)
342 {
343 /* norms_c[nlist] sits right after centroids[nlist * dim]. */
344
2/2
✓ Branch 1 taken 44 times.
✓ Branch 2 taken 86 times.
636 return prism_dsm_centroids(base) + (size_t)nlist * dim;
345 }
346
347 /* ----------------------------------------------------------------
348 * K-means per-worker accumulators in DSM
349 *
350 * Per worker: centroid_sums[nlist * dim] + centroid_cnts[nlist]
351 * + cost (1 float)
352 * ---------------------------------------------------------------- */
353
354 static inline Size
355 6611 prism_dsm_km_worker_size(uint32_t nlist, Dimension dim)
356 {
357 6611 return (Size)nlist * dim * sizeof(float) + (Size)nlist * sizeof(uint32_t) +
358 sizeof(float);
359 }
360
361 static inline Size
362 527 prism_dsm_km_workers_size(int nparticipants, uint32_t nlist, Dimension dim)
363 {
364 572 return (Size)nparticipants * prism_dsm_km_worker_size(nlist, dim);
365 }
366
367 static inline float *
368 7707 prism_dsm_km_worker_sums(
369 char *base, uint32_t nlist, Dimension dim, int worker_id)
370 {
371 13791 return (float *)(base +
372
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
7707 (size_t)worker_id * prism_dsm_km_worker_size(nlist, dim));
373 }
374
375 static inline uint32_t *
376 5679 prism_dsm_km_worker_cnts(
377 char *base, uint32_t nlist, Dimension dim, int worker_id)
378 {
379 /* cnts[nlist] sits right after sums[nlist * dim]. */
380
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
5679 return (uint32_t *)(prism_dsm_km_worker_sums(base, nlist, dim, worker_id) +
381 4056 (size_t)nlist * dim);
382 }
383
384 static inline float *
385 3651 prism_dsm_km_worker_cost(
386 char *base, uint32_t nlist, Dimension dim, int worker_id)
387 {
388 /* cost sits right after cnts[nlist]. */
389
2/2
✓ Branch 0 taken 44 times.
✓ Branch 1 taken 86 times.
3651 return (float *)(prism_dsm_km_worker_cnts(base, nlist, dim, worker_id) +
390 nlist);
391 }
392
393 /* ----------------------------------------------------------------
394 * Per-worker posting output in DSM
395 *
396 * active[nlist] per worker, packed contiguously: a flag per cluster marking
397 * which lists this worker wrote into.
398 * ---------------------------------------------------------------- */
399
400 static inline Size
401 prism_dsm_worker_output_size(uint32_t nlist, int nparticipants)
402 {
403 return (Size)nparticipants * nlist * sizeof(bool);
404 }
405
406 static inline bool *
407 prism_dsm_worker_active(char *base, uint32_t nlist, int worker_id)
408 {
409 return (bool *)(base + (Size)worker_id * nlist * sizeof(bool));
410 }
411
412 /* ----------------------------------------------------------------
413 * Shared callbacks — used by both leader and workers
414 * ---------------------------------------------------------------- */
415
416 typedef struct SampleCbState
417 {
418 float *samples;
419 uint32_t count;
420 uint64_t seen;
421 uint32_t max_samples;
422 uint32_t stride;
423 uint32_t stride_counter;
424 Dimension dim;
425 DistanceMetric metric;
426 } SampleCbState;
427
428 /*
429 * Shared sampling logic, called per live tuple with a raw vector pointer (no
430 * Datum) so the same code serves both back-ends. prism_build_scan feeds it:
431 * the PG scan unwraps each heap tuple's Datum, the standalone scan passes its
432 * in-memory vectors directly.
433 */
434 extern void
435 prism_sample_cb(void *state, ItemPointerData tid, const float *vec);
436
437 extern void prism_km_assign_and_accumulate(
438 const float *samples,
439 uint32_t nsamples,
440 const float *centroids,
441 const float *norms_c,
442 uint32_t nlist,
443 Dimension dim,
444 DistanceMetric metric,
445 float *out_sums,
446 uint32_t *out_cnts,
447 float *out_cost);
448
449 /*
450 * Build ONE root-child's subtree into `slot`. Used by the batched streaming
451 * build, where each participant builds one child per batch into a ring slot
452 * indexed by participant (so only nparticipants subtrees are resident) and the
453 * leader streams each to pages.
454 */
455 extern void prism_build_child_subtree(
456 uint32_t child,
457 uint32_t child_count,
458 int nparticipants,
459 PrismDsmSamples *dsm_samples,
460 PrismDsmRootAssign *dsm_ra,
461 const float *root_cents,
462 uint32_t nlist,
463 uint32_t fan_out,
464 Dimension dim,
465 DistanceMetric metric,
466 uint32_t km_max_iterations,
467 char *slot,
468 uint64_t slot_size);
469
470 /*
471 * Per-batch leader callback for the batched subtree stream. Fired only on the
472 * leader (participant 0), between the two per-batch barriers, so it can read
473 * the batch's finished subtrees from the slots [0, batch_size) before they are
474 * reused. children[s] is the child id whose subtree sits in slot s (the batch
475 * schedule is largest-first, not sequential).
476 */
477 typedef void (*PrismBatchCb)(
478 void *arg,
479 const uint32_t *children,
480 uint32_t batch_size,
481 char *subtrees_base,
482 uint64_t slot_size);
483
484 /*
485 * Batched subtree build shared by the leader (participant 0) and workers:
486 * in ceil(km_k / nparticipants) batches, each participant clusters one root
487 * child's subtree into its ring slot, then (leader only) batch_cb consumes
488 * the batch -- recording its layout and spilling the blobs for the
489 * leader-only streaming write that follows; two barriers per batch keep all
490 * participants in lockstep. Children are scheduled largest-first (LPT, by
491 * root-assigned sample count with an id tie-break): every batch runs at the
492 * pace of its slowest subtree, so the skewed children go into the full
493 * batches and the tail batch gets the small ones. The schedule derives from
494 * shared state, so every participant computes it identically with no
495 * coordination; out_child_order (leader passes a km_k buffer, workers NULL)
496 * returns it so the blob replay can place each subtree at its child's
497 * reserved range. Called identically by leader and workers (workers pass
498 * batch_cb = NULL), so the barrier sequence matches by construction; one
499 * invocation per build.
500 */
501 extern void prism_pbuild_stream_subtrees(
502 int participant_id,
503 int nparticipants,
504 PrismDsmSamples *dsm_samples,
505 PrismDsmRootAssign *dsm_ra,
506 const float *root_cents,
507 uint32_t km_k,
508 uint32_t nlist,
509 uint32_t fan_out,
510 Dimension dim,
511 DistanceMetric metric,
512 uint32_t km_max_iterations,
513 char *subtrees_base,
514 uint64_t slot_size,
515 Barrier *barrier,
516 PrismBatchCb batch_cb,
517 void *cb_arg,
518 uint32_t *out_child_order);
519
520 /*
521 * Per-participant execution of phases 1, 2, and 2b, shared by the leader
522 * (participant 0) and the workers — see the contract on the definitions in
523 * parallel_build_worker.c. Each runs its phase's per-participant work and the
524 * phase barrier(s); the leader-only seed/reduce in the k-means phase is gated
525 * on participant_id == 0. prism_pbuild_exec_kmeans returns the iteration count
526 * (for the leader's timing log).
527 */
528 extern void prism_pbuild_exec_sampling(
529 int participant_id,
530 Relation heap,
531 Relation index,
532 struct IndexInfo *index_info,
533 PrismBuildShared *shared,
534 PrismDsmSamples *dsm_samples,
535 Barrier *barrier);
536
537 extern uint32_t prism_pbuild_exec_kmeans(
538 int participant_id,
539 PrismBuildShared *shared,
540 PrismDsmSamples *dsm_samples,
541 char *centroids_base,
542 char *km_workers_base,
543 Barrier *barrier);
544
545 extern void prism_pbuild_exec_root_assign(
546 int participant_id,
547 PrismBuildShared *shared,
548 PrismDsmSamples *dsm_samples,
549 PrismDsmRootAssign *dsm_ra,
550 char *centroids_base,
551 Barrier *barrier);
552
553 /* ----------------------------------------------------------------
554 * Full-table leaf-centroid refinement (parallel)
555 *
556 * A single shared accumulator (sums[nleaves*dim] + counts[nleaves]) holds the
557 * per-leaf running mean for all participants -- one copy, independent of the
558 * worker count, so it scales to fine nlist where per-worker accumulators would
559 * not. Concurrent updates are guarded by a striped lock array owned by the
560 * back-end (prism_pbuild_accum_lock/unlock); contention is low because rows
561 * spread across nleaves leaves. sums accumulate in double so the refined
562 * means match the serial build bit-for-bit at any table size.
563 * ---------------------------------------------------------------- */
564
565 #define PRISM_REFINE_LOCK_STRIPES 256
566
567 typedef struct PrismDsmRefineAccum
568 {
569 uint32_t nleaves;
570 Dimension dim;
571 /* double sums[nleaves * dim], then uint64 counts[nleaves], packed after.
572 */
573 double sums[FLEXIBLE_ARRAY_MEMBER];
574 } PrismDsmRefineAccum;
575
576 static inline double *
577 6 prism_dsm_refine_sums(PrismDsmRefineAccum *a)
578 {
579 6 return a->sums;
580 }
581
582 static inline uint64_t *
583 6 prism_dsm_refine_counts(PrismDsmRefineAccum *a)
584 {
585 6 return (uint64_t *)(a->sums + (size_t)a->nleaves * a->dim);
586 }
587
588 static inline Size
589 prism_dsm_refine_accum_size(uint32_t nleaves, Dimension dim)
590 {
591 Size sz = offsetof(PrismDsmRefineAccum, sums);
592 sz += (Size)nleaves * dim * sizeof(double);
593 sz += (Size)nleaves * sizeof(uint64_t);
594 return sz;
595 }
596
597 /*
598 * Leaves per refine tile: the accumulator holds at most this many leaves, so
599 * it is a bounded constant (cap_bytes, derived from maintenance_work_mem and
600 * MaxAllocSize by the caller) rather than O(nleaves). nleaves above it just
601 * means more re-scanned tiles, not a bigger allocation. The DSM region is
602 * sized for this (capacity = accum->nleaves); the leader and workers derive
603 * the tile count from it identically, so they stay in barrier lockstep.
604 */
605 static inline uint32_t
606 48 prism_refine_tile_leaves(uint32_t nleaves, Dimension dim, uint64_t cap_bytes)
607 {
608 48 uint64_t per_leaf = (uint64_t)dim * sizeof(double) + sizeof(uint64_t);
609 48 uint64_t t = cap_bytes / (per_leaf ? per_leaf : 1);
610
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 48 times.
48 if (t < 1)
611 ✗ t = 1;
612 48 if (t > nleaves)
613 t = nleaves;
614 48 return (uint32_t)t;
615 }
616
617 /*
618 * The per-refined-leaf head writer is the shared PrismLeafWriteFn from
619 * posting_build.h: the leader rewrites each leaf's posting-list head with
620 * the full-table pt_centroid; workers pass NULL (they never divide/write).
621 */
622
623 /*
624 * Refine the leaf encode references on the whole table, page-backed (one
625 * pass; routing reads only the centroid pages, which refine never rewrites,
626 * so the assignment is a fixed point): every row is routed exactly as the
627 * query/insert do (prism_query_route k=1 over the centroid pages, head ->
628 * leaf), per-leaf means
629 * accumulate into the tiled DSM accumulator, and the leader rewrites each
630 * leaf's head-page pt_centroid to the full-table mean. Both leader
631 * (participant 0) and workers call it; gated by shared->refine so they
632 * run the same barriers. The workers route with their own page-backed qs; the
633 * leader passes qs == NULL (it does not scan) and a write_head callback.
634 * Leaves are processed in tiles of accum->nleaves, so the accumulator stays
635 * bounded regardless of nlist.
636 */
637 extern void prism_pbuild_exec_refine_paged(
638 int participant_id,
639 Relation heap,
640 Relation index,
641 struct IndexInfo *index_info,
642 PrismBuildShared *shared,
643 struct PrismQueryState *qs,
644 BlockNumber first_posting,
645 PrismDsmRefineAccum *accum,
646 Barrier *barrier,
647 PrismLeafWriteFn write_head,
648 void *write_head_ctx);
649
650 /* Striped lock seam for the refine accumulator (back-end owns the locks). */
651 extern void prism_pbuild_accum_lock(PrismBuildShared *shared, uint32_t stripe);
652 extern void
653 prism_pbuild_accum_unlock(PrismBuildShared *shared, uint32_t stripe);
654
655 /*
656 * Scan every vector cooperatively, invoking cb per live tuple. Back-end seam:
657 * the PG implementation (parallel_backend_pg.c) drives a parallel heap scan
658 * and unwraps each tuple; the standalone implementation iterates its in-memory
659 * vector array. allow_sync/progress are PG table_index_build_scan flags,
660 * ignored in standalone. Returns the number of heap tuples this participant
661 * scanned; the posting pass accumulates those into shared->reltuples so the
662 * build reports the table's true row count (index_update_stats writes it to
663 * pg_class.reltuples).
664 */
665 extern double prism_build_scan(
666 Relation heap,
667 Relation index,
668 struct IndexInfo *indexInfo,
669 PrismBuildShared *shared,
670 bool allow_sync,
671 bool progress,
672 PrismBuildScanCb cb,
673 void *state);
674
675 /* ----------------------------------------------------------------
676 * Worker entry point — registered with CreateParallelContext
677 * ---------------------------------------------------------------- */
678
679 extern void prism_parallel_build_main(dsm_segment *seg, shm_toc *toc);
680
681 /* ----------------------------------------------------------------
682 * Worker lifecycle seam — back-end-specific (parallel_backend.c for PG)
683 *
684 * Runtime handles for one participant. prism_pbuild_worker_attach fills it
685 * when the worker joins the build; prism_pbuild_worker_detach tears it down.
686 * The PG version opens the heap/index relations and reports instrumentation;
687 * the standalone version takes the shared state and vectors directly and joins
688 * a thread barrier.
689 * ---------------------------------------------------------------- */
690
691 typedef struct PrismPBuildWorker
692 {
693 PrismBuildShared *shared;
694 Barrier *barrier;
695 Relation heapRel;
696 Relation indexRel;
697 int worker_id;
698 Dimension dim;
699 } PrismPBuildWorker;
700
701 extern void prism_pbuild_worker_attach(shm_toc *toc, PrismPBuildWorker *w);
702 extern void prism_pbuild_worker_detach(shm_toc *toc, PrismPBuildWorker *w);
703
704 /*
705 * Page-backed phase-3 routing storage seam — back-end-specific. Workers read
706 * centroid + posting head pages while routing, so each needs a VsStorage over
707 * the index. PG (separate process) opens one on the worker's indexRel;
708 * standalone (threads) returns the leader's shared in-memory store published
709 * via prism_pbuild_publish_storage before launch. Release is a no-op for
710 * standalone.
711 */
712 extern void
713 prism_pbuild_publish_storage(PrismBuildShared *shared, VsStorage *s);
714 extern VsStorage *prism_pbuild_worker_storage(PrismPBuildWorker *w);
715 extern void prism_pbuild_worker_storage_release(VsStorage *s);
716
717 /*
718 * Accumulate one worker's tuple counts into the shared state under the
719 * back-end's lock (the lock lives in the back-end's derived shared struct).
720 * heap_tuples is the participant's posting-scan share of the heap; the sum
721 * across participants is the table's row count, returned to PostgreSQL as
722 * IndexBuildResult.heap_tuples.
723 */
724 extern void prism_pbuild_worker_add_counts(
725 PrismBuildShared *shared,
726 double indtuples,
727 double soar_dupes,
728 double heap_tuples);
729
730 /* ----------------------------------------------------------------
731 * Leader launch/teardown seam — back-end-specific (parallel_backend.c for PG)
732 *
733 * prism_pbuild_launch starts the workers and blocks until the whole party has
734 * attached to the phase barrier (returning false, after teardown, if none
735 * started). Fewer workers can start than were planned (the launch competes
736 * for the max_parallel_workers pool), so launch narrows
737 * shared->nparticipants to the party that actually attached. Every phase
738 * that partitions work by participant reads the count after at least one
739 * barrier, which orders the narrowing write ahead of the read.
740 * prism_pbuild_teardown frees the parallel context. The standalone versions
741 * spawn/join threads and free the shared arena. (struct ParallelContext is
742 * PostgreSQL's; standalone provides its own definition.)
743 * ---------------------------------------------------------------- */
744
745 struct ParallelContext;
746 struct WalUsage;
747 struct BufferUsage;
748
749 extern void prism_pbuild_teardown(struct ParallelContext *pcxt);
750
751 /*
752 * Sample-region seam. The k-means sample is the build's largest working set
753 * (up to the whole memory budget), but it is dead once the subtrees are
754 * clustered -- long before the posting sort claims its own budget. The PG
755 * back-end therefore keeps it in a dedicated DSM segment handed back through
756 * this seam right after the refine pass (the refine accumulator overlays the
757 * then-dead sample region, so it rides along for free), keeping the build's
758 * peak at one budget instead of stacking sample + accumulator + sort.
759 * Workers attach at startup and every participant releases independently;
760 * the segment is destroyed with the last detach. The standalone back-end
761 * keeps the samples in its arena (it does not bound memory) and treats
762 * release as a no-op.
763 */
764 extern PrismDsmSamples *prism_pbuild_samples_attach(
765 shm_toc *toc, PrismBuildShared *shared, void **seg_out);
766 extern void prism_pbuild_samples_release(PrismDsmSamples *samples, void *seg);
767
768 /* Subtree-ring seam: leader creates the ring (own segment) after root
769 * assignment and publishes the slot size; workers attach after the ring
770 * barrier; everyone releases when the subtree batches are done. */
771 extern char *prism_pbuild_subtree_ring_create(
772 PrismBuildShared *shared,
773 int nparticipants,
774 uint64_t slot_size,
775 void **seg_out);
776 extern char *
777 prism_pbuild_subtree_ring_attach(PrismBuildShared *shared, void **seg_out);
778 extern void prism_pbuild_subtree_ring_release(void *seg);
779
780 /* Exact-internal-centroids seam: the leader serializes the exact
781 * centroid collection (prism_exact_centroid_collection_write, a few MB
782 * even at very large nlist) it gathered while streaming the tree, and
783 * publishes it before the tree-ready barrier; workers attach after that
784 * barrier and hook it into their phase-2.5/3 routing base
785 * (PrismIndexBase.exact_internal). PG uses a dedicated DSM segment whose
786 * handle travels in the back-end shared struct; standalone shares the
787 * leader's allocation. Every participant releases when its routing is
788 * done; the segment dies with the last detach. */
789 extern char *prism_pbuild_exact_centroids_create(
790 PrismBuildShared *shared, uint64_t nbytes, void **seg_out);
791 extern char *
792 prism_pbuild_exact_centroids_attach(PrismBuildShared *shared, void **seg_out);
793 extern void prism_pbuild_exact_centroids_release(void *seg);
794
795 /* One histogram over the shared root assignments: out_counts[km_k] = how
796 * many samples landed in each root child. Every participant derives the
797 * same counts (same shared data); the leader additionally sizes the
798 * subtree-ring slot from them. */
799 extern void prism_pbuild_count_children(
800 PrismDsmSamples *dsm_samples,
801 PrismDsmRootAssign *dsm_ra,
802 int nparticipants,
803 uint32_t km_k,
804 uint32_t *out_counts);
805
806 /*
807 * The refine accumulator overlays the (dead) sample region: same base
808 * address, initialized by the leader after the last sample use and before
809 * the tree-ready barrier that workers pass ahead of the refine phase.
810 */
811 static inline PrismDsmRefineAccum *
812 ✗ prism_pbuild_refine_overlay(PrismDsmSamples *samples)
813 {
814 ✗ return (PrismDsmRefineAccum *)samples;
815 }
816
817 extern bool prism_pbuild_launch(
818 struct ParallelContext *pcxt,
819 Barrier *barrier,
820 PrismBuildShared *shared);
821
822 /* ----------------------------------------------------------------
823 * Leader-only subtree blob store — back-end-specific spillable storage
824 *
825 * The PLAN pass produces every subtree exactly once; the leader appends each
826 * blob here and, after computing the block layout, reads them back in the
827 * same order to stream centroid + head pages -- reading a blob back costs
828 * far less than re-running its clustering, and the streaming needs no
829 * worker participation. The PG implementation is a BufFile temp file: small
830 * blob sets stay
831 * in the kernel page cache, large ones spill to pgsql_tmp, so build memory
832 * stays bounded regardless of the partition count. Standalone keeps the
833 * blobs in memory (in-memory engine). Sequential put/rewind/get only.
834 * ---------------------------------------------------------------- */
835
836 typedef struct PrismBlobStore PrismBlobStore;
837
838 extern PrismBlobStore *prism_pbuild_blobstore_begin(void);
839 extern void prism_pbuild_blobstore_put(
840 PrismBlobStore *bs, const void *blob, uint64_t size);
841 extern void prism_pbuild_blobstore_rewind(PrismBlobStore *bs);
842 /* Read the next blob into buf (capacity max_size); returns its size. */
843 extern uint64_t
844 prism_pbuild_blobstore_get(PrismBlobStore *bs, void *buf, uint64_t max_size);
845 extern void prism_pbuild_blobstore_end(PrismBlobStore *bs);
846
847 /* ----------------------------------------------------------------
848 * Build configuration — back-end-neutral input to the setup seam
849 *
850 * The fields the setup seam needs to size and populate the shared state. The
851 * PG path fills this from its opclass-resolved PrismBuildParams; the
852 * standalone path fills it from its own config. Keeping it neutral lets one
853 * setup-seam signature serve both back-ends.
854 * ---------------------------------------------------------------- */
855
856 typedef struct PrismBuildConfig
857 {
858 Dimension dim;
859 DistanceMetric metric;
860 PrismCentroidFormat centroid_format;
861 uint32_t nlist;
862 uint32_t fan_out;
863 double soar_lambda;
864 double boundary_epsilon;
865 bool fastscan;
866 /* CREATE INDEX CONCURRENTLY: take weak locks + an MVCC scan snapshot. */
867 bool concurrent;
868 } PrismBuildConfig;
869
870 /* ----------------------------------------------------------------
871 * Leader setup seam — back-end-specific (parallel_backend.c for PG)
872 *
873 * Leader-side runtime state: the parallel context, the shared regions, and the
874 * derived sizes. prism_pbuild_setup_shared allocates and populates them (a DSM
875 * segment + regions in PG, a heap arena in standalone) and fills this struct;
876 * the rest of the driver consumes it. Returns false (after teardown) if the
877 * back-end could not start. The WalUsage/BufferUsage fields are PG
878 * instrumentation, unused in standalone.
879 * ---------------------------------------------------------------- */
880
881 typedef struct PrismPBuildLeader
882 {
883 struct ParallelContext *pcxt;
884 PrismBuildShared *shared;
885 Barrier *barrier;
886 PrismDsmSamples *dsm_samples;
887 /* Back-end token for releasing the sample region early (PG: the DSM
888 * segment the samples live in; standalone: NULL). */
889 void *sample_seg;
890 char *centroids_base;
891 float *cents;
892 char *km_workers_base;
893 PrismDsmRootAssign *dsm_ra;
894 struct WalUsage *walusage;
895 struct BufferUsage *bufferusage;
896 int nparticipants;
897 uint32_t km_k;
898 uint32_t max_per_worker;
899 Dimension dim;
900 uint32_t nlist;
901 uint64_t rabitq_seed;
902 uint32_t fan_out;
903 Size dsm_total; /* committed DSM chunk bytes (for the
904 planned-allocation introspection line) */
905 } PrismPBuildLeader;
906
907 extern bool prism_pbuild_setup_shared(
908 PrismPBuildLeader *lead,
909 Relation heap,
910 Relation index,
911 const PrismBuildConfig *config,
912 int nworkers);
913
914 /*
915 * Re-initialize the scan for the posting pass (the sampling pass consumed the
916 * first). Back-end seam: resets the PG parallel-scan descriptor or, in
917 * standalone, the work-stealing cursor.
918 */
919 extern void prism_pbuild_rescan(Relation heap, PrismBuildShared *shared);
920
921 /* ----------------------------------------------------------------
922 * Posting sort seam — cluster-keyed external sort of encoded entries.
923 *
924 * Back-end seam (like prism_build_scan): the PG back-end implements it with a
925 * parallel tuplesort whose shared coordinator (Sharedsort) lives in the build
926 * DSM; the standalone back-end with per-participant in-memory arrays merged
927 * and sorted by cluster. The shared phase-3 code
928 * (parallel_build_worker/leader) calls only this seam, so the posting build
929 * stays bounded by maintenance_work_mem without the standalone core depending
930 * on PostgreSQL's tuplesort.
931 *
932 * Lifecycle:
933 * leader, in setup_shared before launch: size the DSM region with
934 * prism_pbuild_sort_shared_size(); after launch:
935 * prism_pbuild_sort_shared_init(). each worker: begin(is_leader=false) ->
936 * put... -> performsort -> end. leader: begin(is_leader=true) ->
937 * performsort -> getnext... -> end. Entries are fixed-size opaque blobs sorted
938 * by the uint32 cluster key. `seg` and `region` are void* to keep PostgreSQL
939 * types out of the shared header (the PG impl casts seg to dsm_segment*);
940 * `region` is the PRISM_DSM_KEY_SORTSHARED chunk.
941 * ---------------------------------------------------------------- */
942 typedef struct PrismSorter PrismSorter;
943
944 extern Size prism_pbuild_sort_shared_size(int nparticipants);
945 extern void
946 prism_pbuild_sort_shared_init(void *region, int nparticipants, void *seg);
947 extern PrismSorter *prism_pbuild_sort_begin(
948 void *region,
949 void *seg,
950 int participant,
951 int nparticipants,
952 bool is_leader,
953 uint32_t entry_size,
954 int work_mem_kb);
955 extern void prism_pbuild_sort_put(
956 PrismSorter *sorter, uint32_t cluster, const void *entry);
957 extern void prism_pbuild_sort_performsort(PrismSorter *sorter);
958 extern bool prism_pbuild_sort_getnext(
959 PrismSorter *sorter, uint32_t *cluster, const void **entry);
960 extern void prism_pbuild_sort_end(PrismSorter *sorter);
961
962 /*
963 * Build every cluster's posting list from a populated (not yet performsorted)
964 * cluster-keyed sorter. Shared by the serial build and the parallel leader:
965 * performsort, then read entries grouped by cluster and write each list with a
966 * single resident page builder. Cluster c's head is the formula first_posting
967 * + c (a pre-extended head region); continuation pages are appended at the end
968 * of the relation and chained, so no O(nlist) reserve is needed. Ends the
969 * sorter.
970 */
971 extern void prism_posting_build_lists(
972 PrismSorter *sorter,
973 VsStorage *storage,
974 uint32_t nlist,
975 Dimension dim,
976 bool fastscan,
977 const RaBitQParams *rq_params,
978 BlockNumber first_posting);
979
980 /* ----------------------------------------------------------------
981 * Parallel build entry — shared driver (parallel_build_leader.c)
982 *
983 * Runs sampling + k-means + the bounded streaming posting build over the heap
984 * (its vectors) into storage, returning the centroid tree and per-list posting
985 * heads. Returns false if parallelism could not start, so the caller falls
986 * back to a serial build. The PG caller fills config from its resolved build
987 * params; the standalone caller fills it from its index config.
988 *
989 * prog is the build-progress reporting seam (phases + % to
990 * pg_stat_progress_create_index, per-phase stats under prism.log_build_stats).
991 * The PG caller passes its PrismBuildProgress so the parallel phases surface
992 * the same way the serial ones do; the standalone caller passes NULL (the seam
993 * is a no-op stub there).
994 * ---------------------------------------------------------------- */
995
996 struct PrismBuildProgress; /* index/build_progress.h */
997
998 extern bool do_parallel_build(
999 Relation heap,
1000 Relation index,
1001 struct IndexInfo *index_info,
1002 const PrismBuildConfig *config,
1003 VsStorage *storage,
1004 struct PrismBuildProgress *prog,
1005 /* No in-RAM tree is materialized; the streamed tree's shape (leaf
1006 * count + depth) comes back through these for the caller's metadata
1007 * write. *out_nlist == 0 means the heap had no indexable rows. */
1008 uint32_t *out_nlist,
1009 uint8_t *out_tree_nlevels,
1010 double *out_heap_tuples,
1011 double *out_indtuples,
1012 double *out_soar_dupes,
1013 /* The build routes page-backed, so it writes the centroid pages +
1014 * heads into `storage` before the scan and computes the global mean;
1015 * *out_global_mean is the (owned) mean the centroid pages encode
1016 * against, so the caller's metadata write matches. May be NULL. */
1017 float **out_global_mean,
1018 /* The posting-head base: cluster c's head is *out_first_posting + c.
1019 * Callers that need to locate head pages after the build (e.g. the
1020 * standalone driver + its tests) capture it; may be NULL. */
1021 BlockNumber *out_first_posting);
1022
1023 #endif /* PRISM_PARALLEL_BUILD_H */
1024