From cb7292c0689d58c528a7c6e0866f907208b25612 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 5 Feb 2026 12:21:01 +0000 Subject: [PATCH 01/16] NewCudaStream and fix memory leak --- go/dlpack.go | 36 ++++++++---------------------------- go/resources.go | 7 +++++++ 2 files changed, 15 insertions(+), 28 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index 6fe619fd35..36e01787e7 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -33,19 +33,10 @@ func NewTensor[T TensorNumberType](data [][]T) (Tensor[T], error) { totalElements := len(data) * len(data[0]) dataPtr := C.malloc(C.size_t(totalElements * int(unsafe.Sizeof(T(0))))) - if dataPtr == nil { - return Tensor[T]{}, errors.New("data memory allocation failed") - } - dataSlice := unsafe.Slice((*T)(dataPtr), totalElements) flattenData(data, dataSlice) shapePtr := C.malloc(C.size_t(2 * int(unsafe.Sizeof(C.int64_t(0))))) - if shapePtr == nil { - C.free(dataPtr) - return Tensor[T]{}, errors.New("shape memory allocation failed") - } - shapeSlice := unsafe.Slice((*C.int64_t)(shapePtr), 2) shapeSlice[0] = C.int64_t(len(data)) shapeSlice[1] = C.int64_t(len(data[0])) @@ -85,18 +76,10 @@ func NewVector[T TensorNumberType](data []T) (Tensor[T], error) { totalElements := len(data) dataPtr := C.malloc(C.size_t(totalElements * int(unsafe.Sizeof(T(0))))) - if dataPtr == nil { - return Tensor[T]{}, errors.New("data memory allocation failed") - } - dataSlice := unsafe.Slice((*T)(dataPtr), totalElements) copy(dataSlice, data) shapePtr := C.malloc(C.size_t(int(unsafe.Sizeof(C.int64_t(0))))) - if shapePtr == nil { - C.free(dataPtr) - return Tensor[T]{}, errors.New("shape memory allocation failed") - } shapeSlice := unsafe.Slice((*C.int64_t)(shapePtr), 1) shapeSlice[0] = C.int64_t(len(data)) @@ -133,19 +116,12 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor } shapePtr := C.malloc(C.size_t(len(shape) * int(unsafe.Sizeof(C.int64_t(0))))) - if shapePtr == nil { - return Tensor[T]{}, errors.New("shape memory allocation failed") - } - shapeSlice := unsafe.Slice((*C.int64_t)(shapePtr), len(shape)) for i, dim := range shape { shapeSlice[i] = C.int64_t(dim) } dlm := (*C.DLManagedTensor)(C.malloc(C.size_t(unsafe.Sizeof(C.DLManagedTensor{})))) - if dlm == nil { - return Tensor[T]{}, errors.New("tensor allocation failed") - } dtype := getDLDataType[T]() var deviceDataPtr unsafe.Pointer @@ -233,6 +209,12 @@ func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { C.cuvsRMMFree(res.Resource, DeviceDataPointer, C.size_t(bytes)) return nil, err } + + if t.C_tensor.dl_tensor.data != nil { + C.free(t.C_tensor.dl_tensor.data) + t.C_tensor.dl_tensor.data = nil + } + t.C_tensor.dl_tensor.device.device_type = C.kDLCUDA t.C_tensor.dl_tensor.data = DeviceDataPointer @@ -325,10 +307,6 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { bytes := t.sizeInBytes() addr := (C.malloc(C.size_t(bytes))) - if addr == nil { - return nil, errors.New("memory allocation failed") - } - err := CheckCuda( C.cudaMemcpy( addr, @@ -337,12 +315,14 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { C.cudaMemcpyDeviceToHost, )) if err != nil { + C.free(addr) return nil, err } err = CheckCuvs(CuvsError( C.cuvsRMMFree(res.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) if err != nil { + C.free(addr) return nil, err } diff --git a/go/resources.go b/go/resources.go index 562aa7b7fc..576e13a3b4 100644 --- a/go/resources.go +++ b/go/resources.go @@ -1,6 +1,7 @@ package cuvs // #include +// #include import "C" type cuvsResource C.cuvsResources_t @@ -12,6 +13,12 @@ type Resource struct { Resource C.cuvsResources_t } +func NewCudaStream() C.cudaStream_t { + var stream C.cudaStream_t + C.cudaStreamCreate(&stream) + return stream +} + // Returns a new Resource object func NewResource(stream C.cudaStream_t) (Resource, error) { res := C.cuvsResources_t(0) From 53da6f5b6434fdce377b7af59e3e62a6d8540771 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 5 Feb 2026 12:25:05 +0000 Subject: [PATCH 02/16] remove unnecessary nil check after C.malloc --- go/dlpack.go | 4 ---- 1 file changed, 4 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index 36e01787e7..bc10d7a03d 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -43,10 +43,6 @@ func NewTensor[T TensorNumberType](data [][]T) (Tensor[T], error) { // Create DLManagedTensor dlm := (*C.DLManagedTensor)(C.malloc(C.size_t(unsafe.Sizeof(C.DLManagedTensor{})))) - if dlm == nil { - return Tensor[T]{}, errors.New("tensor allocation failed") - } - dlm.dl_tensor.data = dataPtr dlm.dl_tensor.device = C.DLDevice{ device_type: C.DLDeviceType(C.kDLCPU), From 5bf83a13064d4ccb9987f549f4ed560c8dd28d2e Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 5 Feb 2026 13:05:28 +0000 Subject: [PATCH 03/16] cuvsRMMFree with same resource --- go/dlpack.go | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index bc10d7a03d..e8866cf7e0 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -20,6 +20,7 @@ type TensorNumberType interface { type Tensor[T any] struct { C_tensor *C.DLManagedTensor shape []int64 + resource *Resource } // Creates a new Tensor on the host and copies the data into it. @@ -148,6 +149,7 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor return Tensor[T]{ C_tensor: dlm, shape: shapeCopy, + resource: res, }, nil } @@ -155,11 +157,10 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor func (t *Tensor[T]) Close() error { if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { bytes := t.sizeInBytes() - res, err := NewResource(nil) - if err != nil { - return err + if t.resource != nil { + return errors.New("resource not found") } - err = CheckCuvs(CuvsError(C.cuvsRMMFree(res.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) + err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) return err } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { @@ -213,6 +214,7 @@ func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { t.C_tensor.dl_tensor.device.device_type = C.kDLCUDA t.C_tensor.dl_tensor.data = DeviceDataPointer + t.resource = res return t, nil } @@ -294,6 +296,7 @@ func (t *Tensor[T]) Expand(res *Resource, newData [][]T) (*Tensor[T], error) { t.C_tensor.dl_tensor.data = NewDeviceDataPointer t.C_tensor.dl_tensor.shape = (*C.int64_t)(unsafe.Pointer(&shape[0])) + t.resource = res return t, nil } @@ -324,6 +327,7 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { t.C_tensor.dl_tensor.device.device_type = C.kDLCPU t.C_tensor.dl_tensor.data = addr + t.resource = res return t, nil } From 779b3cfa6c4e607b697cda8af134567ef4a50aaf Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 5 Feb 2026 15:12:11 +0000 Subject: [PATCH 04/16] fix memory leak in Close --- go/dlpack.go | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index e8866cf7e0..b1a4f9f5a7 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -157,12 +157,14 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor func (t *Tensor[T]) Close() error { if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { bytes := t.sizeInBytes() - if t.resource != nil { - return errors.New("resource not found") + if t.resource == nil { + return errors.New("resource not found") } err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) - - return err + if err != nil { + return err + } + t.C_tensor.dl_tensor.data = nil } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { if t.C_tensor.dl_tensor.data != nil { C.free(t.C_tensor.dl_tensor.data) From 4791d8d051b0e4f14c6067ea709327aac6c98343 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 5 Feb 2026 18:03:01 +0000 Subject: [PATCH 05/16] gemini --- go/brute_force/brute_force.go | 13 +++++- go/brute_force/brute_force_test.go | 13 +++++- go/dlpack.go | 66 ++++++++++++++++-------------- go/dlpack_test.go | 13 +++++- go/ivf_flat/ivf_flat.go | 13 ++++-- go/ivf_flat/ivf_flat_test.go | 13 +++++- go/resources.go | 28 ++++++++++--- 7 files changed, 112 insertions(+), 47 deletions(-) diff --git a/go/brute_force/brute_force.go b/go/brute_force/brute_force.go index c75447b99b..9b1d0f1a21 100644 --- a/go/brute_force/brute_force.go +++ b/go/brute_force/brute_force.go @@ -6,6 +6,7 @@ import "C" import ( "errors" "unsafe" + "runtime" cuvs "github.com/rapidsai/cuvs/go" ) @@ -58,6 +59,9 @@ func BuildIndex[T any](Resources cuvs.Resource, Dataset *cuvs.Tensor[T], metric } index.trained = true + runtime.KeepAlive(index) + runtime.KeepAlive(Dataset) + runtime.KeepAlive(Resources) return nil } @@ -69,7 +73,7 @@ func BuildIndex[T any](Resources cuvs.Resource, Dataset *cuvs.Tensor[T], metric // * `queries` - Tensor in device memory to query for // * `neighbors` - Tensor in device memory that receives the indices of the nearest neighbors // * `distances` - Tensor in device memory that receives the distances of the nearest neighbors -func SearchIndex[T any](resources cuvs.Resource, index BruteForceIndex, queries *cuvs.Tensor[T], neighbors *cuvs.Tensor[int64], distances *cuvs.Tensor[float32]) error { +func SearchIndex[T any](resources cuvs.Resource, index *BruteForceIndex, queries *cuvs.Tensor[T], neighbors *cuvs.Tensor[int64], distances *cuvs.Tensor[float32]) error { if !index.trained { return errors.New("index needs to be built before calling search") } @@ -79,7 +83,12 @@ func SearchIndex[T any](resources cuvs.Resource, index BruteForceIndex, queries _type: C.NO_FILTER, } - err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsBruteForceSearch(C.ulong(resources.Resource), index.index, (*C.DLManagedTensor)(unsafe.Pointer(queries.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(neighbors.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), prefilter))) + err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsBruteForceSearch((C.cuvsResources_t)(resources.Resource), index.index, (*C.DLManagedTensor)(unsafe.Pointer(queries.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(neighbors.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), prefilter))) + runtime.KeepAlive(index) + runtime.KeepAlive(queries) + runtime.KeepAlive(neighbors) + runtime.KeepAlive(distances) + runtime.KeepAlive(resources) return err } diff --git a/go/brute_force/brute_force_test.go b/go/brute_force/brute_force_test.go index ba9ac898f5..33d62bcc39 100644 --- a/go/brute_force/brute_force_test.go +++ b/go/brute_force/brute_force_test.go @@ -16,7 +16,16 @@ func TestBruteForce(t *testing.T) { epsilon = 0.001 ) - resource, _ := cuvs.NewResource(nil) + cudaStream, err := cuvs.NewCudaStream() + if err != nil { + t.Fatal(err) + } + defer cudaStream.Close() + + resource, err := cuvs.NewResource(cudaStream) + if err != nil { + t.Fatal(err) + } defer resource.Close() testDataset := make([][]float32, nDataPoints) @@ -75,7 +84,7 @@ func TestBruteForce(t *testing.T) { t.Fatalf("error moving queries to device: %v", err) } - err = SearchIndex(resource, *index, &queries, &neighbors, &distances) + err = SearchIndex(resource, index, &queries, &neighbors, &distances) if err != nil { t.Fatalf("error searching index: %v", err) } diff --git a/go/dlpack.go b/go/dlpack.go index b1a4f9f5a7..21155298ef 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -155,34 +155,34 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor // Destroys Tensor, freeing the memory it was allocated on. func (t *Tensor[T]) Close() error { - if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { - bytes := t.sizeInBytes() - if t.resource == nil { - return errors.New("resource not found") - } - err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) - if err != nil { - return err - } - t.C_tensor.dl_tensor.data = nil - } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { - if t.C_tensor.dl_tensor.data != nil { - C.free(t.C_tensor.dl_tensor.data) + if t.C_tensor != nil { + if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { + bytes := t.sizeInBytes() + if t.resource == nil { + // We cannot free device memory without the resource, + // so we must return an error here. + return errors.New("resource not found for CUDA tensor") + } + err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) + if err != nil { + return err + } t.C_tensor.dl_tensor.data = nil + } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { + if t.C_tensor.dl_tensor.data != nil { + C.free(t.C_tensor.dl_tensor.data) + t.C_tensor.dl_tensor.data = nil + } } - } - if t.C_tensor.dl_tensor.shape != nil { - C.free(unsafe.Pointer(t.C_tensor.dl_tensor.shape)) - t.C_tensor.dl_tensor.shape = nil - } + if t.C_tensor.dl_tensor.shape != nil { + C.free(unsafe.Pointer(t.C_tensor.dl_tensor.shape)) + t.C_tensor.dl_tensor.shape = nil + } - if t.C_tensor != nil { C.free(unsafe.Pointer(t.C_tensor)) t.C_tensor = nil } - - t.C_tensor = nil return nil } @@ -289,17 +289,22 @@ func (t *Tensor[T]) Expand(res *Resource, newData [][]T) (*Tensor[T], error) { return nil, err } - shape := make([]int64, 2) - shape[0] = int64(*t.C_tensor.dl_tensor.shape) + int64(len(newData)) + oldShapePtr := t.C_tensor.dl_tensor.shape + newShapePtr := C.malloc(C.size_t(2 * int(unsafe.Sizeof(C.int64_t(0))))) + newShapeSlice := unsafe.Slice((*C.int64_t)(newShapePtr), 2) + newShapeSlice[0] = C.int64_t(int64(*t.C_tensor.dl_tensor.shape) + int64(len(newData))) + newShapeSlice[1] = C.int64_t(newShape[1]) - shape[1] = newShape[1] - - t.shape = shape + t.shape = []int64{int64(newShapeSlice[0]), int64(newShapeSlice[1])} t.C_tensor.dl_tensor.data = NewDeviceDataPointer - t.C_tensor.dl_tensor.shape = (*C.int64_t)(unsafe.Pointer(&shape[0])) + t.C_tensor.dl_tensor.shape = (*C.int64_t)(newShapePtr) t.resource = res + if oldShapePtr != nil { + C.free(unsafe.Pointer(oldShapePtr)) + } + return t, nil } @@ -329,7 +334,7 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { t.C_tensor.dl_tensor.device.device_type = C.kDLCPU t.C_tensor.dl_tensor.data = addr - t.resource = res + t.resource = nil return t, nil } @@ -392,11 +397,12 @@ func (t *Tensor[T]) sizeInBytes() int64 { func calculateBytes(shape []int64, dtype C.DLDataType) int64 { bytes := int64(1) - for dim := range shape { - bytes *= (shape[dim]) + for _, dim := range shape { + bytes *= dim } bytes *= int64(dtype.bits) / 8 + return bytes } diff --git a/go/dlpack_test.go b/go/dlpack_test.go index 135363f811..7caf303606 100644 --- a/go/dlpack_test.go +++ b/go/dlpack_test.go @@ -10,7 +10,13 @@ import ( ) func TestDlPack(t *testing.T) { - resource, _ := cuvs.NewResource(nil) + cudaStream, err := cuvs.NewCudaStream() + if err != nil { + t.Fatal(err) + } + defer cudaStream.Close() + + resource, err := cuvs.NewResource(cudaStream) rand.Seed(time.Now().UnixNano()) NDataPoints := 256 NFeatures := 16 @@ -106,10 +112,13 @@ func TestEmptyTensor(t *testing.T) { } func TestDeviceOperations(t *testing.T) { - resource, err := cuvs.NewResource(nil) + cudaStream, err := cuvs.NewCudaStream() if err != nil { t.Fatal(err) } + defer cudaStream.Close() + + resource, err := cuvs.NewResource(cudaStream) // Create test data data := make([][]float32, 10) diff --git a/go/ivf_flat/ivf_flat.go b/go/ivf_flat/ivf_flat.go index 79000978ce..1638113ff0 100644 --- a/go/ivf_flat/ivf_flat.go +++ b/go/ivf_flat/ivf_flat.go @@ -6,6 +6,7 @@ import "C" import ( "errors" "unsafe" + "runtime" cuvs "github.com/rapidsai/cuvs/go" ) @@ -17,14 +18,13 @@ type IvfFlatIndex struct { } // Creates a new empty IvfFlatIndex -func CreateIndex[T any](params *IndexParams, dataset *cuvs.Tensor[T]) (*IvfFlatIndex, error) { +func CreateIndex[T any](params *IndexParams) (*IvfFlatIndex, error) { var index C.cuvsIvfFlatIndex_t err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsIvfFlatIndexCreate(&index))) if err != nil { return nil, err } - - return &IvfFlatIndex{index: index}, nil + return &IvfFlatIndex{index: index, trained: false}, nil } // Builds an IvfFlatIndex from the dataset for efficient search. @@ -41,6 +41,11 @@ func BuildIndex[T any](Resources cuvs.Resource, params *IndexParams, dataset *cu return err } index.trained = true + runtime.KeepAlive(Resources) + runtime.KeepAlive(params) + runtime.KeepAlive(dataset) + runtime.KeepAlive(index) + return nil } @@ -63,7 +68,7 @@ func (index *IvfFlatIndex) Close() error { // * `queries` - A tensor in device memory to query for // * `neighbors` - Tensor in device memory that receives the indices of the nearest neighbors // * `distances` - Tensor in device memory that receives the distances of the nearest neighbors -func SearchIndex[T any](Resources cuvs.Resource, params *SearchParams, index *IvfFlatIndex, queries *cuvs.Tensor[T], neighbors *cuvs.Tensor[int64], distances *cuvs.Tensor[T]) error { +func SearchIndex[T any](Resources cuvs.Resource, params *SearchParams, index *IvfFlatIndex, queries *cuvs.Tensor[T], neighbors *cuvs.Tensor[int64], distances *cuvs.Tensor[float32]) error { if !index.trained { return errors.New("index needs to be built before calling search") } diff --git a/go/ivf_flat/ivf_flat_test.go b/go/ivf_flat/ivf_flat_test.go index 92ffdc715f..78d603a8e9 100644 --- a/go/ivf_flat/ivf_flat_test.go +++ b/go/ivf_flat/ivf_flat_test.go @@ -17,7 +17,16 @@ func TestIvfFlat(t *testing.T) { nList = 512 ) - resource, _ := cuvs.NewResource(nil) + cudaStream, err := cuvs.NewCudaStream() + if err != nil { + t.Fatal(err) + } + defer cudaStream.Close() + + resource, err := cuvs.NewResource(cudaStream) + if err != nil { + t.Fatal(err) + } defer resource.Close() testDataset := make([][]float32, nDataPoints) @@ -42,7 +51,7 @@ func TestIvfFlat(t *testing.T) { indexParams.SetNLists(nList) - index, _ := CreateIndex(indexParams, &dataset) + index, _ := CreateIndex[float32](indexParams) defer index.Close() // use the first 4 points from the dataset as queries : will test that we get them back diff --git a/go/resources.go b/go/resources.go index 576e13a3b4..4c7a348693 100644 --- a/go/resources.go +++ b/go/resources.go @@ -13,14 +13,31 @@ type Resource struct { Resource C.cuvsResources_t } -func NewCudaStream() C.cudaStream_t { +type CudaStream struct { + stream C.cudaStream_t +} + +func (s *CudaStream) Close() error { + err := CheckCuda(C.cudaStreamDestroy(s.stream)) + if err != nil { + return err + } + s.stream = nil + return nil +} + +// Creates a new CUDA stream +func NewCudaStream() (*CudaStream, error) { var stream C.cudaStream_t - C.cudaStreamCreate(&stream) - return stream + err := CheckCuda(C.cudaStreamCreate(&stream)) + if err != nil { + return nil, err + } + return &CudaStream{stream: stream}, nil } // Returns a new Resource object -func NewResource(stream C.cudaStream_t) (Resource, error) { +func NewResource(stream *CudaStream) (Resource, error) { res := C.cuvsResources_t(0) err := CheckCuvs(CuvsError(C.cuvsResourcesCreate(&res))) if err != nil { @@ -28,8 +45,9 @@ func NewResource(stream C.cudaStream_t) (Resource, error) { } if stream != nil { - err := CheckCuvs(CuvsError(C.cuvsStreamSet(res, stream))) + err := CheckCuvs(CuvsError(C.cuvsStreamSet(res, stream.stream))) if err != nil { + C.cuvsResourcesDestroy(res) // Clean up the resource created return Resource{}, err } } From a8cd66ff89321f6aadb042c3389f6f469660c0ad Mon Sep 17 00:00:00 2001 From: Eric Date: Fri, 6 Feb 2026 09:21:22 +0000 Subject: [PATCH 06/16] keepAlive --- go/dlpack.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/go/dlpack.go b/go/dlpack.go index 21155298ef..15a87e8f65 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -7,6 +7,7 @@ import "C" import ( "errors" + "runtime" "strconv" "unsafe" ) @@ -182,6 +183,7 @@ func (t *Tensor[T]) Close() error { C.free(unsafe.Pointer(t.C_tensor)) t.C_tensor = nil + runtime.KeepAlive(t.resource) } return nil } @@ -218,6 +220,7 @@ func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { t.C_tensor.dl_tensor.data = DeviceDataPointer t.resource = res + runtime.KeepAlive(res) return t, nil } @@ -305,6 +308,7 @@ func (t *Tensor[T]) Expand(res *Resource, newData [][]T) (*Tensor[T], error) { C.free(unsafe.Pointer(oldShapePtr)) } + runtime.KeepAlive(res) return t, nil } @@ -332,6 +336,7 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { return nil, err } + runtime.KeepAlive(res) t.C_tensor.dl_tensor.device.device_type = C.kDLCPU t.C_tensor.dl_tensor.data = addr t.resource = nil From 9672d80f408154d6f12494775ecf120f65cb99d4 Mon Sep 17 00:00:00 2001 From: Eric Date: Fri, 6 Feb 2026 09:28:36 +0000 Subject: [PATCH 07/16] keep alive --- go/brute_force/brute_force.go | 3 ++- go/ivf_flat/ivf_flat.go | 16 +++++++++++++--- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/go/brute_force/brute_force.go b/go/brute_force/brute_force.go index 9b1d0f1a21..30e80f85b4 100644 --- a/go/brute_force/brute_force.go +++ b/go/brute_force/brute_force.go @@ -5,8 +5,8 @@ import "C" import ( "errors" - "unsafe" "runtime" + "unsafe" cuvs "github.com/rapidsai/cuvs/go" ) @@ -35,6 +35,7 @@ func (index *BruteForceIndex) Close() error { if err != nil { return err } + runtime.KeepAlive(index) return nil } diff --git a/go/ivf_flat/ivf_flat.go b/go/ivf_flat/ivf_flat.go index 1638113ff0..1ba1988d56 100644 --- a/go/ivf_flat/ivf_flat.go +++ b/go/ivf_flat/ivf_flat.go @@ -5,8 +5,8 @@ import "C" import ( "errors" - "unsafe" "runtime" + "unsafe" cuvs "github.com/rapidsai/cuvs/go" ) @@ -77,6 +77,14 @@ func SearchIndex[T any](Resources cuvs.Resource, params *SearchParams, index *Iv _type: C.NO_FILTER, } + defer func() { + runtime.KeepAlive(Resources) + runtime.KeepAlive(params) + runtime.KeepAlive(index) + runtime.KeepAlive(queries) + runtime.KeepAlive(neighbors) + runtime.KeepAlive(distances) + }() return cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsIvfFlatSearch(C.cuvsResources_t(Resources.Resource), params.params, index.index, (*C.DLManagedTensor)(unsafe.Pointer(queries.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(neighbors.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), prefilter))) } @@ -90,7 +98,7 @@ func GetNLists(index *IvfFlatIndex) (nlist int64, err error) { if err != nil { return } - + runtime.KeepAlive(index) nlist = int64(ret) return } @@ -105,7 +113,7 @@ func GetDim(index *IvfFlatIndex) (dim int64, err error) { if err != nil { return } - + runtime.KeepAlive(index) dim = int64(ret) return } @@ -116,5 +124,7 @@ func GetCenters[T any](index *IvfFlatIndex, centers *cuvs.Tensor[T]) error { } err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsIvfFlatIndexGetCenters(index.index, (*C.DLManagedTensor)(unsafe.Pointer(centers.C_tensor))))) + runtime.KeepAlive(index) + runtime.KeepAlive(centers) return err } From 946a77b3d69f66bae02e8b4b1cb539635dcee9f1 Mon Sep 17 00:00:00 2001 From: Eric Date: Fri, 6 Feb 2026 09:56:17 +0000 Subject: [PATCH 08/16] add test for issue #1622 --- go/test/issue_test.go | 257 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 257 insertions(+) create mode 100644 go/test/issue_test.go diff --git a/go/test/issue_test.go b/go/test/issue_test.go new file mode 100644 index 0000000000..e6ef53baf8 --- /dev/null +++ b/go/test/issue_test.go @@ -0,0 +1,257 @@ +package test + +import ( + //"fmt" + "math/rand/v2" + "sync" + "testing" + //"os" + + "github.com/stretchr/testify/require" + + cuvs "github.com/rapidsai/cuvs/go" + "github.com/rapidsai/cuvs/go/brute_force" + "github.com/rapidsai/cuvs/go/ivf_flat" +) + +func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Distance, maxIterations int) ([][]float32, error) { + stream, err := cuvs.NewCudaStream() + if err != nil { + return nil, err + } + defer stream.Close() + resource, err := cuvs.NewResource(stream) + if err != nil { + return nil, err + } + defer resource.Close() + + indexParams, err := ivf_flat.CreateIndexParams() + if err != nil { + return nil, err + } + defer indexParams.Close() + + indexParams.SetNLists(uint32(clusterCnt)) + indexParams.SetMetric(distanceType) + indexParams.SetKMeansNIters(uint32(maxIterations)) + indexParams.SetKMeansTrainsetFraction(1) // train all sample + + dataset, err := cuvs.NewTensor(vecs) + if err != nil { + return nil, err + } + defer dataset.Close() + + index, err := ivf_flat.CreateIndex[float32](indexParams) + if err != nil { + return nil, err + } + defer index.Close() + + if _, err := dataset.ToDevice(&resource); err != nil { + return nil, err + } + + centers, err := cuvs.NewTensorOnDevice[float32](&resource, []int64{int64(clusterCnt), int64(dim)}) + if err != nil { + return nil, err + } + + if err := ivf_flat.BuildIndex(resource, indexParams, &dataset, index); err != nil { + return nil, err + } + + if err := resource.Sync(); err != nil { + return nil, err + } + + if err := ivf_flat.GetCenters(index, ¢ers); err != nil { + return nil, err + } + + if _, err := centers.ToHost(&resource); err != nil { + return nil, err + } + + if err := resource.Sync(); err != nil { + return nil, err + } + + result, err := centers.Slice() + if err != nil { + return nil, err + } + + return result, nil + +} + +func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distanceType cuvs.Distance) (retkeys any, retdistances []float64, err error) { + stream, err := cuvs.NewCudaStream() + if err != nil { + return + } + defer stream.Close() + + resource, err := cuvs.NewResource(stream) + if err != nil { + return + } + defer resource.Close() + + dataset, err := cuvs.NewTensor(datasetvec) + if err != nil { + return + } + defer dataset.Close() + + index, err := brute_force.CreateIndex() + if err != nil { + return + } + defer index.Close() + + queries, err := cuvs.NewTensor(queriesvec) + if err != nil { + return + } + defer queries.Close() + + neighbors, err := cuvs.NewTensorOnDevice[int64](&resource, []int64{int64(len(queriesvec)), int64(limit)}) + if err != nil { + return + } + defer neighbors.Close() + + distances, err := cuvs.NewTensorOnDevice[float32](&resource, []int64{int64(len(queriesvec)), int64(limit)}) + if err != nil { + return + } + defer distances.Close() + + if _, err = dataset.ToDevice(&resource); err != nil { + return + } + + if err = resource.Sync(); err != nil { + return + } + + err = brute_force.BuildIndex(resource, &dataset, distanceType, 2.0, index) + if err != nil { + //os.Stderr.WriteString(fmt.Sprintf("BruteForceIndex: build index failed %v\n", err)) + //os.Stderr.WriteString(fmt.Sprintf("BruteForceIndex: build index failed centers %v\n", datasetvec)) + return + } + + if err = resource.Sync(); err != nil { + return + } + //os.Stderr.WriteString("built brute force index\n") + + if _, err = queries.ToDevice(&resource); err != nil { + return + } + + //os.Stderr.WriteString("brute force index search Runing....\n") + err = brute_force.SearchIndex(resource, index, &queries, &neighbors, &distances) + if err != nil { + return + } + //os.Stderr.WriteString("brute force index search finished Runing....\n") + + if _, err = neighbors.ToHost(&resource); err != nil { + return + } + //os.Stderr.WriteString("brute force index search neighbour to host done....\n") + + if _, err = distances.ToHost(&resource); err != nil { + return + } + //os.Stderr.WriteString("brute force index search distances to host done....\n") + + if err = resource.Sync(); err != nil { + return + } + + //os.Stderr.WriteString("brute force index search return result....\n") + neighborsSlice, err := neighbors.Slice() + if err != nil { + return + } + + distancesSlice, err := distances.Slice() + if err != nil { + return + } + + //fmt.Printf("flattened %v\n", flatten) + retdistances = make([]float64, len(distancesSlice)*int(limit)) + for i := range distancesSlice { + for j, dist := range distancesSlice[i] { + retdistances[i*int(limit)+j] = float64(dist) + } + } + + keys := make([]int64, len(neighborsSlice)*int(limit)) + for i := range neighborsSlice { + for j, key := range neighborsSlice[i] { + keys[i*int(limit)+j] = int64(key) + } + } + retkeys = keys + return +} + +func TestIvfAndBruteForceForIssue(t *testing.T) { + + dimension := uint(128) + limit := uint(1) + /* + ncpu := uint(1) + elemsz := uint(4) // float32 + */ + + dsize := 100000 + nlist := 128 + vecs := make([][]float32, dsize) + for i := range vecs { + vecs[i] = make([]float32, dimension) + for j := range vecs[i] { + vecs[i][j] = rand.Float32() + } + } + queries := vecs[:8192] + + centers, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) + require.NoError(t, err) + + var wg sync.WaitGroup + + for n := 0; n < 4; n++ { + + wg.Add(1) + go func() { + defer wg.Done() + for i := 0; i < 1000; i++ { + _, _, err := Search(centers, queries, limit, cuvs.DistanceL2) + require.NoError(t, err) + + /* + keys_i64, ok := keys.([]int64) + require.Equal(t, ok, true) + + for j, key := range keys_i64 { + require.Equal(t, key, int64(j)) + require.Equal(t, distances[j], float64(0)) + } + */ + // fmt.Printf("keys %v, dist %v\n", keys, distances) + } + }() + } + + wg.Wait() + +} From 50784b700898bc3ee5a20519f2a5fc3995bbd1a3 Mon Sep 17 00:00:00 2001 From: Eric Date: Sat, 7 Feb 2026 18:57:59 +0000 Subject: [PATCH 09/16] two Tests will IVF index will crash --- go/brute_force/brute_force.go | 3 +++ go/distance.go | 8 ++++++++ go/ivf_flat/ivf_flat.go | 2 ++ go/test/issue_test.go | 34 ++++++++++++++++++++++++++++++++++ 4 files changed, 47 insertions(+) diff --git a/go/brute_force/brute_force.go b/go/brute_force/brute_force.go index 30e80f85b4..426ffd06f4 100644 --- a/go/brute_force/brute_force.go +++ b/go/brute_force/brute_force.go @@ -63,6 +63,9 @@ func BuildIndex[T any](Resources cuvs.Resource, Dataset *cuvs.Tensor[T], metric runtime.KeepAlive(index) runtime.KeepAlive(Dataset) runtime.KeepAlive(Resources) + runtime.KeepAlive(metric_arg) + runtime.KeepAlive(metric) + runtime.KeepAlive(CMetric) return nil } diff --git a/go/distance.go b/go/distance.go index 530e1d61b3..bd5cde51ef 100644 --- a/go/distance.go +++ b/go/distance.go @@ -5,6 +5,7 @@ import "C" import ( "errors" + "runtime" "unsafe" ) @@ -66,5 +67,12 @@ func PairwiseDistance[T any](Resources Resource, x *Tensor[T], y *Tensor[T], dis return errors.New("cuvs: invalid distance metric") } + defer func() { + runtime.KeepAlive(Resources) + runtime.KeepAlive(x) + runtime.KeepAlive(y) + runtime.KeepAlive(distances) + }() + return CheckCuvs(CuvsError(C.cuvsPairwiseDistance(C.cuvsResources_t(Resources.Resource), (*C.DLManagedTensor)(unsafe.Pointer(x.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(y.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), C.cuvsDistanceType(CMetric), C.float(metric_arg)))) } diff --git a/go/ivf_flat/ivf_flat.go b/go/ivf_flat/ivf_flat.go index 1ba1988d56..be5939a660 100644 --- a/go/ivf_flat/ivf_flat.go +++ b/go/ivf_flat/ivf_flat.go @@ -24,6 +24,7 @@ func CreateIndex[T any](params *IndexParams) (*IvfFlatIndex, error) { if err != nil { return nil, err } + defer runtime.KeepAlive(index) return &IvfFlatIndex{index: index, trained: false}, nil } @@ -55,6 +56,7 @@ func (index *IvfFlatIndex) Close() error { if err != nil { return err } + runtime.KeepAlive(index) return nil } diff --git a/go/test/issue_test.go b/go/test/issue_test.go index e6ef53baf8..6f18e866b7 100644 --- a/go/test/issue_test.go +++ b/go/test/issue_test.go @@ -3,6 +3,7 @@ package test import ( //"fmt" "math/rand/v2" + "runtime" "sync" "testing" //"os" @@ -25,6 +26,7 @@ func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Dis return nil, err } defer resource.Close() + defer runtime.KeepAlive(resource) indexParams, err := ivf_flat.CreateIndexParams() if err != nil { @@ -99,6 +101,7 @@ func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distance return } defer resource.Close() + defer runtime.KeepAlive(resource) dataset, err := cuvs.NewTensor(datasetvec) if err != nil { @@ -204,8 +207,36 @@ func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distance return } + +func TestIssueGpu(t *testing.T) { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + + dimension := uint(128) + /* + ncpu := uint(1) + elemsz := uint(4) // float32 + */ + + dsize := 100000 + nlist := 128 + vecs := make([][]float32, dsize) + for i := range vecs { + vecs[i] = make([]float32, dimension) + for j := range vecs[i] { + vecs[i][j] = rand.Float32() + } + } + + _, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) + require.NoError(t, err) +} + func TestIvfAndBruteForceForIssue(t *testing.T) { + runtime.LockOSThread() + defer runtime.LockOSThread() + dimension := uint(128) limit := uint(1) /* @@ -233,6 +264,9 @@ func TestIvfAndBruteForceForIssue(t *testing.T) { wg.Add(1) go func() { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + defer wg.Done() for i := 0; i < 1000; i++ { _, _, err := Search(centers, queries, limit, cuvs.DistanceL2) From e9efeec67484c87ae44ed5038d8e20752c7a5716 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 12 Feb 2026 16:05:15 +0000 Subject: [PATCH 10/16] bug fix deleter is not nil --- go/dlpack.go | 99 +++++++++++++++++++++++++++++++------------ go/resources.go | 2 +- go/test/issue_test.go | 73 +++++++++++++++---------------- 3 files changed, 111 insertions(+), 63 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index 15a87e8f65..8e6e7b836f 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -3,6 +3,13 @@ package cuvs // #include // #include // #include +// +// void tensor_deleter_free(DLManagedTensor *dlm) { +// if (dlm->deleter != NULL) { +// dlm->deleter(dlm); +// dlm->deleter = NULL; +// } +// } import "C" import ( @@ -157,28 +164,35 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor // Destroys Tensor, freeing the memory it was allocated on. func (t *Tensor[T]) Close() error { if t.C_tensor != nil { - if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { - bytes := t.sizeInBytes() - if t.resource == nil { - // We cannot free device memory without the resource, - // so we must return an error here. - return errors.New("resource not found for CUDA tensor") + if t.C_tensor.deleter != nil { + C.tensor_deleter_free((*C.DLManagedTensor)(unsafe.Pointer(t.C_tensor))) + } else { + + if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { + bytes := t.sizeInBytes() + if t.resource == nil { + // We cannot free device memory without the resource, + // so we must return an error here. + return errors.New("resource not found for CUDA tensor") + } + if t.C_tensor.dl_tensor.data != nil { + err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) + if err != nil { + return err + } + t.C_tensor.dl_tensor.data = nil + } + } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { + if t.C_tensor.dl_tensor.data != nil { + C.free(t.C_tensor.dl_tensor.data) + t.C_tensor.dl_tensor.data = nil + } } - err := CheckCuvs(CuvsError(C.cuvsRMMFree(t.resource.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) - if err != nil { - return err - } - t.C_tensor.dl_tensor.data = nil - } else if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { - if t.C_tensor.dl_tensor.data != nil { - C.free(t.C_tensor.dl_tensor.data) - t.C_tensor.dl_tensor.data = nil - } - } - if t.C_tensor.dl_tensor.shape != nil { - C.free(unsafe.Pointer(t.C_tensor.dl_tensor.shape)) - t.C_tensor.dl_tensor.shape = nil + if t.C_tensor.dl_tensor.shape != nil { + C.free(unsafe.Pointer(t.C_tensor.dl_tensor.shape)) + t.C_tensor.dl_tensor.shape = nil + } } C.free(unsafe.Pointer(t.C_tensor)) @@ -329,11 +343,13 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { return nil, err } - err = CheckCuvs(CuvsError( - C.cuvsRMMFree(res.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) - if err != nil { - C.free(addr) - return nil, err + if t.C_tensor.deleter == nil { + err = CheckCuvs(CuvsError( + C.cuvsRMMFree(res.Resource, t.C_tensor.dl_tensor.data, C.size_t(bytes)))) + if err != nil { + C.free(addr) + return nil, err + } } runtime.KeepAlive(res) @@ -344,6 +360,38 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { return t, nil } +// Creates a new Tensor with nil data on the current device. +func NewTensorNoDataOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor[T], error) { + if len(shape) < 2 { + return Tensor[T]{}, errors.New("shape must be at least 2") + } + + dlm := (*C.DLManagedTensor)(C.malloc(C.size_t(unsafe.Sizeof(C.DLManagedTensor{})))) + dtype := getDLDataType[T]() + + dlm.dl_tensor.data = nil + dlm.dl_tensor.device = C.DLDevice{ + device_type: C.DLDeviceType(C.kDLCUDA), + device_id: 0, + } + dlm.dl_tensor.dtype = dtype + dlm.dl_tensor.ndim = C.int(len(shape)) + dlm.dl_tensor.shape = nil + dlm.dl_tensor.strides = nil + dlm.dl_tensor.byte_offset = 0 + dlm.manager_ctx = nil + dlm.deleter = nil + + shapeCopy := make([]int64, len(shape)) + copy(shapeCopy, shape) + + return Tensor[T]{ + C_tensor: dlm, + shape: shapeCopy, + resource: res, + }, nil +} + // Returns a slice of the data in the Tensor. // The Tensor must be on the host. func (t *Tensor[T]) Slice() ([][]T, error) { @@ -408,6 +456,5 @@ func calculateBytes(shape []int64, dtype C.DLDataType) int64 { bytes *= int64(dtype.bits) / 8 - return bytes } diff --git a/go/resources.go b/go/resources.go index 4c7a348693..1337203ada 100644 --- a/go/resources.go +++ b/go/resources.go @@ -29,7 +29,7 @@ func (s *CudaStream) Close() error { // Creates a new CUDA stream func NewCudaStream() (*CudaStream, error) { var stream C.cudaStream_t - err := CheckCuda(C.cudaStreamCreate(&stream)) + err := CheckCuda(C.cudaStreamCreate(&stream)) if err != nil { return nil, err } diff --git a/go/test/issue_test.go b/go/test/issue_test.go index 6f18e866b7..f0aea47aba 100644 --- a/go/test/issue_test.go +++ b/go/test/issue_test.go @@ -55,23 +55,28 @@ func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Dis return nil, err } - centers, err := cuvs.NewTensorOnDevice[float32](&resource, []int64{int64(clusterCnt), int64(dim)}) - if err != nil { + if err := ivf_flat.BuildIndex(resource, indexParams, &dataset, index); err != nil { return nil, err } - if err := ivf_flat.BuildIndex(resource, indexParams, &dataset, index); err != nil { + if err := resource.Sync(); err != nil { return nil, err } - if err := resource.Sync(); err != nil { + centers, err := cuvs.NewTensorNoDataOnDevice[float32](&resource, []int64{int64(clusterCnt), int64(dim)}) + if err != nil { return nil, err } + defer centers.Close() if err := ivf_flat.GetCenters(index, ¢ers); err != nil { return nil, err } + if err := resource.Sync(); err != nil { + return nil, err + } + if _, err := centers.ToHost(&resource); err != nil { return nil, err } @@ -86,7 +91,6 @@ func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Dis } return result, nil - } func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distanceType cuvs.Distance) (retkeys any, retdistances []float64, err error) { @@ -207,29 +211,28 @@ func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distance return } - func TestIssueGpu(t *testing.T) { - runtime.LockOSThread() - defer runtime.UnlockOSThread() - - dimension := uint(128) - /* - ncpu := uint(1) - elemsz := uint(4) // float32 - */ - - dsize := 100000 - nlist := 128 - vecs := make([][]float32, dsize) - for i := range vecs { - vecs[i] = make([]float32, dimension) - for j := range vecs[i] { - vecs[i][j] = rand.Float32() - } - } - - _, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) - require.NoError(t, err) + runtime.LockOSThread() + defer runtime.UnlockOSThread() + + dimension := uint(128) + /* + ncpu := uint(1) + elemsz := uint(4) // float32 + */ + + dsize := 100000 + nlist := 128 + vecs := make([][]float32, dsize) + for i := range vecs { + vecs[i] = make([]float32, dimension) + for j := range vecs[i] { + vecs[i][j] = rand.Float32() + } + } + + _, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) + require.NoError(t, err) } func TestIvfAndBruteForceForIssue(t *testing.T) { @@ -272,15 +275,13 @@ func TestIvfAndBruteForceForIssue(t *testing.T) { _, _, err := Search(centers, queries, limit, cuvs.DistanceL2) require.NoError(t, err) - /* - keys_i64, ok := keys.([]int64) - require.Equal(t, ok, true) - - for j, key := range keys_i64 { - require.Equal(t, key, int64(j)) - require.Equal(t, distances[j], float64(0)) - } - */ + // keys_i64, ok := keys.([]int64) + // require.Equal(t, ok, true) + // + // for j, key := range keys_i64 { + // require.Equal(t, key, int64(j)) + // require.Equal(t, distances[j], float64(0)) + // } // fmt.Printf("keys %v, dist %v\n", keys, distances) } }() From 5f59002adc255def144eab74af6c1b89bafa826d Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 12 Feb 2026 16:09:03 +0000 Subject: [PATCH 11/16] fmt --- go/dlpack.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/go/dlpack.go b/go/dlpack.go index 8e6e7b836f..ffeeec79cb 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -7,7 +7,7 @@ package cuvs // void tensor_deleter_free(DLManagedTensor *dlm) { // if (dlm->deleter != NULL) { // dlm->deleter(dlm); -// dlm->deleter = NULL; +// dlm->deleter = NULL; // } // } import "C" From 651a1c79822ecca864f2ed538946bb93871094e1 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 12 Feb 2026 17:56:16 +0000 Subject: [PATCH 12/16] add tests --- go/dlpack.go | 5 -- go/ivf_flat/ivf_flat.go | 19 ------- go/ivf_flat/ivf_flat_test.go | 97 ++++++++++++++++++++++++++++++++++++ 3 files changed, 97 insertions(+), 24 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index ffeeec79cb..48480accf9 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -14,7 +14,6 @@ import "C" import ( "errors" - "runtime" "strconv" "unsafe" ) @@ -197,7 +196,6 @@ func (t *Tensor[T]) Close() error { C.free(unsafe.Pointer(t.C_tensor)) t.C_tensor = nil - runtime.KeepAlive(t.resource) } return nil } @@ -234,7 +232,6 @@ func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { t.C_tensor.dl_tensor.data = DeviceDataPointer t.resource = res - runtime.KeepAlive(res) return t, nil } @@ -322,7 +319,6 @@ func (t *Tensor[T]) Expand(res *Resource, newData [][]T) (*Tensor[T], error) { C.free(unsafe.Pointer(oldShapePtr)) } - runtime.KeepAlive(res) return t, nil } @@ -352,7 +348,6 @@ func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { } } - runtime.KeepAlive(res) t.C_tensor.dl_tensor.device.device_type = C.kDLCPU t.C_tensor.dl_tensor.data = addr t.resource = nil diff --git a/go/ivf_flat/ivf_flat.go b/go/ivf_flat/ivf_flat.go index be5939a660..61040a50d9 100644 --- a/go/ivf_flat/ivf_flat.go +++ b/go/ivf_flat/ivf_flat.go @@ -5,7 +5,6 @@ import "C" import ( "errors" - "runtime" "unsafe" cuvs "github.com/rapidsai/cuvs/go" @@ -24,7 +23,6 @@ func CreateIndex[T any](params *IndexParams) (*IvfFlatIndex, error) { if err != nil { return nil, err } - defer runtime.KeepAlive(index) return &IvfFlatIndex{index: index, trained: false}, nil } @@ -42,10 +40,6 @@ func BuildIndex[T any](Resources cuvs.Resource, params *IndexParams, dataset *cu return err } index.trained = true - runtime.KeepAlive(Resources) - runtime.KeepAlive(params) - runtime.KeepAlive(dataset) - runtime.KeepAlive(index) return nil } @@ -56,7 +50,6 @@ func (index *IvfFlatIndex) Close() error { if err != nil { return err } - runtime.KeepAlive(index) return nil } @@ -79,14 +72,6 @@ func SearchIndex[T any](Resources cuvs.Resource, params *SearchParams, index *Iv _type: C.NO_FILTER, } - defer func() { - runtime.KeepAlive(Resources) - runtime.KeepAlive(params) - runtime.KeepAlive(index) - runtime.KeepAlive(queries) - runtime.KeepAlive(neighbors) - runtime.KeepAlive(distances) - }() return cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsIvfFlatSearch(C.cuvsResources_t(Resources.Resource), params.params, index.index, (*C.DLManagedTensor)(unsafe.Pointer(queries.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(neighbors.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), prefilter))) } @@ -100,7 +85,6 @@ func GetNLists(index *IvfFlatIndex) (nlist int64, err error) { if err != nil { return } - runtime.KeepAlive(index) nlist = int64(ret) return } @@ -115,7 +99,6 @@ func GetDim(index *IvfFlatIndex) (dim int64, err error) { if err != nil { return } - runtime.KeepAlive(index) dim = int64(ret) return } @@ -126,7 +109,5 @@ func GetCenters[T any](index *IvfFlatIndex, centers *cuvs.Tensor[T]) error { } err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsIvfFlatIndexGetCenters(index.index, (*C.DLManagedTensor)(unsafe.Pointer(centers.C_tensor))))) - runtime.KeepAlive(index) - runtime.KeepAlive(centers) return err } diff --git a/go/ivf_flat/ivf_flat_test.go b/go/ivf_flat/ivf_flat_test.go index 78d603a8e9..8a6fdb6b20 100644 --- a/go/ivf_flat/ivf_flat_test.go +++ b/go/ivf_flat/ivf_flat_test.go @@ -1,12 +1,32 @@ package ivf_flat import ( + "github.com/stretchr/testify/require" "math/rand/v2" + "runtime" "testing" cuvs "github.com/rapidsai/cuvs/go" ) +func TestGetCenters(t *testing.T) { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + + dimension := uint(128) + dsize := 100000 + nlist := 128 + vecs := make([][]float32, dsize) + for i := range vecs { + vecs[i] = make([]float32, dimension) + for j := range vecs[i] { + vecs[i][j] = rand.Float32() + } + } + _, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) + require.NoError(t, err) +} + func TestIvfFlat(t *testing.T) { const ( nDataPoints = 1024 @@ -177,3 +197,80 @@ func TestIvfFlat(t *testing.T) { } } } + +func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Distance, maxIterations int) ([][]float32, error) { + stream, err := cuvs.NewCudaStream() + if err != nil { + return nil, err + } + defer stream.Close() + resource, err := cuvs.NewResource(stream) + if err != nil { + return nil, err + } + defer resource.Close() + + indexParams, err := CreateIndexParams() + if err != nil { + return nil, err + } + defer indexParams.Close() + + indexParams.SetNLists(uint32(clusterCnt)) + indexParams.SetMetric(distanceType) + indexParams.SetKMeansNIters(uint32(maxIterations)) + indexParams.SetKMeansTrainsetFraction(1) // train all sample + + dataset, err := cuvs.NewTensor(vecs) + if err != nil { + return nil, err + } + defer dataset.Close() + + index, err := CreateIndex[float32](indexParams) + if err != nil { + return nil, err + } + defer index.Close() + + if _, err := dataset.ToDevice(&resource); err != nil { + return nil, err + } + + if err := BuildIndex(resource, indexParams, &dataset, index); err != nil { + return nil, err + } + + if err := resource.Sync(); err != nil { + return nil, err + } + + centers, err := cuvs.NewTensorNoDataOnDevice[float32](&resource, []int64{int64(clusterCnt), int64(dim)}) + if err != nil { + return nil, err + } + defer centers.Close() + + if err := GetCenters(index, ¢ers); err != nil { + return nil, err + } + + if err := resource.Sync(); err != nil { + return nil, err + } + + if _, err := centers.ToHost(&resource); err != nil { + return nil, err + } + + if err := resource.Sync(); err != nil { + return nil, err + } + + result, err := centers.Slice() + if err != nil { + return nil, err + } + + return result, nil +} From 5120ad81d3c04ba223f045fa8952fe2a7b078944 Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 12 Feb 2026 17:56:55 +0000 Subject: [PATCH 13/16] remove test --- go/test/issue_test.go | 292 ------------------------------------------ 1 file changed, 292 deletions(-) delete mode 100644 go/test/issue_test.go diff --git a/go/test/issue_test.go b/go/test/issue_test.go deleted file mode 100644 index f0aea47aba..0000000000 --- a/go/test/issue_test.go +++ /dev/null @@ -1,292 +0,0 @@ -package test - -import ( - //"fmt" - "math/rand/v2" - "runtime" - "sync" - "testing" - //"os" - - "github.com/stretchr/testify/require" - - cuvs "github.com/rapidsai/cuvs/go" - "github.com/rapidsai/cuvs/go/brute_force" - "github.com/rapidsai/cuvs/go/ivf_flat" -) - -func getCenters(vecs [][]float32, dim int, clusterCnt int, distanceType cuvs.Distance, maxIterations int) ([][]float32, error) { - stream, err := cuvs.NewCudaStream() - if err != nil { - return nil, err - } - defer stream.Close() - resource, err := cuvs.NewResource(stream) - if err != nil { - return nil, err - } - defer resource.Close() - defer runtime.KeepAlive(resource) - - indexParams, err := ivf_flat.CreateIndexParams() - if err != nil { - return nil, err - } - defer indexParams.Close() - - indexParams.SetNLists(uint32(clusterCnt)) - indexParams.SetMetric(distanceType) - indexParams.SetKMeansNIters(uint32(maxIterations)) - indexParams.SetKMeansTrainsetFraction(1) // train all sample - - dataset, err := cuvs.NewTensor(vecs) - if err != nil { - return nil, err - } - defer dataset.Close() - - index, err := ivf_flat.CreateIndex[float32](indexParams) - if err != nil { - return nil, err - } - defer index.Close() - - if _, err := dataset.ToDevice(&resource); err != nil { - return nil, err - } - - if err := ivf_flat.BuildIndex(resource, indexParams, &dataset, index); err != nil { - return nil, err - } - - if err := resource.Sync(); err != nil { - return nil, err - } - - centers, err := cuvs.NewTensorNoDataOnDevice[float32](&resource, []int64{int64(clusterCnt), int64(dim)}) - if err != nil { - return nil, err - } - defer centers.Close() - - if err := ivf_flat.GetCenters(index, ¢ers); err != nil { - return nil, err - } - - if err := resource.Sync(); err != nil { - return nil, err - } - - if _, err := centers.ToHost(&resource); err != nil { - return nil, err - } - - if err := resource.Sync(); err != nil { - return nil, err - } - - result, err := centers.Slice() - if err != nil { - return nil, err - } - - return result, nil -} - -func Search(datasetvec [][]float32, queriesvec [][]float32, limit uint, distanceType cuvs.Distance) (retkeys any, retdistances []float64, err error) { - stream, err := cuvs.NewCudaStream() - if err != nil { - return - } - defer stream.Close() - - resource, err := cuvs.NewResource(stream) - if err != nil { - return - } - defer resource.Close() - defer runtime.KeepAlive(resource) - - dataset, err := cuvs.NewTensor(datasetvec) - if err != nil { - return - } - defer dataset.Close() - - index, err := brute_force.CreateIndex() - if err != nil { - return - } - defer index.Close() - - queries, err := cuvs.NewTensor(queriesvec) - if err != nil { - return - } - defer queries.Close() - - neighbors, err := cuvs.NewTensorOnDevice[int64](&resource, []int64{int64(len(queriesvec)), int64(limit)}) - if err != nil { - return - } - defer neighbors.Close() - - distances, err := cuvs.NewTensorOnDevice[float32](&resource, []int64{int64(len(queriesvec)), int64(limit)}) - if err != nil { - return - } - defer distances.Close() - - if _, err = dataset.ToDevice(&resource); err != nil { - return - } - - if err = resource.Sync(); err != nil { - return - } - - err = brute_force.BuildIndex(resource, &dataset, distanceType, 2.0, index) - if err != nil { - //os.Stderr.WriteString(fmt.Sprintf("BruteForceIndex: build index failed %v\n", err)) - //os.Stderr.WriteString(fmt.Sprintf("BruteForceIndex: build index failed centers %v\n", datasetvec)) - return - } - - if err = resource.Sync(); err != nil { - return - } - //os.Stderr.WriteString("built brute force index\n") - - if _, err = queries.ToDevice(&resource); err != nil { - return - } - - //os.Stderr.WriteString("brute force index search Runing....\n") - err = brute_force.SearchIndex(resource, index, &queries, &neighbors, &distances) - if err != nil { - return - } - //os.Stderr.WriteString("brute force index search finished Runing....\n") - - if _, err = neighbors.ToHost(&resource); err != nil { - return - } - //os.Stderr.WriteString("brute force index search neighbour to host done....\n") - - if _, err = distances.ToHost(&resource); err != nil { - return - } - //os.Stderr.WriteString("brute force index search distances to host done....\n") - - if err = resource.Sync(); err != nil { - return - } - - //os.Stderr.WriteString("brute force index search return result....\n") - neighborsSlice, err := neighbors.Slice() - if err != nil { - return - } - - distancesSlice, err := distances.Slice() - if err != nil { - return - } - - //fmt.Printf("flattened %v\n", flatten) - retdistances = make([]float64, len(distancesSlice)*int(limit)) - for i := range distancesSlice { - for j, dist := range distancesSlice[i] { - retdistances[i*int(limit)+j] = float64(dist) - } - } - - keys := make([]int64, len(neighborsSlice)*int(limit)) - for i := range neighborsSlice { - for j, key := range neighborsSlice[i] { - keys[i*int(limit)+j] = int64(key) - } - } - retkeys = keys - return -} - -func TestIssueGpu(t *testing.T) { - runtime.LockOSThread() - defer runtime.UnlockOSThread() - - dimension := uint(128) - /* - ncpu := uint(1) - elemsz := uint(4) // float32 - */ - - dsize := 100000 - nlist := 128 - vecs := make([][]float32, dsize) - for i := range vecs { - vecs[i] = make([]float32, dimension) - for j := range vecs[i] { - vecs[i][j] = rand.Float32() - } - } - - _, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) - require.NoError(t, err) -} - -func TestIvfAndBruteForceForIssue(t *testing.T) { - - runtime.LockOSThread() - defer runtime.LockOSThread() - - dimension := uint(128) - limit := uint(1) - /* - ncpu := uint(1) - elemsz := uint(4) // float32 - */ - - dsize := 100000 - nlist := 128 - vecs := make([][]float32, dsize) - for i := range vecs { - vecs[i] = make([]float32, dimension) - for j := range vecs[i] { - vecs[i][j] = rand.Float32() - } - } - queries := vecs[:8192] - - centers, err := getCenters(vecs, int(dimension), nlist, cuvs.DistanceL2, 10) - require.NoError(t, err) - - var wg sync.WaitGroup - - for n := 0; n < 4; n++ { - - wg.Add(1) - go func() { - runtime.LockOSThread() - defer runtime.UnlockOSThread() - - defer wg.Done() - for i := 0; i < 1000; i++ { - _, _, err := Search(centers, queries, limit, cuvs.DistanceL2) - require.NoError(t, err) - - // keys_i64, ok := keys.([]int64) - // require.Equal(t, ok, true) - // - // for j, key := range keys_i64 { - // require.Equal(t, key, int64(j)) - // require.Equal(t, distances[j], float64(0)) - // } - // fmt.Printf("keys %v, dist %v\n", keys, distances) - } - }() - } - - wg.Wait() - -} From badddc5c42fe63b5aa7390644aef27ada440498f Mon Sep 17 00:00:00 2001 From: Eric Date: Thu, 12 Feb 2026 18:31:03 +0000 Subject: [PATCH 14/16] cleanup --- go/brute_force/brute_force.go | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/go/brute_force/brute_force.go b/go/brute_force/brute_force.go index 426ffd06f4..47f86018e3 100644 --- a/go/brute_force/brute_force.go +++ b/go/brute_force/brute_force.go @@ -5,7 +5,6 @@ import "C" import ( "errors" - "runtime" "unsafe" cuvs "github.com/rapidsai/cuvs/go" @@ -35,7 +34,6 @@ func (index *BruteForceIndex) Close() error { if err != nil { return err } - runtime.KeepAlive(index) return nil } @@ -60,12 +58,6 @@ func BuildIndex[T any](Resources cuvs.Resource, Dataset *cuvs.Tensor[T], metric } index.trained = true - runtime.KeepAlive(index) - runtime.KeepAlive(Dataset) - runtime.KeepAlive(Resources) - runtime.KeepAlive(metric_arg) - runtime.KeepAlive(metric) - runtime.KeepAlive(CMetric) return nil } @@ -89,10 +81,5 @@ func SearchIndex[T any](resources cuvs.Resource, index *BruteForceIndex, queries err := cuvs.CheckCuvs(cuvs.CuvsError(C.cuvsBruteForceSearch((C.cuvsResources_t)(resources.Resource), index.index, (*C.DLManagedTensor)(unsafe.Pointer(queries.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(neighbors.C_tensor)), (*C.DLManagedTensor)(unsafe.Pointer(distances.C_tensor)), prefilter))) - runtime.KeepAlive(index) - runtime.KeepAlive(queries) - runtime.KeepAlive(neighbors) - runtime.KeepAlive(distances) - runtime.KeepAlive(resources) return err } From c238b00c85a708f14b3b3cbe37e5d438b53c46a1 Mon Sep 17 00:00:00 2001 From: Eric Date: Fri, 13 Feb 2026 09:19:39 +0000 Subject: [PATCH 15/16] bug fix simply return when tensor already in device. double ToDevice() call --- go/dlpack.go | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/go/dlpack.go b/go/dlpack.go index 48480accf9..7de7d978c4 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -202,6 +202,12 @@ func (t *Tensor[T]) Close() error { // Transfers the data in the Tensor to the device. func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { + + if t.C_tensor.dl_tensor.device.device_type == C.kDLCUDA { + // tensor already in device + return t, nil + } + bytes := t.sizeInBytes() var DeviceDataPointer unsafe.Pointer @@ -223,9 +229,14 @@ func (t *Tensor[T]) ToDevice(res *Resource) (*Tensor[T], error) { return nil, err } - if t.C_tensor.dl_tensor.data != nil { - C.free(t.C_tensor.dl_tensor.data) - t.C_tensor.dl_tensor.data = nil + if t.C_tensor.deleter != nil { + C.tensor_deleter_free(t.C_tensor) + + } else { + if t.C_tensor.dl_tensor.data != nil { + C.free(t.C_tensor.dl_tensor.data) + t.C_tensor.dl_tensor.data = nil + } } t.C_tensor.dl_tensor.device.device_type = C.kDLCUDA From ec2d4ea859e5d95ad56b83e8319e37fe56fd8a22 Mon Sep 17 00:00:00 2001 From: Eric Date: Fri, 13 Feb 2026 09:23:06 +0000 Subject: [PATCH 16/16] simply return when tensor already in CPU. double ToHost() --- go/dlpack.go | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/go/dlpack.go b/go/dlpack.go index 7de7d978c4..3a68d7aa1e 100644 --- a/go/dlpack.go +++ b/go/dlpack.go @@ -335,6 +335,11 @@ func (t *Tensor[T]) Expand(res *Resource, newData [][]T) (*Tensor[T], error) { // Transfers the data in the Tensor to the host. func (t *Tensor[T]) ToHost(res *Resource) (*Tensor[T], error) { + if t.C_tensor.dl_tensor.device.device_type == C.kDLCPU { + // tensor is already in CPU + return t, nil + } + bytes := t.sizeInBytes() addr := (C.malloc(C.size_t(bytes)))