From f256e78c6aa64c9d456f4ee0494a25b3b80d8268 Mon Sep 17 00:00:00 2001 From: xuanyili Date: Sat, 5 Sep 2026 21:32:08 +0000 Subject: [PATCH] feat: add table-aware position delete updates --- src/iceberg/CMakeLists.txt | 1 + src/iceberg/arrow/arrow_io.cc | 9 +- .../data/output_file_cleanup_internal.h | 56 ++ src/iceberg/data/position_delete_update.cc | 265 ++++++++ src/iceberg/data/position_delete_update.h | 60 ++ src/iceberg/delete_file_index.cc | 10 + src/iceberg/file_io.h | 3 + src/iceberg/test/CMakeLists.txt | 1 + src/iceberg/test/arrow_io_test.cc | 10 +- src/iceberg/test/delete_file_index_test.cc | 75 +++ .../test/position_delete_update_test.cc | 636 ++++++++++++++++++ src/iceberg/test/std_io.h | 8 +- 12 files changed, 1126 insertions(+), 8 deletions(-) create mode 100644 src/iceberg/data/output_file_cleanup_internal.h create mode 100644 src/iceberg/data/position_delete_update.cc create mode 100644 src/iceberg/data/position_delete_update.h create mode 100644 src/iceberg/test/position_delete_update_test.cc diff --git a/src/iceberg/CMakeLists.txt b/src/iceberg/CMakeLists.txt index 8a98274ff..fbb73a22d 100644 --- a/src/iceberg/CMakeLists.txt +++ b/src/iceberg/CMakeLists.txt @@ -242,6 +242,7 @@ set(ICEBERG_DATA_SOURCES data/delete_loader.cc data/equality_delete_writer.cc data/file_scan_task_reader.cc + data/position_delete_update.cc data/position_delete_writer.cc data/writer.cc) diff --git a/src/iceberg/arrow/arrow_io.cc b/src/iceberg/arrow/arrow_io.cc index 4c795badf..11f905da6 100644 --- a/src/iceberg/arrow/arrow_io.cc +++ b/src/iceberg/arrow/arrow_io.cc @@ -591,7 +591,14 @@ Result> ArrowFileSystemFileIO::NewOutputFile( /// \brief Delete a file at the given location. Status ArrowFileSystemFileIO::DeleteFile(const std::string& file_location) { ICEBERG_ASSIGN_OR_RAISE(auto path, ResolvePath(file_location)); - ICEBERG_ARROW_RETURN_NOT_OK(arrow_fs_->DeleteFile(path)); + auto status = arrow_fs_->DeleteFile(path); + if (!status.ok()) { + auto info = arrow_fs_->GetFileInfo(path); + if (info.ok() && info->type() == ::arrow::fs::FileType::NotFound) { + return {}; + } + } + ICEBERG_ARROW_RETURN_NOT_OK(status); return {}; } diff --git a/src/iceberg/data/output_file_cleanup_internal.h b/src/iceberg/data/output_file_cleanup_internal.h new file mode 100644 index 000000000..c84b5f91e --- /dev/null +++ b/src/iceberg/data/output_file_cleanup_internal.h @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +#include +#include +#include + +#include "iceberg/file_io.h" +#include "iceberg/result.h" + +namespace iceberg::internal { + +// Try every owned output, keeping failed paths for a subsequent cleanup attempt. +inline Status CleanupOutputFiles(FileIO& io, std::set& paths) { + Status first_error; + for (auto it = paths.begin(); it != paths.end();) { + auto status = io.DeleteFile(*it); + if (status.has_value()) { + it = paths.erase(it); + } else { + if (first_error.has_value()) first_error = std::move(status); + ++it; + } + } + return first_error; +} + +// Preserve the operation's error kind while reporting a cleanup failure, if any. +inline Status FailWithOutputCleanup(Error error, FileIO& io, + std::set& paths) { + auto cleanup = CleanupOutputFiles(io, paths); + if (!cleanup.has_value()) { + error.message += "; output cleanup failed: " + cleanup.error().message; + } + return std::unexpected(std::move(error)); +} + +} // namespace iceberg::internal diff --git a/src/iceberg/data/position_delete_update.cc b/src/iceberg/data/position_delete_update.cc new file mode 100644 index 000000000..5b8182361 --- /dev/null +++ b/src/iceberg/data/position_delete_update.cc @@ -0,0 +1,265 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/data/position_delete_update.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "iceberg/data/delete_loader.h" +#include "iceberg/data/output_file_cleanup_internal.h" +#include "iceberg/data/position_delete_writer.h" +#include "iceberg/deletes/dv_writer.h" +#include "iceberg/file_format.h" +#include "iceberg/file_io.h" +#include "iceberg/location_provider.h" +#include "iceberg/manifest/manifest_entry.h" +#include "iceberg/partition_spec.h" +#include "iceberg/schema.h" +#include "iceberg/table.h" +#include "iceberg/table_metadata.h" +#include "iceberg/table_properties.h" +#include "iceberg/table_scan.h" +#include "iceberg/update/row_delta.h" +#include "iceberg/util/macros.h" +#include "iceberg/util/string_util.h" +#include "iceberg/util/uuid.h" + +namespace iceberg { + +namespace { + +struct TargetFile { + std::shared_ptr data_file; + std::shared_ptr spec; + std::vector> position_delete_files; +}; + +} // namespace + +class PositionDeleteUpdate::Impl { + public: + explicit Impl(std::shared_ptr table) : table_(std::move(table)) {} + + Status Delete(std::string_view data_file_path, int64_t pos) { + ICEBERG_PRECHECK(!terminal_, "Position delete update is no longer usable"); + ICEBERG_PRECHECK(!data_file_path.empty(), "Data file path cannot be empty"); + ICEBERG_PRECHECK(pos >= 0, "Position delete must be non-negative: {}", pos); + deletes_[std::string(data_file_path)].push_back(pos); + return {}; + } + + Status Commit() { + ICEBERG_PRECHECK(!terminal_, "Position delete update is no longer usable"); + ICEBERG_PRECHECK(!deletes_.empty(), "Position delete update is empty"); + ICEBERG_PRECHECK(table_->metadata()->format_version >= 2, + "Position deletes require table format version 2 or later"); + ICEBERG_RETURN_UNEXPECTED(internal::CleanupOutputFiles(*table_->io(), output_paths_)); + auto status = CommitDeletes(); + if (!status.has_value() && status.error().kind != ErrorKind::kCommitStateUnknown) { + return internal::FailWithOutputCleanup(std::move(status.error()), *table_->io(), + output_paths_); + } + terminal_ = true; + output_paths_.clear(); + return status; + } + + private: + Status CommitDeletes() { + ICEBERG_ASSIGN_OR_RAISE(auto snapshot, table_->current_snapshot()); + ICEBERG_ASSIGN_OR_RAISE(auto targets, ResolveTargets()); + + ICEBERG_ASSIGN_OR_RAISE(auto written, table_->metadata()->format_version >= 3 + ? WriteDeletionVectors(targets) + : WriteParquetDeletes(targets)); + ICEBERG_ASSIGN_OR_RAISE(auto row_delta, table_->NewRowDelta()); + row_delta->ValidateFromSnapshot(snapshot->snapshot_id) + .ValidateDataFilesExist(written.referenced_data_files) + .ValidateDeletedFiles(); + for (const auto& file : written.data_files) { + row_delta->AddDeletes(file); + } + for (const auto& file : written.rewritten_delete_files) { + row_delta->RemoveDeletes(file); + } + + return row_delta->Commit(); + } + + Result> ResolveTargets() const { + ICEBERG_ASSIGN_OR_RAISE(auto scan_builder, table_->NewScan()); + ICEBERG_ASSIGN_OR_RAISE(auto scan, scan_builder->Build()); + ICEBERG_ASSIGN_OR_RAISE(auto tasks, scan->PlanFilesStream()); + + std::unordered_map targets; + while (targets.size() < deletes_.size()) { + ICEBERG_ASSIGN_OR_RAISE(auto task, tasks->Next()); + if (!task.has_value()) break; + const auto& data_file = (*task)->data_file(); + if (!deletes_.contains(data_file->file_path)) { + continue; + } + + ICEBERG_PRECHECK(data_file->partition_spec_id.has_value(), + "Data file is missing partition spec ID: {}", + data_file->file_path); + ICEBERG_ASSIGN_OR_RAISE(auto spec, table_->metadata()->PartitionSpecById( + *data_file->partition_spec_id)); + + TargetFile target{.data_file = data_file, .spec = std::move(spec)}; + for (const auto& delete_file : (*task)->delete_files()) { + if (delete_file->content == DataFile::Content::kPositionDeletes) { + target.position_delete_files.push_back(delete_file); + } + } + + targets.emplace(data_file->file_path, std::move(target)); + } + + for (const auto& [path, _] : deletes_) { + ICEBERG_PRECHECK(targets.contains(path), "Cannot find live data file: {}", path); + } + return targets; + } + + Result WriteDeletionVectors( + const std::unordered_map& targets) { + ICEBERG_ASSIGN_OR_RAISE(auto location_provider, table_->location_provider()); + auto output_path = location_provider->NewDataLocation( + std::format("position-deletes-{}.puffin", Uuid::GenerateV7().ToString())); + output_paths_.insert(output_path); + + DeleteLoader loader(table_->io()); + ICEBERG_ASSIGN_OR_RAISE( + auto writer, + DVWriter::Make(DVWriterOptions{ + .path = output_path, + .io = table_->io(), + .load_previous_deletes = [&targets, &loader](std::string_view path) + -> Result> { + const auto& previous = targets.at(std::string(path)).position_delete_files; + if (previous.empty()) return std::nullopt; + ICEBERG_ASSIGN_OR_RAISE(auto index, + loader.LoadPositionDeletes(previous, path)); + return std::optional(std::move(index)); + }, + })); + + for (const auto& [path, positions] : deletes_) { + const auto& target = targets.at(path); + for (int64_t pos : positions) { + ICEBERG_RETURN_UNEXPECTED( + writer->Delete(path, pos, target.spec, target.data_file->partition)); + } + } + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + return writer->Metadata(); + } + + Result WriteParquetDeletes( + const std::unordered_map& targets) { + ICEBERG_ASSIGN_OR_RAISE(auto location_provider, table_->location_provider()); + ICEBERG_ASSIGN_OR_RAISE(auto schema, table_->schema()); + + DeleteWriteResult result; + result.data_files.reserve(deletes_.size()); + result.referenced_data_files.reserve(deletes_.size()); + auto properties = table_->properties().configs(); + properties[TableProperties::kParquetCompression.key()] = + table_->properties().Get(TableProperties::kDeleteParquetCompression); + properties[TableProperties::kParquetCompressionLevel.key()] = + table_->properties().Get(TableProperties::kDeleteParquetCompressionLevel); + const auto write_uuid = Uuid::GenerateV7().ToString(); + size_t file_number = 0; + + for (const auto& [path, positions] : deletes_) { + auto sorted_positions = positions; + std::ranges::sort(sorted_positions); + + const auto& target = targets.at(path); + const auto filename = + std::format("position-deletes-{}-{}.parquet", write_uuid, file_number++); + auto output_path = location_provider->NewDataLocation(filename); + output_paths_.insert(output_path); + + ICEBERG_ASSIGN_OR_RAISE(auto writer, + PositionDeleteWriter::Make(PositionDeleteWriterOptions{ + .path = output_path, + .schema = schema, + .spec = target.spec, + .partition = target.data_file->partition, + .format = FileFormatType::kParquet, + .io = table_->io(), + .properties = properties, + })); + for (int64_t pos : sorted_positions) { + ICEBERG_RETURN_UNEXPECTED(writer->WriteDelete(path, pos)); + } + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + ICEBERG_ASSIGN_OR_RAISE(auto metadata, writer->Metadata()); + result.data_files.insert(result.data_files.end(), + std::make_move_iterator(metadata.data_files.begin()), + std::make_move_iterator(metadata.data_files.end())); + result.referenced_data_files.push_back(path); + } + return result; + } + + std::shared_ptr
table_; + std::map, StringLess> deletes_; + std::set output_paths_; + bool terminal_ = false; +}; + +PositionDeleteUpdate::PositionDeleteUpdate(std::unique_ptr impl) + : impl_(std::move(impl)) {} + +PositionDeleteUpdate::~PositionDeleteUpdate() = default; + +Result> PositionDeleteUpdate::Make( + std::shared_ptr
table) { + ICEBERG_PRECHECK(table != nullptr, + "Cannot create position delete update without table"); + return std::unique_ptr( + new PositionDeleteUpdate(std::make_unique(std::move(table)))); +} + +PositionDeleteUpdate& PositionDeleteUpdate::Delete(std::string_view data_file_path, + int64_t pos) { + ICEBERG_BUILDER_RETURN_IF_ERROR(impl_->Delete(data_file_path, pos)); + return *this; +} + +Status PositionDeleteUpdate::Commit() { + ICEBERG_RETURN_UNEXPECTED(CheckErrors()); + return impl_->Commit(); +} + +} // namespace iceberg diff --git a/src/iceberg/data/position_delete_update.h b/src/iceberg/data/position_delete_update.h new file mode 100644 index 000000000..3b531d8fa --- /dev/null +++ b/src/iceberg/data/position_delete_update.h @@ -0,0 +1,60 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +/// \file iceberg/data/position_delete_update.h +/// Table-aware position delete writing and commit. + +#include +#include +#include + +#include "iceberg/iceberg_data_export.h" +#include "iceberg/result.h" +#include "iceberg/type_fwd.h" +#include "iceberg/util/error_collector.h" + +namespace iceberg { + +/// \brief Writes and commits position deletes using the table format. +/// +/// Format v2 tables produce Parquet position delete files. Format v3 tables +/// produce deletion vectors and merge previous file-scoped position deletes. +class ICEBERG_DATA_EXPORT PositionDeleteUpdate : public ErrorCollector { + public: + ~PositionDeleteUpdate() override; + + /// \brief Create a position delete update for a table. + static Result> Make(std::shared_ptr
table); + + /// \brief Add a deleted row position for a live data file. + PositionDeleteUpdate& Delete(std::string_view data_file_path, int64_t pos); + + /// \brief Write and commit all added position deletes. + Status Commit(); + + private: + class Impl; + std::unique_ptr impl_; + + explicit PositionDeleteUpdate(std::unique_ptr impl); +}; + +} // namespace iceberg diff --git a/src/iceberg/delete_file_index.cc b/src/iceberg/delete_file_index.cc index a8c4ef126..2bd6fc733 100644 --- a/src/iceberg/delete_file_index.cc +++ b/src/iceberg/delete_file_index.cc @@ -41,6 +41,7 @@ #include "iceberg/util/content_file_util.h" #include "iceberg/util/executor_util_internal.h" #include "iceberg/util/macros.h" +#include "iceberg/util/struct_like_set.h" namespace iceberg { @@ -458,6 +459,15 @@ Result> DeleteFileIndex::FindDV( "sequence number {}", it->second.sequence_number.value(), seq); + const auto& dv = *it->second.data_file; + ICEBERG_CHECK(dv.partition_spec_id == data_file.partition_spec_id, + "DV and data file have mismatched partition specs: {}", + data_file.file_path); + ICEBERG_ASSIGN_OR_RAISE(auto partitions_match, + StructLikeEqual(dv.partition, data_file.partition)); + ICEBERG_CHECK(partitions_match, "DV and data file have mismatched partitions: {}", + data_file.file_path); + return it->second.data_file; } diff --git a/src/iceberg/file_io.h b/src/iceberg/file_io.h index e22e5cf21..791133fcf 100644 --- a/src/iceberg/file_io.h +++ b/src/iceberg/file_io.h @@ -162,6 +162,9 @@ class ICEBERG_EXPORT FileIO { /// \brief Delete a file at the given location. /// + /// Deletion is idempotent: implementations must return success when the file does + /// not exist. + /// /// \param file_location The location of the file to delete. /// \return void if the delete succeeded, an error code if the delete failed. virtual Status DeleteFile(const std::string& file_location) { diff --git a/src/iceberg/test/CMakeLists.txt b/src/iceberg/test/CMakeLists.txt index 6f6ff7603..c4bdd0733 100644 --- a/src/iceberg/test/CMakeLists.txt +++ b/src/iceberg/test/CMakeLists.txt @@ -265,6 +265,7 @@ if(ICEBERG_BUILD_BUNDLE) delete_loader_test.cc dv_writer_test.cc file_scan_task_reader_test.cc + position_delete_update_test.cc literal_util_test.cc) endif() diff --git a/src/iceberg/test/arrow_io_test.cc b/src/iceberg/test/arrow_io_test.cc index 7bc9ebba5..6b3f6e4dd 100644 --- a/src/iceberg/test/arrow_io_test.cc +++ b/src/iceberg/test/arrow_io_test.cc @@ -364,8 +364,7 @@ TEST_F(LocalFileIOTest, DeleteFile) { EXPECT_THAT(del_res, IsOk()); del_res = file_io_->DeleteFile(temp_filepath_); - EXPECT_THAT(del_res, IsError(ErrorKind::kIOError)); - EXPECT_THAT(del_res, HasErrorMessage("Cannot delete file")); + EXPECT_THAT(del_res, IsOk()); } TEST_F(LocalFileIOTest, DeleteFiles) { @@ -414,6 +413,13 @@ TEST_F(LocalFileIOTest, StdReadFullyReadsFromAbsolutePosition) { VerifyReadFullyReadsFromAbsolutePosition(file_io, temp_filepath_)); } +TEST_F(LocalFileIOTest, StdDeleteFileIsIdempotent) { + auto file_io = std::make_shared(); + ASSERT_THAT(file_io->WriteFile(temp_filepath_, "abc"), IsOk()); + EXPECT_THAT(file_io->DeleteFile(temp_filepath_), IsOk()); + EXPECT_THAT(file_io->DeleteFile(temp_filepath_), IsOk()); +} + TEST_F(LocalFileIOTest, StdReadKeepsPositionAvailableAtEof) { auto file_io = std::make_shared(); ASSERT_THAT(file_io->WriteFile(temp_filepath_, "abc"), IsOk()); diff --git a/src/iceberg/test/delete_file_index_test.cc b/src/iceberg/test/delete_file_index_test.cc index 75c82bc39..48693a9b4 100644 --- a/src/iceberg/test/delete_file_index_test.cc +++ b/src/iceberg/test/delete_file_index_test.cc @@ -1113,6 +1113,81 @@ TEST_P(DeleteFileIndexTest, TestMultipleDVs) { EXPECT_THAT(index_result, HasErrorMessage(file_a_->file_path)); } +TEST_P(DeleteFileIndexTest, TestDVApplicability) { + auto version = GetParam(); + if (version < 3) { + GTEST_SKIP() << "DVs only supported in V3+"; + } + + const auto null_partition = PartitionValues({Literal::Null(int32())}); + auto null_partition_file = MakeDataFile("/path/to/data-null.parquet", null_partition, + partitioned_spec_->spec_id()); + + struct TestCase { + std::string name; + PartitionValues dv_partition; + std::shared_ptr dv_spec; + std::shared_ptr data_file; + bool applies; + }; + const std::vector cases = { + { + .name = "equal-partition", + .dv_partition = file_a_->partition, + .dv_spec = partitioned_spec_, + .data_file = file_a_, + .applies = true, + }, + { + .name = "different-spec", + .dv_partition = PartitionValues{}, + .dv_spec = unpartitioned_spec_, + .data_file = file_a_, + .applies = false, + }, + { + .name = "different-partition-value", + .dv_partition = file_b_->partition, + .dv_spec = partitioned_spec_, + .data_file = file_a_, + .applies = false, + }, + { + .name = "equal-null-partition", + .dv_partition = null_partition, + .dv_spec = partitioned_spec_, + .data_file = null_partition_file, + .applies = true, + }, + { + .name = "null-partition-mismatch", + .dv_partition = file_a_->partition, + .dv_spec = partitioned_spec_, + .data_file = null_partition_file, + .applies = false, + }, + }; + + for (const auto& test_case : cases) { + SCOPED_TRACE(test_case.name); + auto dv = MakeDV("/path/to/" + test_case.name + ".puffin", test_case.dv_partition, + test_case.dv_spec->spec_id(), test_case.data_file->file_path); + std::vector entries; + entries.push_back(MakeDeleteEntry(/*snapshot_id=*/1000L, /*sequence_number=*/2, dv)); + auto manifest = WriteDeleteManifest(version, /*snapshot_id=*/1000L, + std::move(entries), test_case.dv_spec); + ICEBERG_UNWRAP_OR_FAIL(auto index, BuildIndex({manifest})); + auto result = index->ForDataFile(1, *test_case.data_file); + if (test_case.applies) { + ICEBERG_UNWRAP_OR_FAIL(auto deletes, std::move(result)); + ASSERT_EQ(deletes.size(), 1); + EXPECT_EQ(deletes[0]->file_path, dv->file_path); + } else { + EXPECT_THAT(result, IsError(ErrorKind::kValidationFailed)); + } + } +} + TEST_P(DeleteFileIndexTest, TestInvalidDVSequenceNumber) { auto version = GetParam(); if (version < 3) { diff --git a/src/iceberg/test/position_delete_update_test.cc b/src/iceberg/test/position_delete_update_test.cc new file mode 100644 index 000000000..2be9c459e --- /dev/null +++ b/src/iceberg/test/position_delete_update_test.cc @@ -0,0 +1,636 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/data/position_delete_update.h" + +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#include "iceberg/arrow/arrow_io_internal.h" +#include "iceberg/avro/avro_register.h" +#include "iceberg/data/data_writer.h" +#include "iceberg/data/delete_loader.h" +#include "iceberg/data/file_scan_task_reader.h" +#include "iceberg/data/position_delete_writer.h" +#include "iceberg/deletes/position_delete_index.h" +#include "iceberg/file_reader.h" +#include "iceberg/manifest/manifest_reader.h" +#include "iceberg/metadata_columns.h" +#include "iceberg/parquet/parquet_register.h" +#include "iceberg/partition_spec.h" +#include "iceberg/schema.h" +#include "iceberg/schema_field.h" +#include "iceberg/schema_internal.h" +#include "iceberg/snapshot.h" +#include "iceberg/table.h" +#include "iceberg/table_metadata.h" +#include "iceberg/table_properties.h" +#include "iceberg/table_scan.h" +#include "iceberg/test/matchers.h" +#include "iceberg/test/mock_catalog.h" +#include "iceberg/test/update_test_base.h" +#include "iceberg/update/delete_files.h" +#include "iceberg/update/fast_append.h" +#include "iceberg/update/row_delta.h" +#include "iceberg/update/update_partition_spec.h" +#include "iceberg/update/update_properties.h" +#include "iceberg/util/uuid.h" + +namespace iceberg { + +namespace { + +struct RoutingCase { + int8_t format_version; + bool unpartitioned; + FileFormatType expected_format; +}; + +class FailOncePuffinDeleteFileIO : public arrow::ArrowFileSystemFileIO { + public: + explicit FailOncePuffinDeleteFileIO( + std::shared_ptr<::arrow::fs::FileSystem> file_system) + : ArrowFileSystemFileIO(std::move(file_system)) {} + + Status DeleteFile(const std::string& file_location) override { + if (!file_location.ends_with(".puffin")) { + return ArrowFileSystemFileIO::DeleteFile(file_location); + } + + puffin_delete_attempts.push_back(file_location); + if (fail_next_puffin_delete_) { + fail_next_puffin_delete_ = false; + return IOError("injected cleanup failure for {}", file_location); + } + return ArrowFileSystemFileIO::DeleteFile(file_location); + } + + std::vector puffin_delete_attempts; + + private: + bool fail_next_puffin_delete_ = true; +}; + +class PositionDeleteUpdateTest : public MinimalUpdateTestBase, + public ::testing::WithParamInterface { + protected: + static void SetUpTestSuite() { + avro::RegisterAll(); + parquet::RegisterAll(); + } + + int8_t format_version() const override { return GetParam().format_version; } + + void SetUp() override { + MinimalUpdateTestBase::SetUp(); + if (GetParam().unpartitioned) { + RegisterUnpartitionedTable(); + } + if (GetParam().format_version == 2) { + ICEBERG_UNWRAP_OR_FAIL(auto properties, table_->NewUpdateProperties()); + properties->Set(TableProperties::kDeleteParquetCompression.key(), "uncompressed"); + ASSERT_THAT(properties->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + } + ICEBERG_UNWRAP_OR_FAIL(spec_, table_->spec()); + ICEBERG_UNWRAP_OR_FAIL(data_file_, WriteDataFile()); + AppendDataFile(); + } + + void RegisterUnpartitionedTable() { + ICEBERG_UNWRAP_OR_FAIL( + auto metadata, ReadTableMetadataFromResource("TableMetadataV3ValidMinimal.json")); + metadata->location = table_location_; + metadata->partition_specs = {PartitionSpec::Unpartitioned()}; + metadata->default_spec_id = PartitionSpec::kInitialSpecId; + + const auto metadata_location = + std::format("{}/metadata/00001-{}.metadata.json", table_location_, + Uuid::GenerateV7().ToString()); + ASSERT_THAT(TableMetadataUtil::Write(*file_io_, metadata_location, *metadata), + IsOk()); + ASSERT_THAT(catalog_->DropTable(table_ident_, /*purge=*/false), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(table_, + catalog_->RegisterTable(table_ident_, metadata_location)); + } + + Result> WriteDataFile() { + ICEBERG_ASSIGN_OR_RAISE(auto schema, table_->schema()); + ICEBERG_ASSIGN_OR_RAISE( + auto writer, + DataWriter::Make({ + .path = table_location_ + "/data/file.parquet", + .schema = schema, + .spec = spec_, + .partition = GetParam().unpartitioned ? PartitionValues{} + : PartitionValues({Literal::Long(10)}), + .format = FileFormatType::kParquet, + .io = file_io_, + .properties = {{"write.parquet.compression-codec", "uncompressed"}}, + })); + ArrowSchema c_schema{}; + ICEBERG_RETURN_UNEXPECTED(ToArrowSchema(*schema, &c_schema)); + auto type = ::arrow::ImportType(&c_schema).ValueOrDie(); + std::string json = "["; + for (int64_t pos = 0; pos < 10; ++pos) { + if (pos != 0) json += ","; + json += std::format("[10,{},{}]", pos, pos * 10); + } + json += "]"; + auto array = ::arrow::json::ArrayFromJSONString(type, json).ValueOrDie(); + ArrowArray c_array{}; + auto status = ::arrow::ExportArray(*array, &c_array); + if (!status.ok()) return UnknownError(status.ToString()); + ICEBERG_RETURN_UNEXPECTED(writer->Write(&c_array)); + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + ICEBERG_ASSIGN_OR_RAISE(auto metadata, writer->Metadata()); + return metadata.data_files.front(); + } + + Result> ReadSurvivingRows() { + ICEBERG_ASSIGN_OR_RAISE(table_, catalog_->LoadTable(table_ident_)); + ICEBERG_ASSIGN_OR_RAISE(auto schema, table_->schema()); + ICEBERG_ASSIGN_OR_RAISE(auto task, CurrentTask()); + ICEBERG_ASSIGN_OR_RAISE(auto reader, FileScanTaskReader::Make({ + .io = file_io_, + .table_schema = schema, + .schemas = table_->metadata()->schemas, + .projected_schema = schema, + })); + ICEBERG_ASSIGN_OR_RAISE(auto stream, reader->Open(*task)); + auto batches = ::arrow::ImportRecordBatchReader(&stream).ValueOrDie(); + std::vector rows; + while (auto batch = batches->Next().ValueOrDie()) { + auto values = std::static_pointer_cast<::arrow::Int64Array>(batch->column(1)); + for (int64_t i = 0; i < values->length(); ++i) rows.push_back(values->Value(i)); + } + return rows; + } + + void AppendDataFile() { + ICEBERG_UNWRAP_OR_FAIL(auto append, table_->NewFastAppend()); + append->AppendFile(data_file_); + ASSERT_THAT(append->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + } + + Result> CurrentTask() { + ICEBERG_ASSIGN_OR_RAISE(auto builder, table_->NewScan()); + ICEBERG_ASSIGN_OR_RAISE(auto scan, builder->Build()); + ICEBERG_ASSIGN_OR_RAISE(auto tasks, scan->PlanFiles()); + ICEBERG_CHECK(tasks.size() == 1, "Expected one file scan task, found {}", + tasks.size()); + return tasks.front(); + } + + Result LoadPositions(const FileScanTask& task) { + std::vector> deletes; + std::ranges::copy_if(task.delete_files(), std::back_inserter(deletes), + [](const auto& file) { + return file->content == DataFile::Content::kPositionDeletes; + }); + DeleteLoader loader(file_io_); + return loader.LoadPositionDeletes(deletes, task.data_file()->file_path); + } + + std::shared_ptr spec_; + std::shared_ptr data_file_; +}; + +TEST_P(PositionDeleteUpdateTest, RoutesAndCommitsPositionDeletes) { + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 7).Delete(data_file_->file_path, 2); + ASSERT_THAT(update->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto task, CurrentTask()); + ASSERT_EQ(task->delete_files().size(), 1); + const auto& delete_file = task->delete_files().front(); + EXPECT_EQ(delete_file->file_format, GetParam().expected_format); + EXPECT_EQ(delete_file->partition_spec_id, data_file_->partition_spec_id); + EXPECT_EQ(delete_file->partition, data_file_->partition); + + ICEBERG_UNWRAP_OR_FAIL(auto positions, LoadPositions(*task)); + EXPECT_EQ(positions.Cardinality(), 2); + EXPECT_TRUE(positions.IsDeleted(2)); + EXPECT_TRUE(positions.IsDeleted(7)); + if (GetParam().format_version == 2) { + auto delete_schema = std::make_shared(std::vector{ + MetadataColumns::kDeleteFilePath, MetadataColumns::kDeleteFilePos}); + ICEBERG_UNWRAP_OR_FAIL(auto reader, + ReaderFactoryRegistry::Open(FileFormatType::kParquet, + {.path = delete_file->file_path, + .io = file_io_, + .projection = delete_schema})); + ICEBERG_UNWRAP_OR_FAIL(auto batch, reader->Next()); + ASSERT_TRUE(batch.has_value()); + + ArrowSchema arrow_schema; + ASSERT_THAT(ToArrowSchema(*delete_schema, &arrow_schema), IsOk()); + auto arrow_type = ::arrow::ImportType(&arrow_schema).ValueOrDie(); + auto rows = ::arrow::ImportArray(&batch.value(), arrow_type).ValueOrDie(); + auto struct_rows = std::static_pointer_cast<::arrow::StructArray>(rows); + auto positions = std::static_pointer_cast<::arrow::Int64Array>(struct_rows->field(1)); + ASSERT_EQ(positions->length(), 2); + EXPECT_EQ(positions->Value(0), 2); + EXPECT_EQ(positions->Value(1), 7); + } +} + +TEST_P(PositionDeleteUpdateTest, ReopenedTablePreservesBothDeletes) { + ICEBERG_UNWRAP_OR_FAIL(auto first, PositionDeleteUpdate::Make(table_)); + first->Delete(data_file_->file_path, 2).Delete(data_file_->file_path, 7); + ASSERT_THAT(first->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto rows, ReadSurvivingRows()); + EXPECT_THAT(rows, ::testing::ElementsAre(0, 1, 3, 4, 5, 6, 8, 9)); + + if (GetParam().format_version == 3 && !GetParam().unpartitioned) { + ICEBERG_UNWRAP_OR_FAIL(auto evolution, table_->NewUpdatePartitionSpec()); + evolution->RemoveField("x"); + ASSERT_THAT(evolution->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + } + ICEBERG_UNWRAP_OR_FAIL(auto second, PositionDeleteUpdate::Make(table_)); + second->Delete(data_file_->file_path, 4); + ASSERT_THAT(second->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(rows, ReadSurvivingRows()); + EXPECT_THAT(rows, ::testing::ElementsAre(0, 1, 3, 5, 6, 8, 9)); +} + +TEST_P(PositionDeleteUpdateTest, ConcurrentRemovalRejectsPreparedDeletes) { + ICEBERG_UNWRAP_OR_FAIL(auto properties, table_->NewUpdateProperties()); + properties->Set(TableProperties::kCommitMinRetryWaitMs.key(), "1") + .Set(TableProperties::kCommitMaxRetryWaitMs.key(), "1"); + ASSERT_THAT(properties->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + auto mock = std::make_shared<::testing::NiceMock>(); + ON_CALL(*mock, LoadTable(::testing::_)) + .WillByDefault( + [this](const TableIdentifier& name) { return catalog_->LoadTable(name); }); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce([this](const auto&, const auto&, + const auto&) -> Result> { + ICEBERG_ASSIGN_OR_RAISE(auto removal, table_->NewDeleteFiles()); + removal->DeleteFile(data_file_); + ICEBERG_RETURN_UNEXPECTED(removal->Commit()); + return CommitFailed("concurrent removal committed"); + }); + ICEBERG_UNWRAP_OR_FAIL( + auto competing_table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(competing_table)); + update->Delete(data_file_->file_path, 2); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + + ICEBERG_UNWRAP_OR_FAIL(auto reloaded, catalog_->LoadTable(table_ident_)); + ICEBERG_UNWRAP_OR_FAIL(auto builder, reloaded->NewScan()); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +INSTANTIATE_TEST_SUITE_P( + FormatAndPartitioning, PositionDeleteUpdateTest, + ::testing::Values(RoutingCase{.format_version = 2, + .unpartitioned = false, + .expected_format = FileFormatType::kParquet}, + RoutingCase{.format_version = 3, + .unpartitioned = false, + .expected_format = FileFormatType::kPuffin}, + RoutingCase{.format_version = 3, + .unpartitioned = true, + .expected_format = FileFormatType::kPuffin})); + +class PositionDeleteV3Test : public MinimalUpdateTestBase { + protected: + static void SetUpTestSuite() { + avro::RegisterAll(); + parquet::RegisterAll(); + } + + int8_t format_version() const override { return 3; } + + void SetUp() override { + MinimalUpdateTestBase::SetUp(); + ICEBERG_UNWRAP_OR_FAIL(spec_, table_->spec()); + ICEBERG_UNWRAP_OR_FAIL(schema_, table_->schema()); + data_file_ = MakeDataFile(); + AppendDataFile(); + } + + std::shared_ptr MakeDataFile() const { + auto file = std::make_shared(); + file->content = DataFile::Content::kData; + file->file_path = table_location_ + "/data/file.parquet"; + file->file_format = FileFormatType::kParquet; + file->partition = PartitionValues({Literal::Long(10)}); + file->file_size_in_bytes = 1024; + file->record_count = 10; + file->partition_spec_id = spec_->spec_id(); + return file; + } + + void AppendDataFile() { + ICEBERG_UNWRAP_OR_FAIL(auto append, table_->NewFastAppend()); + append->AppendFile(data_file_); + ASSERT_THAT(append->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + } + + Result> CurrentTask(const std::shared_ptr
& table) { + ICEBERG_ASSIGN_OR_RAISE(auto builder, table->NewScan()); + ICEBERG_ASSIGN_OR_RAISE(auto scan, builder->Build()); + ICEBERG_ASSIGN_OR_RAISE(auto tasks, scan->PlanFiles()); + ICEBERG_CHECK(tasks.size() == 1, "Expected one file scan task, found {}", + tasks.size()); + return tasks.front(); + } + + Result LoadPositions(const FileScanTask& task) { + DeleteLoader loader(file_io_); + return loader.LoadPositionDeletes(task.delete_files(), task.data_file()->file_path); + } + + Result> WritePositionDeletes( + std::span positions) { + const auto path = table_location_ + "/data/existing-position-deletes.parquet"; + ICEBERG_ASSIGN_OR_RAISE( + auto writer, + PositionDeleteWriter::Make(PositionDeleteWriterOptions{ + .path = path, + .schema = schema_, + .spec = spec_, + .partition = data_file_->partition, + .format = FileFormatType::kParquet, + .io = file_io_, + .properties = {{"write.parquet.compression-codec", "uncompressed"}}, + })); + for (int64_t pos : positions) { + ICEBERG_RETURN_UNEXPECTED(writer->WriteDelete(data_file_->file_path, pos)); + } + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + ICEBERG_ASSIGN_OR_RAISE(auto result, writer->Metadata()); + ICEBERG_CHECK(result.data_files.size() == 1, + "Expected one position delete file, found {}", + result.data_files.size()); + return result.data_files.front(); + } + + Result> CurrentDeleteEntries() { + ICEBERG_ASSIGN_OR_RAISE(auto snapshot, table_->current_snapshot()); + SnapshotReader cache(snapshot.get()); + ICEBERG_ASSIGN_OR_RAISE(auto manifests, cache.DeleteManifests(file_io_)); + std::vector entries; + for (const auto& manifest : manifests) { + ICEBERG_ASSIGN_OR_RAISE( + auto spec, table_->metadata()->PartitionSpecById(manifest.partition_spec_id)); + ICEBERG_ASSIGN_OR_RAISE( + auto reader, + ManifestReader::Make(manifest, file_io_, schema_, std::move(spec))); + ICEBERG_ASSIGN_OR_RAISE(auto manifest_entries, reader->Entries()); + entries.insert(entries.end(), std::make_move_iterator(manifest_entries.begin()), + std::make_move_iterator(manifest_entries.end())); + } + return entries; + } + + void ConfigureRetries(int32_t retries) { + ICEBERG_UNWRAP_OR_FAIL(auto properties, table_->NewUpdateProperties()); + properties->Set(TableProperties::kCommitNumRetries.key(), std::to_string(retries)) + .Set(TableProperties::kCommitMinRetryWaitMs.key(), "1") + .Set(TableProperties::kCommitMaxRetryWaitMs.key(), "1") + .Set(TableProperties::kCommitTotalRetryTimeMs.key(), "1000"); + ASSERT_THAT(properties->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + } + + void FailCommit(std::unexpected failure) { + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce(::testing::Return(std::move(failure))); + ICEBERG_UNWRAP_OR_FAIL( + table_, Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, + std::move(mock))); + } + + std::vector PuffinFiles() { + auto arrow_io = std::dynamic_pointer_cast(file_io_); + EXPECT_NE(arrow_io, nullptr); + ::arrow::fs::FileSelector selector; + selector.base_dir = table_location_ + "/data"; + selector.recursive = true; + auto infos = arrow_io->fs()->GetFileInfo(selector); + EXPECT_TRUE(infos.ok()) << infos.status().ToString(); + std::vector paths; + if (!infos.ok()) { + return paths; + } + for (const auto& info : *infos) { + if (info.path().ends_with(".puffin")) { + paths.push_back(info.path()); + } + } + return paths; + } + + std::shared_ptr spec_; + std::shared_ptr schema_; + std::shared_ptr data_file_; +}; + +TEST_F(PositionDeleteV3Test, SecondDeleteMergesAndSupersedesPreviousDV) { + ICEBERG_UNWRAP_OR_FAIL(auto first, PositionDeleteUpdate::Make(table_)); + first->Delete(data_file_->file_path, 1); + ASSERT_THAT(first->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto old_task, CurrentTask(table_)); + const auto old_dv = old_task->delete_files().front(); + + ICEBERG_UNWRAP_OR_FAIL(auto second, PositionDeleteUpdate::Make(table_)); + second->Delete(data_file_->file_path, 3); + ASSERT_THAT(second->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto task, CurrentTask(table_)); + ASSERT_EQ(task->delete_files().size(), 1); + EXPECT_NE(task->delete_files().front()->file_path, old_dv->file_path); + ICEBERG_UNWRAP_OR_FAIL(auto positions, LoadPositions(*task)); + EXPECT_EQ(positions.Cardinality(), 2); + EXPECT_TRUE(positions.IsDeleted(1)); + EXPECT_TRUE(positions.IsDeleted(3)); + + ICEBERG_UNWRAP_OR_FAIL(auto entries, CurrentDeleteEntries()); + EXPECT_TRUE(std::ranges::any_of(entries, [&old_dv](const ManifestEntry& entry) { + return entry.status == ManifestStatus::kDeleted && entry.data_file != nullptr && + entry.data_file->file_path == old_dv->file_path && + entry.data_file->content_offset == old_dv->content_offset; + })); +} + +TEST_F(PositionDeleteV3Test, MergesAndSupersedesFileScopedParquetDelete) { + RegisterTableFromResource("TableMetadataV2ValidMinimal.json"); + ICEBERG_UNWRAP_OR_FAIL(spec_, table_->spec()); + ICEBERG_UNWRAP_OR_FAIL(schema_, table_->schema()); + data_file_ = MakeDataFile(); + AppendDataFile(); + + const std::vector existing_positions{1, 3}; + ICEBERG_UNWRAP_OR_FAIL(auto old_delete, WritePositionDeletes(existing_positions)); + ASSERT_EQ(old_delete->referenced_data_file, data_file_->file_path); + ICEBERG_UNWRAP_OR_FAIL(auto row_delta, table_->NewRowDelta()); + row_delta->AddDeletes(old_delete); + ASSERT_THAT(row_delta->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto properties, table_->NewUpdateProperties()); + properties->Set(TableProperties::kFormatVersion.key(), "3"); + ASSERT_THAT(properties->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 5); + ASSERT_THAT(update->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto task, CurrentTask(table_)); + ASSERT_EQ(task->delete_files().size(), 1); + EXPECT_EQ(task->delete_files().front()->file_format, FileFormatType::kPuffin); + ICEBERG_UNWRAP_OR_FAIL(auto positions, LoadPositions(*task)); + EXPECT_EQ(positions.Cardinality(), 3); + EXPECT_TRUE(positions.IsDeleted(1)); + EXPECT_TRUE(positions.IsDeleted(3)); + EXPECT_TRUE(positions.IsDeleted(5)); + + ICEBERG_UNWRAP_OR_FAIL(auto entries, CurrentDeleteEntries()); + EXPECT_TRUE(std::ranges::any_of(entries, [&old_delete](const ManifestEntry& entry) { + return entry.status == ManifestStatus::kDeleted && entry.data_file != nullptr && + entry.data_file->file_path == old_delete->file_path; + })); +} + +TEST_F(PositionDeleteV3Test, UsesTargetDataFileSpecAfterPartitionEvolution) { + ICEBERG_UNWRAP_OR_FAIL(auto spec_update, table_->NewUpdatePartitionSpec()); + spec_update->RemoveField("x"); + ASSERT_THAT(spec_update->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + ASSERT_NE(table_->metadata()->default_spec_id, data_file_->partition_spec_id); + + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 4); + ASSERT_THAT(update->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto task, CurrentTask(table_)); + const auto& dv = task->delete_files().front(); + EXPECT_EQ(dv->partition_spec_id, data_file_->partition_spec_id); + EXPECT_EQ(dv->partition, data_file_->partition); +} + +TEST_F(PositionDeleteV3Test, CommitRetryReusesWrittenDV) { + ConfigureRetries(1); + FailCommits(1); + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 5); + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_EQ(PuffinFiles().size(), 1); +} + +TEST_F(PositionDeleteV3Test, FailedCommitCleansWrittenDV) { + ConfigureRetries(0); + FailCommit(CommitFailed("injected failure")); + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 6); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kCommitFailed)); + EXPECT_TRUE(PuffinFiles().empty()); +} + +TEST_F(PositionDeleteV3Test, CleanupFailureIsReportedAndRetried) { + ConfigureRetries(0); + auto mock_catalog = std::make_shared<::testing::NiceMock>(); + int update_calls = 0; + ON_CALL(*mock_catalog, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault( + [this, &update_calls]( + const TableIdentifier& identifier, + const std::vector>& requirements, + const std::vector>& updates) + -> Result> { + if (++update_calls == 1) { + return CommitFailed("injected commit failure"); + } + return catalog_->UpdateTable(identifier, requirements, updates); + }); + auto arrow_io = std::dynamic_pointer_cast(file_io_); + ASSERT_NE(arrow_io, nullptr); + auto failing_io = std::make_shared(arrow_io->fs()); + ICEBERG_UNWRAP_OR_FAIL(auto mock_table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), + failing_io, mock_catalog)); + + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(mock_table)); + update->Delete(data_file_->file_path, 6); + auto first_status = update->Commit(); + EXPECT_THAT(first_status, IsError(ErrorKind::kCommitFailed)); + EXPECT_THAT(first_status, HasErrorMessage("injected commit failure")); + EXPECT_THAT(first_status, HasErrorMessage("injected cleanup failure")); + ASSERT_EQ(PuffinFiles().size(), 1); + ASSERT_EQ(failing_io->puffin_delete_attempts.size(), 1); + const auto retained_path = failing_io->puffin_delete_attempts.front(); + + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_EQ(update_calls, 2); + EXPECT_EQ(PuffinFiles().size(), 1); + ASSERT_EQ(failing_io->puffin_delete_attempts.size(), 2); + EXPECT_EQ(std::ranges::count(failing_io->puffin_delete_attempts, retained_path), 2); +} + +TEST_F(PositionDeleteV3Test, CommitStateUnknownRelinquishesOutputOwnership) { + ConfigureRetries(1); + FailCommit(CommitStateUnknown("injected unknown state")); + ICEBERG_UNWRAP_OR_FAIL(auto update, PositionDeleteUpdate::Make(table_)); + update->Delete(data_file_->file_path, 6); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kCommitStateUnknown)); + ASSERT_EQ(PuffinFiles().size(), 1); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kInvalidArgument)); + EXPECT_EQ(PuffinFiles().size(), 1); +} + +} // namespace + +} // namespace iceberg diff --git a/src/iceberg/test/std_io.h b/src/iceberg/test/std_io.h index 725fc7ba5..c7d5b312b 100644 --- a/src/iceberg/test/std_io.h +++ b/src/iceberg/test/std_io.h @@ -319,11 +319,9 @@ class StdFileIO : public FileIO { Status DeleteFile(const std::string& file_location) override { std::error_code ec; - if (!std::filesystem::remove(file_location, ec)) { - if (ec) { - return IOError("Failed to delete file {}: {}", file_location, ec.message()); - } - return IOError("File does not exist: {}", file_location); + std::filesystem::remove(file_location, ec); + if (ec) { + return IOError("Failed to delete file {}: {}", file_location, ec.message()); } return {}; }