Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions go/brute_force/brute_force.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Expand All @@ -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
}
13 changes: 11 additions & 2 deletions go/brute_force/brute_force_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
}
Expand Down
8 changes: 8 additions & 0 deletions go/distance.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import "C"

import (
"errors"
"runtime"
"unsafe"
)

Expand Down Expand Up @@ -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))))
}
183 changes: 117 additions & 66 deletions go/dlpack.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,13 @@ package cuvs
// #include <stdlib.h>
// #include <dlpack/dlpack.h>
// #include <cuvs/core/c_api.h>
//
// void tensor_deleter_free(DLManagedTensor *dlm) {
// if (dlm->deleter != NULL) {
// dlm->deleter(dlm);
// dlm->deleter = NULL;
// }
// }
import "C"

import (
Expand All @@ -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.
Expand All @@ -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")
}

Comment thread
cpegeric marked this conversation as resolved.
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),
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
}
Expand Down Expand Up @@ -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,
Expand All @@ -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) {
Expand Down Expand Up @@ -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
Expand Down
Loading