| Line | Branch | Exec | Source |
|---|---|---|---|
| 1 | /* | ||
| 2 | * Copyright (c) 2026 Tiger Data, Inc. | ||
| 3 | * Licensed under the PostgreSQL License. See LICENSE for details. | ||
| 4 | * | ||
| 5 | * bufstorage.c - PG buffer cache VsStorage implementation | ||
| 6 | * | ||
| 7 | * Read path: ReadBuffer + LockBuffer(SHARE) -> return page | ||
| 8 | * UnlockReleaseBuffer on release | ||
| 9 | * | ||
| 10 | * Write path: ReadBufferExtended(P_NEW) + exclusive lock | ||
| 11 | * GenericXLogStart/RegisterBuffer/Finish for WAL | ||
| 12 | * Commit releases buffer | ||
| 13 | */ | ||
| 14 | |||
| 15 | #include <postgres.h> | ||
| 16 | |||
| 17 | #include <access/generic_xlog.h> | ||
| 18 | #include <access/heapam.h> | ||
| 19 | #include <access/htup_details.h> | ||
| 20 | #include <access/tableam.h> | ||
| 21 | #include <catalog/pg_am_d.h> | ||
| 22 | #include <executor/tuptable.h> | ||
| 23 | #include <pgstat.h> | ||
| 24 | #include <storage/buf_internals.h> | ||
| 25 | #include <storage/bufmgr.h> | ||
| 26 | #include <storage/read_stream.h> | ||
| 27 | #include <utils/memutils.h> | ||
| 28 | #include <utils/snapmgr.h> | ||
| 29 | |||
| 30 | #include "algo/distance.h" | ||
| 31 | #include "algo/topk.h" | ||
| 32 | #include "algo/vecops.h" | ||
| 33 | #include "amcache.h" | ||
| 34 | #include "index/posting_page.h" | ||
| 35 | #include "pg/bufstorage.h" | ||
| 36 | #include "support_pg.h" | ||
| 37 | #include "typeinfo.h" | ||
| 38 | |||
| 39 | /* Downcast from base to concrete type */ | ||
| 40 | #define PG_STORAGE(self) ((VsPgStorage *)(self)) | ||
| 41 | |||
| 42 | /* ---------------------------------------------------------------- | ||
| 43 | * Backend-local buffer-id cache (prism.recent_buffers) | ||
| 44 | * | ||
| 45 | * ~22% of warm query CPU is BufTableLookup hash probes inside | ||
| 46 | * ReadBuffer, for index pages that essentially never leave | ||
| 47 | * shared_buffers. Remember the buffer id per block (per backend) and | ||
| 48 | * re-pin it via ReadRecentBuffer, which validates the tag and pins | ||
| 49 | * without touching the buffer mapping table. A stale id (page evicted | ||
| 50 | * or buffer reused) just fails validation and falls back to | ||
| 51 | * ReadBuffer, which refreshes the cached id — correctness never | ||
| 52 | * depends on the cache. | ||
| 53 | * | ||
| 54 | * The cache stores 4-byte buffer ids, not page data, so it does not | ||
| 55 | * duplicate shared_buffers (~11 MB per backend for a 21 GB index). | ||
| 56 | * One relation is cached at a time per backend; switching indexes | ||
| 57 | * swaps the cache. | ||
| 58 | * ---------------------------------------------------------------- */ | ||
| 59 | static bool g_recent_buffers = true; | ||
| 60 | static RelFileLocator g_bufcache_locator; /* zeroed = invalid */ | ||
| 61 | static Buffer *g_bufcache = NULL; | ||
| 62 | static BlockNumber g_bufcache_len = 0; | ||
| 63 | static MemoryContext g_bufcache_ctx = NULL; | ||
| 64 | |||
| 65 | void | ||
| 66 | 255 | vs_pg_storage_set_recent_buffers(bool enabled) | |
| 67 | { | ||
| 68 | 255 | g_recent_buffers = enabled; | |
| 69 | 255 | } | |
| 70 | |||
| 71 | static inline Buffer * | ||
| 72 | 15613030 | bufcache_slot(Relation index, BlockNumber blkno) | |
| 73 | { | ||
| 74 | 15613030 | const RelFileLocator *loc = &index->rd_locator; | |
| 75 | |||
| 76 |
8/10✓ Branch 0 taken 15612859 times.
✓ Branch 1 taken 171 times.
✓ Branch 2 taken 216 times.
✓ Branch 3 taken 15612643 times.
✗ Branch 4 not taken.
✓ Branch 5 taken 15612643 times.
✗ Branch 6 not taken.
✓ Branch 7 taken 15612643 times.
✓ Branch 8 taken 387 times.
✓ Branch 9 taken 15612643 times.
|
15613417 | if (unlikely( |
| 77 | g_bufcache == NULL || | ||
| 78 | !RelFileLocatorEquals(g_bufcache_locator, *loc))) | ||
| 79 | { | ||
| 80 | 387 | BlockNumber nblocks = RelationGetNumberOfBlocks(index); | |
| 81 | |||
| 82 | /* Headroom so post-build inserts don't invalidate the cache. */ | ||
| 83 | 387 | nblocks += nblocks / 8 + 1024; | |
| 84 | |||
| 85 | /* Backend-lifetime cache data belongs under CacheMemoryContext | ||
| 86 | * (as a named child, so it is attributed in | ||
| 87 | * pg_backend_memory_contexts rather than hiding in the top | ||
| 88 | * context). */ | ||
| 89 |
2/2✓ Branch 0 taken 171 times.
✓ Branch 1 taken 216 times.
|
387 | if (g_bufcache_ctx == NULL) |
| 90 | 171 | g_bufcache_ctx = AllocSetContextCreate( | |
| 91 | CacheMemoryContext, | ||
| 92 | "prism recent-buffers cache", | ||
| 93 | ALLOCSET_START_SMALL_SIZES); | ||
| 94 | |||
| 95 |
2/2✓ Branch 0 taken 216 times.
✓ Branch 1 taken 171 times.
|
387 | if (g_bufcache != NULL) |
| 96 | 216 | pfree(g_bufcache); | |
| 97 | 774 | g_bufcache = MemoryContextAlloc( | |
| 98 | 387 | g_bufcache_ctx, (Size)nblocks * sizeof(Buffer)); | |
| 99 |
2/2✓ Branch 0 taken 418319 times.
✓ Branch 1 taken 387 times.
|
418706 | for (BlockNumber i = 0; i < nblocks; i++) |
| 100 | 418319 | g_bufcache[i] = InvalidBuffer; | |
| 101 | 387 | g_bufcache_len = nblocks; | |
| 102 | 387 | g_bufcache_locator = *loc; | |
| 103 | |||
| 104 | /* | ||
| 105 | * Eagerly seed the cache from the buffer descriptors instead of | ||
| 106 | * populating one miss at a time: measured 38% of reads in a | ||
| 107 | * fresh backend (14% even warmed) fall through to the shared | ||
| 108 | * buffer-mapping hash purely because a slot was never | ||
| 109 | * populated. The tags are read WITHOUT the header lock -- a | ||
| 110 | * torn or stale id is harmless because ReadRecentBuffer | ||
| 111 | * re-validates the tag under its own pin, exactly as it does | ||
| 112 | * for ids that went stale after a normal miss fill. | ||
| 113 | */ | ||
| 114 |
2/2✓ Branch 0 taken 5414016 times.
✓ Branch 1 taken 387 times.
|
5414403 | for (int b = 0; b < NBuffers; b++) |
| 115 | { | ||
| 116 |
2/2✓ Branch 0 taken 708109 times.
✓ Branch 1 taken 4705907 times.
|
5414016 | BufferDesc *hdr = GetBufferDescriptor(b); |
| 117 | |||
| 118 |
2/2✓ Branch 0 taken 708109 times.
✓ Branch 1 taken 4705907 times.
|
6122125 | if (BufTagMatchesRelFileLocator(&hdr->tag, loc) && |
| 119 |
1/2✓ Branch 0 taken 15005 times.
✗ Branch 1 not taken.
|
15005 | BufTagGetForkNum(&hdr->tag) == MAIN_FORKNUM && |
| 120 |
1/2✓ Branch 0 taken 15005 times.
✗ Branch 1 not taken.
|
15005 | hdr->tag.blockNum < g_bufcache_len) |
| 121 | 15005 | g_bufcache[hdr->tag.blockNum] = BufferDescriptorGetBuffer(hdr); | |
| 122 | } | ||
| 123 | } | ||
| 124 | |||
| 125 |
1/2✓ Branch 0 taken 15613030 times.
✗ Branch 1 not taken.
|
15613030 | if (unlikely(blkno >= g_bufcache_len)) |
| 126 | return NULL; | ||
| 127 | 15613030 | return &g_bufcache[blkno]; | |
| 128 | } | ||
| 129 | |||
| 130 | /* ---------------------------------------------------------------- | ||
| 131 | * Read path | ||
| 132 | * ---------------------------------------------------------------- */ | ||
| 133 | |||
| 134 | /* Buffer-id cache effectiveness counters (diagnostic; exposed via | ||
| 135 | * vs_routing_stats). */ | ||
| 136 | uint64_t vs_bufcache_hits; | ||
| 137 | uint64_t vs_bufcache_cold; | ||
| 138 | uint64_t vs_bufcache_stale; | ||
| 139 | |||
| 140 | static Page | ||
| 141 | 15613030 | pg_read_page(VsStorage *self, BlockNumber blkno) | |
| 142 | { | ||
| 143 | 15613030 | VsPgStorage *s = PG_STORAGE(self); | |
| 144 | 15613030 | Buffer buf; | |
| 145 | |||
| 146 |
1/2✓ Branch 0 taken 15613030 times.
✗ Branch 1 not taken.
|
15613030 | if (g_recent_buffers) |
| 147 | { | ||
| 148 | 15613030 | Buffer *slot = bufcache_slot(s->index, blkno); | |
| 149 | |||
| 150 |
5/6✓ Branch 0 taken 15613030 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 15608618 times.
✓ Branch 3 taken 4412 times.
✓ Branch 4 taken 15364021 times.
✓ Branch 5 taken 244597 times.
|
31221648 | if (slot != NULL && *slot != InvalidBuffer && |
| 151 | 15608618 | ReadRecentBuffer(s->index->rd_locator, MAIN_FORKNUM, blkno, *slot)) | |
| 152 | { | ||
| 153 | 15364021 | buf = *slot; | |
| 154 | 15364021 | vs_bufcache_hits++; | |
| 155 | /* | ||
| 156 | * ReadRecentBuffer counts the hit in the backend-wide | ||
| 157 | * pgBufferUsage (EXPLAIN BUFFERS) but, taking a locator rather | ||
| 158 | * than a Relation, cannot attribute it to the index the way | ||
| 159 | * ReadBuffer does. Do that here, or pg_statio_*_indexes shows a | ||
| 160 | * warm prism index with almost no block hits. | ||
| 161 | */ | ||
| 162 |
4/4✓ Branch 0 taken 3355188 times.
✓ Branch 1 taken 12008833 times.
✓ Branch 2 taken 67 times.
✓ Branch 3 taken 3355121 times.
|
15364021 | pgstat_count_buffer_hit(s->index); |
| 163 | } | ||
| 164 | else | ||
| 165 | { | ||
| 166 | /* Diagnostic: distinguish never-populated slots (first | ||
| 167 | * touch) from stale ids (eviction churn). */ | ||
| 168 |
2/2✓ Branch 0 taken 4412 times.
✓ Branch 1 taken 244597 times.
|
249009 | if (slot == NULL || *slot == InvalidBuffer) |
| 169 | 4412 | vs_bufcache_cold++; | |
| 170 | else | ||
| 171 | 244597 | vs_bufcache_stale++; | |
| 172 | 249009 | buf = ReadBuffer(s->index, blkno); | |
| 173 |
1/2✗ Branch 0 not taken.
✓ Branch 1 taken 249009 times.
|
249009 | if (slot != NULL) |
| 174 | 249009 | *slot = buf; | |
| 175 | } | ||
| 176 | } | ||
| 177 | else | ||
| 178 | ✗ | buf = ReadBuffer(s->index, blkno); | |
| 179 | |||
| 180 | 15613030 | LockBuffer(buf, BUFFER_LOCK_SHARE); | |
| 181 | 15613030 | s->cur_buf = buf; | |
| 182 | 15613030 | s->read_count++; | |
| 183 | |||
| 184 | 15613030 | return BufferGetPage(buf); | |
| 185 | } | ||
| 186 | |||
| 187 | static void | ||
| 188 | 15639091 | pg_release_page(VsStorage *self, BlockNumber blkno) | |
| 189 | { | ||
| 190 | 15639091 | VsPgStorage *s = PG_STORAGE(self); | |
| 191 | |||
| 192 | 15639091 | (void)blkno; | |
| 193 | 15639091 | UnlockReleaseBuffer(s->cur_buf); | |
| 194 | 15639091 | s->cur_buf = InvalidBuffer; | |
| 195 | 15639091 | } | |
| 196 | |||
| 197 | /* | ||
| 198 | * Issue an async prefetch (posix_fadvise WILLNEED) for one index page so a | ||
| 199 | * batch of upcoming random reads can overlap instead of stalling serially on | ||
| 200 | * a cold cache. No-op for pages already in shared buffers. Gated on | ||
| 201 | * effective_io_concurrency so an operator can disable prefetching (0). | ||
| 202 | */ | ||
| 203 | static void | ||
| 204 | 10134 | pg_prefetch_page(VsStorage *self, BlockNumber blkno) | |
| 205 | { | ||
| 206 | 10134 | VsPgStorage *s = PG_STORAGE(self); | |
| 207 | |||
| 208 |
1/2✓ Branch 0 taken 10134 times.
✗ Branch 1 not taken.
|
10134 | if (effective_io_concurrency == 0) |
| 209 | return; | ||
| 210 | |||
| 211 | 10134 | PrefetchBuffer(s->index, MAIN_FORKNUM, blkno); | |
| 212 | } | ||
| 213 | |||
| 214 | /* ---------------------------------------------------------------- | ||
| 215 | * Write path | ||
| 216 | * ---------------------------------------------------------------- */ | ||
| 217 | |||
| 218 | static Page | ||
| 219 | 70407 | pg_write_page(VsStorage *self, BlockNumber blkno) | |
| 220 | { | ||
| 221 | 70407 | VsPgStorage *s = PG_STORAGE(self); | |
| 222 | |||
| 223 | 70407 | Buffer buf = ReadBuffer(s->index, blkno); | |
| 224 | 70407 | LockBuffer(buf, BUFFER_LOCK_EXCLUSIVE); | |
| 225 | 70407 | s->cur_buf = buf; | |
| 226 | |||
| 227 | 70407 | return BufferGetPage(buf); | |
| 228 | } | ||
| 229 | |||
| 230 | static Page | ||
| 231 | 1405 | pg_new_page(VsStorage *self, BlockNumber *blkno_out) | |
| 232 | { | ||
| 233 | 1405 | VsPgStorage *s = PG_STORAGE(self); | |
| 234 | |||
| 235 | 1405 | Buffer buf = ReadBufferExtended( | |
| 236 | s->index, MAIN_FORKNUM, P_NEW, RBM_NORMAL, NULL); | ||
| 237 | 1405 | LockBuffer(buf, BUFFER_LOCK_EXCLUSIVE); | |
| 238 | |||
| 239 | 1405 | *blkno_out = BufferGetBlockNumber(buf); | |
| 240 | 1405 | s->cur_buf = buf; | |
| 241 | |||
| 242 | 1405 | return BufferGetPage(buf); | |
| 243 | } | ||
| 244 | |||
| 245 | static void | ||
| 246 | 71812 | pg_commit_page(VsStorage *self, BlockNumber blkno) | |
| 247 | { | ||
| 248 | 71812 | VsPgStorage *s = PG_STORAGE(self); | |
| 249 | |||
| 250 | 71812 | (void)blkno; | |
| 251 | |||
| 252 |
2/2✓ Branch 0 taken 35368 times.
✓ Branch 1 taken 36444 times.
|
71812 | if (s->build_mode) |
| 253 | { | ||
| 254 | /* During index build, skip per-page WAL logging. | ||
| 255 | * log_newpage_range() at the end covers all pages. */ | ||
| 256 | 35368 | MarkBufferDirty(s->cur_buf); | |
| 257 | } | ||
| 258 | else | ||
| 259 | { | ||
| 260 | /* | ||
| 261 | * prism pages store their data in the content area between pd_lower | ||
| 262 | * and pd_upper — the region PostgreSQL treats as the free "hole" and | ||
| 263 | * omits from standard-layout full-page images (GenericXLog assumes the | ||
| 264 | * standard layout). Cover the hole (pd_lower = pd_upper) before | ||
| 265 | * logging so the entire page is preserved; prism never reads | ||
| 266 | * pd_lower back (scans locate data via PageGetContents / | ||
| 267 | * pd_special, and centroid pages derive their metadata cursor from | ||
| 268 | * entry_count -- see prism_centroid_meta_end, which exists because of | ||
| 269 | * this overwrite). | ||
| 270 | */ | ||
| 271 | 36444 | PageHeader ph = (PageHeader)BufferGetPage(s->cur_buf); | |
| 272 | 36444 | ph->pd_lower = ph->pd_upper; | |
| 273 | |||
| 274 | 36444 | GenericXLogState *state = GenericXLogStart(s->index); | |
| 275 | 36444 | GenericXLogRegisterBuffer(state, s->cur_buf, GENERIC_XLOG_FULL_IMAGE); | |
| 276 | 36444 | GenericXLogFinish(state); | |
| 277 | } | ||
| 278 | |||
| 279 | 71812 | UnlockReleaseBuffer(s->cur_buf); | |
| 280 | 71812 | s->cur_buf = InvalidBuffer; | |
| 281 | 71812 | } | |
| 282 | |||
| 283 | /* ---------------------------------------------------------------- | ||
| 284 | * Bulk extend: pre-allocate contiguous pages | ||
| 285 | * ---------------------------------------------------------------- */ | ||
| 286 | |||
| 287 | static BlockNumber | ||
| 288 | 194 | pg_extend(VsStorage *self, uint32_t npages) | |
| 289 | { | ||
| 290 | 194 | VsPgStorage *s = PG_STORAGE(self); | |
| 291 | |||
| 292 | 194 | BlockNumber start = InvalidBlockNumber; | |
| 293 | 194 | uint32_t remaining = npages; | |
| 294 | |||
| 295 |
2/2✓ Branch 0 taken 1007 times.
✓ Branch 1 taken 194 times.
|
1201 | while (remaining > 0) |
| 296 | { | ||
| 297 | 1007 | uint32_t batch = Min(remaining, 512); | |
| 298 | 1007 | Buffer *buffers = palloc(batch * sizeof(Buffer)); | |
| 299 | 1007 | uint32_t got = 0; | |
| 300 | |||
| 301 | 2014 | BlockNumber batch_start = ExtendBufferedRelBy( | |
| 302 | 1007 | BMR_REL(s->index), | |
| 303 | MAIN_FORKNUM, | ||
| 304 | NULL, | ||
| 305 | 0, | ||
| 306 | batch, | ||
| 307 | buffers, | ||
| 308 | &got); | ||
| 309 | |||
| 310 |
2/2✓ Branch 0 taken 194 times.
✓ Branch 1 taken 813 times.
|
1007 | if (start == InvalidBlockNumber) |
| 311 | 194 | start = batch_start; | |
| 312 | |||
| 313 |
2/2✓ Branch 0 taken 11193 times.
✓ Branch 1 taken 1007 times.
|
12200 | for (uint32_t i = 0; i < got; i++) |
| 314 | 11193 | ReleaseBuffer(buffers[i]); | |
| 315 | |||
| 316 | 1007 | pfree(buffers); | |
| 317 | 1007 | remaining -= got; | |
| 318 | } | ||
| 319 | |||
| 320 | 194 | return start; | |
| 321 | } | ||
| 322 | |||
| 323 | /* ---------------------------------------------------------------- | ||
| 324 | * Rerank: fetch vectors from heap, compute exact L2 | ||
| 325 | * ---------------------------------------------------------------- */ | ||
| 326 | |||
| 327 | /* Comparator for sorting candidate indices by TID block order */ | ||
| 328 | static int | ||
| 329 | 61088 | cmp_tid_order(const void *a, const void *b, void *arg) | |
| 330 | { | ||
| 331 | 61088 | const VsTopKEntry *cands = arg; | |
| 332 | |||
| 333 | /* Encoded ids are (block << 16) | offset — strictly monotonic in | ||
| 334 | * (block, offset), so raw id comparison IS TID order; no need to | ||
| 335 | * decode both TIDs on every comparison. */ | ||
| 336 | 61088 | uint64_t id_a = cands[*(const uint32_t *)a].id; | |
| 337 | 61088 | uint64_t id_b = cands[*(const uint32_t *)b].id; | |
| 338 | |||
| 339 |
2/2✓ Branch 0 taken 31289 times.
✓ Branch 1 taken 29799 times.
|
61088 | if (id_a < id_b) |
| 340 | return -1; | ||
| 341 |
1/2✓ Branch 0 taken 31289 times.
✗ Branch 1 not taken.
|
31289 | if (id_a > id_b) |
| 342 | 31289 | return 1; | |
| 343 | return 0; | ||
| 344 | } | ||
| 345 | |||
| 346 | /* | ||
| 347 | * pg_rerank - Rerank candidates with exact L2 distances. | ||
| 348 | * | ||
| 349 | * Fetches full-precision vectors from the heap table and computes | ||
| 350 | * exact L2 distance for candidates with nonzero error. Candidates | ||
| 351 | * are visited in TID block-number order to minimize buffer cache | ||
| 352 | * thrash — sequential block access avoids re-pinning the same | ||
| 353 | * buffer multiple times. | ||
| 354 | * | ||
| 355 | * SnapshotAny: medoid vectors are structural index data that happen | ||
| 356 | * to live in the heap. They must be readable regardless of MVCC | ||
| 357 | * state — an MVCC snapshot could fail if the row was DELETEd but | ||
| 358 | * not yet VACUUMed, breaking centroid reranking. If the tuple is | ||
| 359 | * physically reclaimed by VACUUM, the index is stale and needs | ||
| 360 | * REINDEX. | ||
| 361 | */ | ||
| 362 | static uint32_t | ||
| 363 | ✗ | pg_rerank( | |
| 364 | VsStorage *self, | ||
| 365 | const float *query, | ||
| 366 | Dimension dim, | ||
| 367 | const VsTopKEntry *candidates, | ||
| 368 | uint32_t count, | ||
| 369 | uint32_t keep, | ||
| 370 | uint32_t *out_indices, | ||
| 371 | Distance *out_distances) | ||
| 372 | { | ||
| 373 | ✗ | VsPgStorage *s = PG_STORAGE(self); | |
| 374 | |||
| 375 | ✗ | if (s->rel == NULL || count == 0) | |
| 376 | return 0; | ||
| 377 | |||
| 378 | ✗ | AttrNumber vec_attnum = s->index->rd_index->indkey.values[0]; | |
| 379 | |||
| 380 | /* | ||
| 381 | * Widening buffer for a vec16 heap column, allocated once per rerank | ||
| 382 | * call -- rerank runs once per query over the whole candidate set, so this | ||
| 383 | * is not a per-tuple allocation. | ||
| 384 | */ | ||
| 385 | ✗ | Vec32Access input = vec32_access(s->type_info, dim, CurrentMemoryContext); | |
| 386 | |||
| 387 | /* Sort by TID block order for sequential I/O */ | ||
| 388 | ✗ | uint32_t *order = palloc(count * sizeof(uint32_t)); | |
| 389 | ✗ | for (uint32_t i = 0; i < count; i++) | |
| 390 | ✗ | order[i] = i; | |
| 391 | ✗ | qsort_arg( | |
| 392 | order, count, sizeof(uint32_t), cmp_tid_order, (void *)candidates); | ||
| 393 | |||
| 394 | ✗ | VsTopK topk; | |
| 395 | ✗ | vs_topk_init(&topk, keep); | |
| 396 | |||
| 397 | ✗ | TupleTableSlot *slot = table_slot_create(s->rel, NULL); | |
| 398 | |||
| 399 | ✗ | for (uint32_t i = 0; i < count; i++) | |
| 400 | { | ||
| 401 | ✗ | uint32_t idx = order[i]; | |
| 402 | |||
| 403 | ✗ | Distance d; | |
| 404 | ✗ | if (candidates[idx].error == 0.0f) | |
| 405 | { | ||
| 406 | ✗ | d = candidates[idx].distance; | |
| 407 | } | ||
| 408 | else | ||
| 409 | { | ||
| 410 | ✗ | ItemPointerData tid = prism_posting_decode_tid(candidates[idx].id); | |
| 411 | ✗ | if (table_tuple_fetch_row_version(s->rel, &tid, SnapshotAny, slot)) | |
| 412 | { | ||
| 413 | ✗ | bool isnull; | |
| 414 | ✗ | Datum val = slot_getattr(slot, vec_attnum, &isnull); | |
| 415 | ✗ | if (!isnull) | |
| 416 | { | ||
| 417 | ✗ | Vec32Ref qref = {.data = query, .dim = dim}; | |
| 418 | ✗ | Vec32Ref vref = vec32_read(&input, val); | |
| 419 | ✗ | d = vs_distance(qref, vref, s->metric); | |
| 420 | } | ||
| 421 | else | ||
| 422 | { | ||
| 423 | ✗ | d = candidates[idx].distance; | |
| 424 | } | ||
| 425 | ✗ | ExecClearTuple(slot); | |
| 426 | } | ||
| 427 | else | ||
| 428 | { | ||
| 429 | ✗ | d = candidates[idx].distance; | |
| 430 | } | ||
| 431 | } | ||
| 432 | |||
| 433 | ✗ | vs_topk_insert_unique(&topk, d, 0.0f, (uint64_t)idx); | |
| 434 | } | ||
| 435 | |||
| 436 | ✗ | ExecDropSingleTupleTableSlot(slot); | |
| 437 | |||
| 438 | ✗ | VsTopKEntry *entries = palloc(topk.cand_count * sizeof(VsTopKEntry)); | |
| 439 | ✗ | uint32_t nresults; | |
| 440 | ✗ | vs_topk_extract_sorted_unique(&topk, entries, &nresults); | |
| 441 | |||
| 442 | ✗ | for (uint32_t i = 0; i < nresults; i++) | |
| 443 | { | ||
| 444 | ✗ | out_indices[i] = (uint32_t)entries[i].id; | |
| 445 | ✗ | out_distances[i] = entries[i].distance; | |
| 446 | } | ||
| 447 | |||
| 448 | ✗ | pfree(entries); | |
| 449 | ✗ | vs_topk_cleanup(&topk); | |
| 450 | ✗ | pfree(order); | |
| 451 | |||
| 452 | ✗ | return nresults; | |
| 453 | } | ||
| 454 | |||
| 455 | /* ---------------------------------------------------------------- | ||
| 456 | * Heap AM rerank with read_stream | ||
| 457 | * | ||
| 458 | * Uses PostgreSQL's read_stream API for batched async I/O when | ||
| 459 | * fetching heap tuples for reranking. The stream callback yields | ||
| 460 | * block numbers in candidate order, and each buffer is processed | ||
| 461 | * directly as it arrives. | ||
| 462 | * ---------------------------------------------------------------- */ | ||
| 463 | |||
| 464 | typedef struct RerankStreamState | ||
| 465 | { | ||
| 466 | const VsTopKEntry *candidates; | ||
| 467 | const uint32_t *order; | ||
| 468 | VsTopK *topk; | ||
| 469 | VsPgStorage *storage; | ||
| 470 | uint32_t count; | ||
| 471 | uint32_t pos; | ||
| 472 | } RerankStreamState; | ||
| 473 | |||
| 474 | static BlockNumber | ||
| 475 | 12062 | rerank_stream_cb( | |
| 476 | ReadStream *stream, void *callback_private_data, void *per_buffer_data) | ||
| 477 | { | ||
| 478 | 12062 | RerankStreamState *st = callback_private_data; | |
| 479 | |||
| 480 |
2/2✓ Branch 0 taken 11984 times.
✓ Branch 1 taken 278 times.
|
12262 | while (st->pos < st->count) |
| 481 | { | ||
| 482 | 11984 | uint32_t idx = st->order[st->pos++]; | |
| 483 | |||
| 484 | /* Already exact — insert directly, no heap fetch */ | ||
| 485 |
2/2✓ Branch 0 taken 200 times.
✓ Branch 1 taken 11784 times.
|
11984 | if (st->candidates[idx].error == 0.0f) |
| 486 | { | ||
| 487 | 200 | vs_topk_insert_unique( | |
| 488 | st->topk, | ||
| 489 | 200 | st->candidates[idx].distance, | |
| 490 | 0.0f, | ||
| 491 | (uint64_t)idx); | ||
| 492 | 200 | continue; | |
| 493 | } | ||
| 494 | |||
| 495 | 11784 | *(uint32_t *)per_buffer_data = idx; | |
| 496 | 11784 | ItemPointerData tid = prism_posting_decode_tid(st->candidates[idx].id); | |
| 497 | 11784 | return ItemPointerGetBlockNumber(&tid); | |
| 498 | } | ||
| 499 | |||
| 500 | return InvalidBlockNumber; | ||
| 501 | } | ||
| 502 | |||
| 503 | static uint32_t | ||
| 504 | 278 | pg_rerank_readstream( | |
| 505 | VsStorage *self, | ||
| 506 | const float *query, | ||
| 507 | Dimension dim, | ||
| 508 | const VsTopKEntry *candidates, | ||
| 509 | uint32_t count, | ||
| 510 | uint32_t keep, | ||
| 511 | uint32_t *out_indices, | ||
| 512 | Distance *out_distances) | ||
| 513 | { | ||
| 514 | 278 | VsPgStorage *s = PG_STORAGE(self); | |
| 515 | |||
| 516 |
2/4✓ Branch 0 taken 278 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 278 times.
✗ Branch 3 not taken.
|
278 | if (s->rel == NULL || count == 0) |
| 517 | return 0; | ||
| 518 | |||
| 519 | 278 | AttrNumber vec_attnum = s->index->rd_index->indkey.values[0]; | |
| 520 | |||
| 521 | /* See pg_rerank: one widening buffer per call, not per candidate. */ | ||
| 522 | 278 | Vec32Access input = vec32_access(s->type_info, dim, CurrentMemoryContext); | |
| 523 | |||
| 524 | 278 | uint32_t *order = palloc(count * sizeof(uint32_t)); | |
| 525 |
2/2✓ Branch 1 taken 11984 times.
✓ Branch 2 taken 278 times.
|
12262 | for (uint32_t i = 0; i < count; i++) |
| 526 | 11984 | order[i] = i; | |
| 527 | 278 | qsort_arg( | |
| 528 | order, count, sizeof(uint32_t), cmp_tid_order, (void *)candidates); | ||
| 529 | |||
| 530 | 278 | VsTopK topk; | |
| 531 | 278 | vs_topk_init(&topk, keep); | |
| 532 | |||
| 533 | 278 | RerankStreamState state = { | |
| 534 | .candidates = candidates, | ||
| 535 | .order = order, | ||
| 536 | .topk = &topk, | ||
| 537 | .storage = s, | ||
| 538 | .count = count, | ||
| 539 | .pos = 0, | ||
| 540 | }; | ||
| 541 | |||
| 542 | 556 | TupleTableSlot *slot = MakeSingleTupleTableSlot( | |
| 543 | 278 | RelationGetDescr(s->rel), &TTSOpsBufferHeapTuple); | |
| 544 | |||
| 545 | 278 | ReadStream *stream = read_stream_begin_relation( | |
| 546 | READ_STREAM_DEFAULT, | ||
| 547 | NULL, | ||
| 548 | s->rel, | ||
| 549 | MAIN_FORKNUM, | ||
| 550 | rerank_stream_cb, | ||
| 551 | &state, | ||
| 552 | sizeof(uint32_t)); | ||
| 553 | |||
| 554 | 278 | void *per_buffer_data; | |
| 555 | 278 | Buffer buf; | |
| 556 | |||
| 557 |
2/2✓ Branch 1 taken 11784 times.
✓ Branch 2 taken 278 times.
|
12340 | while (BufferIsValid( |
| 558 | 12062 | buf = read_stream_next_buffer(stream, &per_buffer_data))) | |
| 559 | { | ||
| 560 | 11784 | uint32_t idx = *(uint32_t *)per_buffer_data; | |
| 561 | 11784 | ItemPointerData tid = prism_posting_decode_tid(candidates[idx].id); | |
| 562 | |||
| 563 | 11784 | Page page = BufferGetPage(buf); | |
| 564 | 11784 | OffsetNumber off = ItemPointerGetOffsetNumber(&tid); | |
| 565 |
2/2✓ Branch 0 taken 11783 times.
✓ Branch 1 taken 1 times.
|
11784 | ItemId lp = PageGetItemId(page, off); |
| 566 | |||
| 567 | 11784 | Distance d = candidates[idx].distance; | |
| 568 |
2/2✓ Branch 0 taken 11783 times.
✓ Branch 1 taken 1 times.
|
11784 | if (ItemIdIsNormal(lp)) |
| 569 | { | ||
| 570 | 11783 | HeapTupleData tuple; | |
| 571 | 11783 | tuple.t_tableOid = RelationGetRelid(s->rel); | |
| 572 | 11783 | tuple.t_data = (HeapTupleHeader)PageGetItem(page, lp); | |
| 573 | 11783 | tuple.t_len = ItemIdGetLength(lp); | |
| 574 | 11783 | ItemPointerCopy(&tid, &tuple.t_self); | |
| 575 | |||
| 576 | 11783 | ExecStoreBufferHeapTuple(&tuple, slot, buf); | |
| 577 | |||
| 578 | 11783 | bool isnull; | |
| 579 | 11783 | Datum val = slot_getattr(slot, vec_attnum, &isnull); | |
| 580 |
1/2✓ Branch 0 taken 11783 times.
✗ Branch 1 not taken.
|
11783 | if (!isnull) |
| 581 | { | ||
| 582 | 11783 | Vec32Ref qref = {.data = query, .dim = dim}; | |
| 583 | 11783 | Vec32Ref vref = vec32_read(&input, val); | |
| 584 | 11783 | d = vs_distance(qref, vref, s->metric); | |
| 585 | } | ||
| 586 | 11783 | ExecClearTuple(slot); | |
| 587 | } | ||
| 588 | |||
| 589 | 11784 | vs_topk_insert_unique(&topk, d, 0.0f, (uint64_t)idx); | |
| 590 | 11784 | ReleaseBuffer(buf); | |
| 591 | } | ||
| 592 | |||
| 593 | 278 | read_stream_end(stream); | |
| 594 | 278 | ExecDropSingleTupleTableSlot(slot); | |
| 595 | |||
| 596 | 278 | VsTopKEntry *entries = palloc(topk.cand_count * sizeof(VsTopKEntry)); | |
| 597 | 278 | uint32_t nresults; | |
| 598 | 278 | vs_topk_extract_sorted_unique(&topk, entries, &nresults); | |
| 599 | |||
| 600 |
2/2✓ Branch 1 taken 8097 times.
✓ Branch 2 taken 278 times.
|
8375 | for (uint32_t i = 0; i < nresults; i++) |
| 601 | { | ||
| 602 | 8097 | out_indices[i] = (uint32_t)entries[i].id; | |
| 603 | 8097 | out_distances[i] = entries[i].distance; | |
| 604 | } | ||
| 605 | |||
| 606 | 278 | pfree(entries); | |
| 607 | 278 | vs_topk_cleanup(&topk); | |
| 608 | 278 | pfree(order); | |
| 609 | |||
| 610 | 278 | return nresults; | |
| 611 | } | ||
| 612 | |||
| 613 | /* ---------------------------------------------------------------- | ||
| 614 | * Static vtables | ||
| 615 | * ---------------------------------------------------------------- */ | ||
| 616 | |||
| 617 | static const VsStorageOps pg_storage_ops = { | ||
| 618 | .read_page = pg_read_page, | ||
| 619 | .release_page = pg_release_page, | ||
| 620 | .prefetch = pg_prefetch_page, | ||
| 621 | .write_page = pg_write_page, | ||
| 622 | .new_page = pg_new_page, | ||
| 623 | .commit_page = pg_commit_page, | ||
| 624 | .extend = pg_extend, | ||
| 625 | .rerank = pg_rerank, | ||
| 626 | }; | ||
| 627 | |||
| 628 | /* | ||
| 629 | * Read for a caller that sweeps the index rather than following a query's | ||
| 630 | * access pattern -- inspection. Identical to pg_read_page minus the | ||
| 631 | * recent-buffer cache, which for a sweep is worse than useless: its slots | ||
| 632 | * would fill with blocks no query asks for again, evicting the ones the | ||
| 633 | * query path reuses, and its hit/cold/stale counters would report the sweep | ||
| 634 | * instead of the workload they exist to measure. | ||
| 635 | * | ||
| 636 | * A separate entry in the ops table rather than a flag inside pg_read_page: | ||
| 637 | * the query path reads a page for every posting page it scans, and should | ||
| 638 | * not test a condition that only inspection can change. | ||
| 639 | */ | ||
| 640 | static Page | ||
| 641 | 26061 | pg_read_page_nocache(VsStorage *self, BlockNumber blkno) | |
| 642 | { | ||
| 643 | 26061 | VsPgStorage *s = PG_STORAGE(self); | |
| 644 | 26061 | Buffer buf = ReadBuffer(s->index, blkno); | |
| 645 | |||
| 646 | 26061 | LockBuffer(buf, BUFFER_LOCK_SHARE); | |
| 647 | 26061 | s->cur_buf = buf; | |
| 648 | 26061 | s->read_count++; | |
| 649 | |||
| 650 | 26061 | return BufferGetPage(buf); | |
| 651 | } | ||
| 652 | |||
| 653 | static const VsStorageOps pg_storage_readstream_ops = { | ||
| 654 | .read_page = pg_read_page, | ||
| 655 | .release_page = pg_release_page, | ||
| 656 | .prefetch = pg_prefetch_page, | ||
| 657 | .write_page = pg_write_page, | ||
| 658 | .new_page = pg_new_page, | ||
| 659 | .commit_page = pg_commit_page, | ||
| 660 | .extend = pg_extend, | ||
| 661 | .rerank = pg_rerank_readstream, | ||
| 662 | }; | ||
| 663 | |||
| 664 | /* Inspection: the standard ops with the cache-bypassing read. */ | ||
| 665 | static const VsStorageOps pg_storage_inspect_ops = { | ||
| 666 | .read_page = pg_read_page_nocache, | ||
| 667 | .release_page = pg_release_page, | ||
| 668 | .prefetch = pg_prefetch_page, | ||
| 669 | .write_page = pg_write_page, | ||
| 670 | .new_page = pg_new_page, | ||
| 671 | .commit_page = pg_commit_page, | ||
| 672 | .extend = pg_extend, | ||
| 673 | .rerank = NULL, | ||
| 674 | }; | ||
| 675 | |||
| 676 | void | ||
| 677 | 569 | vs_pg_storage_init_inspect(VsPgStorage *s, Relation index) | |
| 678 | { | ||
| 679 | 569 | vs_pg_storage_init(s, index, NULL, DISTANCE_L2); | |
| 680 | 569 | s->base.ops = &pg_storage_inspect_ops; | |
| 681 | 569 | } | |
| 682 | |||
| 683 | /* ---------------------------------------------------------------- | ||
| 684 | * Initialization | ||
| 685 | * ---------------------------------------------------------------- */ | ||
| 686 | |||
| 687 | void | ||
| 688 | 18958 | vs_pg_storage_init( | |
| 689 | VsPgStorage *s, Relation index, Relation rel, DistanceMetric metric) | ||
| 690 | { | ||
| 691 |
1/4✗ Branch 0 not taken.
✓ Branch 1 taken 18958 times.
✗ Branch 2 not taken.
✗ Branch 3 not taken.
|
18958 | if (rel != NULL && RelationGetForm(rel)->relam == HEAP_TABLE_AM_OID) |
| 692 | ✗ | s->base.ops = &pg_storage_readstream_ops; | |
| 693 | else | ||
| 694 | 18958 | s->base.ops = &pg_storage_ops; | |
| 695 | 18958 | s->index = index; | |
| 696 | 18958 | s->rel = rel; | |
| 697 | 18958 | s->build_mode = false; | |
| 698 | 18958 | s->cur_buf = InvalidBuffer; | |
| 699 | 18958 | s->metric = metric; | |
| 700 | 18958 | s->read_count = 0; | |
| 701 | /* | ||
| 702 | * Left NULL: every caller inits with rel = NULL, and rerank -- the only | ||
| 703 | * consumer -- returns early without a heap relation. It is resolved in | ||
| 704 | * vs_pg_storage_set_rel, which is where the heap arrives and therefore | ||
| 705 | * where rerank becomes possible. | ||
| 706 | */ | ||
| 707 | 18958 | s->type_info = NULL; | |
| 708 | 18958 | } | |
| 709 | |||
| 710 | void | ||
| 711 | 267 | vs_pg_storage_set_rel(VsPgStorage *s, Relation rel) | |
| 712 | { | ||
| 713 | 267 | s->rel = rel; | |
| 714 |
2/4✓ Branch 0 taken 267 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 267 times.
✗ Branch 3 not taken.
|
267 | if (rel != NULL && RelationGetForm(rel)->relam == HEAP_TABLE_AM_OID) |
| 715 | 267 | s->base.ops = &pg_storage_readstream_ops; | |
| 716 | |||
| 717 | /* | ||
| 718 | * Rerank reads the heap attribute, which has the indexed column's type. | ||
| 719 | * From the per-backend cache: this runs on the scan path, so the index | ||
| 720 | * has a metadata page and the cache is populated. | ||
| 721 | */ | ||
| 722 | 267 | if (rel != NULL) | |
| 723 | 267 | s->type_info = prism_cache_type_info(s->index); | |
| 724 | 267 | } | |
| 725 |