From c263a0b9e6bdc1d4d73d62249fab4290afae7dfd Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 07:57:43 +0000 Subject: [PATCH 01/11] refactor(serde): move DataFile JSON serde into core --- src/iceberg/catalog/rest/json_serde.cc | 255 +---------------------- src/iceberg/json_serde.cc | 269 +++++++++++++++++++++++++ src/iceberg/json_serde_internal.h | 19 ++ src/iceberg/test/json_serde_test.cc | 36 ++++ 4 files changed, 326 insertions(+), 253 deletions(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..83a2dec09 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -105,35 +105,9 @@ constexpr std::string_view kEndSnapshotId = "end-snapshot-id"; constexpr std::string_view kStatsFields = "stats-fields"; constexpr std::string_view kMinRowsRequested = "min-rows-requested"; constexpr std::string_view kPlanTask = "plan-task"; -constexpr std::string_view kContent = "content"; -constexpr std::string_view kContentData = "data"; -constexpr std::string_view kContentPositionDeletes = "position-deletes"; -constexpr std::string_view kContentEqualityDeletes = "equality-deletes"; -constexpr std::string_view kFilePath = "file-path"; -constexpr std::string_view kFileFormat = "file-format"; -constexpr std::string_view kSpecId = "spec-id"; -constexpr std::string_view kPartition = "partition"; -constexpr std::string_view kRecordCount = "record-count"; -constexpr std::string_view kFileSizeInBytes = "file-size-in-bytes"; -constexpr std::string_view kColumnSizes = "column-sizes"; -constexpr std::string_view kValueCounts = "value-counts"; -constexpr std::string_view kNullValueCounts = "null-value-counts"; -constexpr std::string_view kNanValueCounts = "nan-value-counts"; -constexpr std::string_view kLowerBounds = "lower-bounds"; -constexpr std::string_view kUpperBounds = "upper-bounds"; -constexpr std::string_view kKeyMetadata = "key-metadata"; -constexpr std::string_view kSplitOffsets = "split-offsets"; -constexpr std::string_view kEqualityIds = "equality-ids"; -constexpr std::string_view kSortOrderId = "sort-order-id"; -constexpr std::string_view kFirstRowId = "first-row-id"; -constexpr std::string_view kReferencedDataFile = "referenced-data-file"; -constexpr std::string_view kContentOffset = "content-offset"; -constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; constexpr std::string_view kDataFile = "data-file"; constexpr std::string_view kDeleteFileReferences = "delete-file-references"; constexpr std::string_view kResidualFilter = "residual-filter"; -constexpr std::string_view kMapKeys = "keys"; -constexpr std::string_view kMapValues = "values"; Result StorageCredentialToJson(const StorageCredential& credential) { ICEBERG_RETURN_UNEXPECTED(credential.Validate()); @@ -152,47 +126,6 @@ Result StorageCredentialFromJson(const nlohmann::json& json) return credential; } -template -Result> KeyValueMapFromJson(const nlohmann::json& json, - std::string_view key) { - std::map result; - if (!json.contains(key) || json.at(key).is_null()) { - return result; - } - - ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); - ICEBERG_ASSIGN_OR_RAISE(auto keys, - GetJsonValue>(map_json, kMapKeys)); - ICEBERG_ASSIGN_OR_RAISE(auto values, - GetJsonValue>(map_json, kMapValues)); - if (keys.size() != values.size()) { - return JsonParseError("'{}' map keys and values have different lengths", key); - } - - for (size_t i = 0; i < keys.size(); ++i) { - result[keys[i]] = std::move(values[i]); - } - return result; -} - -template -void SetKeyValueMap(nlohmann::json& json, std::string_view key, - const std::map& map) { - if (map.empty()) { - return; - } - - std::vector keys; - std::vector values; - keys.reserve(map.size()); - values.reserve(map.size()); - for (const auto& [field_id, value] : map) { - keys.push_back(field_id); - values.push_back(value); - } - json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; -} - } // namespace Result DataFileFromJson( @@ -200,106 +133,7 @@ Result DataFileFromJson( const std::unordered_map>& partition_spec_by_id, const Schema& schema) { - if (!json.is_object()) { - return JsonParseError("DataFile must be a JSON object: {}", SafeDumpJson(json)); - } - DataFile data_file; - - ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); - if (content_str == kContentData) { - data_file.content = DataFile::Content::kData; - } else if (content_str == kContentPositionDeletes) { - data_file.content = DataFile::Content::kPositionDeletes; - } else if (content_str == kContentEqualityDeletes) { - data_file.content = DataFile::Content::kEqualityDeletes; - } else { - return JsonParseError("Unknown data file content: {}", content_str); - } - - ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, - GetJsonValue(json, kFilePath)); - ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); - ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); - - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); - data_file.partition_spec_id = spec_id; - - ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, - GetJsonValue(json, kPartition)); - if (!partition_vals.is_array()) { - return JsonParseError("PartitionValues must be a JSON array: {}", - SafeDumpJson(partition_vals)); - } - std::vector literals; - auto it = partition_spec_by_id.find(spec_id); - if (it == partition_spec_by_id.end()) { - return JsonParseError("Invalid partition spec id: {}", spec_id); - } - ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); - auto fields = struct_type->fields(); - if (partition_vals.size() != fields.size()) { - return JsonParseError("Invalid partition data size: expected = {}, actual = {}", - fields.size(), partition_vals.size()); - } - for (size_t pos = 0; pos < fields.size(); ++pos) { - ICEBERG_ASSIGN_OR_RAISE( - auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get())); - literals.push_back(std::move(literal)); - } - data_file.partition = PartitionValues(std::move(literals)); - - ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, - GetJsonValue(json, kRecordCount)); - ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, - GetJsonValue(json, kFileSizeInBytes)); - - ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, - KeyValueMapFromJson(json, kColumnSizes)); - ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts, - KeyValueMapFromJson(json, kValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts, - KeyValueMapFromJson(json, kNullValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts, - KeyValueMapFromJson(json, kNanValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, - KeyValueMapFromJson>(json, kLowerBounds)); - ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, - KeyValueMapFromJson>(json, kUpperBounds)); - - if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, - GetJsonValue>(json, kKeyMetadata)); - } - if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, - GetJsonValue>(json, kSplitOffsets)); - } - if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, - GetJsonValue>(json, kEqualityIds)); - } - if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, - GetJsonValue(json, kSortOrderId)); - } - if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, - GetJsonValue(json, kFirstRowId)); - } - if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, - GetJsonValue(json, kReferencedDataFile)); - } - if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, - GetJsonValue(json, kContentOffset)); - } - if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { - ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, - GetJsonValue(json, kContentSizeInBytes)); - } - - return data_file; + return iceberg::DataFileFromJson(json, partition_spec_by_id, schema); } Result>> FileScanTasksFromJson( @@ -357,97 +191,12 @@ Result>> FileScanTasksFromJson( return file_scan_tasks; } -Result DataFileToJsonUnchecked(const DataFile& data_file) { - nlohmann::json json; - switch (data_file.content) { - case DataFile::Content::kData: - json[kContent] = kContentData; - break; - case DataFile::Content::kPositionDeletes: - json[kContent] = kContentPositionDeletes; - break; - case DataFile::Content::kEqualityDeletes: - json[kContent] = kContentEqualityDeletes; - break; - } - json[kFilePath] = data_file.file_path; - json[kFileFormat] = ToString(data_file.file_format); - - if (!data_file.partition_spec_id.has_value()) { - return ValidationFailed("Cannot serialize REST content file without 'spec-id'"); - } - json[kSpecId] = data_file.partition_spec_id.value(); - - nlohmann::json partition_json = nlohmann::json::array(); - for (const auto& literal : data_file.partition.values()) { - ICEBERG_ASSIGN_OR_RAISE(auto lit_json, iceberg::ToJson(literal)); - partition_json.push_back(std::move(lit_json)); - } - json[kPartition] = std::move(partition_json); - - json[kRecordCount] = data_file.record_count; - json[kFileSizeInBytes] = data_file.file_size_in_bytes; - - SetKeyValueMap(json, kColumnSizes, data_file.column_sizes); - SetKeyValueMap(json, kValueCounts, data_file.value_counts); - SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts); - SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts); - SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds); - SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds); - - if (!data_file.key_metadata.empty()) { - json[kKeyMetadata] = data_file.key_metadata; - } - if (!data_file.split_offsets.empty()) { - json[kSplitOffsets] = data_file.split_offsets; - } - if (!data_file.equality_ids.empty()) { - json[kEqualityIds] = data_file.equality_ids; - } - if (data_file.sort_order_id.has_value()) { - json[kSortOrderId] = data_file.sort_order_id.value(); - } - if (data_file.first_row_id.has_value()) { - json[kFirstRowId] = data_file.first_row_id.value(); - } - if (data_file.referenced_data_file.has_value()) { - json[kReferencedDataFile] = data_file.referenced_data_file.value(); - } - if (data_file.content_offset.has_value()) { - json[kContentOffset] = data_file.content_offset.value(); - } - if (data_file.content_size_in_bytes.has_value()) { - json[kContentSizeInBytes] = data_file.content_size_in_bytes.value(); - } - - return json; -} - Result ToJson( const DataFile& data_file, const std::unordered_map>& partition_specs_by_id, const Schema& schema) { - if (!data_file.partition_spec_id.has_value()) { - return ValidationFailed("Invalid partition spec id from content file: null"); - } - auto it = partition_specs_by_id.find(data_file.partition_spec_id.value()); - if (it == partition_specs_by_id.end() || !it->second) { - return ValidationFailed("Invalid partition spec: null"); - } - if (data_file.partition_spec_id.value() != it->second->spec_id()) { - return ValidationFailed( - "Invalid partition spec id from content file: expected = {}, actual = {}", - it->second->spec_id(), data_file.partition_spec_id.value()); - } - ICEBERG_ASSIGN_OR_RAISE(auto partition_type, it->second->PartitionType(schema)); - if (data_file.partition.num_fields() != partition_type->fields().size()) { - return ValidationFailed( - "Invalid partition data from content file: expected = {}, actual = {}", - partition_type->fields().empty() ? "unpartitioned" : "partitioned", - data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); - } - return DataFileToJsonUnchecked(data_file); + return iceberg::ToJson(data_file, partition_specs_by_id, schema); } namespace { diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index f322996ec..c2e685ef4 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -20,6 +20,7 @@ #include #include #include +#include #include #include #include @@ -29,12 +30,16 @@ #include "iceberg/constants.h" #include "iceberg/expression/json_serde_internal.h" #include "iceberg/expression/literal.h" +#include "iceberg/file_format.h" #include "iceberg/json_serde_internal.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/name_mapping.h" #include "iceberg/partition_field.h" #include "iceberg/partition_spec.h" #include "iceberg/result.h" +#include "iceberg/row/partition_values.h" #include "iceberg/schema.h" +#include "iceberg/schema_field.h" #include "iceberg/snapshot.h" #include "iceberg/sort_order.h" #include "iceberg/statistics_file.h" @@ -240,6 +245,137 @@ constexpr std::string_view kRequirementAssertDefaultSortOrderID = constexpr std::string_view kLastAssignedFieldId = "last-assigned-field-id"; constexpr std::string_view kLastAssignedPartitionId = "last-assigned-partition-id"; +// DataFile JSON (Iceberg REST ContentFile field names) +constexpr std::string_view kContent = "content"; +constexpr std::string_view kContentData = "data"; +constexpr std::string_view kContentPositionDeletes = "position-deletes"; +constexpr std::string_view kContentEqualityDeletes = "equality-deletes"; +constexpr std::string_view kFilePath = "file-path"; +constexpr std::string_view kFileFormat = "file-format"; +constexpr std::string_view kPartition = "partition"; +constexpr std::string_view kRecordCount = "record-count"; +constexpr std::string_view kColumnSizes = "column-sizes"; +constexpr std::string_view kValueCounts = "value-counts"; +constexpr std::string_view kNullValueCounts = "null-value-counts"; +constexpr std::string_view kNanValueCounts = "nan-value-counts"; +constexpr std::string_view kLowerBounds = "lower-bounds"; +constexpr std::string_view kUpperBounds = "upper-bounds"; +constexpr std::string_view kKeyMetadata = "key-metadata"; +constexpr std::string_view kSplitOffsets = "split-offsets"; +constexpr std::string_view kEqualityIds = "equality-ids"; +constexpr std::string_view kReferencedDataFile = "referenced-data-file"; +constexpr std::string_view kContentOffset = "content-offset"; +constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; +constexpr std::string_view kMapKeys = "keys"; +constexpr std::string_view kMapValues = "values"; + +template +Result> KeyValueMapFromJson(const nlohmann::json& json, + std::string_view key) { + std::map result; + if (!json.contains(key) || json.at(key).is_null()) { + return result; + } + + ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); + ICEBERG_ASSIGN_OR_RAISE(auto keys, + GetJsonValue>(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values, + GetJsonValue>(map_json, kMapValues)); + if (keys.size() != values.size()) { + return JsonParseError("'{}' map keys and values have different lengths", key); + } + + for (size_t i = 0; i < keys.size(); ++i) { + result[keys[i]] = std::move(values[i]); + } + return result; +} + +template +void SetKeyValueMap(nlohmann::json& json, std::string_view key, + const std::map& map) { + if (map.empty()) { + return; + } + + std::vector keys; + std::vector values; + keys.reserve(map.size()); + values.reserve(map.size()); + for (const auto& [field_id, value] : map) { + keys.push_back(field_id); + values.push_back(value); + } + json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; +} + +Result DataFileToJsonUnchecked(const DataFile& data_file) { + nlohmann::json json; + switch (data_file.content) { + case DataFile::Content::kData: + json[kContent] = kContentData; + break; + case DataFile::Content::kPositionDeletes: + json[kContent] = kContentPositionDeletes; + break; + case DataFile::Content::kEqualityDeletes: + json[kContent] = kContentEqualityDeletes; + break; + } + json[kFilePath] = data_file.file_path; + json[kFileFormat] = ToString(data_file.file_format); + + if (!data_file.partition_spec_id.has_value()) { + return ValidationFailed("Cannot serialize REST content file without 'spec-id'"); + } + json[kSpecId] = data_file.partition_spec_id.value(); + + nlohmann::json partition_json = nlohmann::json::array(); + for (const auto& literal : data_file.partition.values()) { + ICEBERG_ASSIGN_OR_RAISE(auto lit_json, ToJson(literal)); + partition_json.push_back(std::move(lit_json)); + } + json[kPartition] = std::move(partition_json); + + json[kRecordCount] = data_file.record_count; + json[kFileSizeInBytes] = data_file.file_size_in_bytes; + + SetKeyValueMap(json, kColumnSizes, data_file.column_sizes); + SetKeyValueMap(json, kValueCounts, data_file.value_counts); + SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts); + SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts); + SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds); + SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds); + + if (!data_file.key_metadata.empty()) { + json[kKeyMetadata] = data_file.key_metadata; + } + if (!data_file.split_offsets.empty()) { + json[kSplitOffsets] = data_file.split_offsets; + } + if (!data_file.equality_ids.empty()) { + json[kEqualityIds] = data_file.equality_ids; + } + if (data_file.sort_order_id.has_value()) { + json[kSortOrderId] = data_file.sort_order_id.value(); + } + if (data_file.first_row_id.has_value()) { + json[kFirstRowId] = data_file.first_row_id.value(); + } + if (data_file.referenced_data_file.has_value()) { + json[kReferencedDataFile] = data_file.referenced_data_file.value(); + } + if (data_file.content_offset.has_value()) { + json[kContentOffset] = data_file.content_offset.value(); + } + if (data_file.content_size_in_bytes.has_value()) { + json[kContentSizeInBytes] = data_file.content_size_in_bytes.value(); + } + + return json; +} + } // namespace nlohmann::json ToJson(const SortField& sort_field) { @@ -1985,4 +2121,137 @@ Result> TableRequirementFromJson( return JsonParseError("Unknown table requirement type: {}", type); } +Result DataFileFromJson( + const nlohmann::json& json, + const std::unordered_map>& + partition_spec_by_id, + const Schema& schema) { + if (!json.is_object()) { + return JsonParseError("DataFile must be a JSON object: {}", SafeDumpJson(json)); + } + DataFile data_file; + + ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); + if (content_str == kContentData) { + data_file.content = DataFile::Content::kData; + } else if (content_str == kContentPositionDeletes) { + data_file.content = DataFile::Content::kPositionDeletes; + } else if (content_str == kContentEqualityDeletes) { + data_file.content = DataFile::Content::kEqualityDeletes; + } else { + return JsonParseError("Unknown data file content: {}", content_str); + } + + ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, GetJsonValue(json, kFilePath)); + ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); + + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + data_file.partition_spec_id = spec_id; + + ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, + GetJsonValue(json, kPartition)); + if (!partition_vals.is_array()) { + return JsonParseError("PartitionValues must be a JSON array: {}", + SafeDumpJson(partition_vals)); + } + std::vector literals; + auto it = partition_spec_by_id.find(spec_id); + if (it == partition_spec_by_id.end()) { + return JsonParseError("Invalid partition spec id: {}", spec_id); + } + ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); + auto fields = struct_type->fields(); + if (partition_vals.size() != fields.size()) { + return JsonParseError("Invalid partition data size: expected = {}, actual = {}", + fields.size(), partition_vals.size()); + } + for (size_t pos = 0; pos < fields.size(); ++pos) { + ICEBERG_ASSIGN_OR_RAISE( + auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get())); + literals.push_back(std::move(literal)); + } + data_file.partition = PartitionValues(std::move(literals)); + + ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, + GetJsonValue(json, kRecordCount)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, + GetJsonValue(json, kFileSizeInBytes)); + + ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, + KeyValueMapFromJson(json, kColumnSizes)); + ICEBERG_ASSIGN_OR_RAISE(data_file.value_counts, + KeyValueMapFromJson(json, kValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.null_value_counts, + KeyValueMapFromJson(json, kNullValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts, + KeyValueMapFromJson(json, kNanValueCounts)); + ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, + KeyValueMapFromJson>(json, kLowerBounds)); + ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, + KeyValueMapFromJson>(json, kUpperBounds)); + + if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, + GetJsonValue>(json, kKeyMetadata)); + } + if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, + GetJsonValue>(json, kSplitOffsets)); + } + if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, + GetJsonValue>(json, kEqualityIds)); + } + if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, + GetJsonValue(json, kSortOrderId)); + } + if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, + GetJsonValue(json, kFirstRowId)); + } + if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, + GetJsonValue(json, kReferencedDataFile)); + } + if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, + GetJsonValue(json, kContentOffset)); + } + if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, + GetJsonValue(json, kContentSizeInBytes)); + } + + return data_file; +} + +Result ToJson( + const DataFile& data_file, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema) { + if (!data_file.partition_spec_id.has_value()) { + return ValidationFailed("Invalid partition spec id from content file: null"); + } + auto it = partition_specs_by_id.find(data_file.partition_spec_id.value()); + if (it == partition_specs_by_id.end() || !it->second) { + return ValidationFailed("Invalid partition spec: null"); + } + if (data_file.partition_spec_id.value() != it->second->spec_id()) { + return ValidationFailed( + "Invalid partition spec id from content file: expected = {}, actual = {}", + it->second->spec_id(), data_file.partition_spec_id.value()); + } + ICEBERG_ASSIGN_OR_RAISE(auto partition_type, it->second->PartitionType(schema)); + if (data_file.partition.num_fields() != partition_type->fields().size()) { + return ValidationFailed( + "Invalid partition data from content file: expected = {}, actual = {}", + partition_type->fields().empty() ? "unpartitioned" : "partitioned", + data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); + } + return DataFileToJsonUnchecked(data_file); +} + } // namespace iceberg diff --git a/src/iceberg/json_serde_internal.h b/src/iceberg/json_serde_internal.h index 1a30e8e8f..3aaced5c9 100644 --- a/src/iceberg/json_serde_internal.h +++ b/src/iceberg/json_serde_internal.h @@ -29,6 +29,7 @@ #include #include "iceberg/encryption/encrypted_key.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/result.h" #include "iceberg/statistics_file.h" #include "iceberg/table_metadata.h" @@ -426,4 +427,22 @@ ICEBERG_EXPORT nlohmann::json ToJson(const TableRequirement& requirement); ICEBERG_EXPORT Result> TableRequirementFromJson( const nlohmann::json& json); +/// \brief Serializes a `DataFile` (content file) to JSON. +/// +/// The JSON object uses the Iceberg REST ContentFile field names (`spec-id`, +/// `file-path`, `file-format`, maps as `{keys, values}`, and so on). `spec-id` is +/// required. Partition values are encoded with the given partition spec and schema. +ICEBERG_EXPORT Result ToJson( + const DataFile& data_file, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema); + +/// \brief Deserializes a JSON object into a `DataFile`. +ICEBERG_EXPORT Result DataFileFromJson( + const nlohmann::json& json, + const std::unordered_map>& + partition_spec_by_id, + const Schema& schema); + } // namespace iceberg diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 22b0a08de..b15d08790 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,15 +18,19 @@ */ #include +#include #include #include #include #include "iceberg/expression/literal.h" +#include "iceberg/file_format.h" #include "iceberg/json_serde_internal.h" +#include "iceberg/manifest/manifest_entry.h" #include "iceberg/name_mapping.h" #include "iceberg/partition_spec.h" +#include "iceberg/row/partition_values.h" #include "iceberg/schema.h" #include "iceberg/schema_field.h" #include "iceberg/snapshot.h" @@ -1044,4 +1048,36 @@ TEST(TableRequirementJsonTest, TableRequirementUnknownType) { EXPECT_THAT(result, HasErrorMessage("Unknown table requirement type")); } +std::unordered_map> UnpartitionedSpecs() { + return {{PartitionSpec::kInitialSpecId, PartitionSpec::Unpartitioned()}}; +} + +DataFile MakeUnpartitionedDataFile(std::string path, int64_t record_count = 100, + int64_t file_size = 12345) { + DataFile data_file; + data_file.content = DataFile::Content::kData; + data_file.file_path = std::move(path); + data_file.file_format = FileFormatType::kParquet; + data_file.partition_spec_id = PartitionSpec::kInitialSpecId; + data_file.partition = PartitionValues{}; + data_file.record_count = record_count; + data_file.file_size_in_bytes = file_size; + return data_file; +} + +TEST(DataFileJsonTest, RoundTripRequiredFields) { + auto data_file = MakeUnpartitionedDataFile("s3://bucket/data/file.parquet"); + Schema schema({}, 0); + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(data_file, UnpartitionedSpecs(), schema)); + EXPECT_EQ(json["content"], "data"); + EXPECT_EQ(json["file-path"], "s3://bucket/data/file.parquet"); + EXPECT_EQ(json["spec-id"], 0); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, UnpartitionedSpecs(), schema)); + EXPECT_EQ(parsed.file_path, data_file.file_path); + EXPECT_EQ(parsed.record_count, data_file.record_count); + EXPECT_EQ(parsed.file_size_in_bytes, data_file.file_size_in_bytes); + EXPECT_EQ(parsed.partition_spec_id, data_file.partition_spec_id); +} + } // namespace iceberg From ece593036953b7ed15f02ae2fc28a97b068e4ec4 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 08:54:58 +0000 Subject: [PATCH 02/11] fix(serde): disambiguate REST DataFile calls --- src/iceberg/catalog/rest/json_serde.cc | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 83a2dec09..46593f0b5 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -156,7 +156,8 @@ Result>> FileScanTasksFromJson( ICEBERG_ASSIGN_OR_RAISE(auto data_file_json, GetJsonValue(task_json, kDataFile)); ICEBERG_ASSIGN_OR_RAISE( - auto data_file, DataFileFromJson(data_file_json, partition_spec_by_id, schema)); + auto data_file, + iceberg::rest::DataFileFromJson(data_file_json, partition_spec_by_id, schema)); // FIXME: REST scan-task DataFile JSON currently carries first-row-id, // but not the manifest-entry data sequence number. Until the REST API exposes // it, REST-planned tasks cannot inherit _last_updated_sequence_number. @@ -244,7 +245,7 @@ Result ScanTaskFieldsToJson( if (task->data_file()) { ICEBERG_ASSIGN_OR_RAISE( auto data_file_json, - ToJson(*task->data_file(), partition_specs_by_id, schema)); + iceberg::rest::ToJson(*task->data_file(), partition_specs_by_id, schema)); task_json[kDataFile] = std::move(data_file_json); } if (!task->delete_files().empty()) { @@ -270,7 +271,8 @@ Result ScanTaskFieldsToJson( } nlohmann::json delete_files_json = nlohmann::json::array(); for (const auto& file : delete_files) { - ICEBERG_ASSIGN_OR_RAISE(auto df_json, ToJson(*file, partition_specs_by_id, schema)); + ICEBERG_ASSIGN_OR_RAISE(auto df_json, + iceberg::rest::ToJson(*file, partition_specs_by_id, schema)); delete_files_json.push_back(std::move(df_json)); } if (!delete_files_json.empty()) { @@ -306,8 +308,9 @@ Status ScanTaskFieldsFromJson( SafeDumpJson(delete_files_json)); } for (const auto& entry_json : delete_files_json) { - ICEBERG_ASSIGN_OR_RAISE(auto delete_file, - DataFileFromJson(entry_json, partition_specs_by_id, schema)); + ICEBERG_ASSIGN_OR_RAISE( + auto delete_file, + iceberg::rest::DataFileFromJson(entry_json, partition_specs_by_id, schema)); response.delete_files.push_back(std::make_shared(std::move(delete_file))); } From 5304e282d76b70e84366f5efd32defa71bdb7f4e Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 09:16:33 +0000 Subject: [PATCH 03/11] fix clang error --- src/iceberg/json_serde.cc | 3 ++- src/iceberg/test/json_serde_test.cc | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index c2e685ef4..1f5caf8f0 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -2142,7 +2142,8 @@ Result DataFileFromJson( return JsonParseError("Unknown data file content: {}", content_str); } - ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, GetJsonValue(json, kFilePath)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_path, + GetJsonValue(json, kFilePath)); ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index b15d08790..c7779d703 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1073,7 +1073,8 @@ TEST(DataFileJsonTest, RoundTripRequiredFields) { EXPECT_EQ(json["file-path"], "s3://bucket/data/file.parquet"); EXPECT_EQ(json["spec-id"], 0); - ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, UnpartitionedSpecs(), schema)); + ICEBERG_UNWRAP_OR_FAIL(auto parsed, + DataFileFromJson(json, UnpartitionedSpecs(), schema)); EXPECT_EQ(parsed.file_path, data_file.file_path); EXPECT_EQ(parsed.record_count, data_file.record_count); EXPECT_EQ(parsed.file_size_in_bytes, data_file.file_size_in_bytes); From 2c48f026a40f384187e698f80247f1c7ed355a77 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 08:02:10 +0000 Subject: [PATCH 04/11] fix(serde): align DataFile JSON with Java --- src/iceberg/json_serde.cc | 273 +++++++++++++++++++---- src/iceberg/test/json_serde_test.cc | 108 +++++++++ src/iceberg/test/rest_json_serde_test.cc | 10 +- 3 files changed, 339 insertions(+), 52 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index 1f5caf8f0..45dfa29f5 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -18,10 +18,13 @@ */ #include +#include #include #include +#include #include #include +#include #include #include @@ -269,7 +272,58 @@ constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; constexpr std::string_view kMapKeys = "keys"; constexpr std::string_view kMapValues = "values"; -template +template +Result IntegerFromJson(const nlohmann::json& json, + std::string_view description) { + if (!json.is_number_integer()) { + return JsonParseError("{} must be an integer, but is {}", description, + json.type_name()); + } + if (json.is_number_unsigned()) { + const auto value = json.get(); + if (value > static_cast(std::numeric_limits::max())) { + return JsonParseError("{} integer is out of range: {}", description, + SafeDumpJson(json)); + } + return static_cast(value); + } + + const auto value = json.get(); + if (value < static_cast(std::numeric_limits::min()) || + value > static_cast(std::numeric_limits::max())) { + return JsonParseError("{} integer is out of range: {}", description, + SafeDumpJson(json)); + } + return static_cast(value); +} + +template +Result GetJsonInteger(const nlohmann::json& json, std::string_view key) { + ICEBERG_ASSIGN_OR_RAISE(auto value, GetJsonValue(json, key)); + return IntegerFromJson(value, std::format("'{}'", key)); +} + +template +Result> IntegerVectorFromJson(const nlohmann::json& json, + std::string_view key) { + ICEBERG_ASSIGN_OR_RAISE(auto values_json, GetJsonValue(json, key)); + if (!values_json.is_array()) { + return JsonParseError("'{}' must be an array, but is {}", key, + values_json.type_name()); + } + + std::vector values; + values.reserve(values_json.size()); + for (const auto& value_json : values_json) { + ICEBERG_ASSIGN_OR_RAISE( + auto value, + IntegerFromJson(value_json, std::format("'{}' element", key))); + values.push_back(value); + } + return values; +} + +template Result> KeyValueMapFromJson(const nlohmann::json& json, std::string_view key) { std::map result; @@ -278,16 +332,27 @@ Result> KeyValueMapFromJson(const nlohmann::json& json, } ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); - ICEBERG_ASSIGN_OR_RAISE(auto keys, - GetJsonValue>(map_json, kMapKeys)); - ICEBERG_ASSIGN_OR_RAISE(auto values, - GetJsonValue>(map_json, kMapValues)); - if (keys.size() != values.size()) { + ICEBERG_ASSIGN_OR_RAISE(auto keys_json, + GetJsonValue(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values_json, + GetJsonValue(map_json, kMapValues)); + if (!keys_json.is_array() || !values_json.is_array()) { + return JsonParseError("'{}' map keys and values must be arrays", key); + } + if (keys_json.size() != values_json.size()) { return JsonParseError("'{}' map keys and values have different lengths", key); } - for (size_t i = 0; i < keys.size(); ++i) { - result[keys[i]] = std::move(values[i]); + for (size_t i = 0; i < keys_json.size(); ++i) { + ICEBERG_ASSIGN_OR_RAISE( + auto field_id, + IntegerFromJson(keys_json[i], std::format("'{}' map key", key))); + ICEBERG_ASSIGN_OR_RAISE( + auto value, + IntegerFromJson(values_json[i], std::format("'{}' map value", key))); + if (!result.emplace(field_id, value).second) { + return JsonParseError("'{}' map contains duplicate key {}", key, field_id); + } } return result; } @@ -310,6 +375,118 @@ void SetKeyValueMap(nlohmann::json& json, std::string_view key, json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; } +std::string BytesToHex(const std::vector& bytes) { + std::string hex; + hex.reserve(bytes.size() * 2); + for (uint8_t byte : bytes) { + hex += std::format("{:02X}", byte); + } + return hex; +} + +Result>> BytesMapFromJson( + const nlohmann::json& json, std::string_view key) { + std::map> result; + if (!json.contains(key) || json.at(key).is_null()) { + return result; + } + + ICEBERG_ASSIGN_OR_RAISE(auto map_json, GetJsonValue(json, key)); + ICEBERG_ASSIGN_OR_RAISE(auto keys_json, + GetJsonValue(map_json, kMapKeys)); + ICEBERG_ASSIGN_OR_RAISE(auto values_json, + GetJsonValue(map_json, kMapValues)); + if (!keys_json.is_array() || !values_json.is_array()) { + return JsonParseError("'{}' map keys and values must be arrays", key); + } + if (keys_json.size() != values_json.size()) { + return JsonParseError("'{}' map keys and values have different lengths", key); + } + + for (size_t i = 0; i < keys_json.size(); ++i) { + ICEBERG_ASSIGN_OR_RAISE( + auto field_id, + IntegerFromJson(keys_json[i], std::format("'{}' map key", key))); + if (!values_json[i].is_string()) { + return JsonParseError("'{}' map value must be a string, but is {}", key, + values_json[i].type_name()); + } + ICEBERG_ASSIGN_OR_RAISE( + auto bytes, StringUtils::HexStringToBytes(values_json[i].get())); + if (!result.emplace(field_id, std::move(bytes)).second) { + return JsonParseError("'{}' map contains duplicate key {}", key, field_id); + } + } + return result; +} + +void SetBytesMap(nlohmann::json& json, std::string_view key, + const std::map>& map) { + if (map.empty()) { + return; + } + + std::vector keys; + std::vector values; + keys.reserve(map.size()); + values.reserve(map.size()); + for (const auto& [field_id, value] : map) { + keys.push_back(field_id); + values.push_back(BytesToHex(value)); + } + json[key] = {{kMapKeys, std::move(keys)}, {kMapValues, std::move(values)}}; +} + +Result PartitionLiteralFromJson(const nlohmann::json* value, + const SchemaField& field) { + if (!field.type() || !field.type()->is_primitive()) { + return InvalidSchema("Partition field {} must have a primitive type", + field.field_id()); + } + + auto primitive_type = internal::checked_pointer_cast(field.type()); + if (value == nullptr || value->is_null()) { + return Literal::Null(std::move(primitive_type)); + } + return LiteralFromJson(*value, primitive_type.get()); +} + +Result PartitionValuesFromJson(const nlohmann::json& json, + const StructType& partition_type) { + const auto& fields = partition_type.fields(); + std::vector values; + values.reserve(fields.size()); + + if (json.is_array()) { + if (json.size() != fields.size()) { + return JsonParseError("Invalid partition data size: expected = {}, actual = {}", + fields.size(), json.size()); + } + for (size_t pos = 0; pos < fields.size(); ++pos) { + ICEBERG_ASSIGN_OR_RAISE(auto value, + PartitionLiteralFromJson(&json[pos], fields[pos])); + values.push_back(std::move(value)); + } + } else if (json.is_object()) { + if (json.size() > fields.size()) { + return JsonParseError("Invalid partition data size: expected <= {}, actual = {}", + fields.size(), json.size()); + } + for (const auto& field : fields) { + const auto it = json.find(std::to_string(field.field_id())); + const nlohmann::json* value = it == json.end() ? nullptr : &it.value(); + ICEBERG_ASSIGN_OR_RAISE(auto literal, PartitionLiteralFromJson(value, field)); + values.push_back(std::move(literal)); + } + } else { + return JsonParseError( + "Invalid partition data for content file: expected array or object ({})", + SafeDumpJson(json)); + } + + return PartitionValues(std::move(values)); +} + Result DataFileToJsonUnchecked(const DataFile& data_file) { nlohmann::json json; switch (data_file.content) { @@ -327,7 +504,7 @@ Result DataFileToJsonUnchecked(const DataFile& data_file) { json[kFileFormat] = ToString(data_file.file_format); if (!data_file.partition_spec_id.has_value()) { - return ValidationFailed("Cannot serialize REST content file without 'spec-id'"); + return ValidationFailed("Cannot serialize content file without 'spec-id'"); } json[kSpecId] = data_file.partition_spec_id.value(); @@ -345,11 +522,11 @@ Result DataFileToJsonUnchecked(const DataFile& data_file) { SetKeyValueMap(json, kValueCounts, data_file.value_counts); SetKeyValueMap(json, kNullValueCounts, data_file.null_value_counts); SetKeyValueMap(json, kNanValueCounts, data_file.nan_value_counts); - SetKeyValueMap(json, kLowerBounds, data_file.lower_bounds); - SetKeyValueMap(json, kUpperBounds, data_file.upper_bounds); + SetBytesMap(json, kLowerBounds, data_file.lower_bounds); + SetBytesMap(json, kUpperBounds, data_file.upper_bounds); if (!data_file.key_metadata.empty()) { - json[kKeyMetadata] = data_file.key_metadata; + json[kKeyMetadata] = BytesToHex(data_file.key_metadata); } if (!data_file.split_offsets.empty()) { json[kSplitOffsets] = data_file.split_offsets; @@ -2132,11 +2309,13 @@ Result DataFileFromJson( DataFile data_file; ICEBERG_ASSIGN_OR_RAISE(auto content_str, GetJsonValue(json, kContent)); - if (content_str == kContentData) { + if (content_str == kContentData || content_str == "DATA") { data_file.content = DataFile::Content::kData; - } else if (content_str == kContentPositionDeletes) { + } else if (content_str == kContentPositionDeletes || + content_str == "POSITION_DELETES") { data_file.content = DataFile::Content::kPositionDeletes; - } else if (content_str == kContentEqualityDeletes) { + } else if (content_str == kContentEqualityDeletes || + content_str == "EQUALITY_DELETES") { data_file.content = DataFile::Content::kEqualityDeletes; } else { return JsonParseError("Unknown data file content: {}", content_str); @@ -2147,37 +2326,25 @@ Result DataFileFromJson( ICEBERG_ASSIGN_OR_RAISE(auto format_str, GetJsonValue(json, kFileFormat)); ICEBERG_ASSIGN_OR_RAISE(data_file.file_format, FileFormatTypeFromString(format_str)); - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(json, kSpecId)); data_file.partition_spec_id = spec_id; - ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, - GetJsonValue(json, kPartition)); - if (!partition_vals.is_array()) { - return JsonParseError("PartitionValues must be a JSON array: {}", - SafeDumpJson(partition_vals)); - } - std::vector literals; auto it = partition_spec_by_id.find(spec_id); - if (it == partition_spec_by_id.end()) { + if (it == partition_spec_by_id.end() || !it->second) { return JsonParseError("Invalid partition spec id: {}", spec_id); } - ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); - auto fields = struct_type->fields(); - if (partition_vals.size() != fields.size()) { - return JsonParseError("Invalid partition data size: expected = {}, actual = {}", - fields.size(), partition_vals.size()); - } - for (size_t pos = 0; pos < fields.size(); ++pos) { - ICEBERG_ASSIGN_OR_RAISE( - auto literal, LiteralFromJson(partition_vals[pos], fields[pos].type().get())); - literals.push_back(std::move(literal)); + if (json.contains(kPartition)) { + ICEBERG_ASSIGN_OR_RAISE(auto partition_vals, + GetJsonValue(json, kPartition)); + ICEBERG_ASSIGN_OR_RAISE(auto struct_type, it->second->PartitionType(schema)); + ICEBERG_ASSIGN_OR_RAISE(data_file.partition, + PartitionValuesFromJson(partition_vals, *struct_type)); } - data_file.partition = PartitionValues(std::move(literals)); ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, - GetJsonValue(json, kRecordCount)); + GetJsonInteger(json, kRecordCount)); ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, - GetJsonValue(json, kFileSizeInBytes)); + GetJsonInteger(json, kFileSizeInBytes)); ICEBERG_ASSIGN_OR_RAISE(data_file.column_sizes, KeyValueMapFromJson(json, kColumnSizes)); @@ -2187,30 +2354,30 @@ Result DataFileFromJson( KeyValueMapFromJson(json, kNullValueCounts)); ICEBERG_ASSIGN_OR_RAISE(data_file.nan_value_counts, KeyValueMapFromJson(json, kNanValueCounts)); - ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, - KeyValueMapFromJson>(json, kLowerBounds)); - ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, - KeyValueMapFromJson>(json, kUpperBounds)); + ICEBERG_ASSIGN_OR_RAISE(data_file.lower_bounds, BytesMapFromJson(json, kLowerBounds)); + ICEBERG_ASSIGN_OR_RAISE(data_file.upper_bounds, BytesMapFromJson(json, kUpperBounds)); if (json.contains(kKeyMetadata) && !json.at(kKeyMetadata).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(auto key_metadata, + GetJsonValue(json, kKeyMetadata)); ICEBERG_ASSIGN_OR_RAISE(data_file.key_metadata, - GetJsonValue>(json, kKeyMetadata)); + StringUtils::HexStringToBytes(key_metadata)); } if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, - GetJsonValue>(json, kSplitOffsets)); + IntegerVectorFromJson(json, kSplitOffsets)); } if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, - GetJsonValue>(json, kEqualityIds)); + IntegerVectorFromJson(json, kEqualityIds)); } if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, - GetJsonValue(json, kSortOrderId)); + GetJsonInteger(json, kSortOrderId)); } if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, - GetJsonValue(json, kFirstRowId)); + GetJsonInteger(json, kFirstRowId)); } if (json.contains(kReferencedDataFile) && !json.at(kReferencedDataFile).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.referenced_data_file, @@ -2218,11 +2385,11 @@ Result DataFileFromJson( } if (json.contains(kContentOffset) && !json.at(kContentOffset).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.content_offset, - GetJsonValue(json, kContentOffset)); + GetJsonInteger(json, kContentOffset)); } if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, - GetJsonValue(json, kContentSizeInBytes)); + GetJsonInteger(json, kContentSizeInBytes)); } return data_file; @@ -2252,6 +2419,16 @@ Result ToJson( partition_type->fields().empty() ? "unpartitioned" : "partitioned", data_file.partition.num_fields() == 0 ? "unpartitioned" : "partitioned"); } + for (size_t pos = 0; pos < partition_type->fields().size(); ++pos) { + const auto& literal = data_file.partition.values()[pos]; + const auto& expected_type = partition_type->fields()[pos].type(); + if (!literal.IsNull() && (!literal.type() || *literal.type() != *expected_type)) { + return ValidationFailed( + "Invalid partition value type at position {}: expected = {}, actual = {}", pos, + expected_type->ToString(), + literal.type() ? literal.type()->ToString() : "unknown"); + } + } return DataFileToJsonUnchecked(data_file); } diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index c7779d703..0353f1ced 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1081,4 +1081,112 @@ TEST(DataFileJsonTest, RoundTripRequiredFields) { EXPECT_EQ(parsed.partition_spec_id, data_file.partition_spec_id); } +class DataFileJavaJsonTest : public ::testing::Test { + protected: + void SetUp() override { + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(schema_, /*spec_id=*/0, + {PartitionField(1, 1000, "id", Transform::Identity())}, + /*allow_missing_fields=*/false)); + spec_ = std::shared_ptr(std::move(spec)); + specs_.emplace(spec_->spec_id(), spec_); + } + + DataFile DataFileForTest() const { + DataFile data_file; + data_file.content = DataFile::Content::kData; + data_file.file_path = "/path/to/data.parquet"; + data_file.file_format = FileFormatType::kParquet; + data_file.partition_spec_id = 0; + data_file.partition = PartitionValues({Literal::Int(7)}); + data_file.record_count = 1; + data_file.file_size_in_bytes = 10; + data_file.lower_bounds = {{1, {0x01, 0x00, 0x00, 0x00}}}; + data_file.upper_bounds = {{1, {0x05, 0x00, 0x00, 0x00}}}; + data_file.key_metadata = {0x0A, 0x0B}; + data_file.sort_order_id = 0; + return data_file; + } + + nlohmann::json JavaGolden() const { + return R"({ + "spec-id": 0, + "content": "data", + "file-path": "/path/to/data.parquet", + "file-format": "parquet", + "partition": [7], + "file-size-in-bytes": 10, + "record-count": 1, + "lower-bounds": {"keys": [1], "values": ["01000000"]}, + "upper-bounds": {"keys": [1], "values": ["05000000"]}, + "key-metadata": "0A0B", + "sort-order-id": 0 + })"_json; + } + + Schema schema_{{SchemaField::MakeRequired(1, "id", int32())}, 0}; + std::shared_ptr spec_; + std::unordered_map> specs_; +}; + +TEST_F(DataFileJavaJsonTest, SerializesAndParsesJavaEncoding) { + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(DataFileForTest(), specs_, schema_)); + EXPECT_EQ(json, JavaGolden()); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed, DataFileForTest()); +} + +TEST_F(DataFileJavaJsonTest, ParsesPartitionObjectAndLegacyContentName) { + auto json = JavaGolden(); + json["content"] = "DATA"; + json["partition"] = {{"1000", 7}}; + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.content, DataFile::Content::kData); + ASSERT_EQ(parsed.partition.num_fields(), 1U); + EXPECT_EQ(parsed.partition.values()[0], Literal::Int(7)); +} + +TEST_F(DataFileJavaJsonTest, RejectsInvalidIntegersAndDuplicateMetricKeys) { + for (std::string_view field : + {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json[field] = 0.5; + EXPECT_THAT(DataFileFromJson(json, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + } + + auto duplicate_keys = JavaGolden(); + duplicate_keys["column-sizes"] = {{"keys", {1, 1}}, {"values", {10, 20}}}; + auto result = DataFileFromJson(duplicate_keys, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("duplicate key")); +} + +TEST_F(DataFileJavaJsonTest, AcceptsMissingPartitionAndNullMetricMaps) { + auto json = JavaGolden(); + json.erase("partition"); + for (std::string_view field : {"column-sizes", "value-counts", "null-value-counts", + "nan-value-counts", "lower-bounds", "upper-bounds"}) { + json[field] = nullptr; + } + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.partition.num_fields(), 0U); + EXPECT_TRUE(parsed.column_sizes.empty()); + EXPECT_TRUE(parsed.lower_bounds.empty()); +} + +TEST_F(DataFileJavaJsonTest, RejectsInvalidPartitionTypeWhenSerializing) { + auto data_file = DataFileForTest(); + data_file.partition = PartitionValues({Literal::String("not-an-int")}); + + auto result = ToJson(data_file, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(result, HasErrorMessage("partition value type")); +} + } // namespace iceberg diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index ec41e4a66..e1979dbef 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -19,7 +19,9 @@ #include #include +#include #include +#include #include #include @@ -2119,7 +2121,7 @@ TEST(DataFileFromJsonTest, MissingSpecId) { EXPECT_THAT(result, HasErrorMessage("Missing 'spec-id'")); } -TEST(DataFileFromJsonTest, MissingPartition) { +TEST(DataFileFromJsonTest, MissingPartitionIsAccepted) { auto json = R"({ "content": "data", "file-path": "s3://bucket/data/file.parquet", @@ -2129,9 +2131,9 @@ TEST(DataFileFromJsonTest, MissingPartition) { "record-count": 10 })"_json; - auto result = DataFileFromJson(json, UnpartitionedSpecs(), Schema({}, 0)); - EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); - EXPECT_THAT(result, HasErrorMessage("Missing 'partition'")); + ICEBERG_UNWRAP_OR_FAIL(auto data_file, + DataFileFromJson(json, UnpartitionedSpecs(), Schema({}, 0))); + EXPECT_EQ(data_file.partition.num_fields(), 0U); } TEST(DataFileFromJsonTest, NotAnObject) { From fef0a43d59986b1af2a32c783ad4dd7cfa17079d Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 10:41:15 +0000 Subject: [PATCH 05/11] test(serde): cover DataFile Java JSON edge cases --- src/iceberg/test/json_serde_test.cc | 140 ++++++++++++++++++++++++++-- 1 file changed, 134 insertions(+), 6 deletions(-) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 0353f1ced..3f8dbc165 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,6 +18,8 @@ */ #include +#include +#include #include #include @@ -1138,32 +1140,158 @@ TEST_F(DataFileJavaJsonTest, SerializesAndParsesJavaEncoding) { EXPECT_EQ(parsed, DataFileForTest()); } -TEST_F(DataFileJavaJsonTest, ParsesPartitionObjectAndLegacyContentName) { +TEST_F(DataFileJavaJsonTest, ParsesPartitionObject) { auto json = JavaGolden(); - json["content"] = "DATA"; json["partition"] = {{"1000", 7}}; ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); - EXPECT_EQ(parsed.content, DataFile::Content::kData); ASSERT_EQ(parsed.partition.num_fields(), 1U); EXPECT_EQ(parsed.partition.values()[0], Literal::Int(7)); } +TEST_F(DataFileJavaJsonTest, ParsesAllLegacyContentNames) { + struct Case { + std::string_view legacy; + DataFile::Content content; + std::string_view canonical; + }; + for (const auto& test_case : + {Case{"DATA", DataFile::Content::kData, "data"}, + Case{"POSITION_DELETES", DataFile::Content::kPositionDeletes, "position-deletes"}, + Case{"EQUALITY_DELETES", DataFile::Content::kEqualityDeletes, + "equality-deletes"}}) { + SCOPED_TRACE(test_case.legacy); + auto json = JavaGolden(); + json["content"] = std::string(test_case.legacy); + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + EXPECT_EQ(parsed.content, test_case.content); + ICEBERG_UNWRAP_OR_FAIL(auto serialized, ToJson(parsed, specs_, schema_)); + EXPECT_EQ(serialized["content"], std::string(test_case.canonical)); + } +} + +TEST_F(DataFileJavaJsonTest, ParsesFieldIdPartitionsInSpecOrderAndTypedNulls) { + Schema schema{{SchemaField::MakeRequired(1, "id", int32()), + SchemaField::MakeOptional(2, "region", string())}, + 0}; + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(schema, /*spec_id=*/0, + {PartitionField(2, 1001, "region", Transform::Identity()), + PartitionField(1, 1000, "id", Transform::Identity())}, + /*allow_missing_fields=*/false)); + std::unordered_map> specs{ + {0, std::shared_ptr(std::move(spec))}}; + + auto json = JavaGolden(); + json["partition"] = {{"1000", 7}, {"1001", "north"}}; + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(parsed.partition.num_fields(), 2U); + EXPECT_EQ(parsed.partition.values()[0], Literal::String("north")); + EXPECT_EQ(parsed.partition.values()[1], Literal::Int(7)); + + json["partition"] = {{"1000", 7}}; + ICEBERG_UNWRAP_OR_FAIL(auto sparse, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(sparse.partition.num_fields(), 2U); + EXPECT_TRUE(sparse.partition.values()[0].IsNull()); + ASSERT_NE(sparse.partition.values()[0].type(), nullptr); + EXPECT_EQ(sparse.partition.values()[0].type()->type_id(), TypeId::kString); + EXPECT_EQ(sparse.partition.values()[1], Literal::Int(7)); + + json["partition"] = nlohmann::json::array({nullptr, 7}); + ICEBERG_UNWRAP_OR_FAIL(auto positional, DataFileFromJson(json, specs, schema)); + ASSERT_EQ(positional.partition.num_fields(), 2U); + EXPECT_TRUE(positional.partition.values()[0].IsNull()); + ASSERT_NE(positional.partition.values()[0].type(), nullptr); + EXPECT_EQ(positional.partition.values()[0].type()->type_id(), TypeId::kString); + EXPECT_EQ(positional.partition.values()[1], Literal::Int(7)); +} + +TEST_F(DataFileJavaJsonTest, RejectsMalformedPartitionShapes) { + for (const auto& partition : + {nlohmann::json::array(), nlohmann::json::object({{"1000", 7}, {"9999", 8}}), + nlohmann::json("invalid")}) { + SCOPED_TRACE(partition.dump()); + auto json = JavaGolden(); + json["partition"] = partition; + EXPECT_THAT(DataFileFromJson(json, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + } +} + TEST_F(DataFileJavaJsonTest, RejectsInvalidIntegersAndDuplicateMetricKeys) { for (std::string_view field : - {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id"}) { + {"spec-id", "file-size-in-bytes", "record-count", "sort-order-id", "first-row-id", + "content-offset", "content-size-in-bytes"}) { SCOPED_TRACE(field); auto json = JavaGolden(); json[field] = 0.5; - EXPECT_THAT(DataFileFromJson(json, specs_, schema_), - IsError(ErrorKind::kJsonParseError)); + auto result = DataFileFromJson(json, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); } + auto expect_invalid = [&](nlohmann::json json, std::string_view field) { + auto result = DataFileFromJson(json, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + }; + + auto fractional_key = JavaGolden(); + fractional_key["column-sizes"] = {{"keys", {0.5}}, {"values", {10}}}; + expect_invalid(fractional_key, "column-sizes"); + + auto fractional_value = JavaGolden(); + fractional_value["column-sizes"] = {{"keys", {1}}, {"values", {0.5}}}; + expect_invalid(fractional_value, "column-sizes"); + + auto fractional_bound_key = JavaGolden(); + fractional_bound_key["lower-bounds"]["keys"] = {0.5}; + expect_invalid(fractional_bound_key, "lower-bounds"); + + for (std::string_view field : {"split-offsets", "equality-ids"}) { + auto json = JavaGolden(); + json[field] = {0.5}; + expect_invalid(std::move(json), field); + } + + auto oversized_spec_id = JavaGolden(); + oversized_spec_id["spec-id"] = 2147483648ULL; + expect_invalid(oversized_spec_id, "spec-id"); + + auto oversized_record_count = JavaGolden(); + oversized_record_count["record-count"] = 18446744073709551615ULL; + expect_invalid(oversized_record_count, "record-count"); + auto duplicate_keys = JavaGolden(); duplicate_keys["column-sizes"] = {{"keys", {1, 1}}, {"values", {10, 20}}}; auto result = DataFileFromJson(duplicate_keys, specs_, schema_); EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); EXPECT_THAT(result, HasErrorMessage("duplicate key")); + + auto duplicate_bounds = JavaGolden(); + duplicate_bounds["lower-bounds"] = {{"keys", {1, 1}}, {"values", {"01", "02"}}}; + auto bounds_result = DataFileFromJson(duplicate_bounds, specs_, schema_); + EXPECT_THAT(bounds_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(bounds_result, HasErrorMessage("duplicate key")); +} + +TEST_F(DataFileJavaJsonTest, RejectsMalformedHexMetadata) { + auto bad_bound = JavaGolden(); + bad_bound["lower-bounds"]["values"] = {"0G"}; + EXPECT_THAT(DataFileFromJson(bad_bound, specs_, schema_), + IsError(ErrorKind::kInvalidArgument)); + + auto non_string_bound = JavaGolden(); + non_string_bound["lower-bounds"]["values"] = {7}; + EXPECT_THAT(DataFileFromJson(non_string_bound, specs_, schema_), + IsError(ErrorKind::kJsonParseError)); + + auto bad_key_metadata = JavaGolden(); + bad_key_metadata["key-metadata"] = "0G"; + EXPECT_THAT(DataFileFromJson(bad_key_metadata, specs_, schema_), + IsError(ErrorKind::kInvalidArgument)); } TEST_F(DataFileJavaJsonTest, AcceptsMissingPartitionAndNullMetricMaps) { From d4a1c887eaec1debd33e9bace944ef8fc2fd4e94 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 20:31:27 +0000 Subject: [PATCH 06/11] fix(serde): use designated initializers in DataFile tests --- src/iceberg/test/json_serde_test.cc | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 3f8dbc165..7bb867290 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1156,10 +1156,13 @@ TEST_F(DataFileJavaJsonTest, ParsesAllLegacyContentNames) { std::string_view canonical; }; for (const auto& test_case : - {Case{"DATA", DataFile::Content::kData, "data"}, - Case{"POSITION_DELETES", DataFile::Content::kPositionDeletes, "position-deletes"}, - Case{"EQUALITY_DELETES", DataFile::Content::kEqualityDeletes, - "equality-deletes"}}) { + {Case{.legacy = "DATA", .content = DataFile::Content::kData, .canonical = "data"}, + Case{.legacy = "POSITION_DELETES", + .content = DataFile::Content::kPositionDeletes, + .canonical = "position-deletes"}, + Case{.legacy = "EQUALITY_DELETES", + .content = DataFile::Content::kEqualityDeletes, + .canonical = "equality-deletes"}}) { SCOPED_TRACE(test_case.legacy); auto json = JavaGolden(); json["content"] = std::string(test_case.legacy); From d182ea5fb1f2b71f01962d146d87bd609169d21b Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 08:05:40 +0000 Subject: [PATCH 07/11] feat(serde): add core FileScanTask JSON support --- src/iceberg/catalog/rest/json_serde.cc | 26 ++ src/iceberg/json_serde.cc | 151 ++++++++- src/iceberg/json_serde_internal.h | 32 ++ src/iceberg/test/json_serde_test.cc | 383 +++++++++++++++++++++++ src/iceberg/test/rest_json_serde_test.cc | 47 +++ 5 files changed, 638 insertions(+), 1 deletion(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 46593f0b5..c1530b520 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -18,6 +18,7 @@ */ #include +#include #include #include #include @@ -108,6 +109,25 @@ constexpr std::string_view kPlanTask = "plan-task"; constexpr std::string_view kDataFile = "data-file"; constexpr std::string_view kDeleteFileReferences = "delete-file-references"; constexpr std::string_view kResidualFilter = "residual-filter"; +constexpr std::string_view kStart = "start"; +constexpr std::string_view kLength = "length"; + +Result GetRangeValueOrDefault(const nlohmann::json& json, std::string_view key, + int64_t default_value) { + if (!json.contains(key) || json.at(key).is_null()) { + return default_value; + } + const auto& value = json.at(key); + if (!value.is_number_integer()) { + return JsonParseError("'{}' must be an integer, but is {}", key, value.type_name()); + } + if (value.is_number_unsigned() && + value.get() > + static_cast(std::numeric_limits::max())) { + return JsonParseError("'{}' integer is out of range: {}", key, SafeDumpJson(value)); + } + return value.get(); +} Result StorageCredentialToJson(const StorageCredential& credential) { ICEBERG_RETURN_UNEXPECTED(credential.Validate()); @@ -158,6 +178,12 @@ Result>> FileScanTasksFromJson( ICEBERG_ASSIGN_OR_RAISE( auto data_file, iceberg::rest::DataFileFromJson(data_file_json, partition_spec_by_id, schema)); + ICEBERG_ASSIGN_OR_RAISE(auto start, GetRangeValueOrDefault(task_json, kStart, 0)); + ICEBERG_ASSIGN_OR_RAISE( + auto length, + GetRangeValueOrDefault(task_json, kLength, data_file.file_size_in_bytes)); + ICEBERG_RETURN_UNEXPECTED( + CheckFileScanTaskNotSplit(start, length, data_file.file_size_in_bytes)); // FIXME: REST scan-task DataFile JSON currently carries first-row-id, // but not the manifest-entry data sequence number. Until the REST API exposes // it, REST-planned tasks cannot inherit _last_updated_sequence_number. diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index 45dfa29f5..add5e50f1 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -31,6 +31,7 @@ #include #include "iceberg/constants.h" +#include "iceberg/expression/expressions.h" #include "iceberg/expression/json_serde_internal.h" #include "iceberg/expression/literal.h" #include "iceberg/file_format.h" @@ -50,6 +51,7 @@ #include "iceberg/table_metadata.h" #include "iceberg/table_properties.h" #include "iceberg/table_requirement.h" +#include "iceberg/table_scan.h" #include "iceberg/table_update.h" #include "iceberg/transform.h" #include "iceberg/type.h" @@ -248,7 +250,7 @@ constexpr std::string_view kRequirementAssertDefaultSortOrderID = constexpr std::string_view kLastAssignedFieldId = "last-assigned-field-id"; constexpr std::string_view kLastAssignedPartitionId = "last-assigned-partition-id"; -// DataFile JSON (Iceberg REST ContentFile field names) +// DataFile / FileScanTask JSON (Iceberg REST ContentFile field names) constexpr std::string_view kContent = "content"; constexpr std::string_view kContentData = "data"; constexpr std::string_view kContentPositionDeletes = "position-deletes"; @@ -269,6 +271,14 @@ constexpr std::string_view kEqualityIds = "equality-ids"; constexpr std::string_view kReferencedDataFile = "referenced-data-file"; constexpr std::string_view kContentOffset = "content-offset"; constexpr std::string_view kContentSizeInBytes = "content-size-in-bytes"; +constexpr std::string_view kDataFile = "data-file"; +constexpr std::string_view kDeleteFiles = "delete-files"; +constexpr std::string_view kDeleteFileReferences = "delete-file-references"; +constexpr std::string_view kResidualFilter = "residual-filter"; +constexpr std::string_view kTaskType = "task-type"; +constexpr std::string_view kFileScanTaskType = "file-scan-task"; +constexpr std::string_view kStart = "start"; +constexpr std::string_view kLength = "length"; constexpr std::string_view kMapKeys = "keys"; constexpr std::string_view kMapValues = "values"; @@ -2432,4 +2442,143 @@ Result ToJson( return DataFileToJsonUnchecked(data_file); } +Result ToJson( + const FileScanTask& task, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema) { + if (!task.data_file()) { + return ValidationFailed("Cannot serialize FileScanTask without data-file"); + } + if (task.data_file()->content != DataFile::Content::kData) { + return ValidationFailed("FileScanTask data-file must have data content"); + } + + nlohmann::json json; + json[kTaskType] = kFileScanTaskType; + ICEBERG_ASSIGN_OR_RAISE(auto schema_json, ToJson(schema)); + json[kSchema] = std::move(schema_json); + + ICEBERG_ASSIGN_OR_RAISE(auto data_file_json, + ToJson(*task.data_file(), partition_specs_by_id, schema)); + const auto spec_id = task.data_file()->partition_spec_id.value(); + json[kSpec] = ToJson(*partition_specs_by_id.at(spec_id)); + json[kDataFile] = std::move(data_file_json); + json[kStart] = 0; + json[kLength] = task.data_file()->file_size_in_bytes; + + nlohmann::json delete_files_json = nlohmann::json::array(); + for (const auto& delete_file : task.delete_files()) { + if (!delete_file) { + return ValidationFailed("FileScanTask delete-files must not contain null"); + } + if (delete_file->content == DataFile::Content::kData) { + return ValidationFailed("FileScanTask delete-file must have delete content"); + } + if (delete_file->partition_spec_id != task.data_file()->partition_spec_id) { + return ValidationFailed( + "Invalid partition spec id from content file: expected = {}, actual = {}", + spec_id, + delete_file->partition_spec_id.has_value() + ? std::to_string(delete_file->partition_spec_id.value()) + : "null"); + } + ICEBERG_ASSIGN_OR_RAISE(auto delete_file_json, + ToJson(*delete_file, partition_specs_by_id, schema)); + delete_files_json.push_back(std::move(delete_file_json)); + } + json[kDeleteFiles] = std::move(delete_files_json); + + if (task.residual_filter()) { + ICEBERG_ASSIGN_OR_RAISE(auto residual_json, ToJson(*task.residual_filter())); + json[kResidualFilter] = std::move(residual_json); + } + return json; +} + +Status CheckFileScanTaskNotSplit(int64_t start, int64_t length, + int64_t file_size_in_bytes) { + if (start != 0 || length != file_size_in_bytes) { + return NotSupported( + "Split FileScanTask is not supported: start={}, length={}, " + "file-size-in-bytes={}", + start, length, file_size_in_bytes); + } + return {}; +} + +Result> FileScanTaskFromJson(const nlohmann::json& json) { + if (!json.is_object()) { + return JsonParseError("Cannot parse file scan task from a non-object: {}", + SafeDumpJson(json)); + } + ICEBERG_ASSIGN_OR_RAISE(auto task_type, + GetJsonValueOptional(json, kTaskType)); + if (task_type.has_value() && + !StringUtils::EqualsIgnoreCase(*task_type, kFileScanTaskType)) { + return JsonParseError("Unsupported scan task type: {}", *task_type); + } + if (json.contains(kDeleteFileReferences) && !json.at(kDeleteFileReferences).is_null()) { + return JsonParseError( + "Cannot parse FileScanTask with 'delete-file-references'; use " + "rest::FileScanTasksFromJson for REST scan responses"); + } + + ICEBERG_ASSIGN_OR_RAISE(auto schema_json, GetJsonValue(json, kSchema)); + ICEBERG_ASSIGN_OR_RAISE(auto schema, SchemaFromJson(schema_json)); + auto shared_schema = std::shared_ptr(std::move(schema)); + + ICEBERG_ASSIGN_OR_RAISE(auto spec_json, GetJsonValue(json, kSpec)); + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(spec_json, kSpecId)); + ICEBERG_ASSIGN_OR_RAISE(auto spec, + PartitionSpecFromJson(shared_schema, spec_json, spec_id)); + auto shared_spec = std::shared_ptr(std::move(spec)); + const std::unordered_map> partition_spec_by_id{ + {shared_spec->spec_id(), shared_spec}}; + + ICEBERG_ASSIGN_OR_RAISE(auto data_file_json, + GetJsonValue(json, kDataFile)); + ICEBERG_ASSIGN_OR_RAISE( + auto data_file, + DataFileFromJson(data_file_json, partition_spec_by_id, *shared_schema)); + if (data_file.content != DataFile::Content::kData) { + return JsonParseError("FileScanTask data-file must have data content"); + } + + ICEBERG_ASSIGN_OR_RAISE(auto start, GetJsonInteger(json, kStart)); + ICEBERG_ASSIGN_OR_RAISE(auto length, GetJsonInteger(json, kLength)); + ICEBERG_RETURN_UNEXPECTED( + CheckFileScanTaskNotSplit(start, length, data_file.file_size_in_bytes)); + + std::vector> delete_files; + if (json.contains(kDeleteFiles)) { + ICEBERG_ASSIGN_OR_RAISE(auto delete_files_json, + GetJsonValue(json, kDeleteFiles)); + if (!delete_files_json.is_array()) { + return JsonParseError("Cannot parse delete files from non-array: {}", + SafeDumpJson(delete_files_json)); + } + for (const auto& delete_file_json : delete_files_json) { + ICEBERG_ASSIGN_OR_RAISE( + auto delete_file, + DataFileFromJson(delete_file_json, partition_spec_by_id, *shared_schema)); + if (delete_file.content == DataFile::Content::kData) { + return JsonParseError("FileScanTask delete-file must have delete content"); + } + delete_files.push_back(std::make_shared(std::move(delete_file))); + } + } + + std::shared_ptr residual_filter = Expressions::AlwaysTrue(); + if (json.contains(kResidualFilter)) { + ICEBERG_ASSIGN_OR_RAISE(auto filter_json, + GetJsonValue(json, kResidualFilter)); + ICEBERG_ASSIGN_OR_RAISE(residual_filter, ExpressionFromJson(filter_json)); + } + + return std::make_shared(std::make_shared(std::move(data_file)), + std::move(delete_files), + std::move(residual_filter)); +} + } // namespace iceberg diff --git a/src/iceberg/json_serde_internal.h b/src/iceberg/json_serde_internal.h index 3aaced5c9..8686076f8 100644 --- a/src/iceberg/json_serde_internal.h +++ b/src/iceberg/json_serde_internal.h @@ -445,4 +445,36 @@ ICEBERG_EXPORT Result DataFileFromJson( partition_spec_by_id, const Schema& schema); +/// \brief Serializes a `FileScanTask` to a self-contained JSON object. +/// +/// Unlike REST scan responses, delete files are inlined under `delete-files` rather +/// than encoded as `delete-file-references` into a sibling array. +/// +/// JSON fields: +/// - `task-type`: `file-scan-task` +/// - `schema` (required): schema used by this task +/// - `spec` (required): partition spec used by this task +/// - `data-file` (required): ContentFile JSON +/// - `start` (required): always 0 because split tasks are unsupported +/// - `length` (required): always the data file size +/// - `delete-files` (required): array of ContentFile JSON +/// - `residual-filter` (optional): Expression JSON +ICEBERG_EXPORT Result ToJson( + const FileScanTask& task, + const std::unordered_map>& + partition_specs_by_id, + const Schema& schema); + +/// \brief Deserializes a self-contained FileScanTask JSON object. +/// +/// Reads the embedded schema and partition spec. Rejects REST +/// `delete-file-references` and byte-range splits (`start` != 0 or `length` != file +/// size). Both range fields are required, matching Java core. +ICEBERG_EXPORT Result> FileScanTaskFromJson( + const nlohmann::json& json); + +/// Rejects `start != 0 || length != file_size_in_bytes`. +ICEBERG_EXPORT Status CheckFileScanTaskNotSplit(int64_t start, int64_t length, + int64_t file_size_in_bytes); + } // namespace iceberg diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 7bb867290..2ca1d0883 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -26,6 +26,8 @@ #include #include +#include "iceberg/expression/expressions.h" +#include "iceberg/expression/json_serde_internal.h" #include "iceberg/expression/literal.h" #include "iceberg/file_format.h" #include "iceberg/json_serde_internal.h" @@ -40,6 +42,7 @@ #include "iceberg/sort_order.h" #include "iceberg/statistics_file.h" #include "iceberg/table_requirement.h" +#include "iceberg/table_scan.h" #include "iceberg/table_update.h" #include "iceberg/test/matchers.h" #include "iceberg/transform.h" @@ -1320,4 +1323,384 @@ TEST_F(DataFileJavaJsonTest, RejectsInvalidPartitionTypeWhenSerializing) { EXPECT_THAT(result, HasErrorMessage("partition value type")); } +class FileScanTaskJsonTest : public ::testing::Test { + protected: + void SetUp() override { + ICEBERG_UNWRAP_OR_FAIL( + auto spec, + PartitionSpec::Make(schema_, /*spec_id=*/0, + {PartitionField(1, 1000, "id", Transform::Identity())}, + /*allow_missing_fields=*/false)); + spec_ = std::shared_ptr(std::move(spec)); + specs_.emplace(spec_->spec_id(), spec_); + } + + DataFile DataFileForTask() const { + DataFile data_file; + data_file.content = DataFile::Content::kData; + data_file.file_path = "/path/to/data.parquet"; + data_file.file_format = FileFormatType::kParquet; + data_file.partition_spec_id = 0; + data_file.partition = PartitionValues({Literal::Int(7)}); + data_file.record_count = 1; + data_file.file_size_in_bytes = 10; + data_file.lower_bounds = {{1, {0x01, 0x00, 0x00, 0x00}}}; + data_file.upper_bounds = {{1, {0x05, 0x00, 0x00, 0x00}}}; + data_file.key_metadata = {0x0A, 0x0B}; + data_file.sort_order_id = 0; + return data_file; + } + + DataFile DeleteFileForTask() const { + DataFile delete_file; + delete_file.content = DataFile::Content::kPositionDeletes; + delete_file.file_path = "/path/to/delete.parquet"; + delete_file.file_format = FileFormatType::kParquet; + delete_file.partition_spec_id = 0; + delete_file.partition = PartitionValues({Literal::Int(7)}); + delete_file.record_count = 1; + delete_file.file_size_in_bytes = 4; + return delete_file; + } + + nlohmann::json JavaGolden() const { + return R"({ + "task-type": "file-scan-task", + "schema": { + "type": "struct", + "schema-id": 0, + "fields": [ + {"id": 1, "name": "id", "required": true, "type": "int"} + ] + }, + "spec": { + "spec-id": 0, + "fields": [ + {"source-id": 1, "field-id": 1000, "name": "id", "transform": "identity"} + ] + }, + "data-file": { + "spec-id": 0, + "content": "data", + "file-path": "/path/to/data.parquet", + "file-format": "parquet", + "partition": [7], + "file-size-in-bytes": 10, + "record-count": 1, + "lower-bounds": {"keys": [1], "values": ["01000000"]}, + "upper-bounds": {"keys": [1], "values": ["05000000"]}, + "key-metadata": "0A0B", + "sort-order-id": 0 + }, + "start": 0, + "length": 10, + "delete-files": [ + { + "spec-id": 0, + "content": "position-deletes", + "file-path": "/path/to/delete.parquet", + "file-format": "parquet", + "partition": [7], + "file-size-in-bytes": 4, + "record-count": 1 + } + ], + "residual-filter": true + })"_json; + } + + Schema schema_{{SchemaField::MakeRequired(1, "id", int32())}, 0}; + std::shared_ptr spec_; + std::unordered_map> specs_; +}; + +TEST_F(FileScanTaskJsonTest, SerializesJavaGolden) { + FileScanTask task(std::make_shared(DataFileForTask()), + {std::make_shared(DeleteFileForTask())}, + Expressions::AlwaysTrue()); + + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(task, specs_, schema_)); + EXPECT_EQ(json, JavaGolden()); +} + +TEST_F(FileScanTaskJsonTest, ParsesJavaGoldenWithoutExternalSchemaOrSpec) { + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(JavaGolden())); + + ASSERT_NE(task->data_file(), nullptr); + EXPECT_EQ(*task->data_file(), DataFileForTask()); + ASSERT_EQ(task->delete_files().size(), 1U); + ASSERT_NE(task->delete_files()[0], nullptr); + EXPECT_EQ(*task->delete_files()[0], DeleteFileForTask()); + ASSERT_NE(task->residual_filter(), nullptr); + ICEBERG_UNWRAP_OR_FAIL(auto residual, ToJson(*task->residual_filter())); + EXPECT_EQ(residual, true); +} + +TEST_F(FileScanTaskJsonTest, SerializesEmptyDeleteFilesArray) { + FileScanTask task(std::make_shared(DataFileForTask())); + + ICEBERG_UNWRAP_OR_FAIL(auto json, ToJson(task, specs_, schema_)); + EXPECT_EQ(json["delete-files"], nlohmann::json::array()); + EXPECT_EQ(json["start"], 0); + EXPECT_EQ(json["length"], 10); +} + +TEST_F(FileScanTaskJsonTest, RequiresJavaCoreFields) { + for (std::string_view field : {"schema", "spec", "start", "length"}) { + for (bool use_null : {false, true}) { + SCOPED_TRACE(std::string(field) + (use_null ? " null" : " missing")); + auto json = JavaGolden(); + if (use_null) { + json[field] = nullptr; + } else { + json.erase(field); + } + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + } + } +} + +TEST_F(FileScanTaskJsonTest, RejectsSplitTasksUsingExactOrCondition) { + struct Case { + int64_t start; + int64_t length; + bool supported; + }; + for (const auto& test_case : + {Case{.start = 0, .length = 10, .supported = true}, + Case{.start = 1, .length = 10, .supported = false}, + Case{.start = 0, .length = 9, .supported = false}}) { + SCOPED_TRACE(testing::Message() + << "start=" << test_case.start << ", length=" << test_case.length); + auto json = JavaGolden(); + json["start"] = test_case.start; + json["length"] = test_case.length; + + auto result = FileScanTaskFromJson(json); + if (test_case.supported) { + EXPECT_THAT(result, IsOk()); + } else { + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); + EXPECT_THAT(result, HasErrorMessage("Split FileScanTask is not supported")); + } + } +} + +TEST_F(FileScanTaskJsonTest, RejectsFractionalIntegerFields) { + for (std::string_view field : + {"start", "length", "spec-id", "file-size-in-bytes", "record-count", + "sort-order-id", "first-row-id", "content-offset", "content-size-in-bytes"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + if (field == "start" || field == "length") { + json[field] = 0.5; + } else { + json["data-file"][field] = 0.5; + } + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + } + + for (std::string_view field : {"split-offsets", "equality-ids"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json["data-file"][field] = {0.5}; + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + } + + auto scalar_map = JavaGolden(); + scalar_map["data-file"]["column-sizes"] = {{"keys", {1}}, {"values", {0.5}}}; + EXPECT_THAT(FileScanTaskFromJson(scalar_map), IsError(ErrorKind::kJsonParseError)); + + auto bytes_map = JavaGolden(); + bytes_map["data-file"]["lower-bounds"] = {{"keys", {0.5}}, {"values", {"00"}}}; + EXPECT_THAT(FileScanTaskFromJson(bytes_map), IsError(ErrorKind::kJsonParseError)); +} + +TEST_F(FileScanTaskJsonTest, DoesNotTreatOffsetAsStartAlias) { + auto json = JavaGolden(); + json.erase("start"); + json["offset"] = 0; + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("start")); +} + +TEST_F(FileScanTaskJsonTest, ValidatesTaskTypeWhenPresent) { + auto json = JavaGolden(); + json["task-type"] = "FILE-SCAN-TASK"; + EXPECT_THAT(FileScanTaskFromJson(json), IsOk()); + + json["task-type"] = nullptr; + EXPECT_THAT(FileScanTaskFromJson(json), IsOk()); + + json.erase("task-type"); + EXPECT_THAT(FileScanTaskFromJson(json), IsOk()); + + json["task-type"] = "data-task"; + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("task type")); +} + +TEST_F(FileScanTaskJsonTest, ParsesPartitionObjectsByFieldId) { + auto json = JavaGolden(); + json["data-file"]["partition"] = {{"1000", 7}}; + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + ASSERT_EQ(task->data_file()->partition.num_fields(), 1U); + EXPECT_EQ(task->data_file()->partition.values()[0], Literal::Int(7)); +} + +TEST_F(FileScanTaskJsonTest, ParsesNullPartitionValues) { + const std::vector null_partitions = { + nlohmann::json::object(), {{"1000", nullptr}}, nlohmann::json::array({nullptr})}; + for (const auto& partition : null_partitions) { + SCOPED_TRACE(partition.dump()); + auto json = JavaGolden(); + json["data-file"]["partition"] = partition; + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + ASSERT_EQ(task->data_file()->partition.num_fields(), 1U); + const auto& value = task->data_file()->partition.values()[0]; + EXPECT_TRUE(value.IsNull()); + ASSERT_NE(value.type(), nullptr); + EXPECT_EQ(value.type()->type_id(), TypeId::kInt); + } +} + +TEST_F(FileScanTaskJsonTest, AcceptsMissingPartitionLikeJava) { + auto json = JavaGolden(); + json["data-file"].erase("partition"); + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + EXPECT_EQ(task->data_file()->partition.num_fields(), 0U); +} + +TEST_F(FileScanTaskJsonTest, AcceptsNullMetricMapsLikeJava) { + auto json = JavaGolden(); + for (std::string_view field : {"column-sizes", "value-counts", "null-value-counts", + "nan-value-counts", "lower-bounds", "upper-bounds"}) { + json["data-file"][field] = nullptr; + } + + EXPECT_THAT(FileScanTaskFromJson(json), IsOk()); +} + +TEST_F(FileScanTaskJsonTest, RejectsDuplicateMetricMapKeys) { + for (std::string_view field : {"column-sizes", "lower-bounds"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json["data-file"][field] = { + {"keys", {1, 1}}, + {"values", field == "lower-bounds" ? nlohmann::json({"00", "01"}) + : nlohmann::json({1, 2})}}; + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("duplicate key")); + } +} + +TEST_F(FileScanTaskJsonTest, AcceptsLegacyContentEnumNames) { + auto json = JavaGolden(); + json["data-file"]["content"] = "DATA"; + json["delete-files"][0]["content"] = "POSITION_DELETES"; + + EXPECT_THAT(FileScanTaskFromJson(json), IsOk()); +} + +TEST_F(FileScanTaskJsonTest, RejectsInvalidContentRolesWhenParsing) { + auto invalid_data = JavaGolden(); + invalid_data["data-file"]["content"] = "position-deletes"; + auto data_result = FileScanTaskFromJson(invalid_data); + EXPECT_THAT(data_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(data_result, HasErrorMessage("data-file")); + + auto invalid_delete = JavaGolden(); + invalid_delete["delete-files"][0]["content"] = "data"; + auto delete_result = FileScanTaskFromJson(invalid_delete); + EXPECT_THAT(delete_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(delete_result, HasErrorMessage("delete-file")); +} + +TEST_F(FileScanTaskJsonTest, RejectsInvalidContentRolesWhenSerializing) { + auto invalid_data = DataFileForTask(); + invalid_data.content = DataFile::Content::kPositionDeletes; + FileScanTask data_task(std::make_shared(std::move(invalid_data))); + auto data_result = ToJson(data_task, specs_, schema_); + EXPECT_THAT(data_result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(data_result, HasErrorMessage("data-file")); + + auto invalid_delete = DeleteFileForTask(); + invalid_delete.content = DataFile::Content::kData; + FileScanTask delete_task(std::make_shared(DataFileForTask()), + {std::make_shared(std::move(invalid_delete))}); + auto delete_result = ToJson(delete_task, specs_, schema_); + EXPECT_THAT(delete_result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(delete_result, HasErrorMessage("delete-file")); +} + +TEST_F(FileScanTaskJsonTest, RejectsInvalidPartitionValueTypeWhenSerializing) { + auto data_file = DataFileForTask(); + data_file.partition = PartitionValues({Literal::String("not-an-int")}); + FileScanTask task(std::make_shared(std::move(data_file))); + + auto result = ToJson(task, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(result, HasErrorMessage("partition value type")); +} + +TEST_F(FileScanTaskJsonTest, RejectsNullDeleteEntries) { + auto json = JavaGolden(); + json["delete-files"][0] = nullptr; + auto parse_result = FileScanTaskFromJson(json); + EXPECT_THAT(parse_result, IsError(ErrorKind::kJsonParseError)); + + FileScanTask task(std::make_shared(DataFileForTask()), {nullptr}); + auto serialize_result = ToJson(task, specs_, schema_); + EXPECT_THAT(serialize_result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(serialize_result, HasErrorMessage("null")); +} + +TEST_F(FileScanTaskJsonTest, RejectsExplicitNullDeleteFilesAndResidual) { + for (std::string_view field : {"delete-files", "residual-filter"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json[field] = nullptr; + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + } +} + +TEST_F(FileScanTaskJsonTest, DefaultsMissingResidualToAlwaysTrue) { + auto json = JavaGolden(); + json.erase("residual-filter"); + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + ASSERT_NE(task->residual_filter(), nullptr); + ICEBERG_UNWRAP_OR_FAIL(auto residual, ToJson(*task->residual_filter())); + EXPECT_EQ(residual, true); +} + +TEST_F(FileScanTaskJsonTest, RejectsRestDeleteFileReferences) { + auto json = JavaGolden(); + json["delete-file-references"] = nlohmann::json::array({0}); + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("delete-file-references")); +} + } // namespace iceberg diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index e1979dbef..b6956d4cc 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -2286,6 +2286,53 @@ TEST(FileScanTasksFromJsonTest, SingleTaskNoDeleteFiles) { EXPECT_EQ(task->residual_filter(), nullptr); } +TEST(FileScanTasksFromJsonTest, AcceptsWholeFileStartAndLength) { + auto json = R"([{ + "data-file": { + "content": "data", + "file-path": "s3://bucket/data/file.parquet", + "file-format": "PARQUET", + "spec-id": 0, + "partition": [], + "file-size-in-bytes": 12345, + "record-count": 100 + }, + "start": 0, + "length": 12345 + }])"_json; + + EXPECT_THAT(FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)), + IsOk()); +} + +TEST(FileScanTasksFromJsonTest, RejectsSplitTaskWhenEitherRangeFieldDiffers) { + auto json = R"([{ + "data-file": { + "content": "data", + "file-path": "s3://bucket/data/file.parquet", + "file-format": "PARQUET", + "spec-id": 0, + "partition": [], + "file-size-in-bytes": 12345, + "record-count": 100 + }, + "start": 0, + "length": 12345 + }])"_json; + + for (const auto& [field, value] : + {std::pair{"start", 100}, + std::pair{"length", 200}}) { + SCOPED_TRACE(field); + auto split_json = json; + split_json[0][field] = value; + auto result = + FileScanTasksFromJson(split_json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); + EXPECT_THAT(result, HasErrorMessage("Split FileScanTask is not supported")); + } +} + TEST(FileScanTasksFromJsonTest, RowLineageSequence) { GTEST_SKIP() << "REST scan-task JSON does not expose data-sequence-number yet: " << "https://github.com/apache/iceberg-cpp/issues/834"; From ab9c654baf30e399b55d20c5dd59c68c8b4cfd5f Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 10:04:04 +0000 Subject: [PATCH 08/11] fix(serde): accept historical FileScanTask specs --- src/iceberg/json_serde.cc | 36 +++++++++++++---------- src/iceberg/json_serde_internal.h | 2 +- src/iceberg/test/json_serde_test.cc | 45 ++++++++++++++++++++++++++--- 3 files changed, 62 insertions(+), 21 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index add5e50f1..2cf33a4f7 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -1095,10 +1095,11 @@ Result> PartitionFieldFromJson( std::move(transform)); } -Result> PartitionSpecFromJson( - const std::shared_ptr& schema, const nlohmann::json& json, - int32_t default_spec_id) { - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); +namespace { + +Result> ParseBoundPartitionSpec( + const Schema& schema, const nlohmann::json& json, int32_t spec_id, + bool allow_missing_fields) { ICEBERG_ASSIGN_OR_RAISE(auto fields, GetJsonValue(json, kFields)); std::vector partition_fields; @@ -1107,17 +1108,18 @@ Result> PartitionSpecFromJson( partition_fields.push_back(std::move(*partition_field)); } - std::unique_ptr spec; - if (default_spec_id == spec_id) { - ICEBERG_ASSIGN_OR_RAISE( - spec, PartitionSpec::Make(*schema, spec_id, std::move(partition_fields), - /*allow_missing_fields=*/false)); - } else { - ICEBERG_ASSIGN_OR_RAISE( - spec, PartitionSpec::Make(*schema, spec_id, std::move(partition_fields), - /*allow_missing_fields=*/true)); - } - return spec; + return PartitionSpec::Make(schema, spec_id, std::move(partition_fields), + allow_missing_fields); +} + +} // namespace + +Result> PartitionSpecFromJson( + const std::shared_ptr& schema, const nlohmann::json& json, + int32_t default_spec_id) { + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + return ParseBoundPartitionSpec(*schema, json, spec_id, + /*allow_missing_fields=*/spec_id != default_spec_id); } Result> PartitionSpecFromJson(const nlohmann::json& json) { @@ -2530,8 +2532,10 @@ Result> FileScanTaskFromJson(const nlohmann::json& ICEBERG_ASSIGN_OR_RAISE(auto spec_json, GetJsonValue(json, kSpec)); ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(spec_json, kSpecId)); + // An embedded task spec may reference columns dropped from its embedded schema. ICEBERG_ASSIGN_OR_RAISE(auto spec, - PartitionSpecFromJson(shared_schema, spec_json, spec_id)); + ParseBoundPartitionSpec(*shared_schema, spec_json, spec_id, + /*allow_missing_fields=*/true)); auto shared_spec = std::shared_ptr(std::move(spec)); const std::unordered_map> partition_spec_by_id{ {shared_spec->spec_id(), shared_spec}}; diff --git a/src/iceberg/json_serde_internal.h b/src/iceberg/json_serde_internal.h index 8686076f8..2678e3e11 100644 --- a/src/iceberg/json_serde_internal.h +++ b/src/iceberg/json_serde_internal.h @@ -457,7 +457,7 @@ ICEBERG_EXPORT Result DataFileFromJson( /// - `data-file` (required): ContentFile JSON /// - `start` (required): always 0 because split tasks are unsupported /// - `length` (required): always the data file size -/// - `delete-files` (required): array of ContentFile JSON +/// - `delete-files` (optional on input; always emitted): array of ContentFile JSON /// - `residual-filter` (optional): Expression JSON ICEBERG_EXPORT Result ToJson( const FileScanTask& task, diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 2ca1d0883..832d70f81 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -1445,6 +1445,14 @@ TEST_F(FileScanTaskJsonTest, SerializesEmptyDeleteFilesArray) { EXPECT_EQ(json["length"], 10); } +TEST_F(FileScanTaskJsonTest, AcceptsMissingDeleteFiles) { + auto json = JavaGolden(); + json.erase("delete-files"); + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + EXPECT_TRUE(task->delete_files().empty()); +} + TEST_F(FileScanTaskJsonTest, RequiresJavaCoreFields) { for (std::string_view field : {"schema", "spec", "start", "length"}) { for (bool use_null : {false, true}) { @@ -1469,10 +1477,9 @@ TEST_F(FileScanTaskJsonTest, RejectsSplitTasksUsingExactOrCondition) { int64_t length; bool supported; }; - for (const auto& test_case : - {Case{.start = 0, .length = 10, .supported = true}, - Case{.start = 1, .length = 10, .supported = false}, - Case{.start = 0, .length = 9, .supported = false}}) { + for (const auto& test_case : {Case{.start = 0, .length = 10, .supported = true}, + Case{.start = 1, .length = 10, .supported = false}, + Case{.start = 0, .length = 9, .supported = false}}) { SCOPED_TRACE(testing::Message() << "start=" << test_case.start << ", length=" << test_case.length); auto json = JavaGolden(); @@ -1578,6 +1585,36 @@ TEST_F(FileScanTaskJsonTest, ParsesNullPartitionValues) { } } +TEST_F(FileScanTaskJsonTest, RoundTripsHistoricalSpecWithDroppedSourceField) { + auto json = JavaGolden(); + json["schema"]["schema-id"] = 1; + json["schema"]["fields"] = nlohmann::json::array(); + json["data-file"]["partition"] = nlohmann::json::array({nullptr}); + json["delete-files"][0]["partition"] = nlohmann::json::array({nullptr}); + + auto historical_schema = + std::make_shared(std::vector{}, /*schema_id=*/1); + // Parsing a table's default spec remains strict. + auto default_spec_result = PartitionSpecFromJson(historical_schema, json["spec"], 0); + EXPECT_THAT(default_spec_result, IsError(ErrorKind::kInvalidArgument)); + EXPECT_THAT(default_spec_result, + HasErrorMessage("Cannot find source column for partition field")); + + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTaskFromJson(json)); + ASSERT_NE(task->data_file(), nullptr); + ASSERT_EQ(task->delete_files().size(), 1U); + for (const auto& file : {task->data_file(), task->delete_files()[0]}) { + ASSERT_EQ(file->partition.num_fields(), 1U); + const auto& value = file->partition.values()[0]; + EXPECT_TRUE(value.IsNull()); + ASSERT_NE(value.type(), nullptr); + EXPECT_EQ(value.type()->type_id(), TypeId::kUnknown); + } + + ICEBERG_UNWRAP_OR_FAIL(auto serialized, ToJson(*task, specs_, *historical_schema)); + EXPECT_EQ(serialized, json); +} + TEST_F(FileScanTaskJsonTest, AcceptsMissingPartitionLikeJava) { auto json = JavaGolden(); json["data-file"].erase("partition"); From dff1adb770f3a29eb206e451e0c3ae37c4f3b992 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 10:43:24 +0000 Subject: [PATCH 09/11] test(serde): cover FileScanTask range and spec boundaries --- src/iceberg/test/json_serde_test.cc | 44 ++++++++++++++++++++ src/iceberg/test/rest_json_serde_test.cc | 52 ++++++++++++++++++++++++ 2 files changed, 96 insertions(+) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 832d70f81..1085465f8 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,6 +18,7 @@ */ #include +#include #include #include #include @@ -1615,6 +1616,34 @@ TEST_F(FileScanTaskJsonTest, RoundTripsHistoricalSpecWithDroppedSourceField) { EXPECT_EQ(serialized, json); } +TEST_F(FileScanTaskJsonTest, RejectsNonNullPartitionForDroppedSourceField) { + auto json = JavaGolden(); + json["schema"]["fields"] = nlohmann::json::array(); + json["data-file"]["partition"] = {7}; + + // The embedded schema no longer contains the source type, so neither Java nor + // C++ can decode a non-null value for this historical partition field. + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); + EXPECT_THAT(result, HasErrorMessage("Unsupported type for literal JSON parsing")); +} + +TEST_F(FileScanTaskJsonTest, RejectsFilesWithMismatchedEmbeddedSpecId) { + for (bool mismatch_delete_file : {false, true}) { + SCOPED_TRACE(mismatch_delete_file ? "delete-file" : "data-file"); + auto json = JavaGolden(); + if (mismatch_delete_file) { + json["delete-files"][0]["spec-id"] = 1; + } else { + json["data-file"]["spec-id"] = 1; + } + + auto result = FileScanTaskFromJson(json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("Invalid partition spec id")); + } +} + TEST_F(FileScanTaskJsonTest, AcceptsMissingPartitionLikeJava) { auto json = JavaGolden(); json["data-file"].erase("partition"); @@ -1687,6 +1716,21 @@ TEST_F(FileScanTaskJsonTest, RejectsInvalidContentRolesWhenSerializing) { EXPECT_THAT(delete_result, HasErrorMessage("delete-file")); } +TEST_F(FileScanTaskJsonTest, RejectsDeleteFilesWithDifferentSpecIdWhenSerializing) { + for (const auto& spec_id : {std::optional(1), std::optional()}) { + SCOPED_TRACE(spec_id.value_or(-1)); + auto delete_file = DeleteFileForTask(); + delete_file.partition_spec_id = spec_id; + FileScanTask task(std::make_shared(DataFileForTask()), + {std::make_shared(std::move(delete_file))}); + + auto result = ToJson(task, specs_, schema_); + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(result, HasErrorMessage("Invalid partition spec id")); + EXPECT_THAT(result, HasErrorMessage(spec_id ? "actual = 1" : "actual = null")); + } +} + TEST_F(FileScanTaskJsonTest, RejectsInvalidPartitionValueTypeWhenSerializing) { auto data_file = DataFileForTask(); data_file.partition = PartitionValues({Literal::String("not-an-int")}); diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index b6956d4cc..31fc9ef44 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -2305,6 +2305,58 @@ TEST(FileScanTasksFromJsonTest, AcceptsWholeFileStartAndLength) { IsOk()); } +TEST(FileScanTasksFromJsonTest, NullStartAndLengthUseWholeFileDefaults) { + auto json = R"([{ + "data-file": { + "content": "data", + "file-path": "s3://bucket/data/file.parquet", + "file-format": "PARQUET", + "spec-id": 0, + "partition": [], + "file-size-in-bytes": 12345, + "record-count": 100 + }, + "start": null, + "length": null + }])"_json; + + EXPECT_THAT(FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)), + IsOk()); +} + +TEST(FileScanTasksFromJsonTest, RejectsInvalidRangeIntegers) { + auto valid_json = R"([{ + "data-file": { + "content": "data", + "file-path": "s3://bucket/data/file.parquet", + "file-format": "PARQUET", + "spec-id": 0, + "partition": [], + "file-size-in-bytes": 12345, + "record-count": 100 + }, + "start": 0, + "length": 12345 + }])"_json; + + for (std::string_view field : {"start", "length"}) { + for (const auto& [invalid_value, expected_error] : + {std::pair{0.5, "must be an integer"}, + {true, "must be an integer"}, + {"0", "must be an integer"}, + {18446744073709551615ULL, "out of range"}}) { + SCOPED_TRACE(testing::Message() << field << "=" << invalid_value.dump()); + auto json = valid_json; + json[0][field] = invalid_value; + + auto result = FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(field)); + EXPECT_THAT(result, HasErrorMessage(expected_error)); + } + } +} + TEST(FileScanTasksFromJsonTest, RejectsSplitTaskWhenEitherRangeFieldDiffers) { auto json = R"([{ "data-file": { From 8db7cf177ec9794f13ba6d132cb2cd7877ecf219 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Thu, 17 Sep 2026 08:46:27 +0000 Subject: [PATCH 10/11] fix(serde): validate schema and partition IDs as integers --- src/iceberg/json_serde.cc | 47 +++++--- src/iceberg/test/json_serde_test.cc | 136 ++++++++++++++++++++++++ src/iceberg/test/metadata_serde_test.cc | 18 ++++ 3 files changed, 184 insertions(+), 17 deletions(-) diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index 2cf33a4f7..97f034421 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -23,6 +23,7 @@ #include #include #include +#include #include #include #include @@ -313,6 +314,16 @@ Result GetJsonInteger(const nlohmann::json& json, std::string_view key) return IntegerFromJson(value, std::format("'{}'", key)); } +template +Result> GetJsonIntegerOptional(const nlohmann::json& json, + std::string_view key) { + if (!json.contains(key) || json.at(key).is_null()) { + return std::nullopt; + } + ICEBERG_ASSIGN_OR_RAISE(auto value, GetJsonInteger(json, key)); + return value; +} + template Result> IntegerVectorFromJson(const nlohmann::json& json, std::string_view key) { @@ -818,7 +829,7 @@ Result> StructTypeFromJson(const nlohmann::json& json) { Result> ListTypeFromJson(const nlohmann::json& json) { ICEBERG_ASSIGN_OR_RAISE(auto element_type, TypeFromJson(json[kElement])); - ICEBERG_ASSIGN_OR_RAISE(auto element_id, GetJsonValue(json, kElementId)); + ICEBERG_ASSIGN_OR_RAISE(auto element_id, GetJsonInteger(json, kElementId)); ICEBERG_ASSIGN_OR_RAISE(auto element_required, GetJsonValue(json, kElementRequired)); @@ -834,8 +845,8 @@ Result> MapTypeFromJson(const nlohmann::json& json) { ICEBERG_ASSIGN_OR_RAISE( auto value_type, GetJsonValue(json, kValue).and_then(TypeFromJson)); - ICEBERG_ASSIGN_OR_RAISE(auto key_id, GetJsonValue(json, kKeyId)); - ICEBERG_ASSIGN_OR_RAISE(auto value_id, GetJsonValue(json, kValueId)); + ICEBERG_ASSIGN_OR_RAISE(auto key_id, GetJsonInteger(json, kKeyId)); + ICEBERG_ASSIGN_OR_RAISE(auto value_id, GetJsonInteger(json, kValueId)); ICEBERG_ASSIGN_OR_RAISE(auto value_required, GetJsonValue(json, kValueRequired)); SchemaField key_field(key_id, std::string(MapType::kKeyName), std::move(key_type), @@ -996,7 +1007,7 @@ Status ValidateTimestamptzDefaultIsUtc(const Type& type, const nlohmann::json& v Result> FieldFromJson(const nlohmann::json& json) { ICEBERG_ASSIGN_OR_RAISE( auto type, GetJsonValue(json, kType).and_then(TypeFromJson)); - ICEBERG_ASSIGN_OR_RAISE(auto field_id, GetJsonValue(json, kId)); + ICEBERG_ASSIGN_OR_RAISE(auto field_id, GetJsonInteger(json, kId)); ICEBERG_ASSIGN_OR_RAISE(auto name, GetJsonValue(json, kName)); ICEBERG_ASSIGN_OR_RAISE(auto required, GetJsonValue(json, kRequired)); ICEBERG_ASSIGN_OR_RAISE(auto doc, GetJsonValueOrDefault(json, kDoc)); @@ -1029,7 +1040,7 @@ Result> FieldFromJson(const nlohmann::json& json) { Result> SchemaFromJson(const nlohmann::json& json) { ICEBERG_ASSIGN_OR_RAISE(auto schema_id_opt, - GetJsonValueOptional(json, kSchemaId)); + GetJsonIntegerOptional(json, kSchemaId)); ICEBERG_ASSIGN_OR_RAISE(auto type, TypeFromJson(json)); if (type->type_id() != TypeId::kStruct) [[unlikely]] { @@ -1043,9 +1054,11 @@ Result> SchemaFromJson(const nlohmann::json& json) { fields.emplace_back(std::move(field)); } - ICEBERG_ASSIGN_OR_RAISE( - auto identifier_field_ids, - GetJsonValueOrDefault>(json, kIdentifierFieldIds)); + std::vector identifier_field_ids; + if (json.contains(kIdentifierFieldIds) && !json.at(kIdentifierFieldIds).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(identifier_field_ids, + IntegerVectorFromJson(json, kIdentifierFieldIds)); + } return Schema::Make(std::move(fields), schema_id_opt.value_or(Schema::kInitialSchemaId), std::move(identifier_field_ids)); @@ -1078,14 +1091,13 @@ Result ToJsonString(const PartitionSpec& partition_spec) { Result> PartitionFieldFromJson( const nlohmann::json& json, bool allow_field_id_missing) { - ICEBERG_ASSIGN_OR_RAISE(auto source_id, GetJsonValue(json, kSourceId)); + ICEBERG_ASSIGN_OR_RAISE(auto source_id, GetJsonInteger(json, kSourceId)); int32_t field_id; - if (allow_field_id_missing) { + if (allow_field_id_missing && !json.contains(kFieldId)) { // Partition field id in v1 is not tracked, so we use -1 to indicate that. - ICEBERG_ASSIGN_OR_RAISE(field_id, GetJsonValueOrDefault( - json, kFieldId, SchemaField::kInvalidFieldId)); + field_id = SchemaField::kInvalidFieldId; } else { - ICEBERG_ASSIGN_OR_RAISE(field_id, GetJsonValue(json, kFieldId)); + ICEBERG_ASSIGN_OR_RAISE(field_id, GetJsonInteger(json, kFieldId)); } ICEBERG_ASSIGN_OR_RAISE( auto transform, @@ -1117,13 +1129,13 @@ Result> ParseBoundPartitionSpec( Result> PartitionSpecFromJson( const std::shared_ptr& schema, const nlohmann::json& json, int32_t default_spec_id) { - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(json, kSpecId)); return ParseBoundPartitionSpec(*schema, json, spec_id, /*allow_missing_fields=*/spec_id != default_spec_id); } Result> PartitionSpecFromJson(const nlohmann::json& json) { - ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonValue(json, kSpecId)); + ICEBERG_ASSIGN_OR_RAISE(auto spec_id, GetJsonInteger(json, kSpecId)); ICEBERG_ASSIGN_OR_RAISE(auto fields, GetJsonValue(json, kFields)); std::vector partition_fields; @@ -1488,7 +1500,7 @@ Result> ParseSchemas( } ICEBERG_ASSIGN_OR_RAISE(current_schema_id, - GetJsonValue(json, kCurrentSchemaId)); + GetJsonInteger(json, kCurrentSchemaId)); for (const auto& schema_json : schema_array) { ICEBERG_ASSIGN_OR_RAISE(std::shared_ptr schema, SchemaFromJson(schema_json)); @@ -1534,7 +1546,8 @@ Status ParsePartitionSpecs(const nlohmann::json& json, int8_t format_version, return JsonParseError("Cannot parse partition specs from non-array: {}", SafeDumpJson(spec_array)); } - ICEBERG_ASSIGN_OR_RAISE(default_spec_id, GetJsonValue(json, kDefaultSpecId)); + ICEBERG_ASSIGN_OR_RAISE(default_spec_id, + GetJsonInteger(json, kDefaultSpecId)); for (const auto& spec_json : spec_array) { ICEBERG_ASSIGN_OR_RAISE( diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index 1085465f8..d5a6f2fd5 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -17,6 +17,7 @@ * under the License. */ +#include #include #include #include @@ -199,6 +200,70 @@ TEST(JsonInternalTest, SchemaFieldRejectsNonUtcTimestamptzDefault) { EXPECT_TRUE(FieldFromJson(utc).has_value()); } +TEST(JsonInternalTest, SchemaIdsRejectInvalidJsonIntegers) { + const nlohmann::json valid_schema = R"({ + "type": "struct", + "schema-id": 10, + "identifier-field-ids": [20], + "fields": [ + {"id": 20, "name": "id", "required": true, "type": "int"}, + { + "id": 30, + "name": "items", + "required": false, + "type": { + "type": "list", + "element-id": 40, + "element-required": false, + "element": { + "type": "map", + "key-id": 50, + "key": "string", + "value-id": 60, + "value": "long", + "value-required": false + } + } + } + ] + })"_json; + + struct IdPath { + nlohmann::json::json_pointer pointer; + std::string_view key; + }; + const std::array id_paths = {{ + {.pointer = nlohmann::json::json_pointer("/schema-id"), .key = "schema-id"}, + {.pointer = nlohmann::json::json_pointer("/identifier-field-ids/0"), + .key = "identifier-field-ids"}, + {.pointer = nlohmann::json::json_pointer("/fields/0/id"), .key = "id"}, + {.pointer = nlohmann::json::json_pointer("/fields/1/type/element-id"), + .key = "element-id"}, + {.pointer = nlohmann::json::json_pointer("/fields/1/type/element/key-id"), + .key = "key-id"}, + {.pointer = nlohmann::json::json_pointer("/fields/1/type/element/value-id"), + .key = "value-id"}, + }}; + const std::array, 3> invalid_ids = {{ + {nlohmann::json(1.5), "must be an integer"}, + {nlohmann::json(true), "must be an integer"}, + {nlohmann::json(2147483648ULL), "out of range"}, + }}; + + for (const auto& id_path : id_paths) { + for (const auto& [invalid_id, expected_error] : invalid_ids) { + SCOPED_TRACE(id_path.pointer.to_string() + "=" + invalid_id.dump()); + auto invalid_schema = valid_schema; + invalid_schema[id_path.pointer] = invalid_id; + + auto result = SchemaFromJson(invalid_schema); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(id_path.key)); + EXPECT_THAT(result, HasErrorMessage(expected_error)); + } + } +} + TEST(JsonInternalTest, SortField) { auto identity_transform = Transform::Identity(); @@ -312,6 +377,77 @@ TEST(JsonInternalTest, PartitionSpecFromJson) { EXPECT_EQ(*spec, *parsed); } +TEST(JsonInternalTest, PartitionIdsRejectInvalidJsonIntegers) { + const nlohmann::json valid_spec = R"({ + "spec-id": 10, + "fields": [{ + "source-id": 20, + "field-id": 1000, + "transform": "identity", + "name": "id" + }] + })"_json; + auto schema = std::make_shared( + std::vector{SchemaField(20, "id", int32(), false)}, + /*schema_id=*/10); + + struct IdPath { + nlohmann::json::json_pointer pointer; + std::string_view key; + }; + const std::array id_paths = {{ + {.pointer = nlohmann::json::json_pointer("/spec-id"), .key = "spec-id"}, + {.pointer = nlohmann::json::json_pointer("/fields/0/source-id"), + .key = "source-id"}, + {.pointer = nlohmann::json::json_pointer("/fields/0/field-id"), .key = "field-id"}, + }}; + const std::array, 3> invalid_ids = {{ + {nlohmann::json(1.5), "must be an integer"}, + {nlohmann::json(true), "must be an integer"}, + {nlohmann::json(2147483648ULL), "out of range"}, + }}; + + for (const auto& id_path : id_paths) { + for (const auto& [invalid_id, expected_error] : invalid_ids) { + SCOPED_TRACE(id_path.pointer.to_string() + "=" + invalid_id.dump()); + auto invalid_spec = valid_spec; + invalid_spec[id_path.pointer] = invalid_id; + + auto unbound_result = PartitionSpecFromJson(invalid_spec); + EXPECT_THAT(unbound_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(unbound_result, HasErrorMessage(id_path.key)); + EXPECT_THAT(unbound_result, HasErrorMessage(expected_error)); + + auto bound_result = PartitionSpecFromJson(schema, invalid_spec, 10); + EXPECT_THAT(bound_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(bound_result, HasErrorMessage(id_path.key)); + EXPECT_THAT(bound_result, HasErrorMessage(expected_error)); + } + } + + auto legacy_field = valid_spec["fields"][0]; + legacy_field.erase("field-id"); + ICEBERG_UNWRAP_OR_FAIL( + auto parsed_legacy_field, + PartitionFieldFromJson(legacy_field, /*allow_field_id_missing=*/true)); + EXPECT_EQ(parsed_legacy_field->field_id(), SchemaField::kInvalidFieldId); + + for (const auto& [invalid_id, expected_error] : invalid_ids) { + SCOPED_TRACE("legacy field-id=" + invalid_id.dump()); + legacy_field["field-id"] = invalid_id; + auto result = PartitionFieldFromJson(legacy_field, /*allow_field_id_missing=*/true); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("field-id")); + EXPECT_THAT(result, HasErrorMessage(expected_error)); + } + + legacy_field["field-id"] = nullptr; + auto null_result = + PartitionFieldFromJson(legacy_field, /*allow_field_id_missing=*/true); + EXPECT_THAT(null_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(null_result, HasErrorMessage("field-id")); +} + TEST(JsonInternalTest, SnapshotRefBranch) { SnapshotRef ref(1234567890, SnapshotRef::Branch{.min_snapshots_to_keep = 10, .max_snapshot_age_ms = 123456789, diff --git a/src/iceberg/test/metadata_serde_test.cc b/src/iceberg/test/metadata_serde_test.cc index 10529d128..88272866a 100644 --- a/src/iceberg/test/metadata_serde_test.cc +++ b/src/iceberg/test/metadata_serde_test.cc @@ -327,6 +327,24 @@ TEST(MetadataSerdeTest, DeserializeV2ValidMinimal) { ASSERT_FALSE(metadata->Snapshot().has_value()); } +TEST(MetadataSerdeTest, RejectsFractionalSchemaAndDefaultSpecIds) { + ICEBERG_UNWRAP_OR_FAIL( + auto metadata, ReadTableMetadataFromResource("TableMetadataV2ValidMinimal.json")); + ICEBERG_UNWRAP_OR_FAIL(auto valid_json, ToJson(*metadata)); + ASSERT_THAT(TableMetadataFromJson(valid_json), IsOk()); + + for (const auto* key : {"current-schema-id", "default-spec-id"}) { + SCOPED_TRACE(key); + auto invalid_json = valid_json; + invalid_json[key] = 0.5; + + auto result = TableMetadataFromJson(invalid_json); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage(key)); + EXPECT_THAT(result, HasErrorMessage("must be an integer")); + } +} + TEST(MetadataSerdeTest, DeserializeStatisticsFiles) { ICEBERG_UNWRAP_OR_FAIL( auto metadata, ReadTableMetadataFromResource("TableMetadataStatisticsFiles.json")); From 28851ed5d036153d1ea3de27b969ac069904b6af Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Wed, 30 Sep 2026 10:44:57 +0000 Subject: [PATCH 11/11] test(serde): cover signed schema and partition ID bounds --- src/iceberg/test/json_serde_test.cc | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/src/iceberg/test/json_serde_test.cc b/src/iceberg/test/json_serde_test.cc index d5a6f2fd5..79028d380 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -244,9 +244,11 @@ TEST(JsonInternalTest, SchemaIdsRejectInvalidJsonIntegers) { {.pointer = nlohmann::json::json_pointer("/fields/1/type/element/value-id"), .key = "value-id"}, }}; - const std::array, 3> invalid_ids = {{ + const std::array, 5> invalid_ids = {{ {nlohmann::json(1.5), "must be an integer"}, {nlohmann::json(true), "must be an integer"}, + {nlohmann::json(-2147483649LL), "out of range"}, + {nlohmann::json(2147483648LL), "out of range"}, {nlohmann::json(2147483648ULL), "out of range"}, }}; @@ -401,9 +403,11 @@ TEST(JsonInternalTest, PartitionIdsRejectInvalidJsonIntegers) { .key = "source-id"}, {.pointer = nlohmann::json::json_pointer("/fields/0/field-id"), .key = "field-id"}, }}; - const std::array, 3> invalid_ids = {{ + const std::array, 5> invalid_ids = {{ {nlohmann::json(1.5), "must be an integer"}, {nlohmann::json(true), "must be an integer"}, + {nlohmann::json(-2147483649LL), "out of range"}, + {nlohmann::json(2147483648LL), "out of range"}, {nlohmann::json(2147483648ULL), "out of range"}, }};