diff --git a/go/brute_force/brute_force.go b/go/brute_force/brute_force.go index c75447b99b..47f86018e3 100644 --- a/go/brute_force/brute_force.go +++ b/go/brute_force/brute_force.go @@ -69,7 +69,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 +79,7 @@ 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))) 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/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/dlpack.go b/go/dlpack.go index 6fe619fd35..3a68d7aa1e 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 ( @@ -20,6 +27,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. @@ -33,29 +41,16 @@ 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])) // 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), @@ -85,18 +80,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 +120,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 @@ -176,43 +156,58 @@ func NewTensorOnDevice[T TensorNumberType](res *Resource, shape []int64) (Tensor return Tensor[T]{ C_tensor: dlm, shape: shapeCopy, + resource: res, }, nil } // 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() - res, err := NewResource(nil) - if err != nil { - return err - } - err = CheckCuvs(CuvsError(C.cuvsRMMFree(res.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 { - 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 != nil { + 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 + } + } + + 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 } // 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 @@ -233,8 +228,20 @@ 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.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 t.C_tensor.dl_tensor.data = DeviceDataPointer + t.resource = res return t, nil } @@ -307,28 +314,35 @@ 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 } // 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))) - if addr == nil { - return nil, errors.New("memory allocation failed") - } - err := CheckCuda( C.cudaMemcpy( addr, @@ -337,21 +351,58 @@ 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 { - 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 + } } t.C_tensor.dl_tensor.device.device_type = C.kDLCPU t.C_tensor.dl_tensor.data = addr + t.resource = nil 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) { @@ -410,8 +461,8 @@ 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 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..61040a50d9 100644 --- a/go/ivf_flat/ivf_flat.go +++ b/go/ivf_flat/ivf_flat.go @@ -17,14 +17,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 +40,7 @@ func BuildIndex[T any](Resources cuvs.Resource, params *IndexParams, dataset *cu return err } index.trained = true + return nil } @@ -63,7 +63,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") } @@ -85,7 +85,6 @@ func GetNLists(index *IvfFlatIndex) (nlist int64, err error) { if err != nil { return } - nlist = int64(ret) return } @@ -100,7 +99,6 @@ func GetDim(index *IvfFlatIndex) (dim int64, err error) { if err != nil { return } - dim = int64(ret) return } diff --git a/go/ivf_flat/ivf_flat_test.go b/go/ivf_flat/ivf_flat_test.go index 92ffdc715f..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 @@ -17,7 +37,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 +71,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 @@ -168,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 +} diff --git a/go/resources.go b/go/resources.go index 562aa7b7fc..1337203ada 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,8 +13,31 @@ type Resource struct { Resource C.cuvsResources_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 + 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 { @@ -21,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 } }