| 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 |