GCC Code Coverage Report


Directory: src/
File: src/standalone/parallel_backend.c
Date: 2026-09-30 11:11:31
Exec Total Coverage
Lines: 328 337 97.3%
Functions: 35 37 94.6%
Branches: 49 54 90.7%

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_backend.c - Standalone implementations of the build's back-end
6 * seams.
7 *
8 * The build driver (parallel_build_leader.c + parallel_build_worker.c) is
9 * shared with the PG extension; the parts that genuinely differ by back-end
10 * are expressed as same-named seam functions. This file is the thread-based
11 * mirror of src/pg/parallel_backend.c: the shared state lives in a heap arena
12 * instead of a DSM segment, workers run on the thread pool instead of
13 * background processes, and the scan walks the in-memory vector array instead
14 * of a heap relation. Memory bounding is deliberately out of scope here: the
15 * standalone engine holds its vectors, samples, and subtree blobs in RAM (it
16 * exists to isolate and benchmark the shared build logic), so the
17 * maintenance_work_mem-style budgets the PG back-end enforces have no
18 * standalone equivalent -- bounded-memory behavior (sample caps, spill,
19 * refine tiling) is exercised only through the PG build.
20 */
21
22 #ifdef VS_STANDALONE
23
24 #include "vs_config.h"
25
26 #include <pthread.h>
27 #include <stdbool.h>
28 #include <stdlib.h>
29 #include <string.h>
30
31 #include "algo/hkmeans.h"
32 #include "core/log.h"
33 #include "core/memory.h"
34 #include "index/index_build.h" /* prism_auto_fan_out */
35 #include "index/parallel_build.h"
36 #include "quant/rabitq.h"
37 #include "standalone/parallel_ctx.h"
38 #include "standalone/parallel_scan.h" /* PrismParallelScan */
39
40 /*
41 * Standalone shared build state: the neutral PrismBuildShared plus a mutex
42 * guarding its counters, the heap/index relations the workers attach to, and
43 * the work-stealing scan cursor (the analog of PG's appended
44 * ParallelTableScanDesc). The base is the first member, so the
45 * PrismBuildShared
46 * * the workers look up out of the toc is recovered here as a
47 * PrismBuildSharedStandalone *.
48 */
49 typedef struct PrismBuildSharedStandalone
50 {
51 PrismBuildShared base;
52 pthread_mutex_t mutex;
53 Relation heap;
54 Relation index;
55 PrismParallelScan scan;
56 /* Leader's in-memory page store, published for phase-3 page-backed
57 * routing; workers are threads so they share the pointer directly. */
58 VsStorage *storage;
59 /* Per-child subtree ring, allocated by the leader post-assign; workers
60 * are threads and read the pointer directly. */
61 char *subtree_ring;
62 /* Exact centroid collection, allocated by the leader after the
63 * streaming tree write; workers are threads and read the pointer
64 * directly. */
65 char *exact_centroids;
66 } PrismBuildSharedStandalone;
67
68 /*
69 * Scan every vector, invoking cb per vector. The shared work-stealing cursor
70 * lives in the shared state (initialized by setup_shared from the heap's
71 * vector array); every participant runs it cooperatively. The PG flags
72 * (allow_sync/progress) and the relation/index-info handles are unused here.
73 */
74 double
75 294 prism_build_scan(
76 Relation heap,
77 Relation index,
78 struct IndexInfo *indexInfo,
79 PrismBuildShared *shared,
80 bool allow_sync,
81 bool progress,
82 PrismBuildScanCb cb,
83 void *state)
84 {
85 294 PrismBuildSharedStandalone *sh = (PrismBuildSharedStandalone *)shared;
86
87 (void)heap;
88 (void)index;
89 (void)indexInfo;
90 (void)allow_sync;
91 (void)progress;
92
93 294 return prism_parallel_scan_run(&sh->scan, cb, state);
94 }
95
96 /*
97 * Join the parallel build: look up the shared state and barrier, take the
98 * heap/index relations the leader stored, and attach to the phase barrier. No
99 * relations to open and no instrumentation to start (single process).
100 */
101 void
102 106 prism_pbuild_worker_attach(shm_toc *toc, PrismPBuildWorker *w)
103 {
104 PrismBuildShared *shared =
105 106 shm_toc_lookup(toc, PRISM_DSM_KEY_SHARED, false);
106 106 PrismBuildSharedStandalone *sh = (PrismBuildSharedStandalone *)shared;
107 106 Barrier *barrier = shm_toc_lookup(toc, PRISM_DSM_KEY_BARRIER, false);
108
109 106 w->shared = shared;
110 106 w->barrier = barrier;
111 106 w->heapRel = sh->heap;
112 106 w->indexRel = sh->index;
113 106 w->worker_id = ParallelWorkerNumber + 1;
114 106 w->dim = shared->dim;
115
116 106 BarrierAttach(barrier);
117 106 }
118
119 /*
120 * Leave the parallel build. The barrier was already detached in the worker
121 * body at the phase 2->3 transition, the relations are not owned here, and
122 * there is no instrumentation to report, so this is a no-op.
123 */
124 void
125 106 prism_pbuild_worker_detach(shm_toc *toc, PrismPBuildWorker *w)
126 {
127 (void)toc;
128 (void)w;
129 106 }
130
131 /*
132 * Page-backed routing storage seam (see parallel_build.h). Workers are threads
133 * in the leader's process, so they share the leader's in-memory page store:
134 * the leader publishes it and workers return the same pointer. Release is a
135 * no-op.
136 */
137 void
138 82 prism_pbuild_publish_storage(PrismBuildShared *shared, VsStorage *s)
139 {
140 82 ((PrismBuildSharedStandalone *)shared)->storage = s;
141 82 }
142
143 VsStorage *
144 106 prism_pbuild_worker_storage(PrismPBuildWorker *w)
145 {
146 106 return ((PrismBuildSharedStandalone *)w->shared)->storage;
147 }
148
149 void
150 106 prism_pbuild_worker_storage_release(VsStorage *s)
151 {
152 (void)s; /* shared with the leader; not owned by the worker */
153 106 }
154
155 /*
156 * Tear the parallel context down (joins the worker threads and frees the
157 * arena) and leave parallel mode.
158 */
159 void
160 82 prism_pbuild_teardown(ParallelContext *pcxt)
161 {
162 82 DestroyParallelContext(pcxt);
163 82 ExitParallelMode();
164 82 }
165
166 /*
167 * Leader-only subtree blob store (see parallel_build.h). Standalone is the
168 * in-memory engine, so the store is a growing byte buffer with a read
169 * cursor; the PG back-end spills through a BufFile temp file instead.
170 */
171 struct PrismBlobStore
172 {
173 char *data;
174 uint64_t size;
175 uint64_t cap;
176 uint64_t rpos;
177 /* The store outlives pass-scoped scratch contexts (the plan pass fills
178 * it, the write pass drains it), so growth allocates in the context the
179 * store was created in, not the caller's current one. */
180 VsMemCtx ctx;
181 };
182
183 PrismBlobStore *
184 56 prism_pbuild_blobstore_begin(void)
185 {
186 56 PrismBlobStore *bs = vs_alloc0(sizeof(PrismBlobStore));
187 56 bs->ctx = vs_current_memctx; /* the creating context */
188 56 return bs;
189 }
190
191 void
192 1114 prism_pbuild_blobstore_put(PrismBlobStore *bs, const void *blob, uint64_t size)
193 {
194 1114 uint64_t need = bs->size + sizeof(size) + size;
195
2/2
✓ Branch 0 taken 62 times.
✓ Branch 1 taken 1052 times.
1114 if (need > bs->cap)
196 {
197
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 56 times.
62 uint64_t cap = bs->cap ? bs->cap : (uint64_t)1 << 20;
198
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 62 times.
68 while (cap < need)
199 6 cap *= 2;
200 62 VsMemCtx old = vs_memctx_switch(bs->ctx);
201 62 char *grown = vs_alloc(cap);
202 62 vs_memctx_switch(old);
203
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 56 times.
62 if (bs->data != NULL)
204 {
205 6 memcpy(grown, bs->data, bs->size);
206 6 vs_free(bs->data);
207 }
208 62 bs->data = grown;
209 62 bs->cap = cap;
210 }
211 1114 memcpy(bs->data + bs->size, &size, sizeof(size));
212 1114 bs->size += sizeof(size);
213 1114 memcpy(bs->data + bs->size, blob, size);
214 1114 bs->size += size;
215 1114 }
216
217 void
218 56 prism_pbuild_blobstore_rewind(PrismBlobStore *bs)
219 {
220 56 bs->rpos = 0;
221 56 }
222
223 uint64_t
224 1114 prism_pbuild_blobstore_get(PrismBlobStore *bs, void *buf, uint64_t max_size)
225 {
226 uint64_t size;
227 1114 memcpy(&size, bs->data + bs->rpos, sizeof(size));
228 1114 bs->rpos += sizeof(size);
229
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1114 times.
1114 if (size > max_size)
230 ✗ vs_error("subtree blob larger than its slot");
231 1114 memcpy(buf, bs->data + bs->rpos, size);
232 1114 bs->rpos += size;
233 1114 return size;
234 }
235
236 void
237 56 prism_pbuild_blobstore_end(PrismBlobStore *bs)
238 {
239 56 vs_free(bs->data);
240 56 vs_free(bs);
241 56 }
242
243 /*
244 * Launch the workers on the thread pool and wait until the whole party has
245 * attached to the phase barrier. Returns false (after teardown) if none
246 * started. The poll uses a 1ms latch timeout rather than a wakeup, since no
247 * one sets the leader's latch during attach.
248 */
249 bool
250 82 prism_pbuild_launch(
251 ParallelContext *pcxt, Barrier *barrier, PrismBuildShared *shared)
252 {
253 82 LaunchParallelWorkers(pcxt);
254
255
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 82 times.
82 if (pcxt->nworkers_launched == 0)
256 {
257 ✗ WaitForParallelWorkersToFinish(pcxt);
258 ✗ prism_pbuild_teardown(pcxt);
259 ✗ return false;
260 }
261
262 82 WaitForParallelWorkersToAttach(pcxt);
263
2/2
✓ Branch 1 taken 74 times.
✓ Branch 2 taken 82 times.
156 while (BarrierParticipants(barrier) < pcxt->nworkers_launched + 1)
264 {
265 CHECK_FOR_INTERRUPTS();
266 74 (void)WaitLatch(
267 MyLatch,
268 WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH,
269 1L,
270 WAIT_EVENT_PARALLEL_CREATE_INDEX_SCAN);
271 74 ResetLatch(MyLatch);
272 }
273
274 /* Same contract as the PG back-end: the participant count reflects the
275 * party that attached. The thread pool always starts every planned
276 * worker, so this is the planned count. */
277 82 shared->nparticipants = pcxt->nworkers_launched + 1;
278 82 return true;
279 }
280
281 /*
282 * Allocate and populate the parallel build's shared state over a heap arena
283 * (the standalone shm_toc): the shared header, barrier, sample/centroid/
284 * assignment slots, the tree blob, the per-worker page queues, and the partial
285 * pages. The WAL/buffer-usage regions PG carries are replaced by tiny dummies
286 * so the leader's instrumentation loop has valid (ignored) pointers. Always
287 * succeeds (the arena is plain memory), so it returns true.
288 */
289 bool
290 82 prism_pbuild_setup_shared(
291 PrismPBuildLeader *lead,
292 Relation heap,
293 Relation index,
294 const PrismBuildConfig *config,
295 int nworkers)
296 {
297 82 Dimension dim = config->dim;
298 82 uint32_t nlist = config->nlist;
299 82 int nparticipants = nworkers + 1;
300 82 uint64_t rabitq_seed = VS_RABITQ_BUILD_SEED;
301 164 uint32_t fan_out = config->fan_out > 0 ? config->fan_out
302
1/2
✓ Branch 0 taken 82 times.
✗ Branch 1 not taken.
82 : prism_auto_fan_out(0, nlist, 0);
303 82 uint32_t km_k = fan_out < nlist ? fan_out : nlist;
304 82 const Size vec_nbytes = (Size)dim * sizeof(float);
305
306 82 uint32_t total_samples = nlist * 256;
307
2/2
✓ Branch 0 taken 78 times.
✓ Branch 1 taken 4 times.
82 if (total_samples < 10000)
308 78 total_samples = 10000;
309 82 uint32_t max_per_worker = (total_samples + nparticipants - 1) /
310 nparticipants;
311
312 /* The worker entry is resolved by name at launch; register it once. */
313 static bool registered = false;
314
2/2
✓ Branch 0 taken 3 times.
✓ Branch 1 taken 79 times.
82 if (!registered)
315 {
316 3 vs_parallel_register_worker(
317 "prism_parallel_build_main", prism_parallel_build_main);
318 3 registered = true;
319 }
320
321 82 EnterParallelMode();
322 82 ParallelContext *pcxt = CreateParallelContext(
323 VS_MODULE_NAME, "prism_parallel_build_main", nworkers);
324
325 82 int nw_usage = nworkers > 0 ? nworkers : 1;
326 82 Size usage_sz = (Size)nw_usage * sizeof(WalUsage);
327 82 Size bufuse_sz = (Size)nw_usage * sizeof(BufferUsage);
328 82 Size est_shared = BUFFERALIGN(sizeof(PrismBuildSharedStandalone));
329
330 82 shm_toc_estimate_chunk(&pcxt->estimator, est_shared);
331 82 shm_toc_estimate_chunk(&pcxt->estimator, sizeof(Barrier));
332 82 shm_toc_estimate_chunk(
333 &pcxt->estimator,
334 prism_dsm_samples_size(nparticipants, max_per_worker, dim));
335 82 shm_toc_estimate_chunk(
336 &pcxt->estimator, prism_dsm_centroids_size(km_k, dim));
337 82 shm_toc_estimate_chunk(
338 &pcxt->estimator,
339 prism_dsm_km_workers_size(nparticipants, km_k, dim));
340 82 shm_toc_estimate_chunk(
341 &pcxt->estimator,
342 prism_dsm_root_assign_size(nparticipants, max_per_worker));
343 82 shm_toc_estimate_chunk(
344 &pcxt->estimator, prism_pbuild_sort_shared_size(nparticipants));
345 82 shm_toc_estimate_chunk(&pcxt->estimator, usage_sz);
346 82 shm_toc_estimate_chunk(&pcxt->estimator, bufuse_sz);
347 /* Page-backed routing region (leader fills before the tree-ready barrier):
348 * the global mean. The posting-head base is a scalar in PrismBuildShared.
349 */
350 82 shm_toc_estimate_chunk(&pcxt->estimator, vec_nbytes);
351 /* Keyed regions: shared, barrier, samples, centroids, km_workers,
352 * root_assign, sortshared, global_mean (the subtree ring is allocated
353 * by the leader post-assign, outside the toc). */
354 82 shm_toc_estimate_keys(&pcxt->estimator, 8);
355
356 82 InitializeParallelDSM(pcxt);
357
358 /* ---- Populate shared state ---- */
359 82 PrismBuildSharedStandalone *sh = shm_toc_allocate(pcxt->toc, est_shared);
360 82 PrismBuildShared *shared = &sh->base;
361 82 pthread_mutex_init(&sh->mutex, NULL);
362 82 sh->heap = heap;
363 82 sh->index = index;
364 82 prism_parallel_scan_init(&sh->scan, heap->vectors, heap->nvecs, dim);
365
366 82 shared->dim = dim;
367 82 shared->metric = config->metric;
368 82 shared->nlist = nlist;
369 82 shared->fan_out = fan_out; /* resolved (auto if config 0) */
370 82 shared->subtree_slot_size = 0; /* leader sizes the ring post-assign */
371 82 shared->soar_lambda = config->soar_lambda;
372 82 shared->boundary_epsilon = config->boundary_epsilon;
373 82 shared->fastscan = config->fastscan;
374 82 shared->centroid_format = config->centroid_format;
375 82 shared->rabitq_seed = rabitq_seed;
376 82 shared->nparticipants = nparticipants;
377 82 shared->concurrent = config->concurrent; /* never set standalone */
378 82 shared->work_mem_kb = 0; /* standalone sorter is in-memory */
379 82 shared->max_samples_per_worker = max_per_worker;
380 82 shared->km_max_iterations = 20;
381 82 shared->km_tolerance = 1e-4f;
382 82 shared->km_k = km_k;
383 82 shared->km_converged = false;
384 82 shared->refine_threshold = 0; /* standalone builds do not refine */
385 82 shared->refine = false;
386 82 shared->refine_tile_cap = 0; /* standalone builds are not mem-bounded */
387 /* Page-backed routing knobs: route the build scan for accuracy, not
388 * query speed, matching the PG build (see PRISM_BUILD_CENTROID_* in
389 * posting_build.h). */
390 82 shared->centroid_error_scale = PRISM_BUILD_CENTROID_ERROR_SCALE;
391 82 shared->centroid_beam_scale = PRISM_BUILD_CENTROID_BEAM_SCALE;
392 82 shared->fastscan_bits = 16;
393 82 shared->reltuples = 0.0;
394 82 shared->indtuples = 0.0;
395 82 shared->soar_dupes = 0.0;
396 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_SHARED, shared);
397
398 /* Dynamic barrier (0 parties): the leader and every worker attach as they
399 * start, matching the PG back-end (see do_parallel_build). */
400 82 Barrier *barrier = shm_toc_allocate(pcxt->toc, sizeof(Barrier));
401 82 BarrierInit(barrier, 0);
402 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_BARRIER, barrier);
403
404 82 Size samp_sz = prism_dsm_samples_size(nparticipants, max_per_worker, dim);
405 82 PrismDsmSamples *dsm_samples = shm_toc_allocate(pcxt->toc, samp_sz);
406 82 memset(dsm_samples, 0, samp_sz);
407 82 dsm_samples->nparticipants = nparticipants;
408 82 dsm_samples->max_per_worker = max_per_worker;
409 82 dsm_samples->dim = dim;
410 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_SAMPLES, dsm_samples);
411
412 82 Size cent_sz = prism_dsm_centroids_size(km_k, dim);
413 82 char *centroids_base = shm_toc_allocate(pcxt->toc, cent_sz);
414 82 memset(centroids_base, 0, cent_sz);
415 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_CENTROIDS, centroids_base);
416 82 float *cents = prism_dsm_centroids(centroids_base);
417
418 82 Size km_sz = prism_dsm_km_workers_size(nparticipants, km_k, dim);
419 82 char *km_workers_base = shm_toc_allocate(pcxt->toc, km_sz);
420 82 memset(km_workers_base, 0, km_sz);
421 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_KM_WORKERS, km_workers_base);
422
423 82 Size ra_sz = prism_dsm_root_assign_size(nparticipants, max_per_worker);
424 82 PrismDsmRootAssign *dsm_ra = shm_toc_allocate(pcxt->toc, ra_sz);
425 82 memset(dsm_ra, 0, ra_sz);
426 82 dsm_ra->nparticipants = nparticipants;
427 82 dsm_ra->max_per_worker = max_per_worker;
428 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_ROOT_ASSIGN, dsm_ra);
429
430 /* Shared coordinator for the cluster-keyed posting sort (sort seam). */
431 82 Size sort_sz = prism_pbuild_sort_shared_size(nparticipants);
432 82 void *sortshared = shm_toc_allocate(pcxt->toc, sort_sz);
433 82 memset(sortshared, 0, sort_sz);
434 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_SORTSHARED, sortshared);
435
436 /* Page-backed routing region (leader fills before the tree-ready barrier):
437 * the global mean. The posting-head base is a scalar in PrismBuildShared.
438 */
439 82 float *dsm_gmean = shm_toc_allocate(pcxt->toc, vec_nbytes);
440 82 memset(dsm_gmean, 0, vec_nbytes);
441 82 shm_toc_insert(pcxt->toc, PRISM_DSM_KEY_GLOBAL_MEAN, dsm_gmean);
442
443 /* Dummy usage regions so the leader's instrumentation loop is safe. */
444 82 WalUsage *walusage = shm_toc_allocate(pcxt->toc, usage_sz);
445 82 BufferUsage *bufferusage = shm_toc_allocate(pcxt->toc, bufuse_sz);
446
447 82 lead->pcxt = pcxt;
448 82 lead->shared = shared;
449 82 lead->barrier = barrier;
450 82 lead->dsm_samples = dsm_samples;
451 82 lead->sample_seg = NULL;
452 82 lead->centroids_base = centroids_base;
453 82 lead->cents = cents;
454 82 lead->km_workers_base = km_workers_base;
455 82 lead->dsm_ra = dsm_ra;
456 82 lead->walusage = walusage;
457 82 lead->bufferusage = bufferusage;
458 82 lead->nparticipants = nparticipants;
459 82 lead->km_k = km_k;
460 82 lead->max_per_worker = max_per_worker;
461 82 lead->dim = dim;
462 82 lead->nlist = nlist;
463 82 lead->rabitq_seed = rabitq_seed;
464 82 lead->fan_out = fan_out;
465 82 return true;
466 }
467
468 /*
469 * Accumulate one worker's tuple counts under the mutex (the analog of PG's
470 * spinlock in the derived shared struct).
471 */
472 /*
473 * Sample-region seam: standalone keeps the samples in the shared arena for
474 * the whole build (it does not bound memory), so attach is a plain lookup
475 * and release is a no-op.
476 */
477 PrismDsmSamples *
478 106 prism_pbuild_samples_attach(
479 shm_toc *toc, PrismBuildShared *shared, void **seg_out)
480 {
481 (void)shared;
482 106 *seg_out = NULL;
483 106 return (PrismDsmSamples *)
484 106 shm_toc_lookup(toc, PRISM_DSM_KEY_SAMPLES, false);
485 }
486
487 void
488 188 prism_pbuild_samples_release(PrismDsmSamples *samples, void *seg)
489 {
490 (void)samples;
491 (void)seg;
492 188 }
493
494 /*
495 * Subtree-ring seam (thread back-end): one heap allocation the leader makes
496 * after root assignment; workers share the pointer. Only the creator gets a
497 * non-NULL seg to free at release.
498 */
499 char *
500 30 prism_pbuild_subtree_ring_create(
501 PrismBuildShared *shared,
502 int nparticipants,
503 uint64_t slot_size,
504 void **seg_out)
505 {
506 30 PrismBuildSharedStandalone *sa = (PrismBuildSharedStandalone *)shared;
507
508 30 sa->subtree_ring = vs_alloc(
509 prism_dsm_child_subtrees_size(nparticipants, slot_size));
510 30 shared->subtree_slot_size = slot_size;
511 30 *seg_out = sa->subtree_ring;
512 30 return sa->subtree_ring;
513 }
514
515 char *
516 54 prism_pbuild_subtree_ring_attach(PrismBuildShared *shared, void **seg_out)
517 {
518 54 PrismBuildSharedStandalone *sa = (PrismBuildSharedStandalone *)shared;
519
520 54 *seg_out = NULL;
521 54 return sa->subtree_ring;
522 }
523
524 void
525 84 prism_pbuild_subtree_ring_release(void *seg)
526 {
527
2/2
✓ Branch 0 taken 30 times.
✓ Branch 1 taken 54 times.
84 if (seg != NULL)
528 30 vs_free(seg);
529 84 }
530
531 /*
532 * Exact-centroid seam (thread back-end): one heap allocation the leader makes
533 * after the streaming tree write; workers share the pointer. Only the
534 * creator gets a non-NULL seg to free at release, so the collection outlives
535 * every worker's routing (the leader releases last, after the merge).
536 */
537 char *
538 82 prism_pbuild_exact_centroids_create(
539 PrismBuildShared *shared, uint64_t nbytes, void **seg_out)
540 {
541 82 PrismBuildSharedStandalone *sa = (PrismBuildSharedStandalone *)shared;
542
543 82 sa->exact_centroids = vs_alloc(nbytes);
544 82 *seg_out = sa->exact_centroids;
545 82 return sa->exact_centroids;
546 }
547
548 char *
549 106 prism_pbuild_exact_centroids_attach(PrismBuildShared *shared, void **seg_out)
550 {
551 106 PrismBuildSharedStandalone *sa = (PrismBuildSharedStandalone *)shared;
552
553 106 *seg_out = NULL;
554
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 106 times.
106 if (sa->exact_centroids == NULL)
555 ✗ vs_error(
556 "exact centroid collection attached before the leader "
557 "published it");
558 106 return sa->exact_centroids;
559 }
560
561 void
562 188 prism_pbuild_exact_centroids_release(void *seg)
563 {
564
2/2
✓ Branch 0 taken 82 times.
✓ Branch 1 taken 106 times.
188 if (seg != NULL)
565 82 vs_free(seg);
566 188 }
567
568 void
569 106 prism_pbuild_worker_add_counts(
570 PrismBuildShared *shared,
571 double indtuples,
572 double soar_dupes,
573 double heap_tuples)
574 {
575 106 PrismBuildSharedStandalone *sh = (PrismBuildSharedStandalone *)shared;
576
577 106 pthread_mutex_lock(&sh->mutex);
578 106 shared->indtuples += indtuples;
579 106 shared->soar_dupes += soar_dupes;
580 106 shared->reltuples += heap_tuples;
581 106 pthread_mutex_unlock(&sh->mutex);
582 106 }
583
584 /*
585 * Re-initialize the scan for the posting pass; the sampling pass consumed the
586 * first one. Resets the work-stealing cursor.
587 */
588 void
589 82 prism_pbuild_rescan(Relation heap, PrismBuildShared *shared)
590 {
591 82 PrismBuildSharedStandalone *sh = (PrismBuildSharedStandalone *)shared;
592
593 (void)heap;
594 82 prism_parallel_scan_init(
595 82 &sh->scan, sh->scan.vectors, sh->scan.nvecs, sh->scan.dim);
596 82 }
597
598 /*
599 * Leaf-refinement accumulator lock seam. Standalone never refines
600 * (standalone never refines), so these are unused stubs to satisfy the link.
601 */
602 void
603 ✗ prism_pbuild_accum_lock(PrismBuildShared *shared, uint32_t stripe)
604 {
605 (void)shared;
606 (void)stripe;
607 ✗ }
608
609 void
610 ✗ prism_pbuild_accum_unlock(PrismBuildShared *shared, uint32_t stripe)
611 {
612 (void)shared;
613 (void)stripe;
614 ✗ }
615
616 /* ----------------------------------------------------------------
617 * Posting sort seam (standalone back-end) — in-memory cluster sort
618 *
619 * Each worker collects fixed-size records ([uint32 cluster][entry]) into a
620 * malloc'd buffer and, at performsort, hands ownership to the shared
621 * coordinator (one slot per worker). After the phase barrier the leader
622 * gathers every worker's records into one buffer, qsorts by cluster, and
623 * streams them back via getnext. Standalone is in-memory by design, so this
624 * mirrors the PG parallel-tuplesort seam without bounding memory.
625 * ---------------------------------------------------------------- */
626 typedef struct PrismSortShared
627 {
628 int nparticipants;
629 /*
630 * Three per-participant arrays follow, in order:
631 * char *bufs[nparticipants]; worker record buffer (transferred)
632 * size_t counts[nparticipants]; worker record count
633 * VsMemCtx arenas[nparticipants]; worker arena (ownership transferred)
634 * Each worker allocates its records from its own dedicated arena and, at
635 * performsort, hands the buffer + arena to the leader, which gathers,
636 * sorts, and deletes every worker arena (plus its own) in one go at
637 * sort_end.
638 */
639 } PrismSortShared;
640
641 /*
642 * Offset of the per-participant arrays. PrismSortShared is only int-sized, so
643 * its size is not a multiple of the pointer alignment; round up so bufs[] (and
644 * the size_t/pointer arrays after it) start 8-byte aligned.
645 */
646 #define PRISM_SORTSHARED_HDR (((sizeof(PrismSortShared)) + 7u) & ~(size_t)7u)
647
648 static char **
649 1414 ss_bufs(PrismSortShared *s)
650 {
651 1414 return (char **)((char *)s + PRISM_SORTSHARED_HDR);
652 }
653
654 static size_t *
655 984 ss_counts(PrismSortShared *s)
656 {
657 984 return (size_t *)(ss_bufs(s) + s->nparticipants);
658 }
659
660 static VsMemCtx *
661 550 ss_arenas(PrismSortShared *s)
662 {
663 550 return (VsMemCtx *)(ss_counts(s) + s->nparticipants);
664 }
665
666 struct PrismSorter
667 {
668 bool is_leader;
669 uint32_t entry_size;
670 size_t stride; /* align4(4 + entry_size) */
671 char *buf; /* worker: own records; leader: merged */
672 size_t n; /* record count */
673 size_t cap; /* capacity (records) */
674 size_t cursor; /* leader getnext position */
675 PrismSortShared *sh;
676 int participant;
677 VsMemCtx arena; /* dedicated; holds buf (and merged buf for leader) */
678 };
679
680 static int
681 1518988 prism_sort_cluster_cmp(const void *a, const void *b)
682 {
683 1518988 uint32_t ca = *(const uint32_t *)a;
684 1518988 uint32_t cb = *(const uint32_t *)b;
685 1518988 return (ca > cb) - (ca < cb);
686 }
687
688 Size
689 168 prism_pbuild_sort_shared_size(int nparticipants)
690 {
691 336 return PRISM_SORTSHARED_HDR +
692 168 (size_t)nparticipants *
693 (sizeof(char *) + sizeof(size_t) + sizeof(VsMemCtx));
694 }
695
696 void
697 86 prism_pbuild_sort_shared_init(void *region, int nparticipants, void *seg)
698 {
699 (void)seg;
700 86 PrismSortShared *s = (PrismSortShared *)region;
701 86 s->nparticipants = nparticipants;
702 86 memset(ss_bufs(s), 0, (size_t)nparticipants * sizeof(char *));
703 86 memset(ss_counts(s), 0, (size_t)nparticipants * sizeof(size_t));
704 86 memset(ss_arenas(s), 0, (size_t)nparticipants * sizeof(VsMemCtx));
705 86 }
706
707 PrismSorter *
708 202 prism_pbuild_sort_begin(
709 void *region,
710 void *seg,
711 int participant,
712 int nparticipants,
713 bool is_leader,
714 uint32_t entry_size,
715 int work_mem_kb)
716 {
717 (void)seg;
718 (void)nparticipants;
719 (void)work_mem_kb;
720 /* The sorter struct itself stays a plain malloc: it is tiny and freed
721 * deterministically by its own thread at sort_end. Its records go in a
722 * dedicated arena (created here, not the worker's thread-local context,
723 * which is deleted when the worker returns) whose ownership transfers to
724 * the leader at performsort. */
725 202 PrismSorter *s = calloc(1, sizeof(PrismSorter));
726 202 s->is_leader = is_leader;
727 202 s->entry_size = entry_size;
728 202 s->stride = (sizeof(uint32_t) + entry_size + 3u) & ~(size_t)3u;
729 202 s->sh = (PrismSortShared *)region;
730 202 s->participant = participant;
731 202 s->arena = vs_memctx_create(NULL, "vs_sort");
732 202 return s;
733 }
734
735 void
736 171082 prism_pbuild_sort_put(PrismSorter *s, uint32_t cluster, const void *entry)
737 {
738
2/2
✓ Branch 0 taken 124 times.
✓ Branch 1 taken 170958 times.
171082 if (s->n == s->cap)
739 {
740 /* Grow by doubling. The arena cannot free or grow in place, so the old
741 * buffer is left behind and reclaimed when the arena is deleted — the
742 * cost of arena-backed growth in the in-memory standalone path. */
743
2/2
✓ Branch 0 taken 12 times.
✓ Branch 1 taken 112 times.
124 size_t newcap = s->cap ? s->cap * 2 : 4096;
744 124 char *nbuf = vs_memctx_alloc(s->arena, newcap * s->stride);
745
2/2
✓ Branch 0 taken 12 times.
✓ Branch 1 taken 112 times.
124 if (s->buf)
746 12 memcpy(nbuf, s->buf, s->n * s->stride);
747 124 s->buf = nbuf;
748 124 s->cap = newcap;
749 }
750 171082 char *rec = s->buf + s->n * s->stride;
751 171082 memcpy(rec, &cluster, sizeof(uint32_t));
752 171082 memcpy(rec + sizeof(uint32_t), entry, s->entry_size);
753 171082 s->n++;
754 171082 }
755
756 void
757 202 prism_pbuild_sort_performsort(PrismSorter *s)
758 {
759
2/2
✓ Branch 0 taken 116 times.
✓ Branch 1 taken 86 times.
202 if (!s->is_leader)
760 {
761 /* Hand the buffer and its arena to the coordinator; the leader reads
762 * the records and deletes the arena. Drop our references so sort_end
763 * below does not delete the arena out from under the leader. */
764 116 ss_bufs(s->sh)[s->participant] = s->buf;
765 116 ss_counts(s->sh)[s->participant] = s->n;
766 116 ss_arenas(s->sh)[s->participant] = s->arena;
767 116 s->buf = NULL;
768 116 s->arena = NULL;
769 116 return;
770 }
771
772 /* Leader: gather every worker's records, then sort by cluster. */
773 86 PrismSortShared *sh = s->sh;
774 86 size_t total = 0;
775
2/2
✓ Branch 0 taken 116 times.
✓ Branch 1 taken 86 times.
202 for (int i = 0; i < sh->nparticipants; i++)
776 116 total += ss_counts(sh)[i];
777
2/2
✓ Branch 0 taken 84 times.
✓ Branch 1 taken 2 times.
86 s->buf = total ? vs_memctx_alloc(s->arena, total * s->stride) : NULL;
778 86 size_t off = 0;
779
2/2
✓ Branch 0 taken 116 times.
✓ Branch 1 taken 86 times.
202 for (int i = 0; i < sh->nparticipants; i++)
780 {
781 116 size_t c = ss_counts(sh)[i];
782
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 112 times.
116 if (c == 0)
783 4 continue;
784 /* Unreachable
785 * unless the standalone arena is exhausted, which no
786 * caller handles. */
787 /* NOLINTNEXTLINE(clang-analyzer-core.NonNullParamChecker) */
788 112 memcpy(s->buf + off * s->stride, ss_bufs(sh)[i], c * s->stride);
789 112 off += c;
790 }
791 86 s->n = total;
792 86 s->cursor = 0;
793
2/2
✓ Branch 0 taken 84 times.
✓ Branch 1 taken 2 times.
86 if (total)
794 84 qsort(s->buf, total, s->stride, prism_sort_cluster_cmp);
795 }
796
797 bool
798 171168 prism_pbuild_sort_getnext(
799 PrismSorter *s, uint32_t *cluster, const void **entry)
800 {
801
2/2
✓ Branch 0 taken 86 times.
✓ Branch 1 taken 171082 times.
171168 if (s->cursor >= s->n)
802 86 return false;
803 171082 char *rec = s->buf + s->cursor * s->stride;
804 171082 memcpy(cluster, rec, sizeof(uint32_t));
805 171082 *entry = rec + sizeof(uint32_t);
806 171082 s->cursor++;
807 171082 return true;
808 }
809
810 void
811 202 prism_pbuild_sort_end(PrismSorter *s)
812 {
813
2/2
✓ Branch 0 taken 86 times.
✓ Branch 1 taken 116 times.
202 if (s->is_leader)
814 {
815 /* Delete the worker arenas handed over at performsort (this frees
816 * their record buffers); then our own arena frees the merged buffer
817 * below. */
818 86 PrismSortShared *sh = s->sh;
819
2/2
✓ Branch 0 taken 116 times.
✓ Branch 1 taken 86 times.
202 for (int i = 0; i < sh->nparticipants; i++)
820 {
821
1/2
✓ Branch 1 taken 116 times.
✗ Branch 2 not taken.
116 if (ss_arenas(sh)[i])
822 116 vs_memctx_delete(ss_arenas(sh)[i]);
823 116 ss_arenas(sh)[i] = NULL;
824 116 ss_bufs(sh)[i] = NULL;
825 }
826 }
827
2/2
✓ Branch 0 taken 86 times.
✓ Branch 1 taken 116 times.
202 if (s->arena) /* NULL on a worker after performsort transferred ownership
828 */
829 86 vs_memctx_delete(s->arena);
830 202 free(s);
831 202 }
832
833 #endif /* VS_STANDALONE */
834