diff --git a/src/iceberg/test/expire_snapshots_test.cc b/src/iceberg/test/expire_snapshots_test.cc index b96754ec0..9695b4a85 100644 --- a/src/iceberg/test/expire_snapshots_test.cc +++ b/src/iceberg/test/expire_snapshots_test.cc @@ -37,9 +37,14 @@ #include "iceberg/snapshot.h" #include "iceberg/statistics_file.h" #include "iceberg/table_metadata.h" +#include "iceberg/table_properties.h" #include "iceberg/test/executor.h" #include "iceberg/test/matchers.h" +#include "iceberg/test/mock_catalog.h" #include "iceberg/test/update_test_base.h" +#include "iceberg/transaction.h" +#include "iceberg/update/fast_append.h" +#include "iceberg/update/set_snapshot.h" namespace iceberg { @@ -291,21 +296,15 @@ TEST_F(ExpireSnapshotsCleanupTest, RetainsUnreferencedSnapshotAtExpireThreshold) testing::Not(testing::Contains(unreferenced_snapshot_id))); } -TEST_F(ExpireSnapshotsTest, FinalizeRequiresCommittedMetadata) { +TEST_F(ExpireSnapshotsTest, ApplyDoesNotDeleteBeforeCommit) { std::vector deleted_files; - ICEBERG_UNWRAP_OR_FAIL(auto update, table_->NewExpireSnapshots()); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewExpireSnapshots()); update->DeleteWith( [&deleted_files](const std::string& path) { deleted_files.push_back(path); }); - - // Apply first so apply_result_ is cached ICEBERG_UNWRAP_OR_FAIL(auto result, update->Apply()); EXPECT_EQ(result.snapshot_ids_to_remove.size(), 1); - - // A successful finalize now requires the committed metadata from the catalog. - auto finalize_status = update->Finalize(static_cast(nullptr)); - EXPECT_THAT(finalize_status, IsError(ErrorKind::kInvalidArgument)); - EXPECT_THAT(finalize_status, - HasErrorMessage("Missing committed table metadata for cleanup")); + EXPECT_THAT(txn->Abort(), IsOk()); EXPECT_TRUE(deleted_files.empty()); } @@ -319,29 +318,103 @@ TEST_F(ExpireSnapshotsTest, CleanupNoneSkipsDeletion) { ICEBERG_UNWRAP_OR_FAIL(auto result, update->Apply()); EXPECT_EQ(result.snapshot_ids_to_remove.size(), 1); - // With kNone cleanup level, Finalize should skip all file deletion - auto finalize_status = update->Finalize(static_cast(nullptr)); - EXPECT_THAT(finalize_status, IsOk()); + // With kNone cleanup level, Commit should skip all file deletion + auto commit_status = update->Commit(); + EXPECT_THAT(commit_status, IsOk()); EXPECT_TRUE(deleted_files.empty()); } -TEST_F(ExpireSnapshotsTest, FinalizeSkippedOnCommitError) { +TEST_F(ExpireSnapshotsTest, AbortSkipsExpirationDeletion) { std::vector deleted_files; - ICEBERG_UNWRAP_OR_FAIL(auto update, table_->NewExpireSnapshots()); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewExpireSnapshots()); update->DeleteWith( [&deleted_files](const std::string& path) { deleted_files.push_back(path); }); + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_TRUE(deleted_files.empty()); +} - ICEBERG_UNWRAP_OR_FAIL(auto result, update->Apply()); - EXPECT_EQ(result.snapshot_ids_to_remove.size(), 1); +TEST_F(ExpireSnapshotsCleanupTest, CommitFailureSkipsExpirationDeletion) { + const auto expired_list = table_location_ + "/metadata/expired-list.avro"; + const auto current_list = table_location_ + "/metadata/current-list.avro"; + WriteManifestList(expired_list, kExpiredSnapshotId, 0, kExpiredSequenceNumber, {}); + WriteManifestList(current_list, kCurrentSnapshotId, kExpiredSnapshotId, + kCurrentSequenceNumber, {}); + RewriteTableWithManifestLists(expired_list, current_list); + + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce(::testing::Return(CommitFailed("injected commit failure"))); + auto metadata = std::make_shared(*table_->metadata()); + metadata->properties.Set(TableProperties::kCommitNumRetries, 0); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), metadata, std::string(table_->metadata_file_location()), + file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewExpireSnapshots()); + std::vector deleted_files; + update->ExpireSnapshotId(kExpiredSnapshotId).DeleteWith([&](const std::string& path) { + deleted_files.push_back(path); + }); + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_FALSE(txn->current().SnapshotById(kExpiredSnapshotId).has_value()); + EXPECT_TRUE(deleted_files.empty()); - // Simulate a commit failure - Finalize should not delete any files - auto finalize_status = update->Finalize(Result(std::unexpected( - Error{.kind = ErrorKind::kCommitFailed, .message = "simulated failure"}))); - EXPECT_THAT(finalize_status, IsOk()); + EXPECT_THAT(txn->Commit(), + ::testing::AllOf(IsError(ErrorKind::kCommitFailed), + HasErrorMessage("injected commit failure"))); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + EXPECT_THAT(txn->Abort(), IsOk()); EXPECT_TRUE(deleted_files.empty()); + EXPECT_TRUE(ReloadMetadata()->SnapshotById(kExpiredSnapshotId).has_value()); + EXPECT_THAT(file_io_->ReadFile(expired_list, std::nullopt), IsOk()); + EXPECT_THAT(file_io_->ReadFile(current_list, std::nullopt), IsOk()); } -TEST_F(ExpireSnapshotsTest, FinalizeSkipsWhenNothingExpired) { +TEST_F(ExpireSnapshotsCleanupTest, CommitStateUnknownPreservesExpiredFiles) { + const auto data_path = table_location_ + "/data/unknown-expired.parquet"; + const auto manifest_path = table_location_ + "/metadata/unknown-expired.avro"; + const auto expired_list = table_location_ + "/metadata/unknown-expired-list.avro"; + const auto current_list = table_location_ + "/metadata/unknown-current-list.avro"; + ASSERT_THAT(file_io_->WriteFile(data_path, "data"), IsOk()); + auto data_file = MakeDataFile(data_path); + data_file->partition_spec_id = DefaultSpec()->spec_id(); + auto manifest = WriteDataManifest(manifest_path, kExpiredSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kExpiredSnapshotId, + kExpiredSequenceNumber, data_file)}); + WriteManifestList(expired_list, kExpiredSnapshotId, 0, kExpiredSequenceNumber, + {manifest}); + WriteManifestList(current_list, kCurrentSnapshotId, kExpiredSnapshotId, + kCurrentSequenceNumber, {}); + RewriteTableWithManifestLists(expired_list, current_list); + + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillOnce(::testing::Return(CommitStateUnknown("unknown commit state"))); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto update, table->NewExpireSnapshots()); + std::vector deleted_files; + update->ExpireSnapshotId(kExpiredSnapshotId).DeleteWith([&](const std::string& path) { + deleted_files.push_back(path); + return file_io_->DeleteFile(path); + }); + + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kCommitStateUnknown)); + EXPECT_TRUE(deleted_files.empty()); + EXPECT_TRUE(ReloadMetadata()->SnapshotById(kExpiredSnapshotId).has_value()); + for (const auto& path : {data_path, manifest_path, expired_list}) { + EXPECT_THAT(file_io_->ReadFile(path, std::nullopt), IsOk()); + } +} + +TEST_F(ExpireSnapshotsTest, CommitSkipsWhenNothingExpired) { std::vector deleted_files; ICEBERG_UNWRAP_OR_FAIL(auto update, table_->NewExpireSnapshots()); update->RetainLast(2); @@ -351,9 +424,9 @@ TEST_F(ExpireSnapshotsTest, FinalizeSkipsWhenNothingExpired) { ICEBERG_UNWRAP_OR_FAIL(auto result, update->Apply()); EXPECT_TRUE(result.snapshot_ids_to_remove.empty()); - // No snapshots expired, so Finalize should not delete any files - auto finalize_status = update->Finalize(static_cast(nullptr)); - EXPECT_THAT(finalize_status, IsOk()); + // No snapshots expired, so Commit should not delete any files + auto commit_status = update->Commit(); + EXPECT_THAT(commit_status, IsOk()); EXPECT_TRUE(deleted_files.empty()); } @@ -361,7 +434,7 @@ TEST_F(ExpireSnapshotsTest, CommitWithCleanupNone) { ICEBERG_UNWRAP_OR_FAIL(auto update, table_->NewExpireSnapshots()); update->CleanupLevel(CleanupLevel::kNone); - // Commit should succeed - Finalize is called internally but skips cleanup + // Commit should succeed and skip cleanup EXPECT_THAT(update->Commit(), IsOk()); // Verify snapshot was removed from metadata @@ -1016,4 +1089,102 @@ TEST_F(ExpireSnapshotsCleanupTest, CommitIgnoresMalformedSourceSnapshotIdCleanup EXPECT_EQ(committed_metadata->snapshots.at(0)->snapshot_id, kCurrentSnapshotId); } +class ExpirationLifecycleTest + : public ExpireSnapshotsCleanupTest, + public ::testing::WithParamInterface> {}; + +TEST_P(ExpirationLifecycleTest, LaterReferencesSuppressPhysicalDeletion) { + auto [cleanup_level, custom_delete, reattach, retry] = GetParam(); + const auto data_path = table_location_ + "/data/reattached.parquet"; + const auto manifest_path = table_location_ + "/metadata/expired.avro"; + const auto expired_list = table_location_ + "/metadata/expired-list.avro"; + const auto current_list = table_location_ + "/metadata/current-list.avro"; + auto data_file = MakeDataFile(data_path); + data_file->partition_spec_id = DefaultSpec()->spec_id(); + ASSERT_THAT(file_io_->WriteFile(data_path, "data"), IsOk()); + auto manifest = WriteDataManifest(manifest_path, kExpiredSnapshotId, + {MakeEntry(ManifestStatus::kAdded, kExpiredSnapshotId, + kExpiredSequenceNumber, data_file)}); + WriteManifestList(expired_list, kExpiredSnapshotId, 0, kExpiredSequenceNumber, + {manifest}); + WriteManifestList(current_list, kCurrentSnapshotId, kExpiredSnapshotId, + kCurrentSequenceNumber, {}); + RewriteTableWithManifestLists(expired_list, current_list); + if (retry) { + FailCommits(1); + } + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto expire, txn->NewExpireSnapshots()); + expire->ExpireSnapshotId(kExpiredSnapshotId).CleanupLevel(cleanup_level); + int deletes = 0; + if (custom_delete) { + expire->DeleteWith([&](const std::string& path) { + ++deletes; + std::ignore = file_io_->DeleteFile(path); + }); + } + ASSERT_THAT(expire->Commit(), IsOk()); + if (reattach) { + // A later update reuses the expired data file, so final cleanup must preserve it. + ICEBERG_UNWRAP_OR_FAIL(auto append, txn->NewFastAppend()); + append->AppendFile(data_file); + ASSERT_THAT(append->Commit(), IsOk()); + } else { + ICEBERG_UNWRAP_OR_FAIL(auto noop, txn->NewSetSnapshot()); + noop->SetCurrentSnapshot(kCurrentSnapshotId); + ASSERT_THAT(noop->Commit(), IsOk()); + } + ASSERT_THAT(txn->Commit(), IsOk()); + EXPECT_FALSE(ReloadMetadata()->SnapshotById(kExpiredSnapshotId).has_value()); + EXPECT_EQ(file_io_->ReadFile(data_path, std::nullopt).has_value(), + reattach || cleanup_level == CleanupLevel::kMetadataOnly); + EXPECT_FALSE(file_io_->ReadFile(manifest_path, std::nullopt).has_value()); + EXPECT_FALSE(file_io_->ReadFile(expired_list, std::nullopt).has_value()); + if (custom_delete) { + EXPECT_GT(deletes, 0); + } + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsError(ErrorKind::kValidationFailed)); +} + +INSTANTIATE_TEST_SUITE_P( + CleanupPolicies, ExpirationLifecycleTest, + ::testing::Combine(::testing::Values(CleanupLevel::kAll, CleanupLevel::kMetadataOnly), + ::testing::Bool(), ::testing::Bool(), ::testing::Bool())); + +TEST_F(ExpireSnapshotsCleanupTest, RetryRecomputesExpirationAgainstRefreshedMetadata) { + const auto expired_list = table_location_ + "/metadata/expire-before-retry.avro"; + const auto current_list = table_location_ + "/metadata/current-before-retry.avro"; + WriteManifestList(expired_list, kExpiredSnapshotId, 0, kExpiredSequenceNumber, {}); + WriteManifestList(current_list, kCurrentSnapshotId, kExpiredSnapshotId, + kCurrentSequenceNumber, {}); + RewriteTableWithManifestLists(expired_list, current_list); + auto data_file = MakeDataFile(table_location_ + "/data/concurrent.parquet"); + data_file->partition_spec_id = DefaultSpec()->spec_id(); + FailCommits(1, [&](int attempt) { + if (attempt == 0) { + ICEBERG_UNWRAP_OR_FAIL(auto latest, catalog_->LoadTable(table_ident_)); + ICEBERG_UNWRAP_OR_FAIL(auto append, latest->NewFastAppend()); + append->AppendFile(data_file); + ASSERT_THAT(append->Commit(), IsOk()); + } + }); + ICEBERG_UNWRAP_OR_FAIL(auto expire, table_->NewExpireSnapshots()); + expire + ->ExpireOlderThan( + (CurrentTimePointMs() + std::chrono::hours(1)).time_since_epoch().count()) + .RetainLast(1); + std::vector deleted; + expire->DeleteWith([&](const std::string& path) { deleted.push_back(path); }); + ICEBERG_UNWRAP_OR_FAIL(auto preview, expire->Apply()); + EXPECT_THAT(preview.snapshot_ids_to_remove, ::testing::ElementsAre(kExpiredSnapshotId)); + ASSERT_THAT(expire->Commit(), IsOk()); + auto metadata = ReloadMetadata(); + EXPECT_FALSE(metadata->SnapshotById(kExpiredSnapshotId).has_value()); + EXPECT_FALSE(metadata->SnapshotById(kCurrentSnapshotId).has_value()); + EXPECT_EQ(metadata->snapshots.size(), 1U); + EXPECT_THAT(deleted, ::testing::Contains(expired_list)); + EXPECT_THAT(deleted, ::testing::Contains(current_list)); +} + } // namespace iceberg diff --git a/src/iceberg/test/fast_append_test.cc b/src/iceberg/test/fast_append_test.cc index f88d2e011..6ae3eef90 100644 --- a/src/iceberg/test/fast_append_test.cc +++ b/src/iceberg/test/fast_append_test.cc @@ -20,9 +20,11 @@ #include "iceberg/update/fast_append.h" #include +#include #include #include #include +#include #include #include #include @@ -34,6 +36,8 @@ #include "iceberg/avro/avro_register.h" #include "iceberg/constants.h" +#include "iceberg/exception.h" +#include "iceberg/expression/expressions.h" #include "iceberg/manifest/manifest_entry.h" #include "iceberg/manifest/manifest_reader.h" #include "iceberg/manifest/manifest_writer.h" @@ -45,12 +49,16 @@ #include "iceberg/snapshot.h" #include "iceberg/table_metadata.h" #include "iceberg/table_properties.h" +#include "iceberg/table_update.h" #include "iceberg/test/executor.h" #include "iceberg/test/matchers.h" #include "iceberg/test/mock_catalog.h" #include "iceberg/test/update_test_base.h" #include "iceberg/transaction.h" +#include "iceberg/update/delete_files.h" #include "iceberg/update/merge_append.h" +#include "iceberg/update/overwrite_files.h" +#include "iceberg/update/rewrite_files.h" #include "iceberg/update/update_properties.h" #include "iceberg/util/uuid.h" @@ -64,11 +72,32 @@ class TestSnapshotUpdate : public SnapshotUpdate { : SnapshotUpdate(std::move(ctx)) {} using SnapshotUpdate::ManifestPath; + using SnapshotUpdate::SnapshotId; - Status CleanUncommitted(const std::unordered_set&) override { return {}; } - std::string operation() override { return "test"; } + std::function apply_callback; + std::function cleanup_callback; + int applies = 0; + int cleanups = 0; + bool write_partial = false; + std::vector partial_paths; + + protected: + Status CleanUncommitted(const std::unordered_set&) override { + ++cleanups; + return cleanup_callback ? cleanup_callback() : Status{}; + } + std::string operation() override { return "append"; } Result> Apply(const TableMetadata&, const std::shared_ptr&) override { + ++applies; + if (write_partial) { + auto path = ManifestPath(); + partial_paths.push_back(path); + ICEBERG_RETURN_UNEXPECTED(ctx_->table->io()->WriteFile(path, "partial manifest")); + } + if (apply_callback) { + ICEBERG_RETURN_UNEXPECTED(apply_callback(applies)); + } return std::vector{}; } std::unordered_map Summary() override { return {}; } @@ -161,6 +190,18 @@ class FastAppendTest : public UpdateTestBase { class SnapshotUpdateTest : public UpdateTestBase {}; +TEST_F(SnapshotUpdateTest, StandaloneCommitRequiresSharedOwnership) { + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table_, TransactionKind::kUpdate)); + auto update = std::make_unique(ctx); + + EXPECT_THAT(update->Commit(), + ::testing::AllOf( + IsError(ErrorKind::kInvalidArgument), + HasErrorMessage("PendingUpdate must be owned by std::shared_ptr"))); + EXPECT_FALSE(ctx->transaction.has_value()); +} + TEST_F(FastAppendTest, AppendDataFile) { std::shared_ptr fast_append; ICEBERG_UNWRAP_OR_FAIL(fast_append, table_->NewFastAppend()); @@ -303,50 +344,155 @@ TEST_F(FastAppendTest, AppendNullFile) { EXPECT_THAT(table_->current_snapshot(), HasErrorMessage("No current snapshot")); } -TEST_F(FastAppendTest, FinalizeIgnoresCleanupDeleteFailure) { - std::shared_ptr fast_append; - ICEBERG_UNWRAP_OR_FAIL(fast_append, table_->NewFastAppend()); +TEST_F(FastAppendTest, AbortIgnoresCleanupDeleteFailure) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto fast_append, txn->NewFastAppend()); + fast_append->AppendFile(file_a_); + std::vector deleted_paths; + fast_append->DeleteWith([&](const std::string& path) { + deleted_paths.push_back(path); + return IOError("delete failed"); + }); + EXPECT_THAT(fast_append->Commit(), IsOk()); + EXPECT_THAT(txn->Abort(), IsOk()); + std::unordered_set unique_deleted_paths(deleted_paths.begin(), + deleted_paths.end()); + EXPECT_THAT(unique_deleted_paths, ::testing::SizeIs(2)); + const auto delete_attempts = deleted_paths.size(); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_EQ(deleted_paths.size(), delete_attempts); +} + +TEST_F(FastAppendTest, CommitFailureIgnoresCleanupDeleteFailure) { + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce(::testing::Return(CommitFailed("injected commit failure"))); + auto metadata = std::make_shared(*table_->metadata()); + metadata->properties.Set(TableProperties::kCommitNumRetries, 0); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), metadata, std::string(table_->metadata_file_location()), + file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto append, txn->NewFastAppend()); + std::vector deleted_paths; + append->AppendFile(file_a_).DeleteWith([&](const std::string& path) { + deleted_paths.push_back(path); + return IOError("delete failed"); + }); + ASSERT_THAT(append->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, txn->current().Snapshot()); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.Manifests(file_io_)); + ASSERT_THAT(manifests, ::testing::SizeIs(1)); + EXPECT_TRUE(deleted_paths.empty()); + + EXPECT_THAT(txn->Commit(), + ::testing::AllOf(IsError(ErrorKind::kCommitFailed), + HasErrorMessage("injected commit failure"))); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + std::unordered_set unique_deleted_paths(deleted_paths.begin(), + deleted_paths.end()); + EXPECT_THAT(unique_deleted_paths, + ::testing::UnorderedElementsAre(manifests[0].manifest_path, + snapshot->manifest_list)); + const auto cleanup_attempts = deleted_paths.size(); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_GT(deleted_paths.size(), cleanup_attempts); + unique_deleted_paths = {deleted_paths.begin(), deleted_paths.end()}; + EXPECT_THAT(unique_deleted_paths, + ::testing::UnorderedElementsAre(manifests[0].manifest_path, + snapshot->manifest_list)); + EXPECT_FALSE(ReloadMetadata()->Snapshot().has_value()); +} + +TEST_F(FastAppendTest, TransactionApplyFailureCleansUpStagedFiles) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto fast_append, txn->NewFastAppend()); + std::vector deleted_paths; + fast_append->DeleteWith([&](const std::string& path) { + deleted_paths.push_back(path); + return file_io_->DeleteFile(path); + }); fast_append->AppendFile(file_a_); - fast_append->DeleteWith([](const std::string&) { return IOError("delete failed"); }); + EXPECT_THAT(fast_append->Commit(), IsOk()); - EXPECT_THAT(static_cast(*fast_append).Apply(), IsOk()); - EXPECT_THAT(fast_append->Finalize(Result( - std::unexpected(CommitFailed("commit failed").error()))), - IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto failed_append, txn->NewFastAppend()); + failed_append->AppendFile(nullptr); + EXPECT_THAT(failed_append->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(deleted_paths, ::testing::SizeIs(2U)); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_THAT(deleted_paths, ::testing::SizeIs(2U)); } -TEST_F(FastAppendTest, RetryCopiesAppendManifestAgain) { +TEST_F(FastAppendTest, RebaseCopiesAppendManifestAgain) { table_->metadata()->format_version = 1; + ASSERT_THAT( + TableMetadataUtil::Write(*file_io_, std::string(table_->metadata_file_location()), + *table_->metadata()), + IsOk()); const auto path = table_location_ + "/metadata/input.avro"; ICEBERG_UNWRAP_OR_FAIL(auto manifest, WriteManifest(path, {file_a_})); - - std::shared_ptr fast_append; - ICEBERG_UNWRAP_OR_FAIL(fast_append, table_->NewFastAppend()); + std::shared_ptr txn; + std::vector attempt_manifests; + std::vector attempt_lists; + FailCommits(2, [&](int attempt) { + SCOPED_TRACE(attempt); + if (attempt < 2) { + ICEBERG_UNWRAP_OR_FAIL(auto concurrent, catalog_->LoadTable(table_ident_)); + ICEBERG_UNWRAP_OR_FAIL(auto properties, concurrent->NewUpdateProperties()); + properties->Set(std::format("conflict.{}", attempt), "true"); + ASSERT_THAT(properties->Commit(), IsOk()); + } + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, txn->current().Snapshot()); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.DataManifests(file_io_)); + ASSERT_THAT(manifests, ::testing::SizeIs(1)); + const auto& manifest_path = manifests[0].manifest_path; + EXPECT_NE(manifest_path, path); + EXPECT_THAT(attempt_manifests, ::testing::Not(::testing::Contains(manifest_path))); + attempt_manifests.push_back(manifest_path); + attempt_lists.push_back(snapshot->manifest_list); + ICEBERG_UNWRAP_OR_FAIL(auto entries, ReadEntries(manifests[0])); + ASSERT_THAT(entries, ::testing::SizeIs(1)); + ASSERT_NE(entries[0].data_file, nullptr); + EXPECT_EQ(entries[0].data_file->file_path, file_a_->file_path); + }); + ICEBERG_UNWRAP_OR_FAIL(txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto fast_append, txn->NewFastAppend()); std::vector deleted_paths; fast_append->DeleteWith([&](const std::string& deleted_path) { deleted_paths.push_back(deleted_path); + if (deleted_paths.size() == 1) { + throw std::runtime_error("cleanup callback threw"); + } + if (deleted_paths.size() == 2) { + return Status(IOError("cleanup returned an error")); + } return file_io_->DeleteFile(deleted_path); }); fast_append->AppendManifest(manifest); - - auto& update = static_cast(*fast_append); - // First Apply() copies the input manifest because v1 cannot inherit snapshot IDs. - ICEBERG_UNWRAP_OR_FAIL(auto first_apply, update.Apply()); - SnapshotCache first_cache(first_apply.snapshot.get()); - ICEBERG_UNWRAP_OR_FAIL(auto first_manifests, first_cache.Manifests(file_io_)); - ASSERT_EQ(first_manifests.size(), 1U); - const auto first_rewritten_path = first_manifests[0].manifest_path; - EXPECT_NE(first_rewritten_path, path); - - // Second Apply() simulates retry cleanup, then copies the original manifest again. - ICEBERG_UNWRAP_OR_FAIL(auto second_apply, update.Apply()); - EXPECT_THAT(deleted_paths, testing::Contains(first_rewritten_path)); - - SnapshotCache second_cache(second_apply.snapshot.get()); - ICEBERG_UNWRAP_OR_FAIL(auto second_manifests, second_cache.Manifests(file_io_)); - ASSERT_EQ(second_manifests.size(), 1U); - EXPECT_NE(second_manifests[0].manifest_path, path); - EXPECT_NE(second_manifests[0].manifest_path, first_rewritten_path); + ASSERT_THAT(fast_append->Commit(), IsOk()); + ASSERT_THAT(txn->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, CurrentDataManifests()); + ASSERT_EQ(manifests.size(), 1U); + ASSERT_THAT(attempt_manifests, ::testing::SizeIs(3)); + ASSERT_THAT(attempt_lists, ::testing::SizeIs(3)); + EXPECT_EQ(manifests[0].manifest_path, attempt_manifests[2]); + std::unordered_set unique_deleted_paths(deleted_paths.begin(), + deleted_paths.end()); + EXPECT_THAT(unique_deleted_paths, + ::testing::UnorderedElementsAre(attempt_manifests[0], attempt_lists[0], + attempt_manifests[1], attempt_lists[1])); + EXPECT_THAT(deleted_paths, ::testing::Not(::testing::Contains(path))); + EXPECT_THAT(deleted_paths, + ::testing::Not(::testing::Contains(manifests[0].manifest_path))); + EXPECT_THAT(file_io_->ReadFile(path, std::nullopt), IsOk()); + EXPECT_THAT(file_io_->ReadFile(attempt_lists[2], std::nullopt), IsOk()); } TEST_F(FastAppendTest, AppendDuplicateFile) { @@ -660,21 +806,19 @@ TEST_F(FastAppendMetricsTest, ReporterOverrideAppliesOnlyToItsOwnUpdate) { EXPECT_EQ(second_report.commit_metrics.attempts->value, 1); } -TEST_F(FastAppendMetricsTest, TransactionRetryReportsOnceAfterSuccess) { +TEST_F(FastAppendMetricsTest, ExplicitRetryWithoutMetadataChangeDoesNotReplay) { auto mock_catalog = std::make_shared<::testing::NiceMock>(); - const std::string refreshed_metadata_location = - table_location_ + "/metadata/refreshed.metadata.json"; ON_CALL(*mock_catalog, LoadTable(::testing::_)) - .WillByDefault([this, &mock_catalog, &refreshed_metadata_location]( + .WillByDefault([this, &mock_catalog]( const TableIdentifier&) -> Result> { ICEBERG_ASSIGN_OR_RAISE( auto metadata, TableMetadataUtil::Read(*table_->io(), std::string(table_->metadata_file_location()))); return Table::Make(table_->name(), std::move(metadata), - refreshed_metadata_location, table_->io(), mock_catalog, - table_->full_name(), reporter_); + std::string(table_->metadata_file_location()), table_->io(), + mock_catalog, table_->full_name(), reporter_); }); int update_call_count = 0; @@ -713,7 +857,17 @@ TEST_F(FastAppendMetricsTest, TransactionRetryReportsOnceAfterSuccess) { ASSERT_TRUE(std::holds_alternative(reporter_->reports()[0])); const auto& report = std::get(reporter_->reports()[0]); ASSERT_TRUE(report.commit_metrics.attempts.has_value()); - EXPECT_EQ(report.commit_metrics.attempts->value, 2); + EXPECT_EQ(report.commit_metrics.attempts->value, 1); +} + +TEST_F(FastAppendTest, StandaloneRetryReplaysWithoutMetadataChange) { + FailCommits(1); + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table_, TransactionKind::kUpdate)); + auto update = std::make_shared(ctx); + + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_EQ(update->applies, 2); } TEST_F(FastAppendMetricsTest, CommitStateUnknownDoesNotReport) { @@ -738,4 +892,456 @@ TEST_F(FastAppendMetricsTest, CommitStateUnknownDoesNotReport) { EXPECT_TRUE(reporter_->reports().empty()); } +TEST_F(FastAppendTest, StandaloneReplayFailureStopsRetryAndCleansPartialFiles) { + struct ReplayFailure { + const char* name; + std::function callback; + ErrorKind kind; + const char* message; + }; + const std::vector failures{ + {.name = "retryable status", + .callback = []() -> Status { + return RetryableValidationFailed("replay validation failed"); + }, + .kind = ErrorKind::kValidationFailed, + .message = "replay validation failed"}, + {.name = "unknown status", + .callback = []() -> Status { + return CommitStateUnknown("replay validation failed"); + }, + .kind = ErrorKind::kValidationFailed, + .message = "replay validation failed"}, + {.name = "standard exception", + .callback = []() -> Status { throw std::runtime_error("replay callback threw"); }, + .kind = ErrorKind::kValidationFailed, + .message = "replay callback threw"}, + {.name = "unknown exception", + .callback = []() -> Status { throw 42; }, + .kind = ErrorKind::kValidationFailed, + .message = "Transaction preparation threw an unknown exception"}}; + for (const auto& failure : failures) { + SCOPED_TRACE(failure.name); + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce(::testing::Return(CommitFailed("first conflict"))); + EXPECT_CALL(*mock, LoadTable(::testing::_)).Times(1).WillOnce([&](const auto& name) { + return catalog_->LoadTable(name); + }); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table, TransactionKind::kUpdate)); + auto update = std::make_shared(ctx); + update->write_partial = true; + update->apply_callback = [&failure](int attempt) -> Status { + if (attempt == 2) { + return failure.callback(); + } + return {}; + }; + std::vector deleted; + update->DeleteWith([&](const std::string& path) { + deleted.push_back(path); + return file_io_->DeleteFile(path); + }); + EXPECT_THAT(update->Commit(), ::testing::AllOf(IsError(failure.kind), + HasErrorMessage(failure.message))); + EXPECT_EQ(update->applies, 2); + EXPECT_THAT(deleted, ::testing::SizeIs(3U)); + for (const auto& path : update->partial_paths) { + EXPECT_THAT(deleted, ::testing::Contains(path)); + EXPECT_FALSE(file_io_->ReadFile(path, std::nullopt).has_value()); + } + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_FALSE(ctx->transaction.has_value()); + EXPECT_THAT(deleted, ::testing::SizeIs(3U)); + } +} + +TEST_F(FastAppendTest, CleanupHookFailureDoesNotSkipStagedFiles) { + struct CleanupFailure { + const char* name; + std::function callback; + }; + const std::vector failures{ + {.name = "error status", + .callback = []() -> Status { return IOError("cleanup failed"); }}, + {.name = "standard exception", + .callback = []() -> Status { throw std::runtime_error("cleanup threw"); }}, + {.name = "unknown exception", .callback = []() -> Status { throw 42; }}}; + for (const auto& failure : failures) { + SCOPED_TRACE(failure.name); + for (bool apply_fails : {false, true}) { + SCOPED_TRACE(apply_fails); + ICEBERG_UNWRAP_OR_FAIL(auto table, catalog_->LoadTable(table_ident_)); + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table, TransactionKind::kUpdate)); + auto update = std::make_shared(ctx); + update->write_partial = true; + update->cleanup_callback = failure.callback; + if (apply_fails) { + update->apply_callback = [](int) -> Status { + return ValidationFailed("original apply failure"); + }; + } + std::vector deleted; + update->DeleteWith([&](const std::string& path) { + deleted.push_back(path); + return file_io_->DeleteFile(path); + }); + + auto result = update->Commit(); + if (apply_fails) { + EXPECT_THAT(result, ::testing::AllOf(IsError(ErrorKind::kValidationFailed), + HasErrorMessage("original apply failure"))); + } else { + ASSERT_THAT(result, IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, ctx->table->current_snapshot()); + EXPECT_THAT(file_io_->ReadFile(snapshot->manifest_list, std::nullopt), IsOk()); + } + EXPECT_EQ(update->applies, 1); + EXPECT_EQ(update->cleanups, 1); + ASSERT_THAT(update->partial_paths, ::testing::SizeIs(1)); + EXPECT_THAT(deleted, ::testing::ElementsAre(update->partial_paths[0])); + EXPECT_FALSE( + file_io_->ReadFile(update->partial_paths[0], std::nullopt).has_value()); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_EQ(update->cleanups, 1); + EXPECT_THAT(deleted, ::testing::SizeIs(1)); + } + } +} + +class CallbackReporter : public MetricsReporter { + public: + explicit CallbackReporter(std::function callback) + : callback_(std::move(callback)) {} + Status Report(const MetricsReport&) override { return callback_(); } + + private: + std::function callback_; +}; + +TEST_F(FastAppendTest, ExplicitSuccessReportsAllUpdatesAndIgnoresReporterFailures) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto first, txn->NewFastAppend()); + std::shared_ptr second; + int reports = 0; + first->AppendFile(file_a_).ReportWith( + std::make_shared([&]() -> Status { + ++reports; + throw std::runtime_error("report failure"); + })); + ASSERT_THAT(first->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(second, txn->NewFastAppend()); + second->AppendFile(file_b_).ReportWith( + std::make_shared([&]() -> Status { + ++reports; + return IOError("another report failure"); + })); + ASSERT_THAT(second->Commit(), IsOk()); + ASSERT_THAT(txn->Commit(), IsOk()); + EXPECT_EQ(txn->state(), TransactionState::kCommitted); + EXPECT_EQ(reports, 2); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_EQ(reports, 2); +} + +TEST_F(FastAppendTest, TransientConflictThenUnknownPreservesInitialAttempt) { + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillOnce(::testing::Return(CommitFailed("conflict"))) + .WillOnce(::testing::Return(CommitStateUnknown("unknown"))); + EXPECT_CALL(*mock, LoadTable(::testing::_)).WillOnce([&](const auto& name) { + return catalog_->LoadTable(name); + }); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto append, txn->NewFastAppend()); + std::vector deleted_paths; + append->AppendFile(file_a_).DeleteWith([&](const std::string& path) { + deleted_paths.push_back(path); + return file_io_->DeleteFile(path); + }); + ASSERT_THAT(append->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kCommitStateUnknown)); + EXPECT_EQ(txn->state(), TransactionState::kCommitStateUnknown); + EXPECT_TRUE(deleted_paths.empty()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, txn->current().Snapshot()); + EXPECT_THAT(file_io_->ReadFile(snapshot->manifest_list, std::nullopt), IsOk()); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.Manifests(file_io_)); + ASSERT_EQ(manifests.size(), 1U); + EXPECT_THAT(file_io_->ReadFile(manifests[0].manifest_path, std::nullopt), IsOk()); + EXPECT_THAT(txn->Abort(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(append->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_TRUE(deleted_paths.empty()); +} + +TEST_F(FastAppendTest, ApplyFailureAfterWritingCleansEveryUpdate) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + int deletes = 0; + auto delete_file = [&](const std::string& path) -> Status { + ++deletes; + if (deletes == 1) { + throw std::runtime_error("delete callback threw"); + } + return IOError("delete failure"); + }; + ICEBERG_UNWRAP_OR_FAIL(auto first, txn->NewFastAppend()); + first->AppendFile(file_a_).DeleteWith(delete_file); + ASSERT_THAT(first->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto failed, txn->NewRewriteFiles()); + failed->DeleteDataFile(file_a_).AddDataFile(file_b_).DeleteWith(delete_file); + EXPECT_THAT(failed->Commit(), HasErrorMessage("Invalid REPLACE operation")); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + EXPECT_GE(deletes, 4); + const int cleanup_attempts = deletes; + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_GT(deletes, cleanup_attempts); + const int abort_cleanup_attempts = deletes; + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_EQ(deletes, abort_cleanup_attempts); +} + +TEST_F(FastAppendTest, ReplayBecomingNoopCleansStagingWithoutReporting) { + for (bool other_update : {false, true}) { + SCOPED_TRACE(other_update); + auto mock = std::make_shared<::testing::NiceMock>(); + std::shared_ptr refreshed; + std::shared_ptr refreshed_metadata; + int attempts = 0; + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(other_update ? 2 : 1) + .WillRepeatedly([&](const auto&, const auto&, + const auto& changes) -> Result> { + const bool first_attempt = ++attempts == 1; + auto builder = TableMetadataBuilder::BuildFrom( + first_attempt ? table_->metadata().get() : refreshed_metadata.get()); + for (const auto& change : changes) { + if (first_attempt && change->kind() == TableUpdate::Kind::kSetProperties) { + continue; + } + if (!first_attempt) { + EXPECT_EQ(change->kind(), TableUpdate::Kind::kSetProperties); + } + change->ApplyTo(*builder); + } + ICEBERG_ASSIGN_OR_RAISE(std::shared_ptr metadata, + builder->Build()); + if (first_attempt) { + // Use independent files for the refreshed snapshot; generated paths are + // owned by the update that created them. + metadata->snapshots.back() = + std::make_shared(*metadata->snapshots.back()); + auto snapshot = metadata->snapshots.back(); + snapshot->manifest_list = table_location_ + "/metadata/concurrent-list.avro"; + ICEBERG_ASSIGN_OR_RAISE( + auto writer, ManifestListWriter::MakeWriter( + metadata->format_version, snapshot->snapshot_id, + snapshot->parent_snapshot_id, snapshot->manifest_list, + file_io_, snapshot->sequence_number)); + ICEBERG_RETURN_UNEXPECTED(writer->Close()); + } + // Simulate the same logical snapshot already becoming current on refresh. + refreshed_metadata = metadata; + ICEBERG_ASSIGN_OR_RAISE( + refreshed, + Table::Make(table_->name(), std::move(metadata), + std::string(table_->metadata_file_location()) + ".refreshed", + file_io_, mock)); + if (first_attempt) { + return CommitFailed("conflict"); + } + return refreshed; + }); + EXPECT_CALL(*mock, LoadTable(::testing::_)) + .WillOnce( + [&](const auto&) -> Result> { return refreshed; }); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + std::shared_ptr txn; + if (other_update) { + ICEBERG_UNWRAP_OR_FAIL(txn, table->NewTransaction()); + } + ICEBERG_UNWRAP_OR_FAIL(auto append, + txn ? txn->NewFastAppend() : table->NewFastAppend()); + int reports = 0; + int deletes = 0; + append->AppendFile(file_a_) + .ReportWith(std::make_shared([&]() -> Status { + ++reports; + return {}; + })) + .DeleteWith([&](const std::string& path) { + ++deletes; + return file_io_->DeleteFile(path); + }); + ASSERT_THAT(append->Commit(), IsOk()); + if (txn) { + ICEBERG_UNWRAP_OR_FAIL(auto props, txn->NewUpdateProperties()); + props->Set("effective", "value"); + ASSERT_THAT(props->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto committed, txn->Commit()); + EXPECT_EQ(committed->properties().configs().at("effective"), "value"); + EXPECT_EQ(txn->state(), TransactionState::kCommitted); + } + EXPECT_EQ(deletes, 4); + EXPECT_EQ(reports, 0); + EXPECT_THAT(append->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_EQ(deletes, 4); + } +} + +TEST_F(FastAppendTest, StagedNoopIsCleanedAndCanBecomeEffectiveOnRebase) { + for (bool rebase : {false, true}) { + SCOPED_TRACE(rebase); + auto mock = std::make_shared<::testing::NiceMock>(); + ON_CALL(*mock, LoadTable(::testing::_)).WillByDefault([&](const auto& name) { + return catalog_->LoadTable(name); + }); + ON_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault( + [&](const auto& name, const auto& requirements, const auto& changes) { + return catalog_->UpdateTable(name, requirements, changes); + }); + EXPECT_CALL(*mock, LoadTable(::testing::_)).Times(rebase ? 1 : 0); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(rebase ? 1 : 0); + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table_, TransactionKind::kUpdate)); + auto update = std::make_shared(ctx); + update->write_partial = true; + auto builder = TableMetadataBuilder::BuildFrom(table_->metadata().get()); + auto existing = std::make_shared(Snapshot{ + .snapshot_id = update->SnapshotId(), + .sequence_number = table_->metadata()->NextSequenceNumber(), + .timestamp_ms = CurrentTimePointMs(), + .manifest_list = table_location_ + "/metadata/existing-noop-list.avro", + .summary = {{SnapshotSummaryFields::kOperation, DataOperation::kAppend}}, + .schema_id = table_->metadata()->current_schema_id, + }); + ICEBERG_UNWRAP_OR_FAIL( + auto writer, ManifestListWriter::MakeWriter(table_->metadata()->format_version, + existing->snapshot_id, std::nullopt, + existing->manifest_list, file_io_, + existing->sequence_number)); + ASSERT_THAT(writer->Close(), IsOk()); + builder->SetBranchSnapshot(existing, std::string(SnapshotRef::kMainBranch)); + ICEBERG_UNWRAP_OR_FAIL(std::shared_ptr metadata, builder->Build()); + // Prepare a base whose current snapshot has this update's logical ID. + ICEBERG_UNWRAP_OR_FAIL( + ctx->table, + Table::Make(table_->name(), metadata, + std::string(table_->metadata_file_location()) + ".synthetic", + file_io_, mock)); + ctx->metadata_builder = TableMetadataBuilder::BuildFrom(metadata.get()); + std::vector deleted; + update->DeleteWith([&](const std::string& path) { + deleted.push_back(path); + return file_io_->DeleteFile(path); + }); + if (rebase) { + // The table refreshes before the update commits its original builder. + // Initial Apply is a no-op; internal replay sees the real empty table. + ASSERT_THAT(ctx->table->Refresh(), IsOk()); + } + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_EQ(update->applies, rebase ? 2 : 1); + EXPECT_THAT(deleted, ::testing::SizeIs(rebase ? 3U : 2U)); + EXPECT_THAT(file_io_->ReadFile(existing->manifest_list, std::nullopt), IsOk()); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + } +} + +TEST_F(FastAppendTest, ReplayFailurePartwayThroughCleansAllUncommittedGenerations) { + file_a_->partition = PartitionValues({Literal::Long(1)}); + file_b_->partition = PartitionValues({Literal::Long(3)}); + ICEBERG_UNWRAP_OR_FAIL(auto initial, table_->NewFastAppend()); + initial->AppendFile(file_a_); + ASSERT_THAT(initial->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto original_snapshot, table_->current_snapshot()); + auto concurrent_file = CreateDataFile("/data/concurrent.parquet", 1, 10, 1); + auto replacement = CreateDataFile("/data/replacement.parquet", 100, 50, 1); + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce( + [&](const auto&, const auto&, const auto&) -> Result> { + ICEBERG_ASSIGN_OR_RAISE(auto latest, catalog_->LoadTable(table_ident_)); + ICEBERG_ASSIGN_OR_RAISE(auto concurrent, latest->NewFastAppend()); + concurrent->AppendFile(concurrent_file); + ICEBERG_RETURN_UNEXPECTED(concurrent->Commit()); + return CommitFailed("concurrent append"); + }); + EXPECT_CALL(*mock, LoadTable(::testing::_)).Times(1).WillOnce([&](const auto& name) { + return catalog_->LoadTable(name); + }); + ICEBERG_UNWRAP_OR_FAIL( + auto table, + Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock)); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table->NewTransaction()); + std::vector deleted; + std::shared_ptr append; + auto delete_file = [&](const std::string& path) { + deleted.push_back(path); + return file_io_->DeleteFile(path); + }; + ICEBERG_UNWRAP_OR_FAIL(append, txn->NewFastAppend()); + append->AppendFile(file_b_).DeleteWith(delete_file); + ASSERT_THAT(append->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto overwrite, txn->NewOverwrite()); + overwrite->DeleteFile(file_a_) + .AddFile(replacement) + .ValidateFromSnapshot(original_snapshot->snapshot_id) + .ConflictDetectionFilter(Expressions::Equal("x", Literal::Long(1))) + .ValidateNoConflictingData() + .DeleteWith(delete_file); + ASSERT_THAT(overwrite->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), HasErrorMessage("Found conflicting files")); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + EXPECT_GE(deleted.size(), 6U); + std::unordered_set unique_deletes(deleted.begin(), deleted.end()); + EXPECT_EQ(unique_deletes.size(), deleted.size()); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_EQ(deleted.size(), unique_deletes.size()); + + std::unordered_set committed_paths; + auto metadata = ReloadMetadata(); + for (const auto& snapshot : metadata->snapshots) { + committed_paths.insert(snapshot->manifest_list); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.Manifests(file_io_)); + for (const auto& manifest : manifests) { + committed_paths.insert(manifest.manifest_path); + } + } + auto& io = static_cast(*file_io_); + ::arrow::fs::FileSelector selector; + selector.base_dir = table_location_ + "/metadata"; + selector.recursive = true; + auto files = io.fs()->GetFileInfo(selector); + ASSERT_TRUE(files.ok()) << files.status(); + std::unordered_set remaining; + for (const auto& file : *files) { + if (file.path().ends_with(".avro")) { + remaining.insert(file.path()); + } + } + EXPECT_EQ(remaining, committed_paths); +} + } // namespace iceberg diff --git a/src/iceberg/test/merging_snapshot_update_test.cc b/src/iceberg/test/merging_snapshot_update_test.cc index 25907c445..84e372952 100644 --- a/src/iceberg/test/merging_snapshot_update_test.cc +++ b/src/iceberg/test/merging_snapshot_update_test.cc @@ -55,6 +55,7 @@ #include "iceberg/test/retry.h" #include "iceberg/test/update_test_base.h" #include "iceberg/transaction.h" +#include "iceberg/update/delete_files.h" #include "iceberg/update/fast_append.h" #include "iceberg/update/merge_append.h" #include "iceberg/update/row_delta.h" @@ -81,20 +82,33 @@ class MergingSnapshotCapturingReporter final : public MetricsReporter { /// \brief Concrete subclass of MergingSnapshotUpdate for testing. class TestMergeAppend : public MergingSnapshotUpdate { public: - static Result> Make(std::string table_name, + static Result> Make(std::string table_name, std::shared_ptr
table) { ICEBERG_ASSIGN_OR_RAISE( auto ctx, TransactionContext::Make(std::move(table), TransactionKind::kUpdate)); - return std::unique_ptr( + return std::shared_ptr( new TestMergeAppend(std::move(table_name), std::move(ctx))); } std::string operation() override { return "append"; } // Expose protected API for test access - using MergingSnapshotUpdate::Apply; - using MergingSnapshotUpdate::CleanUncommitted; - using MergingSnapshotUpdate::Summary; + Result> StagedSnapshot() const { + return ctx_->current().Snapshot(); + } + + Result> ApplyForTest( + const TableMetadata& metadata, const std::shared_ptr& snapshot) { + return MergingSnapshotUpdate::Apply(metadata, snapshot); + } + + Result> CommitManifests() { + ICEBERG_RETURN_UNEXPECTED(Commit()); + ICEBERG_ASSIGN_OR_RAISE(auto snapshot, ctx_->table->current_snapshot()); + SnapshotCache cache(snapshot.get()); + ICEBERG_ASSIGN_OR_RAISE(auto manifests, cache.Manifests(ctx_->table->io())); + return std::vector(manifests.begin(), manifests.end()); + } Status AddFile(std::shared_ptr file) { return AddDataFile(std::move(file)); } Status AddDelete(std::shared_ptr file) { @@ -231,18 +245,21 @@ class TestMergeAppend : public MergingSnapshotUpdate { class TestOverwriteUpdate : public MergingSnapshotUpdate { public: - static Result> Make(std::string table_name, + static Result> Make(std::string table_name, std::shared_ptr
table) { ICEBERG_ASSIGN_OR_RAISE( auto ctx, TransactionContext::Make(std::move(table), TransactionKind::kUpdate)); - return std::unique_ptr( + return std::shared_ptr( new TestOverwriteUpdate(std::move(table_name), std::move(ctx))); } std::string operation() override { return DataOperation::kOverwrite; } int64_t GeneratedSnapshotId() { return SnapshotId(); } - using MergingSnapshotUpdate::Apply; + Result> ApplyForTest( + const TableMetadata& metadata, const std::shared_ptr& snapshot) { + return MergingSnapshotUpdate::Apply(metadata, snapshot); + } Status AddDelete(std::shared_ptr file) { return AddDeleteFile(std::move(file)); @@ -335,11 +352,11 @@ class MergingSnapshotUpdateTest : public MinimalUpdateTestBase { return f; } - Result> NewMergeAppend() { + Result> NewMergeAppend() { return TestMergeAppend::Make(TableName(), table_); } - Result> NewOverwriteUpdate() { + Result> NewOverwriteUpdate() { return TestOverwriteUpdate::Make(TableName(), table_); } @@ -887,13 +904,28 @@ TEST_F(MergingSnapshotUpdateTest, BaseSetCustomSummaryPropertySurvivesApplyRebui // CleanUncommitted test // ------------------------------------------------------------------------- -TEST_F(MergingSnapshotUpdateTest, CleanUncommittedAfterSuccessfulCommitDoesNotCrash) { - ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); - EXPECT_THAT(op->AddFile(file_a_), IsOk()); - EXPECT_THAT(op->Commit(), IsOk()); - - // Cleanup may run from an error handler even after commit success. - EXPECT_THAT(op->CleanUncommitted({}), IsOk()); +TEST_F(MergingSnapshotUpdateTest, CommittedFilesSurviveRejectedCommitAndAbort) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto op, txn->NewMergeAppend()); + std::vector deleted_paths; + op->AppendFile(file_a_).DeleteWith([&](const std::string& path) { + deleted_paths.push_back(path); + return file_io_->DeleteFile(path); + }); + ASSERT_THAT(op->Commit(), IsOk()); + ASSERT_THAT(txn->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, txn->table()->current_snapshot()); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.Manifests(file_io_)); + ASSERT_THAT(manifests, ::testing::SizeIs(1)); + + EXPECT_TRUE(deleted_paths.empty()); + EXPECT_THAT(op->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsError(ErrorKind::kValidationFailed)); + EXPECT_TRUE(deleted_paths.empty()); + EXPECT_THAT(file_io_->ReadFile(manifests[0].manifest_path, std::nullopt), IsOk()); + EXPECT_THAT(file_io_->ReadFile(snapshot->manifest_list, std::nullopt), IsOk()); } TEST_F(MergingSnapshotUpdateTest, @@ -901,23 +933,39 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto initial, NewMergeAppend()); EXPECT_THAT(initial->AddFile(file_a_), IsOk()); EXPECT_THAT(initial->AddFile(file_b_), IsOk()); - EXPECT_THAT(initial->Commit(), IsOk()); - EXPECT_THAT(table_->Refresh(), IsOk()); + ASSERT_THAT(initial->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ICEBERG_UNWRAP_OR_FAIL(auto initial_snapshot, table_->current_snapshot()); + SnapshotCache initial_cache(initial_snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto initial_manifests, initial_cache.Manifests(file_io_)); + ASSERT_THAT(initial_manifests, ::testing::SizeIs(1)); std::vector deleted_paths; - ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto op, txn->NewDeleteFiles()); op->DeleteWith([&deleted_paths](const std::string& path) { deleted_paths.push_back(path); return Status{}; }); - EXPECT_THAT(op->RemoveDataFile(file_a_), IsOk()); - - ICEBERG_UNWRAP_OR_FAIL( - auto manifests, op->Apply(*table_->metadata(), table_->current_snapshot().value())); - EXPECT_THAT(manifests, ::testing::SizeIs(1)); + op->DeleteFile(file_a_->file_path); + ASSERT_THAT(op->Commit(), IsOk()); - EXPECT_THAT(op->CleanUncommitted({}), IsOk()); - EXPECT_THAT(deleted_paths, ::testing::Contains(::testing::HasSubstr("/metadata/"))); + ICEBERG_UNWRAP_OR_FAIL(auto staged_snapshot, txn->current().Snapshot()); + SnapshotCache staged_cache(staged_snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto staged_manifests, staged_cache.Manifests(file_io_)); + ASSERT_THAT(staged_manifests, ::testing::SizeIs(1)); + ASSERT_NE(staged_manifests[0].manifest_path, initial_manifests[0].manifest_path); + EXPECT_TRUE(deleted_paths.empty()); + + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_THAT(deleted_paths, + ::testing::UnorderedElementsAre(staged_manifests[0].manifest_path, + staged_snapshot->manifest_list)); + EXPECT_THAT(deleted_paths, + ::testing::Not(::testing::Contains(initial_manifests[0].manifest_path))); + EXPECT_THAT(deleted_paths, + ::testing::Not(::testing::Contains(initial_snapshot->manifest_list))); } // ------------------------------------------------------------------------- @@ -957,7 +1005,7 @@ TEST_F(MergingSnapshotUpdateTest, AddDeleteFileWithExplicitSequenceWritesSequenc ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); EXPECT_THAT(op->AddDelete(del_file, 17), IsOk()); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), nullptr)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->ApplyForTest(*table_->metadata(), nullptr)); auto delete_manifest_it = std::ranges::find_if(manifests, [](const ManifestFile& manifest) { return manifest.content == ManifestContent::kDeletes; @@ -994,24 +1042,37 @@ TEST_F(MergingSnapshotUpdateTest, WriteDeleteGroups) { } } -TEST_F(MergingSnapshotUpdateTest, ApplyRebuildsDeleteSummaryAfterPreparingDeletes) { +TEST_F(MergingSnapshotUpdateTest, RetryRebuildsDeleteSummary) { auto del_file = MakeDeleteFile("/delete/del_a.parquet", 1L); - - ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); + std::shared_ptr op; + FailCommits(2, [&](int attempt) { + SCOPED_TRACE(attempt); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, op->StagedSnapshot()); + EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kAddedDeleteFiles), "1"); + EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kAddedPosDeleteFiles), "1"); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.DeleteManifests(file_io_)); + ASSERT_THAT(manifests, ::testing::SizeIs(1)); + ICEBERG_UNWRAP_OR_FAIL( + auto entries, + ReadAllEntries(std::vector(manifests.begin(), manifests.end()), + *table_->metadata())); + ASSERT_THAT(entries, ::testing::SizeIs(1)); + ASSERT_NE(entries[0].data_file, nullptr); + EXPECT_EQ(entries[0].status, ManifestStatus::kAdded); + EXPECT_EQ(entries[0].snapshot_id, snapshot->snapshot_id); + auto expected_file = *del_file; + expected_file.data_sequence_number = snapshot->sequence_number; + EXPECT_EQ(*entries[0].data_file, expected_file); + }); + ICEBERG_UNWRAP_OR_FAIL(op, NewMergeAppend()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); - - ICEBERG_UNWRAP_OR_FAIL(auto first_manifests, op->Apply(*table_->metadata(), nullptr)); - EXPECT_THAT(first_manifests, ::testing::Contains(::testing::Field( - &ManifestFile::content, ManifestContent::kDeletes))); - - ICEBERG_UNWRAP_OR_FAIL(auto second_manifests, op->Apply(*table_->metadata(), nullptr)); - EXPECT_THAT(second_manifests, ::testing::Contains(::testing::Field( - &ManifestFile::content, ManifestContent::kDeletes))); - - auto summary = op->Summary(); - EXPECT_EQ(summary.at(SnapshotSummaryFields::kAddedDeleteFiles), "1"); - EXPECT_EQ(summary.at(SnapshotSummaryFields::kAddedPosDeleteFiles), "1"); + EXPECT_THAT(op->Commit(), IsOk()); + EXPECT_THAT(table_->Refresh(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, table_->current_snapshot()); + EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kAddedDeleteFiles), "1"); + EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kAddedPosDeleteFiles), "1"); } // Covers the bug where deleted delete files were not tracked in the snapshot summary. @@ -1101,7 +1162,7 @@ class MergingSnapshotUpdateV1Test : public UpdateTestBase { return f; } - Result> NewMergeAppend() { + Result> NewMergeAppend() { return TestMergeAppend::Make(TableName(), table_); } @@ -1179,7 +1240,7 @@ TEST_F(MergingSnapshotUpdateTest, ApplyMergesDuplicateDeletionVectors) { EXPECT_THAT(op->AddDelete(dv_a2, 7), IsOk()); EXPECT_THAT(op->AddDelete(dv_b, 8), IsOk()); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), nullptr)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->ApplyForTest(*table_->metadata(), nullptr)); auto delete_manifest_it = std::ranges::find_if(manifests, [](const ManifestFile& manifest) { return manifest.content == ManifestContent::kDeletes; @@ -1230,7 +1291,7 @@ TEST_F(MergingSnapshotUpdateTest, ApplyMergesDuplicateDeletionVectorsWithNullPar EXPECT_THAT(op->AddDelete(dv_a1, 7), IsOk()); EXPECT_THAT(op->AddDelete(dv_a2, 7), IsOk()); - EXPECT_THAT(op->Apply(*table_->metadata(), nullptr), IsOk()); + EXPECT_THAT(op->ApplyForTest(*table_->metadata(), nullptr), IsOk()); } TEST_F(MergingSnapshotUpdateTest, ValidateNewDeleteFileRejectsUnsupportedVersion) { @@ -1246,15 +1307,23 @@ TEST_F(MergingSnapshotUpdateTest, ValidateNewDeleteFileRejectsUnsupportedVersion EXPECT_THAT(op->AddDelete(del_file), IsError(ErrorKind::kInvalidArgument)); } -TEST_F(MergingSnapshotUpdateTest, ApplyRejectsV2StagedPositionDeleteAfterV3Upgrade) { +TEST_F(MergingSnapshotUpdateTest, + CommitRejectsV2StagedPositionDeleteAfterV3UpgradeOnReplay) { auto del_file = MakeDeleteFile("/delete/del_a.parquet", 1L); ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); - auto metadata = std::make_shared(*table_->metadata()); - metadata->format_version = 3; - EXPECT_THAT(op->Apply(*metadata, nullptr), IsError(ErrorKind::kInvalidArgument)); + ICEBERG_UNWRAP_OR_FAIL(auto properties, table_->NewUpdateProperties()); + properties->Set("format-version", "3"); + ASSERT_THAT(properties->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + EXPECT_THAT(op->Commit(), + ::testing::AllOf( + IsError(ErrorKind::kValidationFailed), + HasErrorMessage("Transaction replay failed: Must use DVs for position " + "deletes in V3:"), + HasErrorMessage("/delete/del_a.parquet"))); } // ------------------------------------------------------------------------- @@ -1314,29 +1383,41 @@ TEST_F(MergingSnapshotUpdateTest, AddManifestRetryCopiesManifestAgain) { auto path = table_location_ + "/metadata/retry-input.avro"; ICEBERG_UNWRAP_OR_FAIL(auto manifest, WriteManifest(path, {file_a_})); manifest.added_snapshot_id = 12345; - - ICEBERG_UNWRAP_OR_FAIL(auto op, NewMergeAppend()); + std::shared_ptr op; + std::vector attempt_manifests; + std::vector attempt_lists; + FailCommits(2, [&](int attempt) { + SCOPED_TRACE(attempt); + ICEBERG_UNWRAP_OR_FAIL(auto snapshot, op->StagedSnapshot()); + SnapshotCache cache(snapshot.get()); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, cache.DataManifests(file_io_)); + ASSERT_THAT(manifests, ::testing::SizeIs(1)); + const auto& manifest_path = manifests[0].manifest_path; + EXPECT_NE(manifest_path, path); + EXPECT_THAT(attempt_manifests, ::testing::Not(::testing::Contains(manifest_path))); + attempt_manifests.push_back(manifest_path); + attempt_lists.push_back(snapshot->manifest_list); + }); + ICEBERG_UNWRAP_OR_FAIL(op, NewMergeAppend()); + std::vector deleted; + op->DeleteWith([&](const std::string& path) { + deleted.push_back(path); + return file_io_->DeleteFile(path); + }); EXPECT_THAT(op->AppendManifest(manifest), IsOk()); - - ICEBERG_UNWRAP_OR_FAIL(auto first_apply, static_cast(*op).Apply()); - SnapshotCache first_snapshot_cache(first_apply.snapshot.get()); - ICEBERG_UNWRAP_OR_FAIL(auto first_manifests, - first_snapshot_cache.DataManifests(file_io_)); - ASSERT_EQ(first_manifests.size(), 1U); - EXPECT_NE(first_manifests[0].manifest_path, path); - - ICEBERG_UNWRAP_OR_FAIL(auto second_apply, static_cast(*op).Apply()); - SnapshotCache second_snapshot_cache(second_apply.snapshot.get()); - ICEBERG_UNWRAP_OR_FAIL(auto second_manifests, - second_snapshot_cache.DataManifests(file_io_)); - ASSERT_EQ(second_manifests.size(), 1U); - EXPECT_NE(second_manifests[0].manifest_path, path); - EXPECT_NE(second_manifests[0].manifest_path, first_manifests[0].manifest_path); - - std::vector second_manifest_vector(second_manifests.begin(), - second_manifests.end()); - ICEBERG_UNWRAP_OR_FAIL(auto entries, - ReadAllEntries(second_manifest_vector, *table_->metadata())); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->CommitManifests()); + ASSERT_EQ(manifests.size(), 1U); + ASSERT_THAT(attempt_manifests, ::testing::SizeIs(3)); + ASSERT_THAT(attempt_lists, ::testing::SizeIs(3)); + EXPECT_EQ(manifests[0].manifest_path, attempt_manifests[2]); + EXPECT_THAT(deleted, + ::testing::UnorderedElementsAre(attempt_manifests[0], attempt_lists[0], + attempt_manifests[1], attempt_lists[1])); + EXPECT_THAT(deleted, ::testing::Not(::testing::Contains(path))); + EXPECT_THAT(deleted, ::testing::Not(::testing::Contains(manifests[0].manifest_path))); + EXPECT_THAT(file_io_->ReadFile(path, std::nullopt), IsOk()); + EXPECT_THAT(file_io_->ReadFile(attempt_lists[2], std::nullopt), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto entries, ReadAllEntries(manifests, *table_->metadata())); ASSERT_EQ(entries.size(), 1U); ASSERT_NE(entries[0].data_file, nullptr); EXPECT_EQ(entries[0].data_file->file_path, file_a_->file_path); @@ -1722,7 +1803,8 @@ TEST_F(MergingSnapshotUpdateTest, ValidateDataFilesExistUsesRowFilter) { ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->RemoveDataFile(file_a_), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1756,7 +1838,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1784,7 +1867,8 @@ TEST_F(MergingSnapshotUpdateTest, ValidateNoNewDeletesForDataFilesDetectsConflic ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1813,7 +1897,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1848,7 +1933,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1880,7 +1966,7 @@ TEST_F(MergingSnapshotUpdateTest, EXPECT_THAT(overwrite->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = overwrite->GeneratedSnapshotId(); ICEBERG_UNWRAP_OR_FAIL(auto manifests, - overwrite->Apply(*table_->metadata(), first_snapshot)); + overwrite->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1912,7 +1998,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1940,7 +2027,8 @@ TEST_F(MergingSnapshotUpdateTest, ValidateNoNewDeleteFilesWithExpressionDetectsC ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1967,7 +2055,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -1994,7 +2083,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->AddDelete(del_file), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -2138,7 +2228,8 @@ TEST_F(MergingSnapshotUpdateTest, ValidateDeletedDataFilesWithExpressionDetectsC ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->RemoveDataFile(file_a_), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, @@ -2164,7 +2255,8 @@ TEST_F(MergingSnapshotUpdateTest, ICEBERG_UNWRAP_OR_FAIL(auto op, NewOverwriteUpdate()); EXPECT_THAT(op->RemoveDataFile(file_a_), IsOk()); const int64_t second_snapshot_id = op->GeneratedSnapshotId(); - ICEBERG_UNWRAP_OR_FAIL(auto manifests, op->Apply(*table_->metadata(), first_snapshot)); + ICEBERG_UNWRAP_OR_FAIL(auto manifests, + op->ApplyForTest(*table_->metadata(), first_snapshot)); ICEBERG_UNWRAP_OR_FAIL( auto second_snapshot, MakeSyntheticSnapshot(DataOperation::kOverwrite, second_snapshot_id, diff --git a/src/iceberg/test/snapshot_manager_test.cc b/src/iceberg/test/snapshot_manager_test.cc index f4fe68958..8b6c82d46 100644 --- a/src/iceberg/test/snapshot_manager_test.cc +++ b/src/iceberg/test/snapshot_manager_test.cc @@ -387,6 +387,14 @@ TEST_F(SnapshotManagerTest, SnapshotManagerThroughTransaction) { ExpectCurrentSnapshot(oldest_snapshot_id_); } +TEST_F(SnapshotManagerTest, ExternalManagerRejectsTerminalTransaction) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ASSERT_THAT(txn->Commit(), IsOk()); + ICEBERG_UNWRAP_OR_FAIL(auto manager, txn->NewSnapshotManager()); + + EXPECT_THAT(manager->Commit(), HasErrorMessage("Transaction is terminal")); +} + TEST_F(SnapshotManagerTest, SnapshotManagerFromTableAllowsMultipleSnapshotOperations) { ICEBERG_UNWRAP_OR_FAIL(auto manager, table_->NewSnapshotManager()); manager->SetCurrentSnapshot(oldest_snapshot_id_); diff --git a/src/iceberg/test/transaction_test.cc b/src/iceberg/test/transaction_test.cc index 3a13b7bc5..36c6bb514 100644 --- a/src/iceberg/test/transaction_test.cc +++ b/src/iceberg/test/transaction_test.cc @@ -19,14 +19,24 @@ #include "iceberg/transaction.h" +#include +#include +#include +#include + +#include "iceberg/exception.h" #include "iceberg/expression/expressions.h" #include "iceberg/expression/term.h" +#include "iceberg/schema.h" #include "iceberg/sort_order.h" +#include "iceberg/table_metadata.h" #include "iceberg/test/matchers.h" #include "iceberg/test/mock_catalog.h" #include "iceberg/test/update_test_base.h" #include "iceberg/transform.h" #include "iceberg/type.h" +#include "iceberg/update/fast_append.h" +#include "iceberg/update/set_snapshot.h" #include "iceberg/update/update_properties.h" #include "iceberg/update/update_schema.h" #include "iceberg/update/update_sort_order.h" @@ -46,6 +56,14 @@ TEST_F(TransactionTest, CommitEmptyTransaction) { EXPECT_THAT(txn->Commit(), IsOk()); } +TEST_F(TransactionTest, CommitNoOpUpdate) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewSetSnapshot()); + + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), IsOk()); +} + TEST_F(TransactionTest, CommitTransactionWithPropertyUpdate) { ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); @@ -151,6 +169,14 @@ TEST_F(TransactionRetryTest, CommitRetrySucceedsAfterConflict) { EXPECT_EQ(update_call_count, 2); } +TEST_F(TransactionTest, StandaloneCommitRetryReappliesUpdate) { + FailCommits(1); + ICEBERG_UNWRAP_OR_FAIL(auto update, table_->NewUpdateProperties()); + update->Set("retry.test", "value"); + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_EQ(ReloadMetadata()->properties.configs().at("retry.test"), "value"); +} + TEST_F(TransactionRetryTest, CommitRetryExhausted) { int update_call_count = 0; ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) @@ -195,6 +221,32 @@ TEST_F(TransactionRetryTest, CommitNonRetryableErrorStopsImmediately) { EXPECT_EQ(update_call_count, 1); // Should not retry } +TEST_F(TransactionRetryTest, CommitExceptionMakesOutcomeUnknown) { + int update_call_count = 0; + ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault( + [&update_call_count](const TableIdentifier&, + const std::vector>&, + const std::vector>&) + -> Result> { + ++update_call_count; + throw std::runtime_error("injected catalog failure"); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto txn, mock_table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto properties, txn->NewUpdateProperties()); + properties->Set("exception.test", "value"); + EXPECT_THAT(properties->Commit(), IsOk()); + + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kCommitStateUnknown)); + EXPECT_EQ(txn->state(), TransactionState::kCommitStateUnknown); + EXPECT_THAT(txn->NewFastAppend(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(properties->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsError(ErrorKind::kValidationFailed)); + EXPECT_EQ(update_call_count, 1); +} + TEST_F(TransactionRetryTest, CreateTransactionDoesNotRetry) { int update_call_count = 0; ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) @@ -240,4 +292,81 @@ TEST_F(TransactionRetryTest, NonRetryableUpdatePreventsRetry) { EXPECT_EQ(update_call_count, 1); } +TEST_F(TransactionTest, AppliedUpdateCannotCompleteAnotherPendingOperation) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto first, txn->NewUpdateProperties()); + first->Set("first", "1"); + ASSERT_THAT(first->Commit(), IsOk()); + EXPECT_EQ(txn->state(), TransactionState::kReady); + ICEBERG_UNWRAP_OR_FAIL(auto second, txn->NewUpdateProperties()); + EXPECT_THAT(first->Commit(), HasErrorMessage("Update has already been committed")); + EXPECT_EQ(txn->state(), TransactionState::kUpdatePending); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->NewUpdateProperties(), IsError(ErrorKind::kValidationFailed)); + second->Set("second", "2"); + ASSERT_THAT(second->Commit(), IsOk()); + ASSERT_THAT(txn->Commit(), IsOk()); + EXPECT_EQ(txn->state(), TransactionState::kCommitted); + EXPECT_EQ(ReloadMetadata()->properties.configs().at("first"), "1"); + EXPECT_EQ(ReloadMetadata()->properties.configs().at("second"), "2"); + EXPECT_THAT(txn->Abort(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(first->Commit(), IsError(ErrorKind::kValidationFailed)); +} + +TEST_F(TransactionTest, StandaloneCommitIsTerminal) { + ICEBERG_UNWRAP_OR_FAIL(auto ctx, + TransactionContext::Make(table_, TransactionKind::kUpdate)); + ICEBERG_UNWRAP_OR_FAIL(auto update, UpdateProperties::Make(ctx)); + update->Set("once", "value"); + ASSERT_THAT(update->Commit(), IsOk()); + EXPECT_FALSE(ctx->transaction.has_value()); + EXPECT_EQ(ReloadMetadata()->properties.configs().at("once"), "value"); + EXPECT_THAT(update->Commit(), HasErrorMessage("Update has already been committed")); +} + +TEST_F(TransactionTest, AbortReadyAndPendingIsIdempotent) { + for (bool add_update : {false, true}) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + std::shared_ptr update; + if (add_update) { + ICEBERG_UNWRAP_OR_FAIL(update, txn->NewUpdateProperties()); + update->Set("discard", "value"); + } + ASSERT_THAT(txn->Abort(), IsOk()); + EXPECT_EQ(txn->state(), TransactionState::kAborted); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->NewFastAppend(), IsError(ErrorKind::kValidationFailed)); + if (update) { + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + } + } +} + +TEST_F(TransactionTest, ApplyFailureCannotBeCorrected) { + ICEBERG_UNWRAP_OR_FAIL(auto txn, table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); + update->Set("format-version", "100"); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kInvalidArgument)); + EXPECT_EQ(txn->state(), TransactionState::kFailed); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Commit(), + ::testing::AllOf(IsError(ErrorKind::kValidationFailed), + HasErrorMessage("Transaction is not ready"))); + EXPECT_THAT(txn->NewUpdateProperties(), IsError(ErrorKind::kValidationFailed)); + EXPECT_THAT(txn->Abort(), IsOk()); + EXPECT_THAT(txn->Abort(), IsOk()); +} + +TEST_F(TransactionRetryTest, NonRetryableStandaloneUpdateStopsAtFirstConflict) { + EXPECT_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(1) + .WillOnce(::testing::Return(CommitFailed("conflict"))); + EXPECT_CALL(*mock_catalog_, LoadTable(::testing::_)).Times(0); + ICEBERG_UNWRAP_OR_FAIL(auto update, mock_table_->NewUpdateSchema()); + update->AddColumn("new_column", int64()); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kCommitFailed)); + EXPECT_THAT(update->Commit(), IsError(ErrorKind::kValidationFailed)); +} + } // namespace iceberg diff --git a/src/iceberg/test/update_test_base.h b/src/iceberg/test/update_test_base.h index 232060d2e..e0ea53c3b 100644 --- a/src/iceberg/test/update_test_base.h +++ b/src/iceberg/test/update_test_base.h @@ -20,6 +20,7 @@ #pragma once #include +#include #include #include @@ -34,6 +35,7 @@ #include "iceberg/table_identifier.h" #include "iceberg/table_metadata.h" #include "iceberg/test/matchers.h" +#include "iceberg/test/mock_catalog.h" #include "iceberg/test/test_resource.h" #include "iceberg/util/uuid.h" @@ -85,6 +87,41 @@ class UpdateTestBase : public ::testing::Test { catalog_->RegisterTable(table_ident_, metadata_location)); } + // Exercise the public commit retry path against real catalog metadata. + void FailCommits(int conflicts, std::function before_attempt = {}) { + auto mock = std::make_shared<::testing::NiceMock>(); + EXPECT_CALL(*mock, LoadTable(::testing::_)) + .Times(::testing::AtLeast(conflicts)) + .WillRepeatedly([catalog = catalog_](const TableIdentifier& name) { + return catalog->LoadTable(name); + }); + EXPECT_CALL(*mock, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .Times(conflicts + 1) + .WillRepeatedly( + [catalog = catalog_, conflicts, before_attempt, + attempt = std::make_shared(0)]( + const auto& name, const auto& requirements, + const auto& updates) mutable -> Result> { + if (before_attempt) { + before_attempt(*attempt); + } + if ((*attempt)++ < conflicts) { + return CommitFailed("injected conflict"); + } + auto result = catalog->UpdateTable(name, requirements, updates); + if (!result) { + ADD_FAILURE() << result.error().message; + return ValidationFailed("Test catalog failed: {}", + result.error().message); + } + return result; + }); + ICEBERG_UNWRAP_OR_FAIL( + table_, Table::Make(table_->name(), table_->metadata(), + std::string(table_->metadata_file_location()), file_io_, mock, + table_->full_name(), table_->reporter())); + } + /// \brief Reload the table from catalog and return its metadata. std::shared_ptr ReloadMetadata() { auto result = catalog_->LoadTable(table_ident_); diff --git a/src/iceberg/transaction.cc b/src/iceberg/transaction.cc index bacd880de..f5042f59a 100644 --- a/src/iceberg/transaction.cc +++ b/src/iceberg/transaction.cc @@ -53,6 +53,7 @@ #include "iceberg/update/update_sort_order.h" #include "iceberg/update/update_statistics.h" #include "iceberg/util/checked_cast.h" +#include "iceberg/util/error_util_internal.h" #include "iceberg/util/location_util.h" #include "iceberg/util/macros.h" #include "iceberg/util/retry_util.h" @@ -123,9 +124,7 @@ Result> Transaction::Make(std::shared_ptr
ta Result> Transaction::Make( std::shared_ptr ctx) { ICEBERG_PRECHECK(ctx != nullptr, "TransactionContext cannot be null"); - auto txn = std::shared_ptr(new Transaction(ctx)); - ctx->transaction = std::weak_ptr(txn); - return txn; + return std::shared_ptr(new Transaction(std::move(ctx))); } const std::shared_ptr
& Transaction::table() const { return ctx_->table; } @@ -138,16 +137,56 @@ std::string Transaction::MetadataFileLocation(std::string_view filename) const { return ctx_->MetadataFileLocation(filename); } -Status Transaction::AddUpdate(const std::shared_ptr& update) { - ICEBERG_CHECK(last_update_committed_, - "Cannot add update when previous update is not committed"); +Status Transaction::CheckReady() const { + ICEBERG_CHECK(state_ == TransactionState::kReady, "Transaction is not ready (state {})", + static_cast(state_)); + return {}; +} +Status Transaction::AddUpdate(const std::shared_ptr& update) { + ICEBERG_RETURN_UNEXPECTED(CheckReady()); + ICEBERG_PRECHECK(update && update->ctx_.get() == ctx_.get(), + "Update must belong to this transaction context"); + ICEBERG_CHECK(!update->commit_called_, "Update has already been committed"); pending_updates_.push_back(update); - last_update_committed_ = false; + state_ = TransactionState::kUpdatePending; + return {}; +} + +Status Transaction::CommitUpdate(PendingUpdate& update) { + ICEBERG_CHECK(state_ == TransactionState::kUpdatePending, + "Transaction has no pending operation (state {})", + static_cast(state_)); + ICEBERG_CHECK(!pending_updates_.empty() && pending_updates_.back().get() == &update, + "Update is not the current pending operation"); + update.commit_called_ = true; + Status status; + try { + status = ApplyUpdate(update); + } catch (const std::exception& e) { + status = ValidationFailed("Update Apply threw: {}", e.what()); + } catch (...) { + status = ValidationFailed("Update Apply threw an unknown exception"); + } + if (!status) { + state_ = TransactionState::kFailed; + CleanupUpdates(); + return status; + } + state_ = TransactionState::kReady; return {}; } -Status Transaction::Apply(PendingUpdate& update) { +Status Transaction::ReplayUpdates() { + ICEBERG_CHECK(state_ == TransactionState::kReady, + "Replay requires a ready transaction"); + for (const auto& update : pending_updates_) { + ICEBERG_RETURN_UNEXPECTED(ApplyUpdate(*update)); + } + return {}; +} + +Status Transaction::ApplyUpdate(PendingUpdate& update) { switch (update.kind()) { case PendingUpdate::Kind::kExpireSnapshots: ICEBERG_RETURN_UNEXPECTED( @@ -198,9 +237,7 @@ Status Transaction::Apply(PendingUpdate& update) { static_cast(update.kind())); } - last_update_committed_ = true; - - return {}; + return ctx_->metadata_builder->CheckErrors(); } Status Transaction::ApplyExpireSnapshots(ExpireSnapshots& update) { @@ -294,9 +331,10 @@ Status Transaction::ApplyUpdateSnapshot(SnapshotUpdate& update) { ICEBERG_RETURN_UNEXPECTED(temp_update->CheckErrors()); if (temp_update->changes().empty()) { - // Do not commit if the metadata has not changed. for example, this may happen - // when setting the current snapshot to an ID that is already current. note that - // this check uses identity. + // A no-op may still write temporary files. Clean them now; the update remains + // registered so it can be replayed after a refresh. + internal::LogAndIgnoreFailure("Update staging cleanup", + [&update] { return update.CleanStaged(); }); return {}; } @@ -358,78 +396,124 @@ Status Transaction::ApplyUpdatePartitionStatistics(UpdatePartitionStatistics& up } Result> Transaction::Commit() { - ICEBERG_CHECK(!committed_, "Transaction already committed"); - ICEBERG_CHECK(last_update_committed_, - "Cannot commit transaction when previous update is not committed"); - - const auto& updates = ctx_->metadata_builder->changes(); - if (updates.empty()) { - committed_ = true; - return ctx_->table; - } - - const auto& props = ctx_->table->properties(); - int32_t num_retries = - CanRetry() ? static_cast(props.Get(TableProperties::kCommitNumRetries)) - : 0; - int32_t min_wait_ms = props.Get(TableProperties::kCommitMinRetryWaitMs); - int32_t max_wait_ms = props.Get(TableProperties::kCommitMaxRetryWaitMs); - int32_t total_timeout_ms = props.Get(TableProperties::kCommitTotalRetryTimeMs); - - bool is_first_attempt = true; - auto commit_result = - MakeCommitRetryRunner(num_retries, min_wait_ms, max_wait_ms, total_timeout_ms) - .Run([this, &is_first_attempt]() -> Result> { - auto result = CommitOnce(is_first_attempt); - is_first_attempt = false; - return result; - }); - - Result finalize_result = - commit_result.has_value() - ? Result(commit_result.value()->metadata().get()) - : std::unexpected(commit_result.error()); - - for (const auto& update : pending_updates_) { - std::ignore = update->Finalize(finalize_result); + ICEBERG_RETURN_UNEXPECTED(CheckReady()); + Result> commit_result = ctx_->table; + try { + auto builder_status = ctx_->metadata_builder->CheckErrors(); + if (!builder_status) { + commit_result = std::unexpected(builder_status.error()); + } else { + const auto& props = ctx_->table->properties(); + const int32_t num_retries = + CanRetry() ? static_cast(props.Get(TableProperties::kCommitNumRetries)) + : 0; + bool is_first_attempt = true; + commit_result = + MakeCommitRetryRunner(num_retries, + props.Get(TableProperties::kCommitMinRetryWaitMs), + props.Get(TableProperties::kCommitMaxRetryWaitMs), + props.Get(TableProperties::kCommitTotalRetryTimeMs)) + .Run([this, &is_first_attempt]() -> Result> { + auto result = CommitOnce(is_first_attempt); + is_first_attempt = false; + return result; + }); + } + } catch (const std::exception& e) { + // CommitOnce handles catalog exceptions, so this failed before the commit. + commit_result = ValidationFailed("Transaction preparation threw: {}", e.what()); + } catch (...) { + commit_result = + ValidationFailed("Transaction preparation threw an unknown exception"); } - ICEBERG_RETURN_UNEXPECTED(commit_result); + if (!commit_result) { + if (commit_result.error().kind == ErrorKind::kCommitStateUnknown) { + state_ = TransactionState::kCommitStateUnknown; + } else { + state_ = TransactionState::kFailed; + CleanupUpdates(); + } + return commit_result; + } - // Mark as committed and update table reference - committed_ = true; ctx_->table = std::move(commit_result.value()); - + state_ = TransactionState::kCommitted; + FinalizeUpdates(*ctx_->table->metadata()); return ctx_->table; } +Status Transaction::Abort() { + if (state_ == TransactionState::kAborted) { + return {}; + } + ICEBERG_CHECK(state_ == TransactionState::kReady || + state_ == TransactionState::kUpdatePending || + state_ == TransactionState::kFailed, + "Cannot abort a committed or unknown transaction"); + state_ = TransactionState::kAborted; + CleanupUpdates(); + return {}; +} + +void Transaction::CleanupUpdates() noexcept { + for (const auto& update : pending_updates_) { + internal::LogAndIgnoreFailure("Update staging cleanup", + [&update] { return update->CleanStaged(); }); + } +} + +void Transaction::FinalizeUpdates(const TableMetadata& committed) noexcept { + for (const auto& update : pending_updates_) { + internal::LogAndIgnoreFailure("Update finalization", [&update, &committed] { + return update->Finalize(committed); + }); + } +} + Result> Transaction::CommitOnce(bool is_first_attempt) { std::vector> requirements; - - switch (ctx_->kind) { - case TransactionKind::kCreate: { - ICEBERG_ASSIGN_OR_RAISE(requirements, TableRequirements::ForCreateTable( - ctx_->metadata_builder->changes())); - } break; - case TransactionKind::kUpdate: { - if (!is_first_attempt) { - ICEBERG_RETURN_UNEXPECTED(ctx_->table->Refresh()); - } - if (ctx_->metadata_builder->base() != ctx_->table->metadata().get()) { - ctx_->metadata_builder = - TableMetadataBuilder::BuildFrom(ctx_->table->metadata().get()); - for (const auto& update : pending_updates_) { - ICEBERG_RETURN_UNEXPECTED(update->Commit()); - } + if (ctx_->kind == TransactionKind::kUpdate) { + std::shared_ptr metadata_before_refresh; + if (!is_first_attempt) { + // Keep the builder's base alive while Refresh replaces the table metadata. + metadata_before_refresh = ctx_->table->metadata(); + ICEBERG_RETURN_UNEXPECTED(ctx_->table->Refresh()); + } + const bool metadata_changed = + ctx_->metadata_builder->base() != ctx_->table->metadata().get(); + const bool standalone_retry = !is_first_attempt && !ctx_->transaction.has_value(); + if (metadata_changed || standalone_retry) { + ICEBERG_CHECK(CanRetry(), + "Cannot rebase a transaction containing a non-retryable update"); + CleanupUpdates(); + ctx_->metadata_builder = + TableMetadataBuilder::BuildFrom(ctx_->table->metadata().get()); + auto applied = ReplayUpdates(); + if (!applied) { + return ValidationFailed("Transaction replay failed: {}", applied.error().message); } - ICEBERG_ASSIGN_OR_RAISE(requirements, TableRequirements::ForUpdateTable( - *ctx_->metadata_builder->base(), - ctx_->metadata_builder->changes())); - } break; + } + if (ctx_->metadata_builder->changes().empty()) { + return ctx_->table; + } + ICEBERG_ASSIGN_OR_RAISE(requirements, TableRequirements::ForUpdateTable( + *ctx_->metadata_builder->base(), + ctx_->metadata_builder->changes())); + } else { + ICEBERG_ASSIGN_OR_RAISE(requirements, TableRequirements::ForCreateTable( + ctx_->metadata_builder->changes())); } - return ctx_->table->catalog()->UpdateTable(ctx_->table->name(), requirements, - ctx_->metadata_builder->changes()); + // UpdateTable may throw after the catalog has committed the changes. + try { + return ctx_->table->catalog()->UpdateTable(ctx_->table->name(), requirements, + ctx_->metadata_builder->changes()); + } catch (const std::exception& e) { + return CommitStateUnknown("Catalog commit threw: {}", e.what()); + } catch (...) { + return CommitStateUnknown("Catalog commit threw an unknown exception"); + } } bool Transaction::CanRetry() const { diff --git a/src/iceberg/transaction.h b/src/iceberg/transaction.h index 3ee0372ad..75e3c831e 100644 --- a/src/iceberg/transaction.h +++ b/src/iceberg/transaction.h @@ -39,6 +39,16 @@ namespace iceberg { /// \brief Whether a transaction creates a new table or updates an existing one. enum class TransactionKind : uint8_t { kCreate, kUpdate }; +/// \brief Lifecycle outcomes of a transaction. Failed transactions may only be aborted. +enum class TransactionState : uint8_t { + kReady, + kUpdatePending, + kCommitted, + kFailed, + kAborted, + kCommitStateUnknown, +}; + /// \brief A transaction for performing multiple updates to a table class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this { public: @@ -50,8 +60,8 @@ class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this> Make( std::shared_ptr ctx); @@ -77,6 +87,14 @@ class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this> Commit(); + /// \brief Discard staged changes and clean owned files best effort. + /// Repeated Abort succeeds without repeating cleanup. Committed and unknown + /// transactions cannot be aborted. Destructors never perform cleanup. + Status Abort(); + + /// \brief Return the current state of the transaction. + TransactionState state() const { return state_; } + /// \brief Create a new UpdatePartitionSpec to update the partition spec of this table /// and commit the changes. Result> NewUpdatePartitionSpec(); @@ -146,10 +164,12 @@ class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this ctx); + Status CheckReady() const; Status AddUpdate(const std::shared_ptr& update); - /// \brief Apply the pending changes to current table. - Status Apply(PendingUpdate& updates); + Status CommitUpdate(PendingUpdate& update); + Status ApplyUpdate(PendingUpdate& update); + Status ReplayUpdates(); // Helper methods for applying different types of updates Status ApplyExpireSnapshots(ExpireSnapshots& update); @@ -164,12 +184,14 @@ class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this> CommitOnce(bool is_first_attempt); /// \brief Whether this transaction can retry after a commit conflict. bool CanRetry() const; + void CleanupUpdates() noexcept; + void FinalizeUpdates(const TableMetadata& committed) noexcept; + private: friend class PendingUpdate; @@ -178,10 +200,7 @@ class ICEBERG_EXPORT Transaction : public std::enable_shared_from_this> pending_updates_; - // To make the state simple, we require updates are added and committed in order. - bool last_update_committed_ = true; - // Tracks if transaction has been committed to prevent double-commit - bool committed_ = false; + TransactionState state_ = TransactionState::kReady; }; /// \brief Shared context between Transaction and PendingUpdate instances. diff --git a/src/iceberg/update/delete_files.cc b/src/iceberg/update/delete_files.cc index 9759e3eb9..e738108de 100644 --- a/src/iceberg/update/delete_files.cc +++ b/src/iceberg/update/delete_files.cc @@ -30,11 +30,11 @@ namespace iceberg { -Result> DeleteFiles::Make( +Result> DeleteFiles::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create DeleteFiles without a context"); - return std::unique_ptr( + return std::shared_ptr( new DeleteFiles(std::move(table_name), std::move(ctx))); } diff --git a/src/iceberg/update/delete_files.h b/src/iceberg/update/delete_files.h index 7e567830e..18756ff48 100644 --- a/src/iceberg/update/delete_files.h +++ b/src/iceberg/update/delete_files.h @@ -43,7 +43,7 @@ namespace iceberg { /// differently-normalized URIs are not considered matches. class ICEBERG_EXPORT DeleteFiles : public MergingSnapshotUpdate { public: - static Result> Make( + static Result> Make( std::string table_name, std::shared_ptr ctx); /// \brief Delete a file by path from the underlying table. diff --git a/src/iceberg/update/expire_snapshots.cc b/src/iceberg/update/expire_snapshots.cc index 5573efa77..1c9e11c43 100644 --- a/src/iceberg/update/expire_snapshots.cc +++ b/src/iceberg/update/expire_snapshots.cc @@ -31,6 +31,7 @@ #include #include "iceberg/file_io.h" +#include "iceberg/logging/log_macros.h" #include "iceberg/manifest/manifest_entry.h" #include "iceberg/manifest/manifest_reader.h" #include "iceberg/result.h" @@ -115,7 +116,9 @@ class FileCleanupStrategy { group.Submit([this, paths = std::move(path_list)]() -> Status { return file_io_->DeleteFiles(paths); }); - std::ignore = std::move(group).Run(); + if (auto status = std::move(group).Run(); !status) { + ICEBERG_LOG_WARN("Expiration deletion failed: {}", status.error().message); + } return; } @@ -133,7 +136,9 @@ class FileCleanupStrategy { } }); } - std::ignore = std::move(group).Run(); + if (auto status = std::move(group).Run(); !status) { + ICEBERG_LOG_WARN("Expiration deletion failed: {}", status.error().message); + } } bool HasAnyStatisticsFiles(const TableMetadata& metadata) const { @@ -999,39 +1004,31 @@ Result ExpireSnapshots::Apply() { std::ranges::to>(); } - // Cache the result for use during Finalize() apply_result_ = result; - return result; } -Status ExpireSnapshots::Finalize(Result commit_result) { - if (!commit_result.has_value()) { - return {}; - } - - if (cleanup_level_ == CleanupLevel::kNone) { - return {}; - } - - if (!apply_result_.has_value() || apply_result_->snapshot_ids_to_remove.empty()) { +Status ExpireSnapshots::Finalize(const TableMetadata& metadata_after_expiration) { + // The cached Apply result belongs to this generation only; consume it regardless + // of whether any physical cleanup happens. + auto apply_result = std::exchange(apply_result_, std::nullopt); + if (cleanup_level_ == CleanupLevel::kNone || !apply_result.has_value() || + apply_result->snapshot_ids_to_remove.empty()) { return {}; } - ICEBERG_PRECHECK(apply_result_->metadata_before_expiration != nullptr, + ICEBERG_PRECHECK(apply_result->metadata_before_expiration != nullptr, "Missing pre-expiration table metadata for cleanup"); - ICEBERG_PRECHECK(commit_result.value() != nullptr, - "Missing committed table metadata for cleanup"); - auto metadata_before_expiration_ptr = apply_result_->metadata_before_expiration; - const TableMetadata& metadata_before_expiration = *metadata_before_expiration_ptr; - const TableMetadata& metadata_after_expiration = *commit_result.value(); - apply_result_.reset(); + const TableMetadata& metadata_before_expiration = + *apply_result->metadata_before_expiration; // Pick incremental cleanup when the expiration is a simple linear-ancestry walk: // no explicit snapshot IDs, no removed snapshots outside main ancestry, and no // retained snapshots outside main ancestry. + // An explicit transaction may apply later updates that add file references. + // Reachable cleanup evaluates the final committed metadata and preserves those files. const bool can_use_incremental = - !specified_snapshot_id_ && + !ctx_->transaction.has_value() && !specified_snapshot_id_ && !HasRemovedNonMainAncestors(metadata_before_expiration, metadata_after_expiration) && !HasNonMainSnapshots(metadata_after_expiration); diff --git a/src/iceberg/update/expire_snapshots.h b/src/iceberg/update/expire_snapshots.h index 215ad85fd..a2512f0e0 100644 --- a/src/iceberg/update/expire_snapshots.h +++ b/src/iceberg/update/expire_snapshots.h @@ -64,7 +64,7 @@ enum class CleanupLevel : uint8_t { /// that were deleted by snapshots that are expired will be deleted. DeleteWith() can be /// used to pass an alternative deletion method. /// -/// Apply() returns a list of the snapshots that will be removed. +/// Apply() returns the snapshots that will be removed. class ICEBERG_EXPORT ExpireSnapshots : public PendingUpdate { public: static Result> Make( @@ -158,22 +158,16 @@ class ICEBERG_EXPORT ExpireSnapshots : public PendingUpdate { Kind kind() const final { return Kind::kExpireSnapshots; } bool IsRetryable() const override { return true; } - /// \brief Apply the pending changes and return the results - /// \return The results of changes Result Apply(); - /// \brief Finalize the expire snapshots update, cleaning up expired files. - /// - /// After a successful commit, this method deletes manifest files, manifest lists, - /// data files, and statistics files that are no longer referenced by any valid - /// snapshot. The cleanup behavior is controlled by the CleanupLevel setting. - /// - /// \param commit_result The committed table metadata when the commit succeeds, or the - /// commit error when it fails. - /// \return Status indicating success or failure - Status Finalize(Result commit_result) override; - private: + friend class Transaction; + + Status Finalize(const TableMetadata& committed) override; + Status CleanStaged() override { + apply_result_.reset(); + return {}; + } explicit ExpireSnapshots(std::shared_ptr ctx); using SnapshotToRef = std::unordered_map>; diff --git a/src/iceberg/update/fast_append.cc b/src/iceberg/update/fast_append.cc index 4167387c7..eb2d1c20f 100644 --- a/src/iceberg/update/fast_append.cc +++ b/src/iceberg/update/fast_append.cc @@ -35,11 +35,11 @@ namespace iceberg { -Result> FastAppend::Make( +Result> FastAppend::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create FastAppend without a context"); - return std::unique_ptr( + return std::shared_ptr( new FastAppend(std::move(table_name), std::move(ctx))); } @@ -75,14 +75,10 @@ FastAppend& FastAppend::AppendManifest(const ManifestFile& manifest) { "Sequence number must be assigned during commit"); if (can_inherit_snapshot_id() && manifest.added_snapshot_id == kInvalidSnapshotId) { - appended_manifests_summary_.AddedManifest(manifest); append_manifests_.push_back(manifest); } else { // The manifest must be rewritten with this update's snapshot ID - ICEBERG_BUILDER_ASSIGN_OR_RETURN(auto copied_manifest, - CopyManifest(manifest, /*update_summary=*/true)); append_manifests_to_copy_.push_back(manifest); - rewritten_append_manifests_.push_back(std::move(copied_manifest)); } return *this; @@ -93,14 +89,17 @@ std::string FastAppend::operation() { return DataOperation::kAppend; } Result> FastAppend::Apply( const TableMetadata& metadata_to_update, const std::shared_ptr& snapshot) { std::vector manifests; + appended_manifests_summary_.Clear(); + for (const auto& manifest : append_manifests_) { + appended_manifests_summary_.AddedManifest(manifest); + } ICEBERG_ASSIGN_OR_RAISE(auto new_written_manifests, WriteNewManifests()); // A retry cleanup deletes copied append manifests and clears the rewritten // list; rebuild them from the original appended manifests before re-applying. if (rewritten_append_manifests_.empty() && !append_manifests_to_copy_.empty()) { for (const auto& manifest : append_manifests_to_copy_) { - ICEBERG_ASSIGN_OR_RAISE(auto copied_manifest, - CopyManifest(manifest, /*update_summary=*/false)); + ICEBERG_ASSIGN_OR_RAISE(auto copied_manifest, CopyManifest(manifest)); rewritten_append_manifests_.push_back(std::move(copied_manifest)); } } @@ -184,21 +183,11 @@ Status FastAppend::CleanUncommitted(const std::unordered_set& commi return {}; } -bool FastAppend::CleanupAfterCommit() const { - // Cleanup after committing is disabled for FastAppend unless append manifests - // were copied or need to be copied on retry because: - // 1.) Directly appended manifests are never rewritten - // 2.) Manifests which are written out as part of AppendFile are already cleaned - // up between commit attempts in WriteNewManifests - return !rewritten_append_manifests_.empty() || !append_manifests_to_copy_.empty(); -} - Result> FastAppend::Spec(int32_t spec_id) { return base().PartitionSpecById(spec_id); } -Result FastAppend::CopyManifest(const ManifestFile& manifest, - bool update_summary) { +Result FastAppend::CopyManifest(const ManifestFile& manifest) { const TableMetadata& current = base(); ICEBERG_ASSIGN_OR_RAISE(auto schema, current.Schema()); ICEBERG_ASSIGN_OR_RAISE(auto spec, @@ -211,7 +200,7 @@ Result FastAppend::CopyManifest(const ManifestFile& manifest, // Copy the manifest with the new snapshot ID. return CopyAppendManifest(manifest, ctx_->table->io(), schema, spec, snapshot_id, new_manifest_path, current.format_version, - update_summary ? &appended_manifests_summary_ : nullptr); + &appended_manifests_summary_); } Result> FastAppend::WriteNewManifests() { diff --git a/src/iceberg/update/fast_append.h b/src/iceberg/update/fast_append.h index f5e1ceaba..e20576654 100644 --- a/src/iceberg/update/fast_append.h +++ b/src/iceberg/update/fast_append.h @@ -51,7 +51,7 @@ class ICEBERG_EXPORT FastAppend : public SnapshotUpdate { /// \param table_name The name of the table /// \param ctx The transaction context to use for this update /// \return A Result containing the FastAppend instance or an error - static Result> Make( + static Result> Make( std::string table_name, std::shared_ptr ctx); /// \brief Append a DataFile to the table. @@ -73,6 +73,7 @@ class ICEBERG_EXPORT FastAppend : public SnapshotUpdate { /// \return This FastAppend for method chaining. FastAppend& AppendManifest(const ManifestFile& manifest); + protected: std::string operation() override; Result> Apply( @@ -81,7 +82,6 @@ class ICEBERG_EXPORT FastAppend : public SnapshotUpdate { std::unordered_map Summary() override; void SetSummaryProperty(const std::string& property, const std::string& value) override; Status CleanUncommitted(const std::unordered_set& committed) override; - bool CleanupAfterCommit() const override; private: explicit FastAppend(std::string table_name, std::shared_ptr ctx); @@ -92,9 +92,8 @@ class ICEBERG_EXPORT FastAppend : public SnapshotUpdate { /// \brief Copy a manifest file with a new snapshot ID. /// /// \param manifest The manifest to copy - /// \param update_summary Whether to add copied entries to the append summary /// \return The copied manifest file - Result CopyManifest(const ManifestFile& manifest, bool update_summary); + Result CopyManifest(const ManifestFile& manifest); /// \brief Write new manifests for the accumulated data files. /// diff --git a/src/iceberg/update/merge_append.cc b/src/iceberg/update/merge_append.cc index 70cd7b8e7..fd6cc35ae 100644 --- a/src/iceberg/update/merge_append.cc +++ b/src/iceberg/update/merge_append.cc @@ -29,11 +29,11 @@ namespace iceberg { -Result> MergeAppend::Make( +Result> MergeAppend::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create MergeAppend without a context"); - return std::unique_ptr( + return std::shared_ptr( new MergeAppend(std::move(table_name), std::move(ctx))); } diff --git a/src/iceberg/update/merge_append.h b/src/iceberg/update/merge_append.h index cd4f4acbb..18503b3c1 100644 --- a/src/iceberg/update/merge_append.h +++ b/src/iceberg/update/merge_append.h @@ -47,7 +47,7 @@ class ICEBERG_EXPORT MergeAppend : public MergingSnapshotUpdate { /// \param table_name The name of the table /// \param ctx The transaction context to use for this update /// \return A Result containing the MergeAppend instance or an error - static Result> Make( + static Result> Make( std::string table_name, std::shared_ptr ctx); /// \brief Append a DataFile to the table. diff --git a/src/iceberg/update/merging_snapshot_update.cc b/src/iceberg/update/merging_snapshot_update.cc index 7ce577076..48cc8dbb3 100644 --- a/src/iceberg/update/merging_snapshot_update.cc +++ b/src/iceberg/update/merging_snapshot_update.cc @@ -662,26 +662,21 @@ Status MergingSnapshotUpdate::AddManifest(ManifestFile manifest) { return InvalidArgument("Cannot append manifest with assigned first row ID: {}", manifest.manifest_path); } - appended_manifests_summary_.AddedManifest(manifest); append_manifests_.push_back(std::move(manifest)); } else { - ICEBERG_ASSIGN_OR_RAISE(auto copied, CopyManifest(manifest, /*update_summary=*/true)); append_manifests_to_copy_.push_back(std::move(manifest)); - rewritten_append_manifests_.push_back(std::move(copied)); } return {}; } -Result MergingSnapshotUpdate::CopyManifest(const ManifestFile& manifest, - bool update_summary) { +Result MergingSnapshotUpdate::CopyManifest(const ManifestFile& manifest) { const TableMetadata& current = base(); ICEBERG_ASSIGN_OR_RAISE(auto schema, SnapshotUtil::SchemaFor(current, target_branch())); ICEBERG_ASSIGN_OR_RAISE(auto spec, current.PartitionSpecById(manifest.partition_spec_id)); std::string path = ManifestPath(); return CopyAppendManifest(manifest, ctx_->table->io(), schema, spec, SnapshotId(), path, - current.format_version, - update_summary ? &appended_manifests_summary_ : nullptr); + current.format_version, &appended_manifests_summary_); } // ------------------------------------------------------------------------- @@ -824,6 +819,7 @@ MergingSnapshotUpdate::MergeDVs() { auto output_path = location_provider->NewDataLocation( std::format("merged-dvs-{}-{}.puffin", SnapshotId(), ++dv_merge_attempt_)); + RegisterStagedFile(output_path); auto merged_files = DVUtil::MergeAndWriteDVs(groups, output_path, ctx_->table->io()); if (!merged_files) { std::ignore = DeleteFile(output_path); @@ -854,6 +850,9 @@ MergingSnapshotUpdate::MergeDVs() { result.push_back(std::move(merged)); } + // The operation-specific merged DV cache now owns cleanup for this file. + UnregisterStagedFile(output_path); + return result; } @@ -924,6 +923,11 @@ Result> MergingSnapshotUpdate::WriteNewDeleteManifests Result> MergingSnapshotUpdate::Apply( const TableMetadata& metadata_to_update, const std::shared_ptr& snapshot) { + appended_manifests_summary_.Clear(); + for (const auto& manifest : append_manifests_) { + appended_manifests_summary_.AddedManifest(manifest); + } + ICEBERG_RETURN_UNEXPECTED(ManagersReady()); // Re-validate buffered delete files against the current format version. A format @@ -988,8 +992,7 @@ Result> MergingSnapshotUpdate::Apply( ICEBERG_ASSIGN_OR_RAISE(auto written_data_manifests, WriteNewDataManifests()); if (rewritten_append_manifests_.empty() && !append_manifests_to_copy_.empty()) { for (const auto& manifest : append_manifests_to_copy_) { - ICEBERG_ASSIGN_OR_RAISE(auto copied, CopyManifest(manifest, - /*update_summary=*/false)); + ICEBERG_ASSIGN_OR_RAISE(auto copied, CopyManifest(manifest)); rewritten_append_manifests_.push_back(std::move(copied)); } } diff --git a/src/iceberg/update/merging_snapshot_update.h b/src/iceberg/update/merging_snapshot_update.h index 65c3f6d3f..e1cd7303b 100644 --- a/src/iceberg/update/merging_snapshot_update.h +++ b/src/iceberg/update/merging_snapshot_update.h @@ -331,8 +331,7 @@ class ICEBERG_EXPORT MergingSnapshotUpdate : public SnapshotUpdate { /// \brief Copy a manifest with the current snapshot ID, for use when snapshot /// ID inheritance is not possible. - /// \param update_summary Whether to add copied entries to the append summary - Result CopyManifest(const ManifestFile& manifest, bool update_summary); + Result CopyManifest(const ManifestFile& manifest); Status AddDeleteFile(std::shared_ptr file, std::optional data_sequence_number); diff --git a/src/iceberg/update/pending_update.cc b/src/iceberg/update/pending_update.cc index 4b3000652..88cd8907d 100644 --- a/src/iceberg/update/pending_update.cc +++ b/src/iceberg/update/pending_update.cc @@ -31,34 +31,29 @@ PendingUpdate::PendingUpdate(std::shared_ptr ctx) PendingUpdate::~PendingUpdate() = default; -Status PendingUpdate::Commit() { - if (!ctx_->transaction) { - // Table-created path: no transaction exists yet, create a temporary one. - ICEBERG_ASSIGN_OR_RAISE(auto txn, Transaction::Make(ctx_)); - auto apply_status = txn->Apply(*this); - if (!apply_status.has_value()) { - std::ignore = Finalize(std::unexpected(apply_status.error())); - return apply_status; - } - - auto commit_result = txn->Commit(); - if (!commit_result.has_value()) { - std::ignore = Finalize(std::unexpected(commit_result.error())); - return std::unexpected(commit_result.error()); - } +Status PendingUpdate::CheckCommitAllowed() const { + ICEBERG_CHECK(!commit_called_, "Update has already been committed"); + return {}; +} - std::ignore = Finalize(commit_result.value()->metadata().get()); - return {}; - } - auto txn = ctx_->transaction->lock(); - if (!txn) { - return CommitFailed("Transaction has been destroyed"); +Status PendingUpdate::Commit() { + ICEBERG_RETURN_UNEXPECTED(CheckCommitAllowed()); + if (ctx_->transaction) { + auto txn = ctx_->transaction->lock(); + ICEBERG_CHECK(txn != nullptr, "Transaction has been destroyed"); + return txn->CommitUpdate(*this); } - return txn->Apply(*this); + + auto self = weak_from_this().lock(); + ICEBERG_PRECHECK(self != nullptr, "PendingUpdate must be owned by std::shared_ptr"); + ICEBERG_ASSIGN_OR_RAISE(auto txn, Transaction::Make(ctx_)); + ICEBERG_RETURN_UNEXPECTED(txn->AddUpdate(self)); + ICEBERG_RETURN_UNEXPECTED(txn->CommitUpdate(*this)); + ICEBERG_RETURN_UNEXPECTED(txn->Commit()); + return {}; } -Status PendingUpdate::Finalize( - [[maybe_unused]] Result commit_result) { +Status PendingUpdate::Finalize([[maybe_unused]] const TableMetadata& committed) { return {}; } diff --git a/src/iceberg/update/pending_update.h b/src/iceberg/update/pending_update.h index 19998ddb3..81bf4563e 100644 --- a/src/iceberg/update/pending_update.h +++ b/src/iceberg/update/pending_update.h @@ -36,9 +36,16 @@ namespace iceberg { /// Any created `PendingUpdate` instance is tracked by the `Transaction` instance /// and commit is also delegated to the `Transaction` instance. /// +/// Lifecycle: an update is configured through its fluent mutators, then committed +/// exactly once. After that point the update is owned by its transaction until the +/// transaction reaches a terminal state. +/// /// \note Implementations are expected to use builder pattern and errors -/// should be handled by the ErrorCollector base class. -class ICEBERG_EXPORT PendingUpdate : public ErrorCollector { +/// should be handled by the ErrorCollector base class. Configuration errors are +/// collected and surfaced by `Commit()`. Callers must not mutate input objects or +/// reconfigure an update after its first commit. +class ICEBERG_EXPORT PendingUpdate : public ErrorCollector, + public std::enable_shared_from_this { public: enum class Kind : uint8_t { kExpireSnapshots, @@ -66,23 +73,14 @@ class ICEBERG_EXPORT PendingUpdate : public ErrorCollector { /// - ValidationFailed: if it cannot be applied to the current table metadata. /// - CommitFailed: if it cannot be committed due to conflicts. /// - CommitStateUnknown: unknown status, no cleanup should be done. + /// \note The update must be owned by a `std::shared_ptr` before calling Commit(). + /// An Apply failure is terminal and cannot be corrected within the same transaction. virtual Status Commit(); - /// \brief Finalize the pending update. - /// - /// This method is called after the update is committed. - /// Implementations should override this method to clean up any resources. - /// - /// \param commit_result The committed table metadata when the commit succeeds, or the - /// commit error when it fails. - /// \return Status indicating success or failure - virtual Status Finalize(Result commit_result); - - // Non-copyable, movable PendingUpdate(const PendingUpdate&) = delete; PendingUpdate& operator=(const PendingUpdate&) = delete; - PendingUpdate(PendingUpdate&&) noexcept = default; - PendingUpdate& operator=(PendingUpdate&&) noexcept = default; + PendingUpdate(PendingUpdate&&) = delete; + PendingUpdate& operator=(PendingUpdate&&) = delete; ~PendingUpdate() override; @@ -91,7 +89,20 @@ class ICEBERG_EXPORT PendingUpdate : public ErrorCollector { const TableMetadata& base() const; + Status CheckCommitAllowed() const; + + /// \brief Discard this generation's staging, retaining applied operation intent. + virtual Status CleanStaged() { return {}; } + /// \brief Complete a known successful commit. Called only by Transaction with the + /// final committed table metadata. + virtual Status Finalize(const TableMetadata& committed); + std::shared_ptr ctx_; + + private: + friend class Transaction; + + bool commit_called_ = false; }; } // namespace iceberg diff --git a/src/iceberg/update/replace_partitions.cc b/src/iceberg/update/replace_partitions.cc index 3d224375c..4487751bf 100644 --- a/src/iceberg/update/replace_partitions.cc +++ b/src/iceberg/update/replace_partitions.cc @@ -32,11 +32,11 @@ namespace iceberg { -Result> ReplacePartitions::Make( +Result> ReplacePartitions::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create ReplacePartitions without a context"); - return std::unique_ptr( + return std::shared_ptr( new ReplacePartitions(std::move(table_name), std::move(ctx))); } diff --git a/src/iceberg/update/replace_partitions.h b/src/iceberg/update/replace_partitions.h index 40048ba5c..ee57b6124 100644 --- a/src/iceberg/update/replace_partitions.h +++ b/src/iceberg/update/replace_partitions.h @@ -62,7 +62,7 @@ class ICEBERG_EXPORT ReplacePartitions : public MergingSnapshotUpdate { /// \param table_name The name of the table /// \param ctx The transaction context /// \return A Result containing the ReplacePartitions instance or an error - static Result> Make( + static Result> Make( std::string table_name, std::shared_ptr ctx); /// \brief Add a data file to the table. diff --git a/src/iceberg/update/rewrite_files.cc b/src/iceberg/update/rewrite_files.cc index b7fc048d3..3c743c01d 100644 --- a/src/iceberg/update/rewrite_files.cc +++ b/src/iceberg/update/rewrite_files.cc @@ -39,11 +39,11 @@ RewriteFiles::RewriteFiles(std::string table_name, FailMissingDeletePaths(); } -Result> RewriteFiles::Make( +Result> RewriteFiles::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create RewriteFiles without a context"); - return std::unique_ptr( + return std::shared_ptr( new RewriteFiles(std::move(table_name), std::move(ctx))); } diff --git a/src/iceberg/update/rewrite_files.h b/src/iceberg/update/rewrite_files.h index ce219c3ae..283091c13 100644 --- a/src/iceberg/update/rewrite_files.h +++ b/src/iceberg/update/rewrite_files.h @@ -54,8 +54,8 @@ class ICEBERG_EXPORT RewriteFiles : public MergingSnapshotUpdate { /// /// \param table_name The name of the table /// \param ctx The transaction context - /// \return A unique pointer to the new RewriteFiles operation - static Result> Make( + /// \return A shared pointer to the new RewriteFiles operation + static Result> Make( std::string table_name, std::shared_ptr ctx); ~RewriteFiles() override = default; diff --git a/src/iceberg/update/row_delta.cc b/src/iceberg/update/row_delta.cc index dd3f50c58..f239ea495 100644 --- a/src/iceberg/update/row_delta.cc +++ b/src/iceberg/update/row_delta.cc @@ -38,11 +38,11 @@ namespace iceberg { -Result> RowDelta::Make( +Result> RowDelta::Make( std::string table_name, std::shared_ptr ctx) { ICEBERG_PRECHECK(!table_name.empty(), "Table name cannot be empty"); ICEBERG_PRECHECK(ctx != nullptr, "Cannot create RowDelta without a context"); - return std::unique_ptr(new RowDelta(std::move(table_name), std::move(ctx))); + return std::shared_ptr(new RowDelta(std::move(table_name), std::move(ctx))); } RowDelta::RowDelta(std::string table_name, std::shared_ptr ctx) diff --git a/src/iceberg/update/row_delta.h b/src/iceberg/update/row_delta.h index bf58a048c..f6faf8c92 100644 --- a/src/iceberg/update/row_delta.h +++ b/src/iceberg/update/row_delta.h @@ -47,7 +47,7 @@ namespace iceberg { class ICEBERG_EXPORT RowDelta : public MergingSnapshotUpdate { public: /// \brief Create a new RowDelta instance. - static Result> Make(std::string table_name, + static Result> Make(std::string table_name, std::shared_ptr ctx); /// \brief Add a DataFile to the table. diff --git a/src/iceberg/update/snapshot_manager.cc b/src/iceberg/update/snapshot_manager.cc index 5473f3033..405b029f8 100644 --- a/src/iceberg/update/snapshot_manager.cc +++ b/src/iceberg/update/snapshot_manager.cc @@ -181,6 +181,10 @@ SnapshotManager& SnapshotManager::SetMaxRefAgeMs(const std::string& name, Status SnapshotManager::Commit() { ICEBERG_RETURN_UNEXPECTED(CheckErrors()); + const auto state = transaction_->state(); + ICEBERG_CHECK( + state == TransactionState::kReady || state == TransactionState::kUpdatePending, + "Transaction is terminal"); ICEBERG_RETURN_UNEXPECTED(CommitIfRefUpdatesExist()); if (!is_external_transaction_) { ICEBERG_RETURN_UNEXPECTED(transaction_->Commit()); diff --git a/src/iceberg/update/snapshot_update.cc b/src/iceberg/update/snapshot_update.cc index 6e8d28bd5..0d168b7f3 100644 --- a/src/iceberg/update/snapshot_update.cc +++ b/src/iceberg/update/snapshot_update.cc @@ -26,6 +26,7 @@ #include "iceberg/constants.h" #include "iceberg/file_io.h" +#include "iceberg/logging/log_macros.h" #include "iceberg/manifest/manifest_entry.h" #include "iceberg/manifest/manifest_list.h" #include "iceberg/manifest/manifest_reader.h" @@ -37,6 +38,7 @@ #include "iceberg/partition_summary_internal.h" #include "iceberg/table.h" // IWYU pragma: keep #include "iceberg/transaction.h" +#include "iceberg/util/error_util_internal.h" #include "iceberg/util/executor_util_internal.h" #include "iceberg/util/macros.h" #include "iceberg/util/snapshot_util_internal.h" @@ -202,17 +204,17 @@ SnapshotUpdate::SnapshotUpdate(std::shared_ptr ctx) reporter_(ctx_->table->reporter()) {} Status SnapshotUpdate::Commit() { - commit_metrics_->attempts->Increment(); + ICEBERG_RETURN_UNEXPECTED(CheckCommitAllowed()); [[maybe_unused]] auto commit_timer = commit_metrics_->total_duration->Start(); return PendingUpdate::Commit(); } -void SnapshotUpdate::ReportCommit() const { +Status SnapshotUpdate::ReportCommit() const { ICEBERG_DCHECK(staged_snapshot_ != nullptr, "Staged snapshot is null after a successful commit"); if (!reporter_) { - return; + return {}; } const auto operation = staged_snapshot_->Operation(); @@ -225,7 +227,7 @@ void SnapshotUpdate::ReportCommit() const { CommitMetricsResult::From(*commit_metrics_, staged_snapshot_->summary), .metadata = {}, }; - std::ignore = reporter_->Report(report); + return reporter_->Report(report); } void SnapshotUpdate::SetSummaryProperty(const std::string& property, @@ -302,19 +304,9 @@ int64_t SnapshotUpdate::SnapshotId() { } Result SnapshotUpdate::Apply() { + commit_metrics_->attempts->Increment(); ICEBERG_RETURN_UNEXPECTED(CheckErrors()); - if (staged_snapshot_ != nullptr) { - for (const auto& manifest_list : manifest_lists_) { - std::ignore = DeleteFile(manifest_list); - } - manifest_lists_.clear(); - ICEBERG_RETURN_UNEXPECTED(CleanUncommitted(std::unordered_set{})); - - staged_snapshot_ = nullptr; - summary_.Clear(); - } - ICEBERG_ASSIGN_OR_RAISE(auto parent_snapshot, SnapshotUtil::OptionalLatestSnapshot(base(), target_branch_)); @@ -338,7 +330,6 @@ Result SnapshotUpdate::Apply() { ICEBERG_RETURN_UNEXPECTED(std::move(metadata_tasks).Run()); std::string manifest_list_path = ManifestListPath(); - manifest_lists_.push_back(manifest_list_path); ICEBERG_ASSIGN_OR_RAISE( auto writer, ManifestListWriter::MakeWriter(base().format_version, SnapshotId(), parent_snapshot_id, manifest_list_path, @@ -388,39 +379,46 @@ Result SnapshotUpdate::Apply() { .stage_only = stage_only_}; } -Status SnapshotUpdate::Finalize(Result commit_result) { - if (!commit_result.has_value()) { - if (commit_result.error().kind == ErrorKind::kCommitStateUnknown) { - return {}; - } - std::ignore = CleanAll(); +Status SnapshotUpdate::Finalize([[maybe_unused]] const TableMetadata& metadata) { + if (staged_snapshot_ == nullptr) { return {}; } - if (CleanupAfterCommit()) { - ICEBERG_CHECK(staged_snapshot_ != nullptr, - "Staged snapshot is null during finalize after commit"); + // Only files created by this update are tracked, and committed paths are kept below. + internal::LogAndIgnoreFailure("Snapshot cleanup", [this]() -> Status { auto cached_snapshot = SnapshotCache(staged_snapshot_.get()); - if (auto manifests = cached_snapshot.Manifests(ctx_->table->io()); - manifests.has_value()) { - std::ignore = CleanUncommitted(manifests.value() | - std::views::transform([](const auto& manifest) { - return manifest.manifest_path; - }) | - std::ranges::to>()); + ICEBERG_ASSIGN_OR_RAISE(auto manifests, cached_snapshot.Manifests(ctx_->table->io())); + auto committed = manifests | std::views::transform([](const auto& manifest) { + return manifest.manifest_path; + }) | + std::ranges::to>(); + internal::LogAndIgnoreFailure("Snapshot operation cleanup", [this, &committed] { + return CleanUncommitted(committed); + }); + committed.insert(staged_snapshot_->manifest_list); + std::vector unused; + { + std::lock_guard lock(staging_mutex_); + for (const auto& path : staged_files_) { + if (!committed.contains(path)) { + unused.push_back(path); + } + } } - } - - // Also clean up unused manifest lists created by multiple attempts - for (const auto& manifest_list : manifest_lists_) { - if (manifest_list != staged_snapshot_->manifest_list) { - std::ignore = DeleteFile(manifest_list); + for (const auto& path : unused) { + std::ignore = DeleteFile(path); } - } - - ReportCommit(); + return {}; + }); - return {}; + { + std::lock_guard lock(staging_mutex_); + staged_files_.clear(); + } + auto report_status = ReportCommit(); + staged_snapshot_.reset(); + summary_.Clear(); + return report_status; } Result> SnapshotUpdate::ComputeSummary( @@ -473,20 +471,51 @@ Result> SnapshotUpdate::ComputeSumm return summary; } -Status SnapshotUpdate::CleanAll() { - for (const auto& manifest_list : manifest_lists_) { - std::ignore = DeleteFile(manifest_list); +Status SnapshotUpdate::CleanStaged() { + internal::LogAndIgnoreFailure("Snapshot staging cleanup", + [this] { return CleanUncommitted({}); }); + // Retry paths that operation-specific cleanup did not delete. + std::vector paths; + { + std::lock_guard lock(staging_mutex_); + paths.assign(staged_files_.begin(), staged_files_.end()); + } + for (const auto& path : paths) { + std::ignore = DeleteFile(path); } - manifest_lists_.clear(); - std::ignore = CleanUncommitted(std::unordered_set{}); + staged_snapshot_.reset(); + summary_.Clear(); return {}; } -Status SnapshotUpdate::DeleteFile(const std::string& path) { - if (delete_func_) { - return delete_func_(path); +void SnapshotUpdate::RegisterStagedFile(const std::string& path) { + std::lock_guard lock(staging_mutex_); + staged_files_.insert(path); +} + +void SnapshotUpdate::UnregisterStagedFile(const std::string& path) { + std::lock_guard lock(staging_mutex_); + staged_files_.erase(path); +} + +Status SnapshotUpdate::DeleteFile(const std::string& path) noexcept { + try { + auto result = delete_func_ ? delete_func_(path) : ctx_->table->io()->DeleteFile(path); + if (!result && result.error().kind != ErrorKind::kNotFound) { + RegisterStagedFile(path); + ICEBERG_LOG_WARN("Cannot clean staged file {}: {}", path, result.error().message); + return {}; + } + std::lock_guard lock(staging_mutex_); + staged_files_.erase(path); + } catch (const std::exception& e) { + RegisterStagedFile(path); + ICEBERG_LOG_WARN("Cannot clean staged file {}: {}", path, e.what()); + } catch (...) { + RegisterStagedFile(path); + ICEBERG_LOG_WARN("Cannot clean staged file {}: unknown exception", path); } - return ctx_->table->io()->DeleteFile(path); + return {}; } std::string SnapshotUpdate::ManifestListPath() { @@ -496,7 +525,9 @@ std::string SnapshotUpdate::ManifestListPath() { auto attempt = attempt_.fetch_add(1, std::memory_order_relaxed) + 1; std::string filename = std::format("snap-{}-{}-{}.avro", snapshot_id, attempt, commit_uuid_); - return ctx_->MetadataFileLocation(filename); + auto path = ctx_->MetadataFileLocation(filename); + RegisterStagedFile(path); + return path; } SnapshotSummaryBuilder SnapshotUpdate::BuildManifestCountSummary( @@ -526,7 +557,9 @@ std::string SnapshotUpdate::ManifestPath() { // Format: {metadata_location}/{uuid}-m{manifest_count}.avro auto manifest_count = manifest_count_.fetch_add(1, std::memory_order_relaxed); std::string filename = std::format("{}-m{}.avro", commit_uuid_, manifest_count); - return ctx_->MetadataFileLocation(filename); + auto path = ctx_->MetadataFileLocation(filename); + RegisterStagedFile(path); + return path; } } // namespace iceberg diff --git a/src/iceberg/update/snapshot_update.h b/src/iceberg/update/snapshot_update.h index 1812b7072..3537bd7d9 100644 --- a/src/iceberg/update/snapshot_update.h +++ b/src/iceberg/update/snapshot_update.h @@ -25,6 +25,7 @@ #include #include #include +#include #include #include #include @@ -156,20 +157,14 @@ class ICEBERG_EXPORT SnapshotUpdate : public PendingUpdate { return self; } - /// \brief Apply the update's changes to create a new snapshot. - /// - /// This method validates the changes, applies them to the current base - /// metadata, and creates a new snapshot without committing it. Commit retries - /// call Apply() again with refreshed metadata so the same changes can be - /// applied to the new latest snapshot. - /// - /// \return A result containing the new snapshot, or an error. - Result Apply(); + protected: + friend class Transaction; - /// \brief Finalize the snapshot update, cleaning up any uncommitted files. - Status Finalize(Result commit_result) override; + /// Build snapshot state for the transaction's current metadata. + Result Apply(); + Status Finalize(const TableMetadata& committed) override; + Status CleanStaged() override; - protected: struct ContentFileWithSequenceNumber { std::shared_ptr file; std::optional data_sequence_number; @@ -249,19 +244,14 @@ class ICEBERG_EXPORT SnapshotUpdate : public PendingUpdate { /// retry-safe summary rebuilds. virtual void SetSummaryProperty(const std::string& property, const std::string& value); - /// \brief Check if cleanup should happen after commit - /// - /// \return True if cleanup should happen after commit - virtual bool CleanupAfterCommit() const { return true; } - /// \brief Get or generate the snapshot ID for the new snapshot. int64_t SnapshotId(); - /// \brief Delete a file at the given path. - /// - /// \param path The path of the file to delete - /// \return A status indicating the result of the deletion - Status DeleteFile(const std::string& path); + /// Best-effort delete. Failed paths remain registered for later cleanup. + Status DeleteFile(const std::string& path) noexcept; + void RegisterStagedFile(const std::string& path); + /// Transfer cleanup ownership to the derived update. + void UnregisterStagedFile(const std::string& path); std::string ManifestPath(); std::string ManifestListPath(); @@ -274,11 +264,7 @@ class ICEBERG_EXPORT SnapshotUpdate : public PendingUpdate { Result> ComputeSummary( const TableMetadata& previous); - /// \brief Clean up all uncommitted files - Status CleanAll(); - - /// \brief Report metrics for the most recently staged snapshot. - void ReportCommit() const; + Status ReportCommit() const; protected: SnapshotSummaryBuilder summary_; @@ -290,7 +276,9 @@ class ICEBERG_EXPORT SnapshotUpdate : public PendingUpdate { int32_t write_manifest_parallelism_{1}; std::atomic manifest_count_{0}; std::atomic attempt_{0}; - std::vector manifest_lists_; + // All paths are registered before writes, including partial/failed writes. + std::mutex staging_mutex_; + std::unordered_set staged_files_; const int64_t target_manifest_size_bytes_; std::optional snapshot_id_; OptionalExecutor plan_executor_; diff --git a/src/iceberg/update/update_properties.cc b/src/iceberg/update/update_properties.cc index d2e75c301..c97175569 100644 --- a/src/iceberg/update/update_properties.cc +++ b/src/iceberg/update/update_properties.cc @@ -68,6 +68,8 @@ UpdateProperties& UpdateProperties::Remove(const std::string& key) { Result UpdateProperties::Apply() { ICEBERG_RETURN_UNEXPECTED(CheckErrors()); + auto updates = updates_; + std::optional format_version; const auto& current_props = base().properties.configs(); std::unordered_map new_properties; std::vector removals; @@ -91,17 +93,18 @@ Result UpdateProperties::Apply() { "Cannot upgrade table to unsupported format version: v{} (supported: v{})", parsed_version, TableMetadata::kSupportedTableFormatVersion); } - format_version_ = static_cast(parsed_version); + format_version = static_cast(parsed_version); - updates_.erase(TableProperties::kFormatVersion.key()); + updates.erase(TableProperties::kFormatVersion.key()); } if (auto schema = base().Schema(); schema.has_value()) { ICEBERG_RETURN_UNEXPECTED( MetricsConfig::VerifyReferencedColumns(new_properties, *schema.value())); } - return ApplyResult{ - .updates = updates_, .removals = removals_, .format_version = format_version_}; + return ApplyResult{.updates = std::move(updates), + .removals = removals_, + .format_version = format_version}; } } // namespace iceberg diff --git a/src/iceberg/update/update_properties.h b/src/iceberg/update/update_properties.h index 18eba427b..67032f072 100644 --- a/src/iceberg/update/update_properties.h +++ b/src/iceberg/update/update_properties.h @@ -77,7 +77,6 @@ class ICEBERG_EXPORT UpdateProperties : public PendingUpdate { std::unordered_map updates_; std::unordered_set removals_; - std::optional format_version_; }; } // namespace iceberg diff --git a/src/iceberg/update/update_schema.cc b/src/iceberg/update/update_schema.cc index 56e167e10..1d8130fba 100644 --- a/src/iceberg/update/update_schema.cc +++ b/src/iceberg/update/update_schema.cc @@ -654,9 +654,9 @@ Result UpdateSchema::Apply() { } auto new_fields = temp_schema->fields() | std::ranges::to>(); - ICEBERG_ASSIGN_OR_RAISE( - auto new_schema, - Schema::Make(std::move(new_fields), schema_->schema_id(), fresh_identifier_ids)); + ICEBERG_ASSIGN_OR_RAISE(auto new_schema, + Schema::Make(std::move(new_fields), schema_->schema_id(), + std::move(fresh_identifier_ids))); ICEBERG_RETURN_UNEXPECTED(new_schema->Validate(base().format_version)); std::unordered_map updated_props; diff --git a/src/iceberg/util/error_util_internal.h b/src/iceberg/util/error_util_internal.h new file mode 100644 index 000000000..8ecfd3ce1 --- /dev/null +++ b/src/iceberg/util/error_util_internal.h @@ -0,0 +1,42 @@ +/* + * 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 "iceberg/logging/log_macros.h" + +namespace iceberg::internal { + +template +void LogAndIgnoreFailure(std::string_view name, Action&& action) noexcept { + try { + if (auto result = action(); !result) { + ICEBERG_LOG_WARN("{} failed: {}", name, result.error().message); + } + } catch (const std::exception& e) { + ICEBERG_LOG_WARN("{} threw: {}", name, e.what()); + } catch (...) { + ICEBERG_LOG_WARN("{} threw an unknown exception", name); + } +} + +} // namespace iceberg::internal