diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..46593f0b5 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( @@ -322,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. @@ -357,97 +192,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 { @@ -495,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()) { @@ -521,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()) { @@ -557,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))); } diff --git a/src/iceberg/json_serde.cc b/src/iceberg/json_serde.cc index f322996ec..45dfa29f5 100644 --- a/src/iceberg/json_serde.cc +++ b/src/iceberg/json_serde.cc @@ -18,9 +18,13 @@ */ #include +#include #include #include +#include +#include #include +#include #include #include @@ -29,12 +33,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 +248,311 @@ 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 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; + 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))); + 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; +} + +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)}}; +} + +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) { + 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 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); + SetBytesMap(json, kLowerBounds, data_file.lower_bounds); + SetBytesMap(json, kUpperBounds, data_file.upper_bounds); + + if (!data_file.key_metadata.empty()) { + json[kKeyMetadata] = BytesToHex(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 +2298,138 @@ 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 || content_str == "DATA") { + data_file.content = DataFile::Content::kData; + } else if (content_str == kContentPositionDeletes || + content_str == "POSITION_DELETES") { + data_file.content = DataFile::Content::kPositionDeletes; + } else if (content_str == kContentEqualityDeletes || + content_str == "EQUALITY_DELETES") { + 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, GetJsonInteger(json, kSpecId)); + data_file.partition_spec_id = spec_id; + + auto it = partition_spec_by_id.find(spec_id); + if (it == partition_spec_by_id.end() || !it->second) { + return JsonParseError("Invalid partition spec id: {}", spec_id); + } + 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)); + } + + ICEBERG_ASSIGN_OR_RAISE(data_file.record_count, + GetJsonInteger(json, kRecordCount)); + ICEBERG_ASSIGN_OR_RAISE(data_file.file_size_in_bytes, + GetJsonInteger(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, 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, + StringUtils::HexStringToBytes(key_metadata)); + } + if (json.contains(kSplitOffsets) && !json.at(kSplitOffsets).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.split_offsets, + IntegerVectorFromJson(json, kSplitOffsets)); + } + if (json.contains(kEqualityIds) && !json.at(kEqualityIds).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.equality_ids, + IntegerVectorFromJson(json, kEqualityIds)); + } + if (json.contains(kSortOrderId) && !json.at(kSortOrderId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.sort_order_id, + GetJsonInteger(json, kSortOrderId)); + } + if (json.contains(kFirstRowId) && !json.at(kFirstRowId).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.first_row_id, + GetJsonInteger(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, + GetJsonInteger(json, kContentOffset)); + } + if (json.contains(kContentSizeInBytes) && !json.at(kContentSizeInBytes).is_null()) { + ICEBERG_ASSIGN_OR_RAISE(data_file.content_size_in_bytes, + GetJsonInteger(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"); + } + 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); +} + } // 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..7bb867290 100644 --- a/src/iceberg/test/json_serde_test.cc +++ b/src/iceberg/test/json_serde_test.cc @@ -18,15 +18,21 @@ */ #include +#include +#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 +1050,274 @@ 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); +} + +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, ParsesPartitionObject) { + auto json = JavaGolden(); + json["partition"] = {{"1000", 7}}; + + ICEBERG_UNWRAP_OR_FAIL(auto parsed, DataFileFromJson(json, specs_, schema_)); + 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{.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); + + 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", "first-row-id", + "content-offset", "content-size-in-bytes"}) { + SCOPED_TRACE(field); + auto json = JavaGolden(); + json[field] = 0.5; + 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) { + 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) {