| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * scan.c - Index scan for prism | ||
| 6 | * | ||
| 7 | * Uses PrismQueryState (shared with standalone) for the search hot | ||
| 8 | * path. PG-specific concerns: scan iterator protocol, memory | ||
| 9 | * contexts, query vector extraction. | ||
| 10 | * | ||
| 11 | * Memory layout: | ||
| 12 | * scan_ctx — scan lifetime (index base, query state, storage) | ||
| 13 | */ | ||
| 14 | |||
| 15 | #include <postgres.h> | ||
| 16 | |||
| 17 | #include <access/genam.h> | ||
| 18 | #include <access/relscan.h> | ||
| 19 | #include <fmgr.h> | ||
| 20 | #include <miscadmin.h> | ||
| 21 | #include <pgstat.h> | ||
| 22 | #include <utils/builtins.h> | ||
| 23 | #include <utils/memutils.h> | ||
| 24 | #include <utils/rel.h> | ||
| 25 | |||
| 26 | #include "algo/vecops.h" | ||
| 27 | #include "amcache.h" | ||
| 28 | #include "build.h" | ||
| 29 | #include "core/platform.h" | ||
| 30 | #include "index/posting_page.h" | ||
| 31 | #include "index/query_scan.h" | ||
| 32 | #include "pg/bufstorage.h" | ||
| 33 | #include "quant/rabitq.h" | ||
| 34 | #include "scan.h" | ||
| 35 | #include "scan_bound.h" | ||
| 36 | #include "support_pg.h" | ||
| 37 | #include "typeinfo.h" | ||
| 38 | #include "types/vec32.h" | ||
| 39 | |||
| 40 | /* Default nprobe — will become a GUC later */ | ||
| 41 | #define PRISM_DEFAULT_NPROBE 10 | ||
| 42 | /* | ||
| 43 | * Smallest top-k any scan is sized for. Not a default in the sense of "what | ||
| 44 | * you get when you ask for nothing" -- a query with no LIMIT is sized from | ||
| 45 | * work_mem (see resolve_top_k) -- but a floor under every sizing, so that a | ||
| 46 | * LIMIT 1 has slack for a candidate that turns out to be a dead tuple the | ||
| 47 | * heap fetch discards, and a work_mem too small to hold more still answers. | ||
| 48 | */ | ||
| 49 | |||
| 50 | /* ---------------------------------------------------------------- | ||
| 51 | * Process-global per-phase accumulators (diagnostic). | ||
| 52 | * | ||
| 53 | * Summed across every prism index scan in this backend so phase timing | ||
| 54 | * can be measured over a large query set (e.g. a full 10k-query | ||
| 55 | * benchmark run in one session) instead of eyeballing EXPLAIN on a | ||
| 56 | * single query. Exposed via vs_phase_stats() / vs_phase_stats_reset(): | ||
| 57 | * | ||
| 58 | * CREATE FUNCTION vs_phase_stats_reset() RETURNS void | ||
| 59 | * AS '$libdir/pg_vectorsearch','vs_phase_stats_reset' LANGUAGE C; | ||
| 60 | * CREATE FUNCTION vs_phase_stats() RETURNS text | ||
| 61 | * AS '$libdir/pg_vectorsearch','vs_phase_stats' LANGUAGE C; | ||
| 62 | * ---------------------------------------------------------------- */ | ||
| 63 | static uint64_t g_phase_nqueries = 0; | ||
| 64 | static uint64_t g_phase_centroid_ns = 0; | ||
| 65 | static uint64_t g_phase_posting_ns = 0; | ||
| 66 | static uint64_t g_phase_rerank_ns = 0; | ||
| 67 | static uint64_t g_phase_entries = 0; | ||
| 68 | static uint64_t g_phase_rotation_ns = 0; | ||
| 69 | static uint64_t g_phase_clut_ns = 0; | ||
| 70 | static uint64_t g_phase_cpageread_ns = 0; | ||
| 71 | static uint64_t g_phase_cscore_ns = 0; | ||
| 72 | /* Routing-depth histogram: deepest probe rank the final top-k came from. */ | ||
| 73 | static uint64_t g_route_sum = 0; | ||
| 74 | static uint64_t g_route_le[7] = {0}; /* <=8,16,32,64,128,256,>256 */ | ||
| 75 | |||
| 76 | 2 | PG_FUNCTION_INFO_V1(vs_phase_stats_reset); | |
| 77 | |||
| 78 | Datum | ||
| 79 | 1 | vs_phase_stats_reset(PG_FUNCTION_ARGS) | |
| 80 | { | ||
| 81 | 1 | g_phase_nqueries = 0; | |
| 82 | 1 | g_phase_centroid_ns = 0; | |
| 83 | 1 | g_phase_posting_ns = 0; | |
| 84 | 1 | g_phase_rerank_ns = 0; | |
| 85 | 1 | g_phase_entries = 0; | |
| 86 | 1 | g_phase_rotation_ns = 0; | |
| 87 | 1 | g_phase_clut_ns = 0; | |
| 88 | 1 | g_phase_cpageread_ns = 0; | |
| 89 | 1 | g_phase_cscore_ns = 0; | |
| 90 | 1 | g_route_sum = 0; | |
| 91 | 1 | vs_bufcache_hits = 0; | |
| 92 | 1 | vs_bufcache_cold = 0; | |
| 93 | 1 | vs_bufcache_stale = 0; | |
| 94 |
2/2✓ Branch 0 taken 7 times.
✓ Branch 1 taken 1 times.
|
8 | for (int i = 0; i < 7; i++) |
| 95 | 7 | g_route_le[i] = 0; | |
| 96 | 1 | PG_RETURN_VOID(); | |
| 97 | } | ||
| 98 | |||
| 99 | 2 | PG_FUNCTION_INFO_V1(vs_routing_stats); | |
| 100 | |||
| 101 | Datum | ||
| 102 | 1 | vs_routing_stats(PG_FUNCTION_ARGS) | |
| 103 | { | ||
| 104 | 1 | char buf[448]; | |
| 105 | 1 | uint64_t n = g_phase_nqueries ? g_phase_nqueries : 1; | |
| 106 | 1 | snprintf( | |
| 107 | buf, | ||
| 108 | sizeof(buf), | ||
| 109 | "bufcache hits=" UINT64_FORMAT " cold=" UINT64_FORMAT | ||
| 110 | " stale=" UINT64_FORMAT " | " | ||
| 111 | "queries=" UINT64_FORMAT " avg_deepest_contrib_rank=%.1f | " | ||
| 112 | "deepest-rank histogram: <=8:%.1f%% <=16:%.1f%% <=32:%.1f%% " | ||
| 113 | "<=64:%.1f%% <=128:%.1f%% <=256:%.1f%% >256:%.1f%%", | ||
| 114 | vs_bufcache_hits, | ||
| 115 | vs_bufcache_cold, | ||
| 116 | vs_bufcache_stale, | ||
| 117 | g_phase_nqueries, | ||
| 118 | 1 | (double)g_route_sum / n, | |
| 119 | 1 | 100.0 * g_route_le[0] / n, | |
| 120 | 1 | 100.0 * g_route_le[1] / n, | |
| 121 | 1 | 100.0 * g_route_le[2] / n, | |
| 122 | 1 | 100.0 * g_route_le[3] / n, | |
| 123 | 1 | 100.0 * g_route_le[4] / n, | |
| 124 | 1 | 100.0 * g_route_le[5] / n, | |
| 125 | 1 | 100.0 * g_route_le[6] / n); | |
| 126 | 1 | PG_RETURN_TEXT_P(cstring_to_text(buf)); | |
| 127 | } | ||
| 128 | |||
| 129 | ✗ | PG_FUNCTION_INFO_V1(vs_phase_stats); | |
| 130 | |||
| 131 | Datum | ||
| 132 | ✗ | vs_phase_stats(PG_FUNCTION_ARGS) | |
| 133 | { | ||
| 134 | ✗ | char buf[256]; | |
| 135 | ✗ | uint64_t n = g_phase_nqueries ? g_phase_nqueries : 1; | |
| 136 | ✗ | snprintf( | |
| 137 | buf, | ||
| 138 | sizeof(buf), | ||
| 139 | "queries=" UINT64_FORMAT " | per-query us: centroid=%.1f " | ||
| 140 | "[rot=%.1f lut=%.1f pageread=%.1f score=%.1f] " | ||
| 141 | "posting=%.1f rerank=%.1f | entries/q=" UINT64_FORMAT, | ||
| 142 | g_phase_nqueries, | ||
| 143 | ✗ | (double)g_phase_centroid_ns / n / VS_NS_PER_US, | |
| 144 | ✗ | (double)g_phase_rotation_ns / n / VS_NS_PER_US, | |
| 145 | ✗ | (double)g_phase_clut_ns / n / VS_NS_PER_US, | |
| 146 | ✗ | (double)g_phase_cpageread_ns / n / VS_NS_PER_US, | |
| 147 | ✗ | (double)g_phase_cscore_ns / n / VS_NS_PER_US, | |
| 148 | ✗ | (double)g_phase_posting_ns / n / VS_NS_PER_US, | |
| 149 | ✗ | (double)g_phase_rerank_ns / n / VS_NS_PER_US, | |
| 150 | g_phase_entries / n); | ||
| 151 | ✗ | PG_RETURN_TEXT_P(cstring_to_text(buf)); | |
| 152 | } | ||
| 153 | |||
| 154 | /* ---------------------------------------------------------------- | ||
| 155 | * Scan result entry | ||
| 156 | * ---------------------------------------------------------------- */ | ||
| 157 | |||
| 158 | typedef struct PrismScanResult | ||
| 159 | { | ||
| 160 | ItemPointerData tid; | ||
| 161 | Distance distance; | ||
| 162 | } PrismScanResult; | ||
| 163 | |||
| 164 | /* ---------------------------------------------------------------- | ||
| 165 | * Scan state | ||
| 166 | * ---------------------------------------------------------------- */ | ||
| 167 | |||
| 168 | typedef struct PrismScanState | ||
| 169 | { | ||
| 170 | /* Common index descriptor (first for cast compatibility) */ | ||
| 171 | PrismIndexBase index_base; | ||
| 172 | |||
| 173 | /* Result iterator */ | ||
| 174 | PrismScanResult *results; | ||
| 175 | uint32_t results_cap; | ||
| 176 | uint32_t nresults; | ||
| 177 | uint32_t curr; | ||
| 178 | bool first; | ||
| 179 | |||
| 180 | /* Shared query state. Allocated on the first search rather than at | ||
| 181 | * beginscan, so it can be sized for the top-k the query actually | ||
| 182 | * asks for -- the LIMIT hint arrives between the two (see | ||
| 183 | * scan_bound.c). Rebuilt only if a later search needs a larger k. */ | ||
| 184 | PrismQueryState qstate; | ||
| 185 | bool qstate_ready; | ||
| 186 | uint32_t max_nprobe; | ||
| 187 | bool has_fastscan; | ||
| 188 | |||
| 189 | /* Rows the enclosing LIMIT will pull, 0 when unknown */ | ||
| 190 | uint32_t scan_bound; | ||
| 191 | |||
| 192 | /* PG storage (index page I/O) */ | ||
| 193 | VsPgStorage storage; | ||
| 194 | |||
| 195 | /* EXPLAIN ANALYZE stats (accumulated across rescans) */ | ||
| 196 | PrismScanStats stats; | ||
| 197 | |||
| 198 | /* Resource owner the params checkout was registered with (the | ||
| 199 | * CurrentResourceOwner at the prism_index_base_init call below); | ||
| 200 | * the endscan release must name the same owner. */ | ||
| 201 | ResourceOwner params_owner; | ||
| 202 | |||
| 203 | MemoryContext scan_ctx; | ||
| 204 | |||
| 205 | /* | ||
| 206 | * Reads the query vector out of the ORDER BY argument: the opclass input | ||
| 207 | * type bound to this index's dimension, plus the conversion buffer a | ||
| 208 | * float32 view needs. It is the same binding build and insert use for | ||
| 209 | * column values, because the ORDER BY argument has the opclass's input | ||
| 210 | * type -- an access, not the vector itself. | ||
| 211 | */ | ||
| 212 | Vec32Access query_vector_access; | ||
| 213 | } PrismScanState; | ||
| 214 | |||
| 215 | const PrismScanStats * | ||
| 216 | 20 | prism_scan_get_stats(IndexScanDesc scan) | |
| 217 | { | ||
| 218 | 20 | PrismScanState *ss = (PrismScanState *)scan->opaque; | |
| 219 |
1/2✓ Branch 0 taken 20 times.
✗ Branch 1 not taken.
|
20 | return ss ? &ss->stats : NULL; |
| 220 | } | ||
| 221 | |||
| 222 | /* ---------------------------------------------------------------- | ||
| 223 | * beginscan | ||
| 224 | * ---------------------------------------------------------------- */ | ||
| 225 | |||
| 226 | IndexScanDesc | ||
| 227 | 267 | prism_beginscan(Relation index, int nkeys, int norderbys) | |
| 228 | { | ||
| 229 | 267 | IndexScanDesc scan = RelationGetIndexScan(index, nkeys, norderbys); | |
| 230 | |||
| 231 | 267 | MemoryContext scan_ctx = AllocSetContextCreate( | |
| 232 | CurrentMemoryContext, "prism scan", ALLOCSET_DEFAULT_SIZES); | ||
| 233 | 267 | MemoryContext old_ctx = MemoryContextSwitchTo(scan_ctx); | |
| 234 | |||
| 235 | 267 | PrismScanState *ss = palloc0(sizeof(PrismScanState)); | |
| 236 | 267 | ss->scan_ctx = scan_ctx; | |
| 237 | 267 | ss->first = true; | |
| 238 | |||
| 239 | /* Immutable index parameters from the per-backend cache (metapage read at | ||
| 240 | * most once per backend). */ | ||
| 241 | 267 | prism_index_base_init(index, &ss->index_base); | |
| 242 | 267 | ss->params_owner = CurrentResourceOwner; | |
| 243 | 267 | PrismScanInfo info = prism_cache_scan_info(index); | |
| 244 | |||
| 245 | /* From the per-backend cache; the buffer must outlive a rescan. */ | ||
| 246 | 267 | ss->query_vector_access = vec32_access( | |
| 247 | 267 | prism_cache_type_info(index), ss->index_base.dim, scan_ctx); | |
| 248 | |||
| 249 | /* Size the per-scan query buffers to the nprobe actually requested | ||
| 250 | * (the GUC is set before the query runs) rather than the worst-case | ||
| 251 | * ceiling: the centroid-search scratch alone is | ||
| 252 | * O(max_nprobe * entries_per_page) candidates, several MB per query | ||
| 253 | * at the ceiling but a few hundred KB at typical nprobe. Headroom | ||
| 254 | * covers routing more leaf candidates than are scanned (bounded | ||
| 255 | * probe expansion); requests beyond the sizing are clamped by | ||
| 256 | * prism_query_execute exactly as they were against the old ceiling. */ | ||
| 257 | 534 | uint32_t req_nprobe = prism_nprobe > 0 ? (uint32_t)prism_nprobe | |
| 258 |
2/2✓ Branch 0 taken 78 times.
✓ Branch 1 taken 189 times.
|
267 | : prism_auto_nprobe(info.nlist); |
| 259 | 267 | uint32_t max_nprobe = req_nprobe + Min(req_nprobe, 256) + 16; | |
| 260 | 267 | if (max_nprobe > 4096) | |
| 261 | max_nprobe = 4096; | ||
| 262 | 267 | if (max_nprobe > info.nlist) | |
| 263 | max_nprobe = info.nlist; | ||
| 264 | 267 | ss->max_nprobe = max_nprobe; | |
| 265 | 267 | ss->has_fastscan = ss->index_base.fastscan != 0; | |
| 266 | |||
| 267 | /* Initialize PG storage */ | ||
| 268 | 267 | vs_pg_storage_init(&ss->storage, index, NULL, ss->index_base.metric); | |
| 269 | 267 | ss->index_base.centroid_storage = &ss->storage.base; | |
| 270 | 267 | ss->index_base.posting_storage = &ss->storage.base; | |
| 271 | 267 | ss->index_base.page_base = NULL; | |
| 272 | |||
| 273 | /* The query state and result buffer are sized on the first search | ||
| 274 | * (ensure_query_state), once the top-k is known. */ | ||
| 275 | |||
| 276 | /* Order-by arrays */ | ||
| 277 |
1/2✓ Branch 0 taken 267 times.
✗ Branch 1 not taken.
|
267 | if (norderbys > 0) |
| 278 | { | ||
| 279 | 267 | scan->xs_orderbyvals = palloc0(norderbys * sizeof(Datum)); | |
| 280 | 267 | scan->xs_orderbynulls = palloc(norderbys * sizeof(bool)); | |
| 281 | 267 | memset(scan->xs_orderbynulls, true, norderbys * sizeof(bool)); | |
| 282 | } | ||
| 283 | |||
| 284 | 267 | MemoryContextSwitchTo(old_ctx); | |
| 285 | 267 | scan->opaque = ss; | |
| 286 | 267 | return scan; | |
| 287 | } | ||
| 288 | |||
| 289 | /* ---------------------------------------------------------------- | ||
| 290 | * rescan | ||
| 291 | * ---------------------------------------------------------------- */ | ||
| 292 | |||
| 293 | void | ||
| 294 | 282 | prism_rescan( | |
| 295 | IndexScanDesc scan, | ||
| 296 | ScanKey keys, | ||
| 297 | int nkeys, | ||
| 298 | ScanKey orderbys, | ||
| 299 | int norderbys) | ||
| 300 | { | ||
| 301 | 282 | PrismScanState *ss = (PrismScanState *)scan->opaque; | |
| 302 | |||
| 303 |
2/4✓ Branch 0 taken 282 times.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 282 times.
|
282 | if (keys && scan->numberOfKeys > 0) |
| 304 | ✗ | memcpy(scan->keyData, keys, scan->numberOfKeys * sizeof(ScanKeyData)); | |
| 305 |
2/4✓ Branch 0 taken 282 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 282 times.
✗ Branch 3 not taken.
|
282 | if (orderbys && scan->numberOfOrderBys > 0) |
| 306 | 282 | memcpy(scan->orderByData, | |
| 307 | orderbys, | ||
| 308 | 282 | scan->numberOfOrderBys * sizeof(ScanKeyData)); | |
| 309 | |||
| 310 | 282 | ss->first = true; | |
| 311 | 282 | ss->curr = 0; | |
| 312 | 282 | ss->nresults = 0; | |
| 313 | |||
| 314 | /* | ||
| 315 | * Size the top-k from the query's LIMIT, if this scan runs under one. | ||
| 316 | * Here rather than pushed in from the executor: by rescan the executor | ||
| 317 | * node already points at this scan, so scan_bound.c can find the right | ||
| 318 | * node by identity, and nothing has to open the scan descriptor before | ||
| 319 | * the executor would. | ||
| 320 | * | ||
| 321 | * Re-resolved on every rescan, not cached: a correlated LIMIT takes a | ||
| 322 | * new value for each outer row, and a stale one would return too few | ||
| 323 | * rows -- the very failure this exists to prevent. nodeLimit.c | ||
| 324 | * re-derives its own bound per rescan for the same reason ("in case | ||
| 325 | * this is a rescan and the previous time we got a different result"). | ||
| 326 | */ | ||
| 327 | 282 | ss->scan_bound = prism_scan_bound(scan); | |
| 328 | 282 | } | |
| 329 | |||
| 330 | /* | ||
| 331 | * Bytes the scan commits per row its top-k can return. | ||
| 332 | * | ||
| 333 | * Every allocation that scales with k, and none scale with the vector | ||
| 334 | * dimension. Sizing this from anything less makes the work_mem ceiling | ||
| 335 | * below a fiction: the extraction buffer alone is PRISM_QUERY_CAND_PER_K | ||
| 336 | * entries per row, which dominates the rest by an order of magnitude. | ||
| 337 | * | ||
| 338 | * vs_topk_init one upper bound and one id per row, plus a | ||
| 339 | * candidate array of two entries per row | ||
| 340 | * prism_query_state_init PRISM_QUERY_CAND_PER_K candidates per row and an | ||
| 341 | * index and a distance for each of them | ||
| 342 | * execute_search one result slot per row | ||
| 343 | * | ||
| 344 | * Keep in step with those three. PRISM_QUERY_CAND_PER_K is the same constant | ||
| 345 | * the allocator uses. | ||
| 346 | * | ||
| 347 | * This prices the state as initialized. extract_candidates doubles the | ||
| 348 | * extraction arrays if a query admits more candidates than the sizing | ||
| 349 | * allowed for, so a scan can exceed this budget; work_mem bounds what the | ||
| 350 | * scan asks for, not the high-water mark of a pathological query. | ||
| 351 | */ | ||
| 352 | #define PRISM_TOP_K_BYTES_PER_ROW \ | ||
| 353 | (sizeof(Distance) + sizeof(uint64_t) + 2 * sizeof(VsTopKEntry) + \ | ||
| 354 | PRISM_QUERY_CAND_PER_K * \ | ||
| 355 | (sizeof(VsTopKEntry) + sizeof(uint32_t) + sizeof(Distance)) + \ | ||
| 356 | sizeof(PrismScanResult)) | ||
| 357 | |||
| 358 | /* | ||
| 359 | * Rows the top-k may be sized to, from work_mem. | ||
| 360 | * | ||
| 361 | * The ceiling belongs to the memory the administrator granted, not to a row | ||
| 362 | * constant: a session with work_mem raised for a large query should be able | ||
| 363 | * to ask for a large LIMIT, and one with it lowered should not be able to | ||
| 364 | * commit the backend to more. The LIMIT and the relation's row count are | ||
| 365 | * inputs to the sizing rather than limits on it; prism.query_limit lowers it | ||
| 366 | * when set, and this bounds whatever the rest of the sizing arrives at. | ||
| 367 | * | ||
| 368 | * Floored at the built-in default so a query always answers something. A | ||
| 369 | * scan, unlike a split, can always return a few rows -- so a work_mem too | ||
| 370 | * small to hold more is a reason to return fewer, not to raise an error. | ||
| 371 | */ | ||
| 372 | static uint32_t | ||
| 373 | 664 | max_top_k_for_work_mem(void) | |
| 374 | { | ||
| 375 | 664 | uint64 rows = ((uint64)work_mem * 1024) / PRISM_TOP_K_BYTES_PER_ROW; | |
| 376 | |||
| 377 | 664 | if (rows < PRISM_DEFAULT_K) | |
| 378 | return PRISM_DEFAULT_K; | ||
| 379 |
1/2✓ Branch 0 taken 664 times.
✗ Branch 1 not taken.
|
664 | if (rows > PG_UINT32_MAX) |
| 380 | return PG_UINT32_MAX; | ||
| 381 | 664 | return (uint32_t)rows; | |
| 382 | } | ||
| 383 | |||
| 384 | /* | ||
| 385 | * Resolve the top-k for a search. | ||
| 386 | * | ||
| 387 | * The rows the query's LIMIT asks for (resolved by scan_bound.c) are the | ||
| 388 | * primary source. With no usable LIMIT the query has asked for every row in | ||
| 389 | * distance order, so the scan is sized for as many as it could possibly | ||
| 390 | * return: what work_mem affords, or the relation's estimated row count if | ||
| 391 | * that is smaller -- an ordered scan cannot return more rows than exist. | ||
| 392 | * Sizing for a fixed handful instead would silently answer a complete | ||
| 393 | * ordered scan with a fraction of it. | ||
| 394 | * | ||
| 395 | * prism.query_limit then lowers the result if it is set below it, and | ||
| 396 | * work_mem bounds it in every case, so a query asking for more rows than | ||
| 397 | * the backend may hold returns as many as it can. | ||
| 398 | * | ||
| 399 | * The built-in default is the floor. It keeps a little slack under a small | ||
| 400 | * LIMIT for rows the executor's heap fetch discards (deleted but not yet | ||
| 401 | * vacuumed), which would otherwise leave a LIMIT 1 empty when its single | ||
| 402 | * candidate is dead, and it leaves a query something to answer with under a | ||
| 403 | * work_mem too small to hold more. | ||
| 404 | */ | ||
| 405 | static uint32_t | ||
| 406 | 282 | resolve_top_k(const PrismScanState *ss, Relation heap) | |
| 407 | { | ||
| 408 |
1/2✓ Branch 0 taken 282 times.
✗ Branch 1 not taken.
|
282 | double rows = heap != NULL ? prism_estimate_heap_tuples(heap) : -1.0; |
| 409 | |||
| 410 | 282 | return prism_scan_resolve_top_k(ss->scan_bound, rows); | |
| 411 | } | ||
| 412 | |||
| 413 | uint32_t | ||
| 414 | 664 | prism_scan_resolve_top_k(uint32_t scan_bound, double heap_rows) | |
| 415 | { | ||
| 416 |
1/2✓ Branch 0 taken 664 times.
✗ Branch 1 not taken.
|
664 | uint32_t cap = max_top_k_for_work_mem(); |
| 417 | 664 | uint32_t k = scan_bound; | |
| 418 | |||
| 419 |
2/2✓ Branch 0 taken 8 times.
✓ Branch 1 taken 656 times.
|
664 | if (k == 0) |
| 420 | { | ||
| 421 | /* No LIMIT to size from: the query has asked for every row in | ||
| 422 | * order, so size for as many as it could return. */ | ||
| 423 | 8 | k = cap; | |
| 424 |
3/4✓ Branch 0 taken 8 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 7 times.
✓ Branch 3 taken 1 times.
|
8 | if (heap_rows >= 1.0 && heap_rows < (double)cap) |
| 425 | 7 | k = (uint32_t)heap_rows; | |
| 426 | } | ||
| 427 | |||
| 428 | /* | ||
| 429 | * prism.query_limit only ever lowers the sizing. Raising it above what | ||
| 430 | * the query asked for would have the scan rank rows the LIMIT then | ||
| 431 | * throws away; the lever exists to cap a query that asks for too much | ||
| 432 | * -- one with no LIMIT, or with one set far higher than the rows the | ||
| 433 | * caller will read. | ||
| 434 | */ | ||
| 435 |
2/2✓ Branch 0 taken 33 times.
✓ Branch 1 taken 631 times.
|
664 | if (prism_query_limit > 0 && (uint32_t)prism_query_limit < k) |
| 436 | 664 | k = (uint32_t)prism_query_limit; | |
| 437 | |||
| 438 | 664 | if (k < PRISM_DEFAULT_K) | |
| 439 | k = PRISM_DEFAULT_K; | ||
| 440 | |||
| 441 | 664 | return k > cap ? cap : k; | |
| 442 | } | ||
| 443 | |||
| 444 | /* | ||
| 445 | * Allocate (or resize) the shared query state and result buffer for a | ||
| 446 | * top-k of at least k. prism_query_execute clamps k to the allocated | ||
| 447 | * max_k, so undersizing here is what silently truncates results. | ||
| 448 | * | ||
| 449 | * A resize discards the previous state rather than adding to it. Every | ||
| 450 | * buffer in it is sized to max_k, and a rescan can resize -- a correlated | ||
| 451 | * LIMIT resolves afresh for each outer row -- so keeping the old ones | ||
| 452 | * would accumulate a full set per resize for the life of the scan. | ||
| 453 | * prism_query_state_cleanup owns that: the state holds its buffers in a | ||
| 454 | * context of its own. | ||
| 455 | * | ||
| 456 | * The result array is deliberately not part of that state. It grows with | ||
| 457 | * what the rerank returns rather than with the sizing, so it lives in the | ||
| 458 | * scan context where repalloc preserves it across a resize. | ||
| 459 | */ | ||
| 460 | static void | ||
| 461 | 282 | ensure_query_state(PrismScanState *ss, uint32_t k) | |
| 462 | { | ||
| 463 | 282 | uint32_t max_k = Max(k, (uint32_t)PRISM_DEFAULT_K); | |
| 464 | |||
| 465 |
4/4✓ Branch 0 taken 15 times.
✓ Branch 1 taken 267 times.
✓ Branch 2 taken 4 times.
✓ Branch 3 taken 11 times.
|
282 | if (ss->qstate_ready && max_k <= ss->qstate.max_k) |
| 466 | return; | ||
| 467 | |||
| 468 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 267 times.
|
271 | MemoryContext old_ctx = MemoryContextSwitchTo(ss->scan_ctx); |
| 469 | |||
| 470 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 267 times.
|
271 | if (ss->qstate_ready) |
| 471 | { | ||
| 472 | 4 | prism_query_state_cleanup(&ss->qstate); | |
| 473 | 4 | ss->qstate_ready = false; | |
| 474 | } | ||
| 475 | |||
| 476 | 271 | prism_query_state_init( | |
| 477 | &ss->qstate, &ss->index_base, max_k, ss->max_nprobe); | ||
| 478 |
2/2✓ Branch 0 taken 227 times.
✓ Branch 1 taken 44 times.
|
271 | if (ss->has_fastscan) |
| 479 | 227 | prism_posting_scan_enable_fastscan( | |
| 480 | &ss->qstate.pscan, prism_fastscan_bits); | ||
| 481 | 271 | ss->qstate_ready = true; | |
| 482 | |||
| 483 | /* The error-bound rerank can return more than max_k results, so this | ||
| 484 | * is a starting size; execute_search grows it as needed. */ | ||
| 485 |
1/2✓ Branch 0 taken 271 times.
✗ Branch 1 not taken.
|
271 | if (ss->results_cap < max_k) |
| 486 | { | ||
| 487 | 271 | size_t bytes = max_k * sizeof(PrismScanResult); | |
| 488 |
2/2✓ Branch 0 taken 4 times.
✓ Branch 1 taken 267 times.
|
271 | ss->results = ss->results ? repalloc(ss->results, bytes) |
| 489 | 267 | : palloc(bytes); | |
| 490 | 271 | ss->results_cap = max_k; | |
| 491 | } | ||
| 492 | |||
| 493 | 271 | MemoryContextSwitchTo(old_ctx); | |
| 494 | } | ||
| 495 | |||
| 496 | /* ---------------------------------------------------------------- | ||
| 497 | * Search execution (called on first gettuple) | ||
| 498 | * ---------------------------------------------------------------- */ | ||
| 499 | |||
| 500 | static void | ||
| 501 | 282 | execute_search(IndexScanDesc scan) | |
| 502 | { | ||
| 503 | 282 | PrismScanState *ss = (PrismScanState *)scan->opaque; | |
| 504 | |||
| 505 | /* | ||
| 506 | * Count the search where it runs, as every core AM does at the start of | ||
| 507 | * its own search: pg_stat_*_indexes.idx_scan, and the per-scan counter | ||
| 508 | * PostgreSQL 18 prints as EXPLAIN's "Index Searches" (also the one a | ||
| 509 | * parallel scan aggregates across workers). Core's indexam.c maintains | ||
| 510 | * neither; it counts only the tuples the scan returns. | ||
| 511 | */ | ||
| 512 |
3/4✓ Branch 0 taken 3 times.
✓ Branch 1 taken 279 times.
✓ Branch 2 taken 3 times.
✗ Branch 3 not taken.
|
282 | pgstat_count_index_scan(scan->indexRelation); |
| 513 |
1/2✓ Branch 0 taken 282 times.
✗ Branch 1 not taken.
|
282 | if (scan->instrument != NULL) |
| 514 | 282 | scan->instrument->nsearches++; | |
| 515 | |||
| 516 | 282 | uint32_t k = resolve_top_k(ss, scan->heapRelation); | |
| 517 | 282 | ensure_query_state(ss, k); | |
| 518 | |||
| 519 | /* Lazily set heap relation for reranking (rel is NULL at | ||
| 520 | * beginscan time; heapRelation becomes available later) */ | ||
| 521 |
3/4✓ Branch 0 taken 282 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 267 times.
✓ Branch 3 taken 15 times.
|
282 | if (scan->heapRelation != NULL && ss->storage.rel == NULL) |
| 522 | 267 | vs_pg_storage_set_rel(&ss->storage, scan->heapRelation); | |
| 523 | |||
| 524 | /* | ||
| 525 | * Extract query vector. The ORDER BY operator belongs to the opclass, so | ||
| 526 | * its argument has the opclass's input type and the same descriptor | ||
| 527 | * converts it -- into the scan-lifetime buffer, so a rescan does not leak | ||
| 528 | * one per execution. vec32_read rejects a dimension mismatch. | ||
| 529 | */ | ||
| 530 | 282 | Datum query_datum = scan->orderByData[0].sk_argument; | |
| 531 | 282 | Vec32Ref qref = vec32_read(&ss->query_vector_access, query_datum); | |
| 532 | |||
| 533 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 282 times.
|
282 | if (qref.dim != ss->index_base.dim) |
| 534 | ✗ | ereport(ERROR, | |
| 535 | (errcode(ERRCODE_DATA_EXCEPTION), | ||
| 536 | errmsg("query dimension %u does not match " | ||
| 537 | "index dimension %u", | ||
| 538 | qref.dim, | ||
| 539 | ss->index_base.dim))); | ||
| 540 | |||
| 541 | /* Execute shared search */ | ||
| 542 | 564 | uint32_t nprobe = prism_nprobe > 0 | |
| 543 | ? (uint32_t)prism_nprobe | ||
| 544 |
2/2✓ Branch 0 taken 80 times.
✓ Branch 1 taken 202 times.
|
282 | : prism_auto_nprobe(ss->index_base.nlist); |
| 545 | |||
| 546 | 282 | PrismQueryStats qstats = {0}; | |
| 547 | 282 | ss->storage.read_count = 0; | |
| 548 | |||
| 549 | 282 | prism_query_execute( | |
| 550 | &ss->qstate, | ||
| 551 | qref.data, | ||
| 552 | k, | ||
| 553 | nprobe, | ||
| 554 | (VsDistanceMode)prism_distance_mode, | ||
| 555 | prism_rerank, | ||
| 556 | &qstats); | ||
| 557 | |||
| 558 | 282 | ss->stats.clusters_scanned = qstats.clusters_scanned; | |
| 559 | 282 | ss->stats.centroid_pages_read = qstats.centroid_pages_read; | |
| 560 | 282 | ss->stats.posting_pages_read = qstats.posting_pages_read; | |
| 561 | 282 | ss->stats.posting_pages_skipped = qstats.posting_pages_skipped; | |
| 562 | 282 | ss->stats.posting_entries_scanned = qstats.posting_entries_scanned; | |
| 563 | 282 | ss->stats.rerank_candidates = ss->qstate.ncandidates; | |
| 564 | 282 | ss->stats.rerank_results = ss->qstate.nresults; | |
| 565 | 282 | ss->stats.storage_reads = ss->storage.read_count; | |
| 566 | 282 | ss->stats.top_k = k; | |
| 567 | 282 | ss->stats.centroid_ns = qstats.centroid_ns; | |
| 568 | 282 | ss->stats.posting_ns = qstats.posting_ns; | |
| 569 | 282 | ss->stats.rerank_ns = qstats.rerank_ns; | |
| 570 | |||
| 571 | /* Accumulate into the process-global diagnostic counters. */ | ||
| 572 | 282 | g_phase_nqueries++; | |
| 573 | 282 | g_phase_centroid_ns += qstats.centroid_ns; | |
| 574 | 282 | g_phase_posting_ns += qstats.posting_ns; | |
| 575 | 282 | g_phase_rerank_ns += qstats.rerank_ns; | |
| 576 | 282 | g_phase_entries += qstats.posting_entries_scanned; | |
| 577 | 282 | g_phase_rotation_ns += qstats.rotation_ns; | |
| 578 | 282 | g_phase_clut_ns += qstats.centroid_lut_ns; | |
| 579 | 282 | g_phase_cpageread_ns += qstats.centroid_pageread_ns; | |
| 580 | 282 | g_phase_cscore_ns += qstats.centroid_score_ns; | |
| 581 | |||
| 582 | { | ||
| 583 | 282 | uint32_t r = qstats.max_contrib_rank; | |
| 584 | 282 | g_route_sum += r; | |
| 585 |
2/2✓ Branch 0 taken 270 times.
✓ Branch 1 taken 12 times.
|
282 | if (r <= 8) |
| 586 | 270 | g_route_le[0]++; | |
| 587 |
2/2✓ Branch 0 taken 8 times.
✓ Branch 1 taken 4 times.
|
12 | else if (r <= 16) |
| 588 | 8 | g_route_le[1]++; | |
| 589 |
1/2✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
|
4 | else if (r <= 32) |
| 590 | 4 | g_route_le[2]++; | |
| 591 | ✗ | else if (r <= 64) | |
| 592 | ✗ | g_route_le[3]++; | |
| 593 | ✗ | else if (r <= 128) | |
| 594 | ✗ | g_route_le[4]++; | |
| 595 | ✗ | else if (r <= 256) | |
| 596 | ✗ | g_route_le[5]++; | |
| 597 | else | ||
| 598 | ✗ | g_route_le[6]++; | |
| 599 | } | ||
| 600 | |||
| 601 | /* Copy results from result ordering. The error-bound rerank can return | ||
| 602 | * more than the beginscan max_k (the rerank set is inflated beyond k to | ||
| 603 | * guarantee correctness), so grow the result buffer to fit. */ | ||
| 604 | 282 | uint32_t nresults = ss->qstate.nresults; | |
| 605 |
2/2✓ Branch 0 taken 16 times.
✓ Branch 1 taken 266 times.
|
282 | if (nresults > ss->results_cap) |
| 606 | { | ||
| 607 | 32 | ss->results = | |
| 608 | 16 | repalloc(ss->results, nresults * sizeof(PrismScanResult)); | |
| 609 | 16 | ss->results_cap = nresults; | |
| 610 | } | ||
| 611 |
2/2✓ Branch 0 taken 8097 times.
✓ Branch 1 taken 282 times.
|
8379 | for (uint32_t i = 0; i < nresults; i++) |
| 612 | { | ||
| 613 | 8097 | uint32_t ci = ss->qstate.result_order[i]; | |
| 614 | 8097 | ss->results[i].tid = prism_posting_decode_tid( | |
| 615 | 8097 | ss->qstate.candidates[ci].id); | |
| 616 | 8097 | ss->results[i].distance = ss->qstate.result_dists[i]; | |
| 617 | } | ||
| 618 | 282 | ss->nresults = nresults; | |
| 619 | 282 | ss->curr = 0; | |
| 620 | 282 | } | |
| 621 | |||
| 622 | /* ---------------------------------------------------------------- | ||
| 623 | * gettuple | ||
| 624 | * ---------------------------------------------------------------- */ | ||
| 625 | |||
| 626 | bool | ||
| 627 | 6607 | prism_gettuple(IndexScanDesc scan, ScanDirection direction) | |
| 628 | { | ||
| 629 | 6607 | PrismScanState *ss = (PrismScanState *)scan->opaque; | |
| 630 | |||
| 631 | 6607 | (void)direction; | |
| 632 | |||
| 633 |
2/2✓ Branch 0 taken 282 times.
✓ Branch 1 taken 6325 times.
|
6607 | if (ss->first) |
| 634 | { | ||
| 635 | 282 | ss->first = false; | |
| 636 | |||
| 637 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 282 times.
|
282 | if (scan->numberOfOrderBys == 0) |
| 638 | return false; | ||
| 639 | |||
| 640 | 282 | execute_search(scan); | |
| 641 | } | ||
| 642 | |||
| 643 |
2/2✓ Branch 0 taken 19 times.
✓ Branch 1 taken 6588 times.
|
6607 | if (ss->curr >= ss->nresults) |
| 644 | return false; | ||
| 645 | |||
| 646 | 6588 | PrismScanResult *entry = &ss->results[ss->curr]; | |
| 647 | |||
| 648 | 6588 | scan->xs_heaptid = entry->tid; | |
| 649 | |||
| 650 | /* prism never returns a lossy match: the heap tuple always satisfies | ||
| 651 | * the original qual, so no recheck is ever needed. */ | ||
| 652 | 6588 | scan->xs_recheck = false; | |
| 653 | 6588 | scan->xs_recheckorderby = false; | |
| 654 | 6588 | scan->xs_orderbyvals[0] = Float8GetDatum((double)entry->distance); | |
| 655 | 6588 | scan->xs_orderbynulls[0] = false; | |
| 656 | |||
| 657 | 6588 | ss->curr++; | |
| 658 | 6588 | return true; | |
| 659 | } | ||
| 660 | |||
| 661 | /* ---------------------------------------------------------------- | ||
| 662 | * endscan | ||
| 663 | * ---------------------------------------------------------------- */ | ||
| 664 | |||
| 665 | void | ||
| 666 | 265 | prism_endscan(IndexScanDesc scan) | |
| 667 | { | ||
| 668 | 265 | PrismScanState *ss = (PrismScanState *)scan->opaque; | |
| 669 | |||
| 670 |
1/2✓ Branch 0 taken 265 times.
✗ Branch 1 not taken.
|
265 | if (ss != NULL) |
| 671 | { | ||
| 672 |
1/2✓ Branch 0 taken 265 times.
✗ Branch 1 not taken.
|
265 | if (ss->qstate_ready) |
| 673 | 265 | prism_query_state_cleanup(&ss->qstate); | |
| 674 | /* Check the RaBitQParams checkout back in before the scan's own | ||
| 675 | * memory goes away — see prism_index_base_init / the beginscan | ||
| 676 | * call above. */ | ||
| 677 | 265 | prism_release_params( | |
| 678 | 265 | ss->index_base.dim, | |
| 679 | ss->index_base.rabitq_seed, | ||
| 680 | ss->params_owner); | ||
| 681 | 265 | MemoryContextDelete(ss->scan_ctx); | |
| 682 | 265 | scan->opaque = NULL; | |
| 683 | } | ||
| 684 | 265 | } | |
| 685 |