diff --git a/c/vendor/nanoarrow/nanoarrow.c b/c/vendor/nanoarrow/nanoarrow.c index 123e8e8451..cfe9c113be 100644 --- a/c/vendor/nanoarrow/nanoarrow.c +++ b/c/vendor/nanoarrow/nanoarrow.c @@ -22,6 +22,29 @@ #include #include +// For thread safe shared buffers we need C11 + stdatomic.h +// Can compile with -DNANOARROW_USE_STDATOMIC=0 or 1 to override +// automatic detection +#if !defined(NANOARROW_USE_STDATOMIC) +#define NANOARROW_USE_STDATOMIC 0 + +// Check for C11 +#if defined(__STDC_VERSION__) && __STDC_VERSION__ >= 201112L + +// Check for GCC 4.8, which doesn't include stdatomic.h but does +// not define __STDC_NO_ATOMICS__ +#if defined(__clang__) || !defined(__GNUC__) || __GNUC__ >= 5 + +#if !defined(__STDC_NO_ATOMICS__) +#include +#undef NANOARROW_USE_STDATOMIC +#define NANOARROW_USE_STDATOMIC 1 +#endif +#endif +#endif + +#endif + #include "nanoarrow.h" const char* ArrowNanoarrowVersion(void) { return NANOARROW_VERSION; } @@ -276,6 +299,225 @@ struct ArrowBufferAllocator ArrowBufferDeallocator( return allocator; } +#if NANOARROW_USE_STDATOMIC +struct ArrowSharedBufferPrivate { + struct ArrowBuffer src; + atomic_long reference_count; +}; + +static int64_t ArrowSharedBufferUpdate(struct ArrowSharedBufferPrivate* private_data, + int delta) { + int64_t old_count = atomic_fetch_add(&private_data->reference_count, delta); + return old_count + delta; +} + +static void ArrowSharedBufferSet(struct ArrowSharedBufferPrivate* private_data, + int64_t count) { + atomic_store(&private_data->reference_count, count); +} + +int ArrowSharedBufferIsThreadSafe(void) { return 1; } + +struct ArrowSharedArrayPrivate { + struct ArrowArray src; + atomic_long reference_count; +}; + +static int64_t ArrowSharedArrayUpdate(struct ArrowSharedArrayPrivate* private_data, + int delta) { + int64_t old_count = atomic_fetch_add(&private_data->reference_count, delta); + return old_count + delta; +} + +static void ArrowSharedArraySet(struct ArrowSharedArrayPrivate* private_data, + int64_t count) { + atomic_store(&private_data->reference_count, count); +} +#else +struct ArrowSharedBufferPrivate { + struct ArrowBuffer src; + int64_t reference_count; +}; + +static int64_t ArrowSharedBufferUpdate(struct ArrowSharedBufferPrivate* private_data, + int delta) { + private_data->reference_count += delta; + return private_data->reference_count; +} + +static void ArrowSharedBufferSet(struct ArrowSharedBufferPrivate* private_data, + int64_t count) { + private_data->reference_count = count; +} + +int ArrowSharedBufferIsThreadSafe(void) { return 0; } + +struct ArrowSharedArrayPrivate { + struct ArrowArray src; + int64_t reference_count; +}; + +static int64_t ArrowSharedArrayUpdate(struct ArrowSharedArrayPrivate* private_data, + int delta) { + private_data->reference_count += delta; + return private_data->reference_count; +} + +static void ArrowSharedArraySet(struct ArrowSharedArrayPrivate* private_data, + int64_t count) { + private_data->reference_count = count; +} +#endif + +static void ArrowSharedBufferFree(struct ArrowBufferAllocator* allocator, uint8_t* ptr, + int64_t size) { + NANOARROW_UNUSED(allocator); + NANOARROW_UNUSED(ptr); + NANOARROW_UNUSED(size); + + struct ArrowSharedBufferPrivate* private_data = + (struct ArrowSharedBufferPrivate*)allocator->private_data; + + if (ArrowSharedBufferUpdate(private_data, -1) == 0) { + ArrowBufferReset(&private_data->src); + ArrowFree(private_data); + } +} + +static void ArrowSharedArrayBufferFree(struct ArrowBufferAllocator* allocator, + uint8_t* ptr, int64_t size) { + NANOARROW_UNUSED(ptr); + NANOARROW_UNUSED(size); + + struct ArrowSharedArrayPrivate* private_data = + (struct ArrowSharedArrayPrivate*)allocator->private_data; + + if (ArrowSharedArrayUpdate(private_data, -1) == 0) { + if (private_data->src.release != NULL) { + ArrowArrayRelease(&private_data->src); + } + ArrowFree(private_data); + } +} + +ArrowErrorCode ArrowSharedBufferInit(struct ArrowBuffer* shared, + struct ArrowBuffer* src) { + if (src->data == NULL) { + ArrowBufferMove(src, shared); + return NANOARROW_OK; + } + + struct ArrowSharedBufferPrivate* private_data = + (struct ArrowSharedBufferPrivate*)ArrowMalloc( + sizeof(struct ArrowSharedBufferPrivate)); + if (private_data == NULL) { + return ENOMEM; + } + + ArrowBufferMove(src, &private_data->src); + ArrowSharedBufferSet(private_data, 1); + + ArrowBufferInit(shared); + shared->data = private_data->src.data; + shared->size_bytes = private_data->src.size_bytes; + // Don't expose any extra capcity from src so that any calls to ArrowBufferAppend + // on this buffer will fail. + shared->capacity_bytes = private_data->src.size_bytes; + shared->allocator = ArrowBufferDeallocator(&ArrowSharedBufferFree, private_data); + return NANOARROW_OK; +} + +int ArrowIsSharedBuffer(struct ArrowBuffer* buffer) { + return buffer->allocator.free == &ArrowSharedBufferFree || + buffer->allocator.free == &ArrowSharedArrayBufferFree; +} + +ArrowErrorCode ArrowSharedBufferClone(struct ArrowBuffer* shared, + struct ArrowBuffer* shared_out) { + // If the buffer has no data, initialize an empty buffer + if (shared->data == NULL) { + ArrowBufferInit(shared_out); + return NANOARROW_OK; + } + + if (shared->allocator.free == &ArrowSharedBufferFree) { + struct ArrowSharedBufferPrivate* private_data = + (struct ArrowSharedBufferPrivate*)shared->allocator.private_data; + ArrowSharedBufferUpdate(private_data, 1); + memcpy(shared_out, shared, sizeof(struct ArrowBuffer)); + return NANOARROW_OK; + } + + if (shared->allocator.free == &ArrowSharedArrayBufferFree) { + struct ArrowSharedArrayPrivate* private_data = + (struct ArrowSharedArrayPrivate*)shared->allocator.private_data; + ArrowSharedArrayUpdate(private_data, 1); + memcpy(shared_out, shared, sizeof(struct ArrowBuffer)); + return NANOARROW_OK; + } + + return EINVAL; +} + +ArrowErrorCode ArrowSharedArrayInit(struct ArrowSharedArray* shared, + struct ArrowArray* src) { + struct ArrowSharedArrayPrivate* private_data = + (struct ArrowSharedArrayPrivate*)ArrowMalloc( + sizeof(struct ArrowSharedArrayPrivate)); + if (private_data == NULL) { + return ENOMEM; + } + + ArrowArrayMove(src, &private_data->src); + ArrowSharedArraySet(private_data, 1); + shared->private_data = private_data; + return NANOARROW_OK; +} + +void ArrowSharedArrayRelease(struct ArrowSharedArray* shared) { + if (shared->private_data == NULL) { + return; + } + + struct ArrowSharedArrayPrivate* private_data = + (struct ArrowSharedArrayPrivate*)shared->private_data; + + if (ArrowSharedArrayUpdate(private_data, -1) == 0) { + if (private_data->src.release != NULL) { + ArrowArrayRelease(&private_data->src); + } + ArrowFree(private_data); + } + + shared->private_data = NULL; +} + +ArrowErrorCode ArrowSharedArrayBuffer(struct ArrowSharedArray* shared, int64_t i, + struct ArrowBuffer* out) { + struct ArrowSharedArrayPrivate* private_data = + (struct ArrowSharedArrayPrivate*)shared->private_data; + NANOARROW_DCHECK(i >= 0 && i < private_data->src.n_buffers); + + ArrowSharedArrayUpdate(private_data, 1); + ArrowBufferInit(out); + + if (ArrowArrayIsInternal(&private_data->src)) { + // The source array was built with nanoarrow, so we can get buffer size info + struct ArrowBuffer* src = ArrowArrayBuffer(&private_data->src, i); + out->data = src->data; + out->size_bytes = src->size_bytes; + out->capacity_bytes = src->size_bytes; + } else { + // Generic C Data Interface array: buffer sizes are not known + out->data = (uint8_t*)private_data->src.buffers[i]; + out->size_bytes = 0; + out->capacity_bytes = 0; + } + + out->allocator = ArrowBufferDeallocator(&ArrowSharedArrayBufferFree, private_data); + return NANOARROW_OK; +} + static const int kInt32DecimalDigits = 9; static const uint64_t kUInt32PowersOfTen[] = { @@ -578,6 +820,7 @@ ArrowErrorCode ArrowDecimalAppendStringToBuffer(const struct ArrowDecimal* decim #include #include +#include #include #include #include @@ -1931,47 +2174,34 @@ ArrowErrorCode ArrowSchemaViewInit(struct ArrowSchemaView* schema_view, return NANOARROW_OK; } -static int64_t ArrowSchemaTypeToStringInternal(struct ArrowSchemaView* schema_view, - char* out, int64_t n) { - const char* type_string = ArrowTypeString(schema_view->type); - switch (schema_view->type) { - case NANOARROW_TYPE_DECIMAL32: - case NANOARROW_TYPE_DECIMAL64: - case NANOARROW_TYPE_DECIMAL128: - case NANOARROW_TYPE_DECIMAL256: - return snprintf(out, n, "%s(%" PRId32 ", %" PRId32 ")", type_string, - schema_view->decimal_precision, schema_view->decimal_scale); - case NANOARROW_TYPE_TIMESTAMP: - return snprintf(out, n, "%s('%s', '%s')", type_string, - ArrowTimeUnitString(schema_view->time_unit), schema_view->timezone); - case NANOARROW_TYPE_TIME32: - case NANOARROW_TYPE_TIME64: - case NANOARROW_TYPE_DURATION: - return snprintf(out, n, "%s('%s')", type_string, - ArrowTimeUnitString(schema_view->time_unit)); - case NANOARROW_TYPE_FIXED_SIZE_BINARY: - case NANOARROW_TYPE_FIXED_SIZE_LIST: - return snprintf(out, n, "%s(%" PRId32 ")", type_string, schema_view->fixed_size); - case NANOARROW_TYPE_SPARSE_UNION: - case NANOARROW_TYPE_DENSE_UNION: - return snprintf(out, n, "%s([%s])", type_string, schema_view->union_type_ids); - default: - return snprintf(out, n, "%s", type_string); +static inline void ArrowSchemaToStringSNPrintF(char** out, int64_t* n_remaining, + int64_t* n_chars, const char* fmt, ...) { + // With _FORTIFY_SOURCE=3, vsnprintf on a null output warns, so ensure we never pass + // a NULL to vsnprintf + char tmp_not_null[1]; + char* out_not_null; + size_t out_not_null_size; + if (*out == NULL) { + out_not_null = tmp_not_null; + out_not_null_size = sizeof(tmp_not_null); + } else { + out_not_null = *out; + out_not_null_size = (size_t)*n_remaining; } -} -// Helper for bookkeeping to emulate sprintf()-like behaviour spread -// among multiple sprintf calls. -static inline void ArrowToStringLogChars(char** out, int64_t n_chars_last, - int64_t* n_remaining, int64_t* n_chars) { + va_list args; + va_start(args, fmt); + int n_chars_last = vsnprintf(out_not_null, out_not_null_size, fmt, args); + va_end(args); + // In the unlikely snprintf() returning a negative value (encoding error), // ensure the result won't cause an out-of-bounds access. if (n_chars_last < 0) { n_chars_last = 0; } - *n_chars += n_chars_last; *n_remaining -= n_chars_last; + *n_chars += n_chars_last; // n_remaining is never less than 0 if (*n_remaining < 0) { @@ -1984,93 +2214,134 @@ static inline void ArrowToStringLogChars(char** out, int64_t n_chars_last, } } -int64_t ArrowSchemaToString(const struct ArrowSchema* schema, char* out, int64_t n, - char recursive) { +static void ArrowSchemaTypeToStringInternal(struct ArrowSchemaView* schema_view, + char** out, int64_t* n_remaining, + int64_t* n_chars) { + const char* type_string = ArrowTypeString(schema_view->type); + switch (schema_view->type) { + case NANOARROW_TYPE_DECIMAL32: + case NANOARROW_TYPE_DECIMAL64: + case NANOARROW_TYPE_DECIMAL128: + case NANOARROW_TYPE_DECIMAL256: + ArrowSchemaToStringSNPrintF( + out, n_remaining, n_chars, "%s(%" PRId32 ", %" PRId32 ")", type_string, + schema_view->decimal_precision, schema_view->decimal_scale); + return; + case NANOARROW_TYPE_TIMESTAMP: + ArrowSchemaToStringSNPrintF( + out, n_remaining, n_chars, "%s('%s', '%s')", type_string, + ArrowTimeUnitString(schema_view->time_unit), schema_view->timezone); + return; + case NANOARROW_TYPE_TIME32: + case NANOARROW_TYPE_TIME64: + case NANOARROW_TYPE_DURATION: + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%s('%s')", type_string, + ArrowTimeUnitString(schema_view->time_unit)); + return; + case NANOARROW_TYPE_FIXED_SIZE_BINARY: + case NANOARROW_TYPE_FIXED_SIZE_LIST: + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%s(%" PRId32 ")", + type_string, schema_view->fixed_size); + return; + case NANOARROW_TYPE_SPARSE_UNION: + case NANOARROW_TYPE_DENSE_UNION: + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%s([%s])", type_string, + schema_view->union_type_ids); + return; + default: + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%s", type_string); + return; + } +} + +static void ArrowSchemaToStringInternal(const struct ArrowSchema* schema, char** out, + int64_t* n_remaining, int64_t* n_chars, + char recursive) { if (schema == NULL) { - return snprintf(out, n, "[invalid: pointer is null]"); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "[invalid: pointer is null]"); + return; } if (schema->release == NULL) { - return snprintf(out, n, "[invalid: schema is released]"); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, + "[invalid: schema is released]"); + return; } struct ArrowSchemaView schema_view; struct ArrowError error; if (ArrowSchemaViewInit(&schema_view, schema, &error) != NANOARROW_OK) { - return snprintf(out, n, "[invalid: %s]", ArrowErrorMessage(&error)); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "[invalid: %s]", + ArrowErrorMessage(&error)); + return; } // Extension type and dictionary should include both the top-level type // and the storage type. int is_extension = schema_view.extension_name.size_bytes > 0; int is_dictionary = schema->dictionary != NULL; - int64_t n_chars = 0; - int64_t n_chars_last = 0; // Uncommon but not technically impossible that both are true if (is_extension && is_dictionary) { - n_chars_last = snprintf( - out, n, "%.*s{dictionary(%s)<", (int)schema_view.extension_name.size_bytes, - schema_view.extension_name.data, ArrowTypeString(schema_view.storage_type)); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%.*s{dictionary(%s)<", + (int)schema_view.extension_name.size_bytes, + schema_view.extension_name.data, + ArrowTypeString(schema_view.storage_type)); } else if (is_extension) { - n_chars_last = snprintf(out, n, "%.*s{", (int)schema_view.extension_name.size_bytes, - schema_view.extension_name.data); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "%.*s{", + (int)schema_view.extension_name.size_bytes, + schema_view.extension_name.data); } else if (is_dictionary) { - n_chars_last = - snprintf(out, n, "dictionary(%s)<", ArrowTypeString(schema_view.storage_type)); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "dictionary(%s)<", + ArrowTypeString(schema_view.storage_type)); } - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); - if (!is_dictionary) { - n_chars_last = ArrowSchemaTypeToStringInternal(&schema_view, out, n); + ArrowSchemaTypeToStringInternal(&schema_view, out, n_remaining, n_chars); } else { - n_chars_last = ArrowSchemaToString(schema->dictionary, out, n, recursive); + ArrowSchemaToStringInternal(schema->dictionary, out, n_remaining, n_chars, recursive); } - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); - if (recursive && schema->format[0] == '+') { - n_chars_last = snprintf(out, n, "<"); - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "<"); for (int64_t i = 0; i < schema->n_children; i++) { if (i > 0) { - n_chars_last = snprintf(out, n, ", "); - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, ", "); } // ArrowSchemaToStringInternal() will validate the child and print the error, // but we need the name first if (schema->children[i] != NULL && schema->children[i]->release != NULL && schema->children[i]->name != NULL) { - n_chars_last = snprintf(out, n, "%s: ", schema->children[i]->name); - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, + "%s: ", schema->children[i]->name); } - n_chars_last = ArrowSchemaToString(schema->children[i], out, n, recursive); - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); + ArrowSchemaToStringInternal(schema->children[i], out, n_remaining, n_chars, + recursive); } - n_chars_last = snprintf(out, n, ">"); - ArrowToStringLogChars(&out, n_chars_last, &n, &n_chars); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, ">"); } if (is_extension && is_dictionary) { - n_chars += snprintf(out, n, ">}"); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, ">}"); } else if (is_extension) { - n_chars += snprintf(out, n, "}"); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, "}"); } else if (is_dictionary) { - n_chars += snprintf(out, n, ">"); + ArrowSchemaToStringSNPrintF(out, n_remaining, n_chars, ">"); } +} - // Ensure that we always return a positive result - if (n_chars > 0) { - return n_chars; - } else { - return 0; - } +int64_t ArrowSchemaToString(const struct ArrowSchema* schema, char* out, int64_t n, + char recursive) { + NANOARROW_DCHECK(n >= 0); + NANOARROW_DCHECK(out != NULL || n == 0); + int64_t n_chars = 0; + ArrowSchemaToStringInternal(schema, &out, &n, &n_chars, recursive); + return n_chars; } ArrowErrorCode ArrowMetadataReaderInit(struct ArrowMetadataReader* reader, @@ -2363,7 +2634,7 @@ static void ArrowArrayReleaseInternal(struct ArrowArray* array) { array->release = NULL; } -static int ArrowArrayIsInternal(struct ArrowArray* array) { +int ArrowArrayIsInternal(struct ArrowArray* array) { return array->release == &ArrowArrayReleaseInternal; } @@ -2858,6 +3129,172 @@ ArrowErrorCode ArrowArrayFinishBuildingDefault(struct ArrowArray* array, return ArrowArrayFinishBuilding(array, NANOARROW_VALIDATION_LEVEL_DEFAULT, error); } +static int ArrowArrayIsShared(struct ArrowArray* array) { + if (!ArrowArrayIsInternal(array)) { + return 0; + } + + for (int64_t i = 0; i < array->n_buffers; i++) { + struct ArrowBuffer* buffer = ArrowArrayBuffer(array, i); + if (buffer->data != NULL && !ArrowIsSharedBuffer(buffer)) { + return 0; + } + } + + for (int64_t i = 0; i < array->n_children; i++) { + if (!ArrowArrayIsShared(array->children[i])) { + return 0; + } + } + + if (array->dictionary != NULL && !ArrowArrayIsShared(array->dictionary)) { + return 0; + } + + return 1; +} + +static ArrowErrorCode ArrowArrayMoveSharedInternal(struct ArrowArray* src, + struct ArrowArray* dst) { + if (ArrowArrayIsShared(src)) { + ArrowArrayMove(src, dst); + return NANOARROW_OK; + } + + NANOARROW_RETURN_NOT_OK(ArrowArrayInitFromType(dst, NANOARROW_TYPE_UNINITIALIZED)); + + // Allocate children and move source children to dst children + NANOARROW_RETURN_NOT_OK(ArrowArrayAllocateChildren(dst, src->n_children)); + for (int64_t i = 0; i < src->n_children; i++) { + NANOARROW_RETURN_NOT_OK( + ArrowArrayMoveSharedInternal(src->children[i], dst->children[i])); + } + + // Allocate dictionary if needed and move source dictionary to dst dictionary + if (src->dictionary != NULL) { + NANOARROW_RETURN_NOT_OK(ArrowArrayAllocateDictionary(dst)); + NANOARROW_RETURN_NOT_OK( + ArrowArrayMoveSharedInternal(src->dictionary, dst->dictionary)); + } + + // We might need some more buffers if we are shallowly copying a string/binary view + if (src->n_buffers > 3) { + if (src->n_buffers > INT_MAX) { + return EINVAL; + } + + NANOARROW_RETURN_NOT_OK( + ArrowArrayAddVariadicBuffers(dst, (int32_t)src->n_buffers - 3)); + } + + // Move src into a shared array and set dst's buffers using the ref-counted version + struct ArrowSharedArray shared_array; + NANOARROW_RETURN_NOT_OK(ArrowSharedArrayInit(&shared_array, src)); + + for (int64_t i = 0; i < src->n_buffers; i++) { + struct ArrowBuffer* dst_buffer = ArrowArrayBuffer(dst, i); + ArrowErrorCode result = ArrowSharedArrayBuffer(&shared_array, i, dst_buffer); + if (result != NANOARROW_OK) { + ArrowSharedArrayRelease(&shared_array); + return result; + } + } + + ArrowSharedArrayRelease(&shared_array); + + dst->n_buffers = src->n_buffers; + dst->length = src->length; + dst->null_count = src->null_count; + dst->offset = src->offset; + + // Flush internal buffer pointers to array->buffers + NANOARROW_RETURN_NOT_OK(ArrowArrayFlushInternalPointers(dst)); + + return NANOARROW_OK; +} + +static ArrowErrorCode ArrowArrayCloneSharedInternal(struct ArrowArray* src, + struct ArrowArray* dst) { + NANOARROW_RETURN_NOT_OK(ArrowArrayInitFromType(dst, NANOARROW_TYPE_UNINITIALIZED)); + + // Allocate children and clone source children to dst children + NANOARROW_RETURN_NOT_OK(ArrowArrayAllocateChildren(dst, src->n_children)); + for (int64_t i = 0; i < src->n_children; i++) { + NANOARROW_RETURN_NOT_OK( + ArrowArrayCloneSharedInternal(src->children[i], dst->children[i])); + } + + // Allocate dictionary if needed and clone source dictionary to dst dictionary + if (src->dictionary != NULL) { + NANOARROW_RETURN_NOT_OK(ArrowArrayAllocateDictionary(dst)); + NANOARROW_RETURN_NOT_OK( + ArrowArrayCloneSharedInternal(src->dictionary, dst->dictionary)); + } + + // We might need some more buffers if we are shallowly copying a string/binary view + if (src->n_buffers > 3) { + if (src->n_buffers > INT_MAX) { + return EINVAL; + } + + NANOARROW_RETURN_NOT_OK( + ArrowArrayAddVariadicBuffers(dst, (int32_t)src->n_buffers - 3)); + } + + for (int64_t i = 0; i < src->n_buffers; i++) { + struct ArrowBuffer* src_buffer = ArrowArrayBuffer(src, i); + struct ArrowBuffer* dst_buffer = ArrowArrayBuffer(dst, i); + NANOARROW_RETURN_NOT_OK(ArrowSharedBufferClone(src_buffer, dst_buffer)); + } + + dst->n_buffers = src->n_buffers; + dst->length = src->length; + dst->null_count = src->null_count; + dst->offset = src->offset; + + // Flush internal buffer pointers to array->buffers + NANOARROW_RETURN_NOT_OK(ArrowArrayFlushInternalPointers(dst)); + + return NANOARROW_OK; +} + +ArrowErrorCode ArrowArrayMoveShared(struct ArrowArray* array, struct ArrowArray* shared) { + struct ArrowArray tmp; + tmp.release = NULL; + ArrowErrorCode result = ArrowArrayMoveSharedInternal(array, shared); + if (result != NANOARROW_OK) { + // On failure, release the temporary output + if (tmp.release != NULL) { + ArrowArrayRelease(&tmp); + } + + // Because this operation may have partially moved the input array at this point + // we also have to release it on failure to be predictable. These failures are + // usually failed heap allocations and are difficult to trigger in practice. + if (array->release != NULL) { + ArrowArrayRelease(array); + } + } + + return result; +} + +ArrowErrorCode ArrowArrayCloneShared(struct ArrowArray* shared, + struct ArrowArray* array) { + if (!ArrowArrayIsShared(shared)) { + return EINVAL; + } + + struct ArrowArray tmp; + tmp.release = NULL; + ArrowErrorCode result = ArrowArrayCloneSharedInternal(shared, array); + if (result != NANOARROW_OK && tmp.release != NULL) { + ArrowArrayRelease(&tmp); + } + + return result; +} + void ArrowArrayViewInitFromType(struct ArrowArrayView* array_view, enum ArrowType storage_type) { memset(array_view, 0, sizeof(struct ArrowArrayView)); @@ -3803,10 +4240,26 @@ static int ArrowArrayViewValidateFull(struct ArrowArrayView* array_view, NANOARROW_RETURN_NOT_OK(ArrowArrayViewValidateFull(array_view->children[i], error)); } - // Dictionary validation not implemented + // Dictionary index validation if (array_view->dictionary != NULL) { NANOARROW_RETURN_NOT_OK(ArrowArrayViewValidateFull(array_view->dictionary, error)); - // TODO: validate the indices + + // Validate that all non-null indices are within the dictionary bounds + int64_t dictionary_length = array_view->dictionary->length; + for (int64_t i = 0; i < array_view->length; i++) { + if (ArrowArrayViewIsNull(array_view, i)) { + continue; + } + + int64_t index = ArrowArrayViewGetIntUnsafe(array_view, i); + if (index < 0 || index >= dictionary_length) { + ArrowErrorSet(error, + "[%" PRId64 "] Expected dictionary index >= 0 and < %" PRId64 + " but found value %" PRId64, + i, dictionary_length, index); + return EINVAL; + } + } } return NANOARROW_OK; diff --git a/c/vendor/nanoarrow/nanoarrow.h b/c/vendor/nanoarrow/nanoarrow.h index 8a2e78cb74..47c4195736 100644 --- a/c/vendor/nanoarrow/nanoarrow.h +++ b/c/vendor/nanoarrow/nanoarrow.h @@ -512,7 +512,7 @@ enum ArrowType { /// \brief Get a string value of an enum ArrowType value /// \ingroup nanoarrow-utils /// -/// Returns NULL for invalid values for type +/// Returns "" for invalid values for type static inline const char* ArrowTypeString(enum ArrowType type); static inline const char* ArrowTypeString(enum ArrowType type) { @@ -608,7 +608,7 @@ static inline const char* ArrowTypeString(enum ArrowType type) { case NANOARROW_TYPE_LARGE_LIST_VIEW: return "large_list_view"; default: - return NULL; + return ""; } } @@ -1140,6 +1140,17 @@ static inline void ArrowDecimalSetBytes(struct ArrowDecimal* decimal, NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowBufferAllocatorDefault) #define ArrowBufferDeallocator \ NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowBufferDeallocator) +#define ArrowSharedBufferInit NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedBufferInit) +#define ArrowSharedBufferIsThreadSafe \ + NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedBufferIsThreadSafe) +#define ArrowSharedBufferClone \ + NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedBufferClone) +#define ArrowSharedArrayInit NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedArrayInit) +#define ArrowSharedArrayRelease \ + NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedArrayRelease) +#define ArrowSharedArrayBuffer \ + NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSharedArrayBuffer) +#define ArrowIsSharedBuffer NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowIsSharedBuffer) #define ArrowErrorSet NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowErrorSet) #define ArrowLayoutInit NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowLayoutInit) #define ArrowDecimalSetDigits NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowDecimalSetDigits) @@ -1189,6 +1200,7 @@ static inline void ArrowDecimalSetBytes(struct ArrowDecimal* decimal, NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowMetadataBuilderRemove) #define ArrowSchemaViewInit NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSchemaViewInit) #define ArrowSchemaToString NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowSchemaToString) +#define ArrowArrayIsInternal NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayIsInternal) #define ArrowArrayInitFromType \ NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayInitFromType) #define ArrowArrayInitFromSchema \ @@ -1209,6 +1221,8 @@ static inline void ArrowDecimalSetBytes(struct ArrowDecimal* decimal, NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayFinishBuilding) #define ArrowArrayFinishBuildingDefault \ NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayFinishBuildingDefault) +#define ArrowArrayMoveShared NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayMoveShared) +#define ArrowArrayCloneShared NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayCloneShared) #define ArrowArrayViewInitFromType \ NANOARROW_SYMBOL(NANOARROW_NAMESPACE, ArrowArrayViewInitFromType) #define ArrowArrayViewInitFromSchema \ @@ -1298,6 +1312,85 @@ NANOARROW_DLL struct ArrowBufferAllocator ArrowBufferDeallocator( /// @} +/// \defgroup nanoarrow-shared-buffer Shared buffers and arrays +/// +/// The nanoarrow library implements two kinds of shared buffers: those backed by +/// a reference-counted ArrowBuffer and those backed by a reference-counted ArrowArray. +/// These are useful for building ArrowArrays using buffers without copying the bytes +/// of the source. For example the IPC implementation uses shared buffers derived from an +/// ArrowBuffer to avoid copying buffers from an IPC message, and the Shared buffers +/// derived from arrays are used to make a safe shallow clone of an array. +/// +/// Reference counts are implemented using C11 atomics and not necessarily thread safe +/// on all platforms (notably, MSVC requires `/experimental:c11atomics`). Use +/// ArrowSharedBufferIsThreadSafe() to check for thread safety if this is needed. +/// +/// @{ + +/// \brief Initialize a shared buffer from an ArrowBuffer +/// +/// If NANOARROW_OK is returned, shared is a reference-counted buffer that +/// took ownership of src. The caller should release the shared buffer +/// using ArrowBufferReset() when finished. The underlying data will persist +/// until all clones have also been released. +NANOARROW_DLL ArrowErrorCode ArrowSharedBufferInit(struct ArrowBuffer* shared, + struct ArrowBuffer* src); + +/// \brief Check for shared buffer thread safety +/// +/// Thread-safe shared buffers require C11 and the stdatomic.h header. +/// If either are unavailable, shared buffers are still possible but +/// the resulting arrays must not be passed to other threads to be released. +NANOARROW_DLL int ArrowSharedBufferIsThreadSafe(void); + +/// \brief Check if a buffer is a shared buffer +/// +/// Returns non-zero if buffer was created by ArrowSharedBufferInit() or +/// obtained from ArrowSharedArrayBuffer(). +NANOARROW_DLL int ArrowIsSharedBuffer(struct ArrowBuffer* buffer); + +/// \brief Clone a shared buffer, incrementing its reference count +/// +/// This creates a new ArrowBuffer that shares the same underlying data with the +/// original shared buffer. The reference count is incremented. Returns EINVAL +/// if shared is not a shared buffer (i.e. ArrowIsSharedBuffer() returns false). +NANOARROW_DLL ArrowErrorCode ArrowSharedBufferClone(struct ArrowBuffer* shared, + struct ArrowBuffer* shared_out); + +/// \brief An opaque handle to a shared, reference-counted ArrowArray +struct ArrowSharedArray { + void* private_data; +}; + +/// \brief Initialize a shared array from an ArrowArray +/// +/// Takes ownership of src (sets src->release to NULL). The shared array +/// should be released with ArrowSharedArrayRelease() when all buffers have +/// been extracted from the source. The original array will be released +/// when all borrowed buffers have been released. +NANOARROW_DLL ArrowErrorCode ArrowSharedArrayInit(struct ArrowSharedArray* shared, + struct ArrowArray* src); + +/// \brief Release the caller's reference to a shared array +/// +/// Decrements the reference count. When no more references remain +/// (including any buffers obtained via ArrowSharedArrayBuffer()), +/// the underlying ArrowArray is released. +NANOARROW_DLL void ArrowSharedArrayRelease(struct ArrowSharedArray* shared); + +/// \brief Obtain a buffer from a shared array +/// +/// Returns an ArrowBuffer that views buffer i of the underlying ArrowArray. +/// If the source array was built with nanoarrow (i.e. ArrowArrayIsInternal() +/// returns true), buffer sizes are known and propagated to the output. +/// Otherwise the output buffer's size_bytes will be 0. The shared array's +/// reference count is incremented. The returned buffer is compatible with +/// ArrowSharedBufferClone(). +NANOARROW_DLL ArrowErrorCode ArrowSharedArrayBuffer(struct ArrowSharedArray* shared, + int64_t i, struct ArrowBuffer* out); + +/// @} + /// \brief Move the contents of an src ArrowSchema into dst and set src->release to NULL /// \ingroup nanoarrow-arrow-cdata static inline void ArrowSchemaMove(struct ArrowSchema* src, struct ArrowSchema* dst); @@ -2133,6 +2226,12 @@ static inline ArrowErrorCode ArrowArrayShrinkToFit(struct ArrowArray* array); NANOARROW_DLL ArrowErrorCode ArrowArrayFinishBuildingDefault(struct ArrowArray* array, struct ArrowError* error); +/// \brief Check if an ArrowArray was allocated by nanoarrow +/// +/// Returns non-zero if array was allocated using ArrowArrayInitFromType(), +/// ArrowArrayInitFromSchema(), or ArrowArrayInitFromArrayView(). +NANOARROW_DLL int ArrowArrayIsInternal(struct ArrowArray* array); + /// \brief Finish building an ArrowArray with explicit validation /// /// Finish building with an explicit validation level. This could perform less validation @@ -2144,6 +2243,28 @@ NANOARROW_DLL ArrowErrorCode ArrowArrayFinishBuilding( struct ArrowArray* array, enum ArrowValidationLevel validation_level, struct ArrowError* error); +/// \brief Create a shared copy of an ArrowArray +/// +/// Recursively converts array into a version where all buffers are +/// reference-counted. On success, shared is a new ArrowArray whose buffers +/// are backed by ArrowSharedArray references, and array is consumed +/// (release set to NULL). The resulting shared array can be safely moved +/// or have its buffers cloned via ArrowSharedBufferClone(). On error, +/// (e.g., failure to allocate a copy of a buffer), the input array is +/// released as it may have been partially moved at the point the error +/// occurs. +NANOARROW_DLL ArrowErrorCode ArrowArrayMoveShared(struct ArrowArray* array, + struct ArrowArray* shared); + +/// \brief Clone a shared ArrowArray +/// +/// Creates a new ArrowArray whose buffers share the same underlying data as +/// the source shared array. The source must have been created by +/// ArrowArrayMoveShared() (i.e. all buffers are reference-counted). +/// Returns EINVAL if shared is not a shared array. +NANOARROW_DLL ArrowErrorCode ArrowArrayCloneShared(struct ArrowArray* shared, + struct ArrowArray* array); + /// @} /// \defgroup nanoarrow-array-view Reading arrays diff --git a/c/vendor/nanoarrow/nanoarrow.hpp b/c/vendor/nanoarrow/nanoarrow.hpp index 595971b776..5c39a22f33 100644 --- a/c/vendor/nanoarrow/nanoarrow.hpp +++ b/c/vendor/nanoarrow/nanoarrow.hpp @@ -861,7 +861,7 @@ class ViewArrayAs { using const_iterator = typename internal::RandomAccessRange::const_iterator; const_iterator begin() const { return range_.begin(); } const_iterator end() const { return range_.end(); } - value_type operator[](int64_t i) const { return range_.get(i); } + value_type operator[](int64_t i) const { return range_.get(range_.offset + i); } }; /// \brief A range-for compatible wrapper for ArrowArray of binary or utf8 diff --git a/c/vendor/vendor_nanoarrow.sh b/c/vendor/vendor_nanoarrow.sh index d118992305..324bc8d9be 100755 --- a/c/vendor/vendor_nanoarrow.sh +++ b/c/vendor/vendor_nanoarrow.sh @@ -21,7 +21,7 @@ main() { local -r repo_url="https://github.com/apache/arrow-nanoarrow" # Check releases page: https://github.com/apache/arrow-nanoarrow/releases/ - local -r commit_sha=af347fa31d0e2dca6f6b6f818849a437244b6bf5 + local -r commit_sha=51b2880f8690f6198a6f842fab64291b7661ebea echo "Fetching $commit_sha from $repo_url" SCRATCH=$(mktemp -d)