Compare commits

...

16 Commits

Author SHA1 Message Date
Andrew Kane
f1dd4e3b03 Updated urls [skip ci] 2026-06-29 01:08:14 -07:00
Andrew Kane
7d067d7b83 Updated changelog [skip ci] 2026-06-24 11:44:04 -07:00
Andrew Kane
d4dd73d970 Moved repair confirmation after lock wait to catch more potential issues 2026-06-23 11:39:16 -07:00
Bhagyesh Chaturvedi
ffe28bb954 Fix HNSW insert and vacuum race 2026-06-23 06:00:23 +00:00
Andrew Kane
0dbc1a27c0 Removed redundant check [skip ci] 2026-06-18 12:56:57 -07:00
Andrew Kane
08c4e7ff10 Hardened VectorArrayGet and VectorArraySet [skip ci] 2026-06-18 12:56:12 -07:00
Andrew Kane
6731c49811 Fixed overflow check (should never be hit) [skip ci] 2026-06-18 12:51:38 -07:00
Andrew Kane
f2617f02d1 Updated style to be consistent with latest Postgres [skip ci] 2026-06-18 12:44:11 -07:00
Andrew Kane
bdf19077db Ensure centers and samples fit into maintenance_work_mem before allocating for IVFFlat index builds - closes #986 2026-06-18 12:29:58 -07:00
Andrew Kane
90cd2b4ee5 Improved memory tracking for IVFFlat index builds [skip ci] 2026-06-18 11:57:03 -07:00
Andrew Kane
b44d1b4c5f Added todo [skip ci] 2026-06-18 11:51:37 -07:00
Andrew Kane
cc5b865c33 Added itemsize to IvfflatBuildState [skip ci] 2026-06-18 11:49:29 -07:00
Andrew Kane
4895021088 Hardened NeedsUpdated [skip ci] 2026-06-18 11:19:27 -07:00
Andrew Kane
eda77b3492 DRY normalize code for IVFFlat index builds 2026-06-18 11:16:51 -07:00
Andrew Kane
a2364b1793 Switched to VectorArraySet for NormCenters [skip ci] 2026-06-18 11:05:37 -07:00
Andrew Kane
a0eaf70d17 Hardened VectorArraySet [skip ci] 2026-06-18 11:03:09 -07:00
8 changed files with 89 additions and 64 deletions

View File

@@ -1,3 +1,7 @@
## 0.8.4 (unreleased)
- Fixed possible error with inserts during HNSW vacuuming
## 0.8.3 (2026-06-17) ## 0.8.3 (2026-06-17)
- Fixed possible index corruption with HNSW vacuuming - Fixed possible index corruption with HNSW vacuuming

View File

@@ -7,7 +7,7 @@
"Andrew Kane <andrew@ankane.org>" "Andrew Kane <andrew@ankane.org>"
], ],
"license": { "license": {
"PostgreSQL": "http://www.postgresql.org/about/licence" "PostgreSQL": "https://www.postgresql.org/about/licence"
}, },
"prereqs": { "prereqs": {
"runtime": { "runtime": {
@@ -38,7 +38,7 @@
"generated_by": "Andrew Kane", "generated_by": "Andrew Kane",
"meta-spec": { "meta-spec": {
"version": "1.0.0", "version": "1.0.0",
"url": "http://pgxn.org/meta/spec.txt" "url": "https://pgxn.org/meta/spec.txt"
}, },
"tags": [ "tags": [
"vectors", "vectors",

View File

@@ -719,7 +719,7 @@ InitBuildState(HnswBuildState * buildstate, Relation heap, Relation index, Index
/* Get support functions */ /* Get support functions */
HnswInitSupport(&buildstate->support, index); HnswInitSupport(&buildstate->support, index);
InitGraph(&buildstate->graphData, NULL, (Size) maintenance_work_mem * 1024L); InitGraph(&buildstate->graphData, NULL, maintenance_work_mem * (Size) 1024);
buildstate->graph = &buildstate->graphData; buildstate->graph = &buildstate->graphData;
buildstate->ml = HnswGetMl(buildstate->m); buildstate->ml = HnswGetMl(buildstate->m);
buildstate->maxLevel = HnswGetMaxLevel(buildstate->m); buildstate->maxLevel = HnswGetMaxLevel(buildstate->m);
@@ -956,7 +956,7 @@ HnswBeginParallel(HnswBuildState * buildstate, bool isconcurrent, int request)
/* Leave space for other objects in shared memory */ /* Leave space for other objects in shared memory */
/* Docker has a default limit of 64 MB for shm_size */ /* Docker has a default limit of 64 MB for shm_size */
/* which happens to be the default value of maintenance_work_mem */ /* which happens to be the default value of maintenance_work_mem */
esthnswarea = maintenance_work_mem * 1024L; esthnswarea = maintenance_work_mem * (Size) 1024;
estother = 3 * 1024 * 1024; estother = 3 * 1024 * 1024;
if (esthnswarea > estother) if (esthnswarea > estother)
esthnswarea -= estother; esthnswarea -= estother;

View File

@@ -202,13 +202,9 @@ NeedsUpdated(HnswVacuumState * vacuumstate, HnswElement element)
/* Also update if layer 0 is not full */ /* Also update if layer 0 is not full */
/* This could indicate too many candidates being deleted during insert */ /* This could indicate too many candidates being deleted during insert */
if (!needsUpdated) /* There should always be more than zero indextids, but check for safety */
{ if (!needsUpdated && ntup->count > 0)
/* Keep clang-tidy happy */
Assert(ntup->count > 0);
needsUpdated = !ItemPointerIsValid(&ntup->indextids[ntup->count - 1]); needsUpdated = !ItemPointerIsValid(&ntup->indextids[ntup->count - 1]);
}
UnlockReleaseBuffer(buf); UnlockReleaseBuffer(buf);
@@ -588,10 +584,15 @@ MarkDeleted(HnswVacuumState * vacuumstate)
BufferAccessStrategy bas = vacuumstate->bas; BufferAccessStrategy bas = vacuumstate->bas;
/* /*
* Wait for index scans to complete. Scans before this point may contain * Wait for inserts and index scans to complete. Inserts and scans before
* tuples about to be deleted. Scans after this point will not, since the * this point may visit tuples about to be deleted. Inserts and scans
* graph has been repaired. * after this point will not, since the graph has been repaired.
*/ */
LockPage(index, HNSW_UPDATE_LOCK, ExclusiveLock);
UnlockPage(index, HNSW_UPDATE_LOCK, ExclusiveLock);
ConfirmRepaired(vacuumstate);
LockPage(index, HNSW_SCAN_LOCK, ExclusiveLock); LockPage(index, HNSW_SCAN_LOCK, ExclusiveLock);
UnlockPage(index, HNSW_SCAN_LOCK, ExclusiveLock); UnlockPage(index, HNSW_SCAN_LOCK, ExclusiveLock);
@@ -771,10 +772,7 @@ hnswbulkdelete(IndexVacuumInfo *info, IndexBulkDeleteResult *stats,
/* Pass 2: Repair graph */ /* Pass 2: Repair graph */
HnswBench("RepairGraph", RepairGraph(&vacuumstate)); HnswBench("RepairGraph", RepairGraph(&vacuumstate));
/* Pass 3: Confirm repaired */ /* Passes 3 and 4: Confirm repaired and mark as deleted */
HnswBench("ConfirmRepaired", ConfirmRepaired(&vacuumstate));
/* Pass 4: Mark as deleted */
HnswBench("MarkDeleted", MarkDeleted(&vacuumstate)); HnswBench("MarkDeleted", MarkDeleted(&vacuumstate));
FreeVacuumState(&vacuumstate); FreeVacuumState(&vacuumstate);

View File

@@ -152,18 +152,7 @@ SampleRows(IvfflatBuildState * buildstate)
/* Normalize if needed */ /* Normalize if needed */
if (buildstate->kmeansnormprocinfo != NULL) if (buildstate->kmeansnormprocinfo != NULL)
{ IvfflatNormVectors(buildstate->typeInfo, buildstate->collation, buildstate->samples, buildstate->tmpCtx);
VectorArray samples = buildstate->samples;
for (int i = 0; i < samples->length; i++)
{
Datum value = PointerGetDatum(VectorArrayGet(samples, i));
Datum normValue = IvfflatNormValue(buildstate->typeInfo, buildstate->collation, value);
VectorArraySet(samples, i, DatumGetPointer(normValue));
pfree(DatumGetPointer(normValue));
}
}
} }
/* /*
@@ -399,8 +388,14 @@ InitBuildState(IvfflatBuildState * buildstate, Relation heap, Relation index, In
buildstate->slot = MakeSingleTupleTableSlot(buildstate->sortdesc, &TTSOpsVirtual); buildstate->slot = MakeSingleTupleTableSlot(buildstate->sortdesc, &TTSOpsVirtual);
/* TODO Ensure within maintenance_work_mem */ buildstate->memoryUsed = 0;
buildstate->centers = VectorArrayInit(buildstate->lists, buildstate->dimensions, buildstate->typeInfo->itemSize(buildstate->dimensions)); buildstate->itemsize = buildstate->typeInfo->itemSize(buildstate->dimensions);
buildstate->memoryUsed += VECTOR_ARRAY_SIZE(buildstate->lists, buildstate->itemsize);
IvfflatCheckMemoryUsage(buildstate->memoryUsed);
buildstate->centers = VectorArrayInit(buildstate->lists, buildstate->dimensions, buildstate->itemsize);
/* TODO Move allocation to page creation */
buildstate->listInfo = palloc(sizeof(ListInfo) * buildstate->lists); buildstate->listInfo = palloc(sizeof(ListInfo) * buildstate->lists);
buildstate->tmpCtx = AllocSetContextCreate(CurrentMemoryContext, buildstate->tmpCtx = AllocSetContextCreate(CurrentMemoryContext,
@@ -454,8 +449,9 @@ ComputeCenters(IvfflatBuildState * buildstate)
numSamples = 1; numSamples = 1;
/* Sample rows */ /* Sample rows */
/* TODO Ensure within maintenance_work_mem */ buildstate->memoryUsed += VECTOR_ARRAY_SIZE(numSamples, buildstate->itemsize);
buildstate->samples = VectorArrayInit(numSamples, buildstate->dimensions, buildstate->centers->itemsize); IvfflatCheckMemoryUsage(buildstate->memoryUsed);
buildstate->samples = VectorArrayInit(numSamples, buildstate->dimensions, buildstate->itemsize);
if (buildstate->heap != NULL) if (buildstate->heap != NULL)
{ {
IvfflatBench("sample rows", SampleRows(buildstate)); IvfflatBench("sample rows", SampleRows(buildstate));
@@ -470,7 +466,7 @@ ComputeCenters(IvfflatBuildState * buildstate)
} }
/* Calculate centers */ /* Calculate centers */
IvfflatBench("k-means", IvfflatKmeans(buildstate->index, buildstate->samples, buildstate->centers, buildstate->typeInfo)); IvfflatBench("k-means", IvfflatKmeans(buildstate->index, buildstate->samples, buildstate->centers, buildstate->typeInfo, buildstate->memoryUsed));
/* Free samples before we allocate more memory */ /* Free samples before we allocate more memory */
VectorArrayFree(buildstate->samples); VectorArrayFree(buildstate->samples);

View File

@@ -204,6 +204,7 @@ typedef struct IvfflatBuildState
VectorArray samples; VectorArray samples;
VectorArray centers; VectorArray centers;
ListInfo *listInfo; ListInfo *listInfo;
Size itemsize;
#ifdef IVFFLAT_KMEANS_DEBUG #ifdef IVFFLAT_KMEANS_DEBUG
double inertia; double inertia;
@@ -223,6 +224,7 @@ typedef struct IvfflatBuildState
TupleTableSlot *slot; TupleTableSlot *slot;
/* Memory */ /* Memory */
Size memoryUsed;
MemoryContext tmpCtx; MemoryContext tmpCtx;
/* Parallel builds */ /* Parallel builds */
@@ -303,22 +305,32 @@ typedef IvfflatScanOpaqueData * IvfflatScanOpaque;
static inline Pointer static inline Pointer
VectorArrayGet(VectorArray arr, int offset) VectorArrayGet(VectorArray arr, int offset)
{ {
if (offset >= arr->maxlen)
elog(ERROR, "safety check failed");
return ((char *) arr->items) + (offset * arr->itemsize); return ((char *) arr->items) + (offset * arr->itemsize);
} }
static inline void static inline void
VectorArraySet(VectorArray arr, int offset, Pointer val) VectorArraySet(VectorArray arr, int offset, Pointer val)
{ {
memcpy(VectorArrayGet(arr, offset), val, VARSIZE_ANY(val)); Size size = VARSIZE_ANY(val);
if (size > arr->itemsize)
elog(ERROR, "safety check failed");
memcpy(VectorArrayGet(arr, offset), val, size);
} }
/* Methods */ /* Methods */
VectorArray VectorArrayInit(int maxlen, int dimensions, Size itemsize); VectorArray VectorArrayInit(int maxlen, int dimensions, Size itemsize);
void VectorArrayFree(VectorArray arr); void VectorArrayFree(VectorArray arr);
void IvfflatKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo); void IvfflatKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo, Size memoryUsed);
FmgrInfo *IvfflatOptionalProcInfo(Relation index, uint16 procnum); FmgrInfo *IvfflatOptionalProcInfo(Relation index, uint16 procnum);
Datum IvfflatNormValue(const IvfflatTypeInfo * typeInfo, Oid collation, Datum value); Datum IvfflatNormValue(const IvfflatTypeInfo * typeInfo, Oid collation, Datum value);
bool IvfflatCheckNorm(FmgrInfo *procinfo, Oid collation, Datum value); bool IvfflatCheckNorm(FmgrInfo *procinfo, Oid collation, Datum value);
void IvfflatNormVectors(const IvfflatTypeInfo * typeInfo, Oid collation, VectorArray arr, MemoryContext tmpCtx);
void IvfflatCheckMemoryUsage(Size totalSize);
int IvfflatGetLists(Relation index); int IvfflatGetLists(Relation index);
void IvfflatGetMetaPageInfo(Relation index, int *lists, int *dimensions); void IvfflatGetMetaPageInfo(Relation index, int *lists, int *dimensions);
void IvfflatUpdateList(Relation index, ListInfo listInfo, BlockNumber insertPage, BlockNumber originalInsertPage, BlockNumber startPage, ForkNumber forkNum); void IvfflatUpdateList(Relation index, ListInfo listInfo, BlockNumber insertPage, BlockNumber originalInsertPage, BlockNumber startPage, ForkNumber forkNum);

View File

@@ -99,22 +99,8 @@ NormCenters(const IvfflatTypeInfo * typeInfo, Oid collation, VectorArray centers
MemoryContext normCtx = AllocSetContextCreate(CurrentMemoryContext, MemoryContext normCtx = AllocSetContextCreate(CurrentMemoryContext,
"Ivfflat norm temporary context", "Ivfflat norm temporary context",
ALLOCSET_DEFAULT_SIZES); ALLOCSET_DEFAULT_SIZES);
MemoryContext oldCtx = MemoryContextSwitchTo(normCtx);
for (int j = 0; j < centers->length; j++) IvfflatNormVectors(typeInfo, collation, centers, normCtx);
{
Datum center = PointerGetDatum(VectorArrayGet(centers, j));
Datum newCenter = IvfflatNormValue(typeInfo, collation, center);
Size size = VARSIZE_ANY(DatumGetPointer(newCenter));
if (size > centers->itemsize)
elog(ERROR, "safety check failed");
memcpy(DatumGetPointer(center), DatumGetPointer(newCenter), size);
MemoryContextReset(normCtx);
}
MemoryContextSwitchTo(oldCtx);
MemoryContextDelete(normCtx); MemoryContextDelete(normCtx);
} }
@@ -258,7 +244,7 @@ ComputeNewCenters(VectorArray samples, float *agg, VectorArray newCenters, int *
* https://www.aaai.org/Papers/ICML/2003/ICML03-022.pdf * https://www.aaai.org/Papers/ICML/2003/ICML03-022.pdf
*/ */
static void static void
ElkanKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo) ElkanKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo, Size memoryUsed)
{ {
FmgrInfo *procinfo; FmgrInfo *procinfo;
FmgrInfo *normprocinfo; FmgrInfo *normprocinfo;
@@ -277,8 +263,6 @@ ElkanKmeans(Relation index, VectorArray samples, VectorArray centers, const Ivff
float *newcdist; float *newcdist;
/* Calculate allocation sizes */ /* Calculate allocation sizes */
Size samplesSize = VECTOR_ARRAY_SIZE(samples->maxlen, samples->itemsize);
Size centersSize = VECTOR_ARRAY_SIZE(centers->maxlen, centers->itemsize);
Size newCentersSize = VECTOR_ARRAY_SIZE(numCenters, centers->itemsize); Size newCentersSize = VECTOR_ARRAY_SIZE(numCenters, centers->itemsize);
Size aggSize = sizeof(float) * (int64) numCenters * dimensions; Size aggSize = sizeof(float) * (int64) numCenters * dimensions;
Size centerCountsSize = sizeof(int) * numCenters; Size centerCountsSize = sizeof(int) * numCenters;
@@ -290,18 +274,13 @@ ElkanKmeans(Relation index, VectorArray samples, VectorArray centers, const Ivff
Size newcdistSize = sizeof(float) * numCenters; Size newcdistSize = sizeof(float) * numCenters;
/* Calculate total size */ /* Calculate total size */
Size totalSize = samplesSize + centersSize + newCentersSize + aggSize + centerCountsSize + closestCentersSize + lowerBoundSize + upperBoundSize + sSize + halfcdistSize + newcdistSize; Size totalSize = memoryUsed + newCentersSize + aggSize + centerCountsSize + closestCentersSize + lowerBoundSize + upperBoundSize + sSize + halfcdistSize + newcdistSize;
/* Check memory requirements */ /* Check memory requirements */
/* Add one to error message to ceil */ IvfflatCheckMemoryUsage(totalSize);
if (totalSize > (Size) maintenance_work_mem * 1024L)
ereport(ERROR,
(errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED),
errmsg("memory required is %zu MB, maintenance_work_mem is %d MB",
totalSize / (1024 * 1024) + 1, maintenance_work_mem / 1024)));
/* Ensure indexing does not overflow */ /* Ensure indexing does not overflow */
if (numCenters * numCenters > INT_MAX) if (numCenters > INT_MAX / numCenters)
elog(ERROR, "Indexing overflow detected. Please report a bug."); elog(ERROR, "Indexing overflow detected. Please report a bug.");
/* Set support functions */ /* Set support functions */
@@ -562,7 +541,7 @@ CheckCenters(Relation index, VectorArray centers, const IvfflatTypeInfo * typeIn
* We use spherical k-means for inner product and cosine * We use spherical k-means for inner product and cosine
*/ */
void void
IvfflatKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo) IvfflatKmeans(Relation index, VectorArray samples, VectorArray centers, const IvfflatTypeInfo * typeInfo, Size memoryUsed)
{ {
MemoryContext kmeansCtx = AllocSetContextCreate(CurrentMemoryContext, MemoryContext kmeansCtx = AllocSetContextCreate(CurrentMemoryContext,
"Ivfflat kmeans temporary context", "Ivfflat kmeans temporary context",
@@ -572,7 +551,7 @@ IvfflatKmeans(Relation index, VectorArray samples, VectorArray centers, const Iv
if (samples->length == 0) if (samples->length == 0)
RandomCenters(index, centers, typeInfo); RandomCenters(index, centers, typeInfo);
else else
ElkanKmeans(index, samples, centers, typeInfo); ElkanKmeans(index, samples, centers, typeInfo, memoryUsed);
CheckCenters(index, centers, typeInfo); CheckCenters(index, centers, typeInfo);

View File

@@ -6,7 +6,9 @@
#include "halfutils.h" #include "halfutils.h"
#include "halfvec.h" #include "halfvec.h"
#include "ivfflat.h" #include "ivfflat.h"
#include "miscadmin.h"
#include "storage/bufmgr.h" #include "storage/bufmgr.h"
#include "utils/memutils.h"
#include "utils/relcache.h" #include "utils/relcache.h"
#include "utils/varbit.h" #include "utils/varbit.h"
#include "vector.h" #include "vector.h"
@@ -88,6 +90,40 @@ IvfflatCheckNorm(FmgrInfo *procinfo, Oid collation, Datum value)
return DatumGetFloat8(FunctionCall1Coll(procinfo, collation, value)) > 0; return DatumGetFloat8(FunctionCall1Coll(procinfo, collation, value)) > 0;
} }
/*
* Normalize vectors
*/
void
IvfflatNormVectors(const IvfflatTypeInfo * typeInfo, Oid collation, VectorArray arr, MemoryContext tmpCtx)
{
MemoryContext oldCtx = MemoryContextSwitchTo(tmpCtx);
for (int i = 0; i < arr->length; i++)
{
Datum value = PointerGetDatum(VectorArrayGet(arr, i));
Datum newValue = IvfflatNormValue(typeInfo, collation, value);
VectorArraySet(arr, i, DatumGetPointer(newValue));
MemoryContextReset(tmpCtx);
}
MemoryContextSwitchTo(oldCtx);
}
/*
* Check memory usage
*/
void
IvfflatCheckMemoryUsage(Size totalSize)
{
/* Add one to error message to ceil */
if (totalSize > maintenance_work_mem * (Size) 1024)
ereport(ERROR,
(errcode(ERRCODE_PROGRAM_LIMIT_EXCEEDED),
errmsg("memory required is %zu MB, maintenance_work_mem is %d MB",
totalSize / (1024 * 1024) + 1, maintenance_work_mem / 1024)));
}
/* /*
* New buffer * New buffer
*/ */