diff --git a/src/iceberg/test/fast_append_test.cc b/src/iceberg/test/fast_append_test.cc index f88d2e011..120917dd1 100644 --- a/src/iceberg/test/fast_append_test.cc +++ b/src/iceberg/test/fast_append_test.cc @@ -34,6 +34,7 @@ #include "iceberg/avro/avro_register.h" #include "iceberg/constants.h" +#include "iceberg/logging/log_level.h" #include "iceberg/manifest/manifest_entry.h" #include "iceberg/manifest/manifest_reader.h" #include "iceberg/manifest/manifest_writer.h" @@ -46,6 +47,7 @@ #include "iceberg/table_metadata.h" #include "iceberg/table_properties.h" #include "iceberg/test/executor.h" +#include "iceberg/test/logging_test_helpers.h" #include "iceberg/test/matchers.h" #include "iceberg/test/mock_catalog.h" #include "iceberg/test/update_test_base.h" @@ -178,6 +180,31 @@ TEST_F(FastAppendTest, AppendDataFile) { EXPECT_EQ(snapshot->summary.at(SnapshotSummaryFields::kManifestsReplaced), "0"); } +TEST_F(FastAppendTest, StageOnlyCommitLogNamesAddedSnapshot) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + + ICEBERG_UNWRAP_OR_FAIL(auto fast_append, table_->NewFastAppend()); + fast_append->StageOnly(); + fast_append->AppendFile(file_a_); + ASSERT_THAT(fast_append->Commit(), IsOk()); + ASSERT_THAT(table_->Refresh(), IsOk()); + + ASSERT_FALSE(table_->metadata()->snapshots.empty()); + const auto snapshot_id = table_->metadata()->snapshots.back()->snapshot_id; + bool found = false; + for (const auto& record : capturing->records()) { + if (record.level == LogLevel::kInfo && + record.message.find(std::format("committed snapshot {}", snapshot_id)) != + std::string::npos) { + found = true; + break; + } + } + EXPECT_TRUE(found) << "expected the staged snapshot in the commit success log"; +} + TEST_F(FastAppendTest, AppendMultipleDataFiles) { std::shared_ptr fast_append; ICEBERG_UNWRAP_OR_FAIL(fast_append, table_->NewFastAppend()); diff --git a/src/iceberg/test/transaction_test.cc b/src/iceberg/test/transaction_test.cc index 3a13b7bc5..a8a9f80ca 100644 --- a/src/iceberg/test/transaction_test.cc +++ b/src/iceberg/test/transaction_test.cc @@ -21,7 +21,11 @@ #include "iceberg/expression/expressions.h" #include "iceberg/expression/term.h" +#include "iceberg/logging/log_level.h" +#include "iceberg/snapshot.h" #include "iceberg/sort_order.h" +#include "iceberg/table_metadata.h" +#include "iceberg/test/logging_test_helpers.h" #include "iceberg/test/matchers.h" #include "iceberg/test/mock_catalog.h" #include "iceberg/test/update_test_base.h" @@ -173,7 +177,181 @@ TEST_F(TransactionRetryTest, CommitRetryExhausted) { EXPECT_EQ(update_call_count, 5); } -TEST_F(TransactionRetryTest, CommitNonRetryableErrorStopsImmediately) { +namespace { +// True if any captured record has the given level and a message containing `needle`. +bool HasRecord(const std::vector& records, LogLevel level, + std::string_view needle) { + for (const auto& record : records) { + if (record.level == level && record.message.find(needle) != std::string::npos) { + return true; + } + } + return false; +} +} // namespace + +// A commit that succeeds after one retryable conflict emits a WARN for the retry +// (carrying the prior error) and an INFO for the eventual success. +TEST_F(TransactionRetryTest, CommitRetryEmitsRetryAndSuccessLogs) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + + int update_call_count = 0; + ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault([this, &update_call_count]( + const TableIdentifier&, + const std::vector>&, + const std::vector>&) + -> Result> { + ++update_call_count; + if (update_call_count == 1) { + return CommitFailed("conflict on first attempt"); + } + return Table::Make(mock_table_->name(), mock_table_->metadata(), + std::string(mock_table_->metadata_file_location()), + mock_table_->io(), mock_catalog_); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto txn, mock_table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); + update->Set("retry.test", "value"); + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), IsOk()); + + auto records = capturing->records(); + EXPECT_TRUE( + HasRecord(records, LogLevel::kWarn, "Retrying transaction commit (attempt 2)")) + << "expected a retry WARN"; + EXPECT_TRUE(HasRecord(records, LogLevel::kWarn, "conflict on first attempt")) + << "retry WARN should carry the prior error"; + EXPECT_TRUE(HasRecord(records, LogLevel::kInfo, "succeeded after 2 attempts")) + << "expected a success INFO"; +} + +// A metadata-only retry must not attribute a snapshot committed concurrently by +// another writer to this transaction. +TEST_F(TransactionRetryTest, MetadataOnlyRetryDoesNotLogConcurrentSnapshot) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + + constexpr int64_t kConcurrentSnapshotId = 987654321; + auto metadata_builder = TableMetadataBuilder::BuildFrom(mock_table_->metadata().get()); + auto concurrent_snapshot = std::make_shared(Snapshot{ + .snapshot_id = kConcurrentSnapshotId, + .parent_snapshot_id = mock_table_->metadata()->current_snapshot_id, + .sequence_number = mock_table_->metadata()->last_sequence_number + 1, + .timestamp_ms = TimePointMs{}, + .manifest_list = "concurrent-manifest-list.avro", + .summary = {{SnapshotSummaryFields::kOperation, "append"}}, + }); + metadata_builder->SetBranchSnapshot(concurrent_snapshot, + std::string(SnapshotRef::kMainBranch)); + ICEBERG_UNWRAP_OR_FAIL(auto concurrent_metadata, metadata_builder->Build()); + auto concurrent_metadata_ptr = + std::shared_ptr(std::move(concurrent_metadata)); + const std::string concurrent_metadata_location = "concurrent.metadata.json"; + + ON_CALL(*mock_catalog_, LoadTable(::testing::_)) + .WillByDefault([this, concurrent_metadata_ptr, &concurrent_metadata_location]( + const TableIdentifier&) -> Result> { + return Table::Make(mock_table_->name(), concurrent_metadata_ptr, + concurrent_metadata_location, mock_table_->io(), + mock_catalog_); + }); + + int update_call_count = 0; + ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault( + [this, concurrent_metadata_ptr, &concurrent_metadata_location, + &update_call_count](const TableIdentifier&, + const std::vector>&, + const std::vector>&) + -> Result> { + if (++update_call_count == 1) { + return CommitFailed("conflict on first attempt"); + } + return Table::Make(mock_table_->name(), concurrent_metadata_ptr, + concurrent_metadata_location, mock_table_->io(), + mock_catalog_); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto txn, mock_table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); + update->Set("retry.test", "value"); + ASSERT_THAT(update->Commit(), IsOk()); + ASSERT_THAT(txn->Commit(), IsOk()); + + EXPECT_FALSE(HasRecord(capturing->records(), LogLevel::kInfo, + std::to_string(kConcurrentSnapshotId))) + << "metadata-only commit attributed the concurrent snapshot to itself"; +} + +// A commit that exhausts its retries returns the final error without emitting a +// generic ERROR log. Genuine retry attempts still emit WARN records. +TEST_F(TransactionRetryTest, CommitRetryExhaustedDoesNotEmitErrorLog) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + + ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault([](const TableIdentifier&, + const std::vector>&, + const std::vector>&) + -> Result> { + return CommitFailed("always conflicts"); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto txn, mock_table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); + update->Set("retry.test", "value"); + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), IsError(ErrorKind::kCommitFailed)); + + auto records = capturing->records(); + EXPECT_FALSE(HasRecord(records, LogLevel::kError, "")) + << "the final commit error should be propagated without a generic ERROR log"; + // Retries 2..5 each log a WARN. + EXPECT_TRUE( + HasRecord(records, LogLevel::kWarn, "Retrying transaction commit (attempt 5)")); +} + +// A commit that succeeds on the first attempt emits a plain success INFO (no +// "after N attempts"). This is the single-attempt case that was previously silent. +TEST_F(TransactionRetryTest, CommitSuccessEmitsInfoLog) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + + ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) + .WillByDefault([this](const TableIdentifier&, + const std::vector>&, + const std::vector>&) + -> Result> { + return Table::Make(mock_table_->name(), mock_table_->metadata(), + std::string(mock_table_->metadata_file_location()), + mock_table_->io(), mock_catalog_); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto txn, mock_table_->NewTransaction()); + ICEBERG_UNWRAP_OR_FAIL(auto update, txn->NewUpdateProperties()); + update->Set("retry.test", "value"); + EXPECT_THAT(update->Commit(), IsOk()); + EXPECT_THAT(txn->Commit(), IsOk()); + + auto records = capturing->records(); + EXPECT_TRUE(HasRecord(records, LogLevel::kInfo, "Transaction commit succeeded")) + << "expected a success INFO on a single-attempt commit"; + // No retry happened, so there must be no retry WARN. + EXPECT_FALSE(HasRecord(records, LogLevel::kWarn, "Retrying transaction commit")); +} + +TEST_F(TransactionRetryTest, CommitStateUnknownStopsImmediatelyWithoutErrorLog) { + auto capturing = std::make_shared(); + capturing->SetLevel(LogLevel::kTrace); + ScopedDefaultLogger guard(capturing); + int update_call_count = 0; ON_CALL(*mock_catalog_, UpdateTable(::testing::_, ::testing::_, ::testing::_)) .WillByDefault( @@ -193,6 +371,8 @@ TEST_F(TransactionRetryTest, CommitNonRetryableErrorStopsImmediately) { auto result = txn->Commit(); EXPECT_THAT(result, IsError(ErrorKind::kCommitStateUnknown)); EXPECT_EQ(update_call_count, 1); // Should not retry + EXPECT_FALSE(HasRecord(capturing->records(), LogLevel::kError, "")) + << "an unknown commit state must not be logged as a confirmed failure"; } TEST_F(TransactionRetryTest, CreateTransactionDoesNotRetry) { diff --git a/src/iceberg/transaction.cc b/src/iceberg/transaction.cc index 80d39c8a8..b722197ae 100644 --- a/src/iceberg/transaction.cc +++ b/src/iceberg/transaction.cc @@ -19,10 +19,13 @@ #include "iceberg/transaction.h" #include +#include #include +#include #include "iceberg/catalog.h" #include "iceberg/location_provider.h" +#include "iceberg/logging/log_macros.h" #include "iceberg/schema.h" #include "iceberg/snapshot.h" #include "iceberg/statistics_file.h" @@ -376,14 +379,66 @@ Result> Transaction::Commit() { int32_t total_timeout_ms = props.Get(TableProperties::kCommitTotalRetryTimeMs); bool is_first_attempt = true; + int32_t attempt = 0; + std::string last_error; auto commit_result = MakeCommitRetryRunner(num_retries, min_wait_ms, max_wait_ms, total_timeout_ms) - .Run([this, &is_first_attempt]() -> Result> { + .Run([this, &is_first_attempt, &attempt, + &last_error]() -> Result> { + ++attempt; + // The runner only re-invokes this task when it has decided to retry, so + // attempt > 1 here means a genuine retry after a retryable failure. + if (attempt > 1) { + ICEBERG_LOG_WARN("Retrying transaction commit (attempt {}) after: {}", + attempt, last_error); + } auto result = CommitOnce(is_first_attempt); is_first_attempt = false; + if (!result.has_value()) { + last_error = result.error().message; + } return result; }); + if (commit_result.has_value()) { + // The builder contains only changes made by the successful attempt. Inspecting + // AddSnapshot changes avoids attributing a concurrent writer's snapshot to this + // transaction and also detects snapshots committed with StageOnly or ToBranch. + std::string detail; + const auto& changes = ctx_->metadata_builder->changes(); + size_t added_snapshot_count = 0; + for (const auto& change : changes) { + added_snapshot_count += change->kind() == TableUpdate::Kind::kAddSnapshot; + } + if (added_snapshot_count > 0) { + detail.reserve(32 + added_snapshot_count * 48); + std::format_to(std::back_inserter(detail), ": committed snapshot{} ", + added_snapshot_count == 1 ? "" : "s"); + + size_t appended_snapshot_count = 0; + for (const auto& change : changes) { + if (change->kind() != TableUpdate::Kind::kAddSnapshot) { + continue; + } + const auto& snapshot = + internal::checked_cast(*change).snapshot(); + if (appended_snapshot_count++ > 0) { + detail += ", "; + } + const auto& summary = snapshot->summary; + auto op = summary.find(SnapshotSummaryFields::kOperation); + std::format_to(std::back_inserter(detail), "{} (op={})", snapshot->snapshot_id, + op != summary.end() ? op->second : "unknown"); + } + } + if (attempt > 1) { + ICEBERG_LOG_INFO("Transaction commit succeeded after {} attempts{}", attempt, + detail); + } else { + ICEBERG_LOG_INFO("Transaction commit succeeded{}", detail); + } + } + Result finalize_result = commit_result.has_value() ? Result(commit_result.value()->metadata().get())