From cce43f7214a28a00b92b4ac56de5f21636c8975a Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Thu, 1 Oct 2026 09:00:41 +0000 Subject: [PATCH 1/2] feat(scan): support validated FileScanTask ranges Preserve start and length in core and REST JSON, and reject partial reads until the reader can honor ranges. --- src/iceberg/catalog/rest/json_serde.cc | 45 ++- src/iceberg/catalog/rest/types.cc | 1 + src/iceberg/data/file_scan_task_reader.cc | 6 + src/iceberg/data/file_scan_task_reader.h | 3 + src/iceberg/table_scan.cc | 70 ++++- src/iceberg/table_scan.h | 34 +++ .../test/file_scan_task_reader_test.cc | 19 ++ src/iceberg/test/file_scan_task_test.cc | 115 ++++++++ src/iceberg/test/rest_json_serde_test.cc | 264 ++++++++++++++++++ 9 files changed, 552 insertions(+), 5 deletions(-) diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..da754ff61 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 @@ -132,9 +133,24 @@ 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 kStart = "start"; +constexpr std::string_view kLength = "length"; constexpr std::string_view kMapKeys = "keys"; constexpr std::string_view kMapValues = "values"; +Result GetRangeValue(const nlohmann::json& json, std::string_view key) { + 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()); nlohmann::json json; @@ -323,6 +339,17 @@ Result>> FileScanTasksFromJson( GetJsonValue(task_json, kDataFile)); ICEBERG_ASSIGN_OR_RAISE( auto data_file, DataFileFromJson(data_file_json, partition_spec_by_id, schema)); + const bool has_start = task_json.contains(kStart); + const bool has_length = task_json.contains(kLength); + if (has_start != has_length) { + return JsonParseError("'start' and 'length' must be provided together"); + } + int64_t start = 0; + int64_t length = data_file.file_size_in_bytes; + if (has_start) { + ICEBERG_ASSIGN_OR_RAISE(start, GetRangeValue(task_json, kStart)); + ICEBERG_ASSIGN_OR_RAISE(length, GetRangeValue(task_json, kLength)); + } // 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. @@ -350,9 +377,13 @@ Result>> FileScanTasksFromJson( ICEBERG_ASSIGN_OR_RAISE(residual_filter, ExpressionFromJson(filter_json)); } - file_scan_tasks.push_back(std::make_shared( - std::make_shared(std::move(data_file)), std::move(task_delete_files), - std::move(residual_filter))); + auto task = FileScanTask::MakeSplit(std::make_shared(std::move(data_file)), + start, length, std::move(task_delete_files), + std::move(residual_filter)); + if (!task) { + return JsonParseError("Invalid FileScanTask range: {}", task.error().message); + } + file_scan_tasks.push_back(std::move(task.value())); } return file_scan_tasks; } @@ -493,10 +524,18 @@ Result ScanTaskFieldsToJson( if (!task) continue; nlohmann::json task_json; if (task->data_file()) { + auto range_result = + FileScanTask::MakeSplit(task->data_file(), task->start(), task->length()); + if (!range_result) { + return ValidationFailed("Invalid FileScanTask range: {}", + range_result.error().message); + } ICEBERG_ASSIGN_OR_RAISE( auto data_file_json, ToJson(*task->data_file(), partition_specs_by_id, schema)); task_json[kDataFile] = std::move(data_file_json); + task_json[kStart] = task->start(); + task_json[kLength] = task->length(); } if (!task->delete_files().empty()) { std::vector refs; diff --git a/src/iceberg/catalog/rest/types.cc b/src/iceberg/catalog/rest/types.cc index 84fba9a7c..ebc2c2876 100644 --- a/src/iceberg/catalog/rest/types.cc +++ b/src/iceberg/catalog/rest/types.cc @@ -192,6 +192,7 @@ bool FileScanTaskEqual(const std::shared_ptr& lhs, return false; } return SharedPtrEqual(lhs->data_file(), rhs->data_file()) && + lhs->start() == rhs->start() && lhs->length() == rhs->length() && SharedPtrVectorEqual(lhs->delete_files(), rhs->delete_files()) && ExpressionPtrEqual(lhs->residual_filter(), rhs->residual_filter()); } diff --git a/src/iceberg/data/file_scan_task_reader.cc b/src/iceberg/data/file_scan_task_reader.cc index 0c41f25ed..7d3a8bc7e 100644 --- a/src/iceberg/data/file_scan_task_reader.cc +++ b/src/iceberg/data/file_scan_task_reader.cc @@ -163,6 +163,12 @@ class FileScanTaskReader::Impl { ICEBERG_PRECHECK(data_file->file_size_in_bytes >= 0, "Data file size must not be negative: {}", data_file->file_size_in_bytes); + if (task.is_split()) { + return NotSupported( + "Reading a partial FileScanTask is not supported yet: {} " + "(start={}, length={})", + data_file->file_path, task.start(), task.length()); + } if (task.delete_files().empty()) { auto options = MakeReaderOptions( diff --git a/src/iceberg/data/file_scan_task_reader.h b/src/iceberg/data/file_scan_task_reader.h index a71ef5f84..c34f3e7aa 100644 --- a/src/iceberg/data/file_scan_task_reader.h +++ b/src/iceberg/data/file_scan_task_reader.h @@ -67,6 +67,9 @@ class ICEBERG_DATA_EXPORT FileScanTaskReader { ~FileScanTaskReader(); /// \brief Open a task and return an Arrow C stream for its projected live rows. + /// + /// Partial tasks are rejected until range-aware reading is implemented. The task and + /// its shared data/delete-file metadata must not be mutated while the stream uses them. Result Open(const FileScanTask& task); FileScanTaskReader(const FileScanTaskReader&) = delete; diff --git a/src/iceberg/table_scan.cc b/src/iceberg/table_scan.cc index e4a85e9a6..604a9202a 100644 --- a/src/iceberg/table_scan.cc +++ b/src/iceberg/table_scan.cc @@ -40,6 +40,7 @@ #include "iceberg/table.h" #include "iceberg/table_metadata.h" #include "iceberg/util/content_file_util.h" +#include "iceberg/util/int128.h" #include "iceberg/util/macros.h" #include "iceberg/util/snapshot_util.h" #include "iceberg/util/timepoint.h" @@ -271,11 +272,76 @@ FileScanTask::FileScanTask(std::shared_ptr data_file, ICEBERG_DCHECK(data_file_ != nullptr, "Data file cannot be null for FileScanTask"); } -int64_t FileScanTask::size_bytes() const { return data_file_->file_size_in_bytes; } +FileScanTask::FileScanTask(std::shared_ptr data_file, + std::vector> delete_files, + std::shared_ptr residual_filter, Range range) + : FileScanTask(std::move(data_file), std::move(delete_files), + std::move(residual_filter)) { + range_ = range; +} + +Result> FileScanTask::MakeSplit( + std::shared_ptr data_file, int64_t start, int64_t length, + std::vector> delete_files, + std::shared_ptr filter) { + if (data_file == nullptr) { + return InvalidArgument("Cannot create a file scan task without a data file"); + } + + const int64_t file_size = data_file->file_size_in_bytes; + if (file_size < 0) { + return InvalidArgument("Cannot split file {} with negative size {}", + data_file->file_path, file_size); + } + if (start < 0 || start > file_size) { + return InvalidArgument("Split start {} is outside file {} of size {}", start, + data_file->file_path, file_size); + } + // Subtraction is safe after validating start, even when file_size is INT64_MAX. + if (length < 0 || length > file_size - start) { + return InvalidArgument("Split length {} from start {} is outside file {} of size {}", + length, start, data_file->file_path, file_size); + } + if (start == 0 && length == file_size) { + return std::make_shared(std::move(data_file), std::move(delete_files), + std::move(filter)); + } + if (length == 0) { + return InvalidArgument("A partial split of file {} must have positive length", + data_file->file_path); + } + + return std::shared_ptr( + new FileScanTask(std::move(data_file), std::move(delete_files), std::move(filter), + Range{start, length, file_size})); +} + +int64_t FileScanTask::length() const { + return range_ ? range_->length : data_file_->file_size_in_bytes; +} + +int64_t FileScanTask::size_bytes() const { return length(); } int32_t FileScanTask::files_count() const { return 1; } -int64_t FileScanTask::estimated_row_count() const { return data_file_->record_count; } +int64_t FileScanTask::estimated_row_count() const { + if (!range_) { + return data_file_->record_count; + } + + const int64_t record_count = data_file_->record_count; + if (record_count <= 0) { + return 0; + } + + // Cumulative boundaries make estimates additive across contiguous splits, and + // 128-bit products avoid overflow for any valid int64 file size and row count. + const auto end = static_cast(range_->start) + range_->length; + const auto count_at_end = end * record_count / range_->file_size; + const auto count_at_start = + static_cast(range_->start) * record_count / range_->file_size; + return static_cast(count_at_end - count_at_start); +} // ChangelogScanTask implementation diff --git a/src/iceberg/table_scan.h b/src/iceberg/table_scan.h index 7310d435b..c05b5cdee 100644 --- a/src/iceberg/table_scan.h +++ b/src/iceberg/table_scan.h @@ -22,6 +22,7 @@ /// \file iceberg/table_scan.h /// \brief Define table scan APIs and scan task types. +#include #include #include #include @@ -71,10 +72,32 @@ class ICEBERG_EXPORT FileScanTask : public ScanTask { /// \param data_file The data file to read. /// \param delete_files Delete files that apply to this data file. /// \param filter Optional residual filter to apply after reading. + /// \note The data file, delete files, and residual filter must not be mutated while + /// a planning or reading stream uses this task. explicit FileScanTask(std::shared_ptr data_file, std::vector> delete_files = {}, std::shared_ptr filter = nullptr); + /// \brief Make a task for a validated half-open byte range [start, start + length). + /// + /// A range covering the entire file is represented as a whole-file task. An empty + /// range is valid only for an empty whole file. The data file, delete files, and + /// residual filter are shared with the resulting task and must not be mutated + /// while a planning or reading stream uses it. + static Result> MakeSplit( + std::shared_ptr data_file, int64_t start, int64_t length, + std::vector> delete_files = {}, + std::shared_ptr filter = nullptr); + + /// \brief The inclusive byte offset where this task starts, or zero for a whole file. + int64_t start() const { return range_ ? range_->start : 0; } + + /// \brief The byte length of this task, or the current size of a whole file. + int64_t length() const; + + /// \brief Whether this task reads a proper subset of its data file. + bool is_split() const { return range_.has_value(); } + /// \brief The data file that should be read by this scan task. const std::shared_ptr& data_file() const { return data_file_; } @@ -92,9 +115,20 @@ class ICEBERG_EXPORT FileScanTask : public ScanTask { int64_t estimated_row_count() const override; private: + struct Range { + int64_t start; + int64_t length; + int64_t file_size; + }; + + FileScanTask(std::shared_ptr data_file, + std::vector> delete_files, + std::shared_ptr filter, Range range); + std::shared_ptr data_file_; std::vector> delete_files_; std::shared_ptr residual_filter_; + std::optional range_; }; /// \brief Stream of file scan tasks. diff --git a/src/iceberg/test/file_scan_task_reader_test.cc b/src/iceberg/test/file_scan_task_reader_test.cc index cca13ce0d..11366f72d 100644 --- a/src/iceberg/test/file_scan_task_reader_test.cc +++ b/src/iceberg/test/file_scan_task_reader_test.cc @@ -409,6 +409,25 @@ TEST_F(FileScanTaskReaderTest, OpenWithoutDeletesReadsProjectedSchema) { VerifyStream(&stream, R"([[1, "Foo"], [2, "Bar"], [3, "Baz"]])")); } +TEST_F(FileScanTaskReaderTest, RejectsPartialTaskUntilRangeReadsAreSupported) { + ICEBERG_UNWRAP_OR_FAIL( + auto data_file, + MakeDataFile(table_schema_, + R"([[1, "Foo", "blue"], [2, "Bar", "red"], [3, "Baz", "green"]])")); + ICEBERG_UNWRAP_OR_FAIL(auto task, FileScanTask::MakeSplit(data_file, 0, 1)); + + FileScanTaskReader::Options options{ + .io = file_io_, + .table_schema = table_schema_, + .schemas = {table_schema_}, + .projected_schema = projected_schema_, + }; + ICEBERG_UNWRAP_OR_FAIL(auto reader, FileScanTaskReader::Make(std::move(options))); + auto result = reader->Open(*task); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); + EXPECT_THAT(result, HasErrorMessage("Reading a partial FileScanTask")); +} + TEST_F(FileScanTaskReaderTest, ReadLastUpdatedFromDataSeq) { ICEBERG_UNWRAP_OR_FAIL( auto data_file, diff --git a/src/iceberg/test/file_scan_task_test.cc b/src/iceberg/test/file_scan_task_test.cc index cb945ba1f..6394c96fe 100644 --- a/src/iceberg/test/file_scan_task_test.cc +++ b/src/iceberg/test/file_scan_task_test.cc @@ -17,6 +17,10 @@ * under the License. */ +#include +#include +#include + #include #include #include @@ -29,6 +33,7 @@ #include "iceberg/arrow/arrow_io_internal.h" #include "iceberg/data/file_scan_task_reader.h" +#include "iceberg/expression/expressions.h" #include "iceberg/file_format.h" #include "iceberg/manifest/manifest_entry.h" #include "iceberg/parquet/parquet_register.h" @@ -40,6 +45,116 @@ namespace iceberg { +TEST(FileScanTaskRangeTest, WholeFileAndEmptyFile) { + auto data_file = std::make_shared(); + data_file->file_path = "test.parquet"; + data_file->file_size_in_bytes = 100; + data_file->record_count = 7; + + FileScanTask whole(data_file); + EXPECT_FALSE(whole.is_split()); + EXPECT_EQ(whole.start(), 0); + EXPECT_EQ(whole.length(), 100); + EXPECT_EQ(whole.size_bytes(), 100); + EXPECT_EQ(whole.estimated_row_count(), 7); + + ICEBERG_UNWRAP_OR_FAIL(auto explicit_whole, FileScanTask::MakeSplit(data_file, 0, 100)); + EXPECT_FALSE(explicit_whole->is_split()); + EXPECT_EQ(explicit_whole->start(), 0); + EXPECT_EQ(explicit_whole->length(), 100); + EXPECT_EQ(explicit_whole->estimated_row_count(), 7); + + data_file->file_size_in_bytes = 0; + data_file->record_count = 0; + ICEBERG_UNWRAP_OR_FAIL(auto empty, FileScanTask::MakeSplit(data_file, 0, 0)); + EXPECT_FALSE(empty->is_split()); + EXPECT_EQ(empty->start(), 0); + EXPECT_EQ(empty->length(), 0); + EXPECT_EQ(empty->size_bytes(), 0); + EXPECT_EQ(empty->estimated_row_count(), 0); +} + +TEST(FileScanTaskRangeTest, DistinctSplitsShareMetadataAndEstimateRowsOnce) { + auto data_file = std::make_shared(); + data_file->file_path = "test.parquet"; + data_file->file_size_in_bytes = 100; + data_file->record_count = 7; + data_file->first_row_id = 1000; + data_file->data_sequence_number = 22; + auto delete_file = std::make_shared(); + auto residual = Expressions::AlwaysTrue(); + + ICEBERG_UNWRAP_OR_FAIL( + auto first, FileScanTask::MakeSplit(data_file, 0, 40, {delete_file}, residual)); + ICEBERG_UNWRAP_OR_FAIL( + auto second, FileScanTask::MakeSplit(data_file, 40, 60, {delete_file}, residual)); + + EXPECT_TRUE(first->is_split()); + EXPECT_TRUE(second->is_split()); + EXPECT_EQ(first->start(), 0); + EXPECT_EQ(first->length(), 40); + EXPECT_EQ(second->start(), 40); + EXPECT_EQ(second->length(), 60); + EXPECT_EQ(first->size_bytes() + second->size_bytes(), 100); + EXPECT_EQ(first->estimated_row_count(), 2); + EXPECT_EQ(second->estimated_row_count(), 5); + EXPECT_EQ(first->estimated_row_count() + second->estimated_row_count(), 7); + EXPECT_EQ(first->data_file(), data_file); + EXPECT_EQ(second->data_file(), data_file); + ASSERT_EQ(first->delete_files().size(), 1); + ASSERT_EQ(second->delete_files().size(), 1); + EXPECT_EQ(first->delete_files()[0], delete_file); + EXPECT_EQ(second->delete_files()[0], delete_file); + EXPECT_EQ(first->residual_filter(), residual); + EXPECT_EQ(second->residual_filter(), residual); + EXPECT_EQ(first->data_file()->first_row_id, 1000); + EXPECT_EQ(first->data_file()->data_sequence_number, 22); + EXPECT_EQ(second->data_file()->first_row_id, 1000); + EXPECT_EQ(second->data_file()->data_sequence_number, 22); +} + +TEST(FileScanTaskRangeTest, RejectsInvalidRanges) { + auto data_file = std::make_shared(); + data_file->file_path = "test.parquet"; + data_file->file_size_in_bytes = 100; + + const auto expect_invalid = [&](int64_t start, int64_t length) { + auto result = FileScanTask::MakeSplit(data_file, start, length); + ASSERT_FALSE(result.has_value()); + EXPECT_EQ(result.error().kind, ErrorKind::kInvalidArgument); + }; + expect_invalid(-1, 1); + expect_invalid(0, -1); + expect_invalid(0, 0); + expect_invalid(100, 0); + expect_invalid(101, 1); + expect_invalid(99, 2); + + auto null_file = FileScanTask::MakeSplit(nullptr, 0, 1); + ASSERT_FALSE(null_file.has_value()); + EXPECT_EQ(null_file.error().kind, ErrorKind::kInvalidArgument); + + data_file->file_size_in_bytes = -1; + expect_invalid(0, 1); + + data_file->file_size_in_bytes = std::numeric_limits::max(); + expect_invalid(std::numeric_limits::max() - 1, 2); +} + +TEST(FileScanTaskRangeTest, RowEstimateDoesNotOverflow) { + auto data_file = std::make_shared(); + data_file->file_size_in_bytes = std::numeric_limits::max(); + data_file->record_count = std::numeric_limits::max(); + + ICEBERG_UNWRAP_OR_FAIL( + auto first, + FileScanTask::MakeSplit(data_file, 0, data_file->file_size_in_bytes - 1)); + ICEBERG_UNWRAP_OR_FAIL(auto last, FileScanTask::MakeSplit( + data_file, data_file->file_size_in_bytes - 1, 1)); + EXPECT_EQ(first->estimated_row_count(), data_file->record_count - 1); + EXPECT_EQ(last->estimated_row_count(), 1); +} + class FileScanTaskTest : public TempFileTestBase { protected: static void SetUpTestSuite() { parquet::RegisterAll(); } diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index ec41e4a66..90a85dd89 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -17,9 +17,12 @@ * under the License. */ +#include #include #include +#include #include +#include #include #include @@ -2280,10 +2283,197 @@ TEST(FileScanTasksFromJsonTest, SingleTaskNoDeleteFiles) { const auto& task = result.value()[0]; ASSERT_NE(task->data_file(), nullptr); EXPECT_EQ(task->data_file()->file_path, "s3://bucket/data/file.parquet"); + EXPECT_EQ(task->start(), 0); + EXPECT_EQ(task->length(), 12345); + EXPECT_FALSE(task->is_split()); EXPECT_TRUE(task->delete_files().empty()); 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; + + auto result = FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)); + ASSERT_THAT(result, IsOk()); + ASSERT_EQ(result.value().size(), 1U); + EXPECT_EQ(result.value()[0]->start(), 0); + EXPECT_EQ(result.value()[0]->length(), 12345); + EXPECT_FALSE(result.value()[0]->is_split()); +} + +TEST(FileScanTasksFromJsonTest, RejectsNullOrMissingRangeEndpoint) { + 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 (std::string_view missing_field : {"start", "length"}) { + SCOPED_TRACE(missing_field); + auto incomplete_json = json; + incomplete_json[0].erase(std::string(missing_field)); + auto result = + FileScanTasksFromJson(incomplete_json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("must be provided together")); + } + + json[0]["start"] = nullptr; + json[0]["length"] = nullptr; + auto null_result = FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(null_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(null_result, HasErrorMessage("'start' must be an integer")); +} + +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, AcceptsValidPartialRanges) { + 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& [start, length] : + {std::pair{0, 200}, {100, 200}, {12245, 100}}) { + SCOPED_TRACE(testing::Message() << "start=" << start << ", length=" << length); + auto split_json = json; + split_json[0]["start"] = start; + split_json[0]["length"] = length; + auto result = + FileScanTasksFromJson(split_json, {}, UnpartitionedSpecs(), Schema({}, 0)); + ASSERT_THAT(result, IsOk()); + ASSERT_EQ(result.value().size(), 1U); + EXPECT_EQ(result.value()[0]->start(), start); + EXPECT_EQ(result.value()[0]->length(), length); + EXPECT_TRUE(result.value()[0]->is_split()); + } +} + +TEST(FileScanTasksFromJsonTest, RejectsInvalidRanges) { + 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& [start, length] : {std::pair{-1, 100}, + {0, -1}, + {0, 0}, + {12345, 0}, + {12345, 1}, + {12000, 500}}) { + SCOPED_TRACE(testing::Message() << "start=" << start << ", length=" << length); + auto invalid_json = json; + invalid_json[0]["start"] = start; + invalid_json[0]["length"] = length; + auto result = + FileScanTasksFromJson(invalid_json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(result, HasErrorMessage("Invalid FileScanTask range")); + } + + auto overflow_json = json; + overflow_json[0]["data-file"]["file-size-in-bytes"] = + std::numeric_limits::max(); + overflow_json[0]["start"] = std::numeric_limits::max() - 1; + overflow_json[0]["length"] = 2; + auto overflow_result = + FileScanTasksFromJson(overflow_json, {}, UnpartitionedSpecs(), Schema({}, 0)); + EXPECT_THAT(overflow_result, IsError(ErrorKind::kJsonParseError)); + EXPECT_THAT(overflow_result, HasErrorMessage("Invalid FileScanTask range")); +} + +TEST(FileScanTasksFromJsonTest, AcceptsEmptyWholeFile) { + auto json = R"([{ + "data-file": { + "content": "data", + "file-path": "s3://bucket/data/empty.parquet", + "file-format": "PARQUET", + "spec-id": 0, + "partition": [], + "file-size-in-bytes": 0, + "record-count": 0 + }, + "start": 0, + "length": 0 + }])"_json; + + auto result = FileScanTasksFromJson(json, {}, UnpartitionedSpecs(), Schema({}, 0)); + ASSERT_THAT(result, IsOk()); + ASSERT_EQ(result.value().size(), 1U); + EXPECT_EQ(result.value()[0]->start(), 0); + EXPECT_EQ(result.value()[0]->length(), 0); + EXPECT_FALSE(result.value()[0]->is_split()); +} + TEST(FileScanTasksFromJsonTest, RowLineageSequence) { GTEST_SKIP() << "REST scan-task JSON does not expose data-sequence-number yet: " << "https://github.com/apache/iceberg-cpp/issues/834"; @@ -2452,6 +2642,80 @@ TEST(FetchScanTasksResponseRoundtripTest, WithFileScanTasksAndDeleteFiles) { EXPECT_EQ(*result, *result2); } +TEST(FetchScanTasksResponseRoundtripTest, PreservesDistinctSplitRanges) { + auto json = nlohmann::json::parse(R"({ + "file-scan-tasks": [ + { + "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": 5000 + }, + { + "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": 5000, + "length": 7345 + } + ] + })"); + + auto result = FetchScanTasksResponseFromJson(json, UnpartitionedSpecs(), EmptySchema()); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result.value().file_scan_tasks.has_value()); + ASSERT_EQ(result.value().file_scan_tasks->size(), 2U); + EXPECT_EQ(result.value().file_scan_tasks->at(0)->start(), 0); + EXPECT_EQ(result.value().file_scan_tasks->at(0)->length(), 5000); + EXPECT_EQ(result.value().file_scan_tasks->at(1)->start(), 5000); + EXPECT_EQ(result.value().file_scan_tasks->at(1)->length(), 7345); + + ICEBERG_UNWRAP_OR_FAIL(auto roundtrip_json, + ToJson(result.value(), UnpartitionedSpecs(), EmptySchema())); + EXPECT_EQ(roundtrip_json["file-scan-tasks"][0]["start"], 0); + EXPECT_EQ(roundtrip_json["file-scan-tasks"][0]["length"], 5000); + EXPECT_EQ(roundtrip_json["file-scan-tasks"][1]["start"], 5000); + EXPECT_EQ(roundtrip_json["file-scan-tasks"][1]["length"], 7345); + + auto result2 = + FetchScanTasksResponseFromJson(roundtrip_json, UnpartitionedSpecs(), EmptySchema()); + ASSERT_THAT(result2, IsOk()); + EXPECT_EQ(result.value(), result2.value()); +} + +TEST(FetchScanTasksResponseRoundtripTest, RejectsSplitOutsideChangedDataFile) { + auto data_file = std::make_shared(); + data_file->content = DataFile::Content::kData; + data_file->file_path = "s3://bucket/data/file.parquet"; + data_file->file_format = FileFormatType::kParquet; + data_file->partition_spec_id = PartitionSpec::kInitialSpecId; + data_file->partition = PartitionValues{}; + data_file->file_size_in_bytes = 12345; + data_file->record_count = 100; + ICEBERG_UNWRAP_OR_FAIL(auto split, FileScanTask::MakeSplit(data_file, 5000, 7345)); + + FetchScanTasksResponse response; + response.file_scan_tasks = std::vector>{split}; + data_file->file_size_in_bytes = 12000; + + auto result = ToJson(response, UnpartitionedSpecs(), EmptySchema()); + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(result, HasErrorMessage("Invalid FileScanTask range")); +} + TEST(FetchScanTasksResponseRoundtripTest, PreservesResidualFilter) { auto json = nlohmann::json::parse(R"({ "file-scan-tasks": [ From cf714b3a5da0642108343d5967ac5de0d8375331 Mon Sep 17 00:00:00 2001 From: Kam Cheung Ting Date: Sun, 4 Oct 2026 04:39:50 +0000 Subject: [PATCH 2/2] fix(scan): use designated initializer for task range --- src/iceberg/table_scan.cc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/iceberg/table_scan.cc b/src/iceberg/table_scan.cc index 604a9202a..de6c5677e 100644 --- a/src/iceberg/table_scan.cc +++ b/src/iceberg/table_scan.cc @@ -313,7 +313,7 @@ Result> FileScanTask::MakeSplit( return std::shared_ptr( new FileScanTask(std::move(data_file), std::move(delete_files), std::move(filter), - Range{start, length, file_size})); + Range{.start = start, .length = length, .file_size = file_size})); } int64_t FileScanTask::length() const {