From 079efb4f5d3fd7f470832a02b5acbc48a72dae15 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Mon, 24 Aug 2026 02:23:52 -0400 Subject: [PATCH 1/3] feat(io): refresh vended storage credentials before they expire --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 236 ++++++++++++++- src/iceberg/arrow/s3/s3_properties.h | 3 + src/iceberg/test/arrow_s3_file_io_test.cc | 338 +++++++++++++++++++++- 3 files changed, 573 insertions(+), 4 deletions(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 4fbf33098..212770b68 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -18,6 +18,8 @@ */ #include +#include +#include #include #include #include @@ -178,6 +180,71 @@ std::string CanonicalizeS3Scheme(std::string_view location) { return std::string(location); } +// Lead time before expiry, matching Java's VendedCredentialsProvider. +constexpr auto kRefreshLeadTime = std::chrono::minutes(5); + +// After a failed refresh, how long to keep the current credentials before +// asking again, so an unreachable catalog is not queried per file operation. +constexpr auto kRefreshRetryBackoff = std::chrono::seconds(30); + +// Floor on a backoff shortened to land on the expiry. +constexpr auto kMinRefreshRetryBackoff = std::chrono::seconds(1); + +// How long an operation with expired credentials waits for a refresh already +// under way. Bounded: the catalog request behind it has no deadline of its own. +constexpr auto kExpiredCredentialWait = std::chrono::seconds(10); + +// When the earliest of `credentials` stops being valid, or nullopt if none of +// them does. No session token means static keys, which never expire; a token +// with no usable expiry is reported as already expired so it gets replaced +// rather than used until it fails, as Java does. +std::optional EarliestExpiry( + const std::vector& credentials) { + std::optional earliest; + const auto note = [&earliest](std::chrono::system_clock::time_point expires_at) { + if (!earliest.has_value() || expires_at < *earliest) { + earliest = expires_at; + } + }; + for (const auto& credential : credentials) { + if (!IsS3CredentialPrefix(credential.prefix) || + FindProperty(credential.config, S3Properties::kSessionToken) == nullptr) { + continue; + } + const auto* value = + FindProperty(credential.config, S3Properties::kSessionTokenExpiresAtMs); + if (value == nullptr) { + ICEBERG_LOG_WARN("Credential \"{}\" has a session token but no \"{}\"", + credential.prefix, S3Properties::kSessionTokenExpiresAtMs); + note(std::chrono::system_clock::now()); + continue; + } + auto millis = StringUtils::ParseNumber(*value); + if (!millis.has_value()) { + ICEBERG_LOG_WARN( + "Credential \"{}\" has a session token but an unparseable \"{}\" value \"{}\"", + credential.prefix, S3Properties::kSessionTokenExpiresAtMs, *value); + note(std::chrono::system_clock::now()); + continue; + } + // Beyond what the clock can hold, converting would overflow it. + constexpr auto kMaxMillis = std::chrono::duration_cast( + std::chrono::system_clock::duration::max()) + .count(); + constexpr auto kMinMillis = std::chrono::duration_cast( + std::chrono::system_clock::duration::min()) + .count(); + if (*millis > kMaxMillis || *millis < kMinMillis) { + ICEBERG_LOG_WARN("Credential \"{}\" has an out-of-range \"{}\" value \"{}\"", + credential.prefix, S3Properties::kSessionTokenExpiresAtMs, *value); + note(std::chrono::system_clock::now()); + continue; + } + note(std::chrono::system_clock::time_point(std::chrono::milliseconds(*millis))); + } + return earliest; +} + class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { public: ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs, @@ -204,6 +271,13 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { return storage_credentials_; } + void SetCredentialRefresher(StorageCredentialRefresher refresher) override { + std::unique_lock lock(mutex_); + refresher_ = std::move(refresher); + // A refresh in flight was started for the refresher just replaced. + ++credential_generation_; + } + SupportsStorageCredentials* AsSupportsStorageCredentials() override { return this; } private: @@ -235,12 +309,46 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { void InstallCredentials(std::vector& storage_credentials, DelegatesByPrefix& delegates); + /// \brief Whether the installed credentials are close enough to expiring to + /// be replaced, and no backoff is in effect. + /// + /// Callers must hold `mutex_`, at least shared. + bool RefreshDue() const; + + /// \brief Whether the installed credentials have already stopped being valid. + /// + /// Callers must hold `mutex_`, at least shared. + bool Expired() const; + + /// \brief When the next refresh attempt becomes allowed after a failure. + /// + /// Never past the point the credentials stop being valid. + /// + /// Callers must hold `mutex_`, at least shared. + std::chrono::steady_clock::time_point BackoffUntil() const; + + /// \brief Replace the credentials once they are close to expiring. + /// + /// Called before each handle is created; a handle keeps the delegate it was + /// built from, so I/O on an open one is not re-checked. A failure keeps the + /// current credentials rather than failing the read. + void MaybeRefreshCredentials(); + std::shared_ptr default_file_io_; std::unordered_map default_properties_; // Guards everything below; shared because reads happen per file operation. mutable std::shared_mutex mutex_; std::vector storage_credentials_; DelegatesByPrefix file_io_by_prefix_; + StorageCredentialRefresher refresher_; + std::optional expires_at_; + std::chrono::steady_clock::time_point retry_refresh_at_; + // Bumped whenever the credentials or the refresher change, so a refresh that + // fetched before one of those happened can tell its result is already stale. + uint64_t credential_generation_ = 0; + // Held across a refresh so concurrent operations skip it. Timed, so waiting + // on it is bounded. + std::timed_mutex refresh_mutex_; }; Status ArrowS3FileIO::SetStorageCredentials( @@ -260,7 +368,6 @@ Result ArrowS3FileIO::BuildDelegates( const std::vector& storage_credentials) const { DelegatesByPrefix delegates; delegates.reserve(storage_credentials.size()); - // TODO(gangwu): Refresh vended credentials via credentials.uri before tokens expire. for (const auto& credential : storage_credentials) { ICEBERG_RETURN_UNEXPECTED(credential.Validate()); // A server may vend credentials for several storage systems at once; @@ -291,9 +398,132 @@ Result ArrowS3FileIO::BuildDelegates( void ArrowS3FileIO::InstallCredentials( std::vector& storage_credentials, DelegatesByPrefix& delegates) { file_io_by_prefix_.swap(delegates); + expires_at_ = EarliestExpiry(storage_credentials); + retry_refresh_at_ = {}; + ++credential_generation_; storage_credentials_.swap(storage_credentials); } +bool ArrowS3FileIO::RefreshDue() const { + return expires_at_.has_value() && + std::chrono::system_clock::now() + kRefreshLeadTime >= *expires_at_ && + std::chrono::steady_clock::now() >= retry_refresh_at_; +} + +bool ArrowS3FileIO::Expired() const { + return expires_at_.has_value() && std::chrono::system_clock::now() >= *expires_at_; +} + +std::chrono::steady_clock::time_point ArrowS3FileIO::BackoffUntil() const { + auto delay = + std::chrono::duration_cast(kRefreshRetryBackoff); + if (expires_at_.has_value()) { + const auto remaining = std::chrono::duration_cast( + *expires_at_ - std::chrono::system_clock::now()); + // Worth retrying before they run out; once they have, faster retries only + // hammer a catalog that is already failing. + if (remaining > std::chrono::milliseconds::zero()) { + delay = std::clamp( + remaining, + std::chrono::duration_cast(kMinRefreshRetryBackoff), + delay); + } + } + return std::chrono::steady_clock::now() + delay; +} + +void ArrowS3FileIO::MaybeRefreshCredentials() { + { + // Cheap pre-check, so the common case costs one shared lock and no more. + std::shared_lock lock(mutex_); + if (!refresher_ || !RefreshDue()) { + return; + } + } + + std::unique_lock refresh_lock(refresh_mutex_, std::defer_lock); + if (!refresh_lock.try_lock()) { + // Another operation is already fetching; normally just use what we have. + { + std::shared_lock lock(mutex_); + if (!Expired()) { + return; + } + } + // Expired credentials leave nothing to proceed with, so wait instead. + if (!refresh_lock.try_lock_for(kExpiredCredentialWait)) { + return; + } + } + // Read together: pairing this refresher with a generation bumped by another + // one installed in between would make its result look current. + StorageCredentialRefresher refresher; + uint64_t generation = 0; + { + // Whoever held the lock may also have just finished, leaving nothing to do. + std::shared_lock lock(mutex_); + if (!refresher_ || !RefreshDue()) { + return; + } + refresher = refresher_; + generation = credential_generation_; + } + + // Outside `mutex_`: both are slow and must not block readers. + Status status; + DelegatesByPrefix delegates; + auto refreshed = refresher(); + if (refreshed.has_value()) { + auto built = BuildDelegates(*refreshed); + if (!built.has_value()) { + status = std::unexpected(built.error()); + } else if (built->empty()) { + // Installing this would drop working credentials for whatever ambient + // identity the AWS chain finds. Java refuses an empty list too. + status = NotFound("Refreshed credentials contain no S3-compatible prefix"); + } else { + delegates = std::move(built).value(); + } + } else { + status = std::unexpected(refreshed.error()); + } + + std::unique_lock lock(mutex_); + // Credentials installed meanwhile supersede this refresh: what it fetched is + // by now the older set. + const bool superseded = credential_generation_ != generation; + if (!status.has_value()) { + // Reported either way, so a failing catalog stays visible. + if (superseded) { + ICEBERG_LOG_WARN( + "Failed to refresh vended storage credentials ({}); they have since been " + "replaced", + status.error().message); + return; + } + retry_refresh_at_ = BackoffUntil(); + ICEBERG_LOG_WARN( + "Failed to refresh vended storage credentials ({}); keeping the current " + "ones and retrying in {}ms", + status.error().message, + std::chrono::duration_cast( + retry_refresh_at_ - std::chrono::steady_clock::now()) + .count()); + return; + } + if (superseded) { + return; + } + + // `delegates`/`*refreshed` take the retired generation; both are declared + // before the lock, so it destructs only after the lock releases. + InstallCredentials(*refreshed, delegates); + if (RefreshDue()) { + // Tokens shorter-lived than the lead time come back due again at once. + retry_refresh_at_ = BackoffUntil(); + } +} + std::shared_ptr ArrowS3FileIO::MatchDelegate( const std::shared_ptr& fallback, const DelegatesByPrefix& by_prefix, std::string_view location) { @@ -314,6 +544,8 @@ std::shared_ptr ArrowS3FileIO::MatchDelegate( std::shared_ptr ArrowS3FileIO::FileIOForPath( std::string_view location) { + MaybeRefreshCredentials(); + std::shared_lock lock(mutex_); return MatchDelegate(default_file_io_, file_io_by_prefix_, location); } @@ -338,6 +570,8 @@ Status ArrowS3FileIO::DeleteFile(const std::string& file_location) { } Status ArrowS3FileIO::DeleteFiles(const std::vector& file_locations) { + MaybeRefreshCredentials(); + // One snapshot so the whole batch matches the same delegate generation; only // ever a handful of delegates, so a linear scan beats hashing. std::shared_ptr fallback; diff --git a/src/iceberg/arrow/s3/s3_properties.h b/src/iceberg/arrow/s3/s3_properties.h index 50dafa56f..42599f0a0 100644 --- a/src/iceberg/arrow/s3/s3_properties.h +++ b/src/iceberg/arrow/s3/s3_properties.h @@ -42,6 +42,9 @@ struct S3Properties { static constexpr std::string_view kSecretAccessKey = "s3.secret-access-key"; /// AWS session token (for temporary credentials) static constexpr std::string_view kSessionToken = "s3.session-token"; + /// Epoch milliseconds at which a vended session token stops being valid + static constexpr std::string_view kSessionTokenExpiresAtMs = + "s3.session-token-expires-at-ms"; /// AWS region, standard Iceberg client property. static constexpr std::string_view kClientRegion = "client.region"; /// Custom endpoint override (for S3-compatible object stores) diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 24628949a..9b29212a4 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -19,9 +19,12 @@ #include #include +#include +#include #include #include #include +#include #include #include #include @@ -142,10 +145,26 @@ class ArrowS3FileIOTest : public ::testing::Test { std::optional base_uri_; }; -bool HasWarning(const CapturingLogger& logger) { +bool HasWarning(const CapturingLogger& logger, std::string_view substring = {}) { const auto records = logger.records(); - return std::ranges::any_of( - records, [](const LogMessage& record) { return record.level == LogLevel::kWarn; }); + return std::ranges::any_of(records, [substring](const LogMessage& record) { + return record.level == LogLevel::kWarn && + record.message.find(substring) != std::string::npos; + }); +} + +constexpr auto kOutlastsAShortenedBackoff = std::chrono::milliseconds(1200); + +std::vector ExpiringCredentials(std::chrono::milliseconds valid_for, + std::string_view access_key) { + const auto expires_at = std::chrono::duration_cast( + (std::chrono::system_clock::now() + valid_for).time_since_epoch()); + return {{.prefix = "s3", + .config = {{std::string(S3Properties::kAccessKeyId), std::string(access_key)}, + {std::string(S3Properties::kSecretAccessKey), "secret"}, + {std::string(S3Properties::kSessionToken), "token"}, + {std::string(S3Properties::kSessionTokenExpiresAtMs), + std::to_string(expires_at.count())}}}}; } Status CheckReadWrite(FileIO& io, const std::string& object_uri, @@ -240,6 +259,319 @@ TEST_F(ArrowS3FileIOTest, WarnsWhenNoCredentialApplies) { EXPECT_TRUE(HasWarning(*logger)); } +TEST_F(ArrowS3FileIOTest, RefreshesCredentialsCloseToExpiry) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + const auto refreshed = ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return refreshed; + }); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + IsOk()); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 1); + EXPECT_EQ(credentialed->credentials(), refreshed); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 1); +} + +TEST_F(ArrowS3FileIOTest, DoesNotRefreshCredentialsThatAreNotCloseToExpiry) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return std::vector{}; + }); + + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::hours(1), "access-key")), + IsOk()); + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 0); + + const std::vector static_credentials = { + {.prefix = "s3", + .config = {{std::string(S3Properties::kAccessKeyId), "access-key"}, + {std::string(S3Properties::kSecretAccessKey), "secret"}}}}; + ASSERT_THAT(credentialed->SetStorageCredentials(static_credentials), IsOk()); + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 0); +} + +TEST_F(ArrowS3FileIOTest, RefreshesOnceWhenOperationsRaceForIt) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + std::atomic refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + IsOk()); + + constexpr int kThreads = 8; + std::atomic failures = 0; + std::vector threads; + threads.reserve(kThreads); + for (int i = 0; i < kThreads; ++i) { + threads.emplace_back([&] { + for (int op = 0; op < 4; ++op) { + if (!result.value()->NewInputFile("s3://bucket/key").has_value()) { + ++failures; + } + } + }); + } + for (auto& thread : threads) { + thread.join(); + } + + EXPECT_EQ(failures, 0); + EXPECT_EQ(refresh_calls, 1); +} + +TEST_F(ArrowS3FileIOTest, RefreshesOnceWhenCredentialsHaveExpired) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + std::mutex mutex; + std::condition_variable cv; + bool refresh_started = false; + bool release_refresh = false; + std::atomic refresh_calls = 0; + + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + std::unique_lock lock(mutex); + refresh_started = true; + cv.notify_all(); + cv.wait(lock, [&] { return release_refresh; }); + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(-std::chrono::minutes(1), "expired-key")), + IsOk()); + + std::thread winner( + [&] { EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); }); + { + std::unique_lock lock(mutex); + cv.wait(lock, [&] { return refresh_started; }); + } + + bool loser_ready = false; + std::thread loser([&] { + { + std::lock_guard lock(mutex); + loser_ready = true; + } + cv.notify_all(); + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + }); + { + std::unique_lock lock(mutex); + cv.wait(lock, [&] { return loser_ready; }); + release_refresh = true; + } + cv.notify_all(); + winner.join(); + loser.join(); + + EXPECT_EQ(refresh_calls, 1); + EXPECT_THAT(credentialed->credentials(), ::testing::Not(::testing::IsEmpty())); +} + +TEST_F(ArrowS3FileIOTest, RefreshDoesNotUndoCredentialsInstalledWhileItRan) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + std::mutex mutex; + std::condition_variable cv; + bool refresh_started = false; + bool release_refresh = false; + + credentialed->SetCredentialRefresher([&]() -> Result> { + std::unique_lock lock(mutex); + refresh_started = true; + cv.notify_all(); + cv.wait(lock, [&] { return release_refresh; }); + return ExpiringCredentials(std::chrono::hours(1), "fetched-by-refresh"); + }); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + IsOk()); + + std::thread operation( + [&] { EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); }); + + const auto installed = + ExpiringCredentials(std::chrono::hours(2), "installed-meanwhile"); + { + std::unique_lock lock(mutex); + cv.wait(lock, [&] { return refresh_started; }); + } + ASSERT_THAT(credentialed->SetStorageCredentials(installed), IsOk()); + { + std::lock_guard lock(mutex); + release_refresh = true; + } + cv.notify_all(); + operation.join(); + + EXPECT_EQ(credentialed->credentials(), installed); +} + +TEST_F(ArrowS3FileIOTest, RefreshesSessionCredentialsWithoutAUsableExpiry) { + for (std::string_view expiry : {"", "not-a-number"}) { + SCOPED_TRACE(expiry); + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); + + auto logger = std::make_shared(); + ScopedDefaultLogger scoped(logger); + std::unordered_map config = { + {std::string(S3Properties::kAccessKeyId), "access-key"}, + {std::string(S3Properties::kSecretAccessKey), "secret"}, + {std::string(S3Properties::kSessionToken), "token"}}; + if (!expiry.empty()) { + config[std::string(S3Properties::kSessionTokenExpiresAtMs)] = std::string(expiry); + } + ASSERT_THAT(credentialed->SetStorageCredentials( + {{.prefix = "s3", .config = std::move(config)}}), + IsOk()); + EXPECT_TRUE(HasWarning(*logger, "session token")); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 1); + } +} + +TEST_F(ArrowS3FileIOTest, BacksOffWhenReplacementsAlsoLackAnExpiry) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + const std::vector undatable = { + {.prefix = "s3", + .config = {{std::string(S3Properties::kAccessKeyId), "access-key"}, + {std::string(S3Properties::kSecretAccessKey), "secret"}, + {std::string(S3Properties::kSessionToken), "token"}}}}; + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return undatable; + }); + ASSERT_THAT(credentialed->SetStorageCredentials(undatable), IsOk()); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + ASSERT_EQ(refresh_calls, 1); + + std::this_thread::sleep_for(kOutlastsAShortenedBackoff); + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 1); +} + +TEST_F(ArrowS3FileIOTest, IgnoresUnparseableExpiry) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return std::vector{}; + }); + + auto logger = std::make_shared(); + ScopedDefaultLogger scoped(logger); + std::unordered_map config = { + {std::string(S3Properties::kAccessKeyId), "access-key"}, + {std::string(S3Properties::kSecretAccessKey), "secret"}, + {std::string(S3Properties::kSessionTokenExpiresAtMs), "not-a-number"}}; + const std::vector credentials = { + {.prefix = "s3", .config = std::move(config)}}; + ASSERT_THAT(credentialed->SetStorageCredentials(credentials), IsOk()); + EXPECT_EQ(credentialed->credentials(), credentials); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 0); +} + +TEST_F(ArrowS3FileIOTest, BacksOffWhenTheReplacementIsAlsoCloseToExpiry) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::minutes(1), "short-lived-key"); + }); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + IsOk()); + + for (int i = 0; i < 3; ++i) { + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + } + EXPECT_EQ(refresh_calls, 1); +} + +TEST_F(ArrowS3FileIOTest, KeepsCredentialsWhenRefreshFails) { + auto result = MakeS3FileIO({}); + ASSERT_THAT(result, IsOk()); + auto* credentialed = result.value()->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + int refresh_calls = 0; + credentialed->SetCredentialRefresher([&]() -> Result> { + ++refresh_calls; + return NotFound("catalog unreachable"); + }); + const auto expiring = ExpiringCredentials(std::chrono::minutes(1), "expiring-key"); + ASSERT_THAT(credentialed->SetStorageCredentials(expiring), IsOk()); + + auto logger = std::make_shared(); + ScopedDefaultLogger scoped(logger); + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(credentialed->credentials(), expiring); + EXPECT_TRUE(HasWarning(*logger, "Failed to refresh")); + + EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(refresh_calls, 1); +} + TEST_F(ArrowS3FileIOTest, OperationsSurviveConcurrentCredentialInstalls) { auto result = MakeS3FileIO({}); ASSERT_THAT(result, IsOk()); From 1fdf718e2e468baf057c56e6b8fbf0e0900038aa Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Tue, 29 Sep 2026 02:13:59 -0400 Subject: [PATCH 2/3] docs(io): document s3.session-token-expires-at-ms --- mkdocs/docs/file-io.md | 1 + 1 file changed, 1 insertion(+) diff --git a/mkdocs/docs/file-io.md b/mkdocs/docs/file-io.md index 5b123addf..af1134491 100644 --- a/mkdocs/docs/file-io.md +++ b/mkdocs/docs/file-io.md @@ -63,6 +63,7 @@ each file location's scheme. | `s3.access-key-id` | `admin` | Static access key ID; must be set together with the secret key | | `s3.secret-access-key` | `password` | Static secret access key | | `s3.session-token` | `AQoDYXdzEJr...` | Session token, for temporary credentials. Ignored unless both static keys are set | +| `s3.session-token-expires-at-ms` | `1767225600000` | When the session token expires, in epoch milliseconds. Vended credentials are refreshed from the catalog five minutes before | | `client.region` | `us-east-1` | Region to sign requests for | | `s3.endpoint` | `https://127.0.0.1:9000` | Endpoint to use instead of the AWS one. When absent, the `AWS_ENDPOINT_URL_S3` / `AWS_ENDPOINT_URL` environment variables are consulted | | `s3.path-style-access` | `true` | Address buckets as a path (`endpoint/bucket`) instead of a virtual host (`bucket.endpoint`). Only takes effect together with a custom endpoint | From a1cf5ef40a37152018014ab6998ed1587d08d0f0 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Thu, 8 Oct 2026 22:34:21 -0400 Subject: [PATCH 3/3] fix(io): use credential providers for S3 refresh Adapt the refresh policy to the provider API merged in #899. Install initial credentials and the provider together, and keep the provider across later credential updates. Migrate the refresh tests and cover failed initialization and the REST-to-S3 refresh path with real object storage. AI-Model: gpt-6 AI-Contributed/Feature: 40/40 AI-Contributed/UT: 248/248 --- src/iceberg/arrow/s3/arrow_s3_file_io.cc | 40 ++-- src/iceberg/test/arrow_s3_file_io_test.cc | 196 +++++++++++++------- src/iceberg/test/rest_arrow_file_io_test.cc | 52 ++++++ 3 files changed, 206 insertions(+), 82 deletions(-) diff --git a/src/iceberg/arrow/s3/arrow_s3_file_io.cc b/src/iceberg/arrow/s3/arrow_s3_file_io.cc index 212770b68..44fd409e2 100644 --- a/src/iceberg/arrow/s3/arrow_s3_file_io.cc +++ b/src/iceberg/arrow/s3/arrow_s3_file_io.cc @@ -271,12 +271,9 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { return storage_credentials_; } - void SetCredentialRefresher(StorageCredentialRefresher refresher) override { - std::unique_lock lock(mutex_); - refresher_ = std::move(refresher); - // A refresh in flight was started for the refresher just replaced. - ++credential_generation_; - } + Status InitializeStorageCredentials( + const std::vector& storage_credentials, + std::shared_ptr provider) override; SupportsStorageCredentials* AsSupportsStorageCredentials() override { return this; } @@ -340,17 +337,29 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials { mutable std::shared_mutex mutex_; std::vector storage_credentials_; DelegatesByPrefix file_io_by_prefix_; - StorageCredentialRefresher refresher_; + std::shared_ptr provider_; std::optional expires_at_; std::chrono::steady_clock::time_point retry_refresh_at_; - // Bumped whenever the credentials or the refresher change, so a refresh that - // fetched before one of those happened can tell its result is already stale. + // Bumped on each install so a refresh cannot overwrite newer credentials. uint64_t credential_generation_ = 0; // Held across a refresh so concurrent operations skip it. Timed, so waiting // on it is bounded. std::timed_mutex refresh_mutex_; }; +Status ArrowS3FileIO::InitializeStorageCredentials( + const std::vector& storage_credentials, + std::shared_ptr provider) { + ICEBERG_ASSIGN_OR_RAISE(auto delegates, BuildDelegates(storage_credentials)); + auto credentials = storage_credentials; + { + std::unique_lock lock(mutex_); + provider_.swap(provider); + InstallCredentials(credentials, delegates); + } + return {}; +} + Status ArrowS3FileIO::SetStorageCredentials( const std::vector& storage_credentials) { ICEBERG_ASSIGN_OR_RAISE(auto delegates, BuildDelegates(storage_credentials)); @@ -436,7 +445,7 @@ void ArrowS3FileIO::MaybeRefreshCredentials() { { // Cheap pre-check, so the common case costs one shared lock and no more. std::shared_lock lock(mutex_); - if (!refresher_ || !RefreshDue()) { + if (!provider_ || !RefreshDue()) { return; } } @@ -455,24 +464,23 @@ void ArrowS3FileIO::MaybeRefreshCredentials() { return; } } - // Read together: pairing this refresher with a generation bumped by another - // one installed in between would make its result look current. - StorageCredentialRefresher refresher; + // Snapshot the provider and credentials' generation together. + std::shared_ptr provider; uint64_t generation = 0; { // Whoever held the lock may also have just finished, leaving nothing to do. std::shared_lock lock(mutex_); - if (!refresher_ || !RefreshDue()) { + if (!provider_ || !RefreshDue()) { return; } - refresher = refresher_; + provider = provider_; generation = credential_generation_; } // Outside `mutex_`: both are slow and must not block readers. Status status; DelegatesByPrefix delegates; - auto refreshed = refresher(); + auto refreshed = provider->Load(); if (refreshed.has_value()) { auto built = BuildDelegates(*refreshed); if (!built.has_value()) { diff --git a/src/iceberg/test/arrow_s3_file_io_test.cc b/src/iceberg/test/arrow_s3_file_io_test.cc index 9b29212a4..3f4056a47 100644 --- a/src/iceberg/test/arrow_s3_file_io_test.cc +++ b/src/iceberg/test/arrow_s3_file_io_test.cc @@ -153,6 +153,11 @@ bool HasWarning(const CapturingLogger& logger, std::string_view substring = {}) }); } +class MockStorageCredentialProvider : public StorageCredentialProvider { + public: + MOCK_METHOD((Result>), Load, (), (override)); +}; + constexpr auto kOutlastsAShortenedBackoff = std::chrono::milliseconds(1200); std::vector ExpiringCredentials(std::chrono::milliseconds valid_for, @@ -267,13 +272,16 @@ TEST_F(ArrowS3FileIOTest, RefreshesCredentialsCloseToExpiry) { const auto refreshed = ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return refreshed; - }); - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return refreshed; + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key"), provider), IsOk()); + EXPECT_EQ(refresh_calls, 0); EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); EXPECT_EQ(refresh_calls, 1); @@ -290,13 +298,15 @@ TEST_F(ArrowS3FileIOTest, DoesNotRefreshCredentialsThatAreNotCloseToExpiry) { ASSERT_NE(credentialed, nullptr); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return std::vector{}; - }); - - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(std::chrono::hours(1), "access-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return std::vector{}; + }); + + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::hours(1), "access-key"), provider), IsOk()); EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); EXPECT_EQ(refresh_calls, 0); @@ -310,6 +320,43 @@ TEST_F(ArrowS3FileIOTest, DoesNotRefreshCredentialsThatAreNotCloseToExpiry) { EXPECT_EQ(refresh_calls, 0); } +TEST_F(ArrowS3FileIOTest, KeepsProviderWhenCredentialsAreUpdated) { + ICEBERG_UNWRAP_OR_FAIL(auto io, MakeS3FileIO({})); + auto* credentialed = io->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + const auto refreshed = ExpiringCredentials(std::chrono::hours(2), "refreshed-key"); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()).WillOnce(::testing::Return(refreshed)); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::hours(1), "initial-key"), provider), + IsOk()); + ASSERT_THAT(credentialed->SetStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + IsOk()); + + EXPECT_THAT(io->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(credentialed->credentials(), refreshed); +} + +TEST_F(ArrowS3FileIOTest, FailedInitializationDoesNotInstallProvider) { + ICEBERG_UNWRAP_OR_FAIL(auto io, MakeS3FileIO({})); + auto* credentialed = io->AsSupportsStorageCredentials(); + ASSERT_NE(credentialed, nullptr); + + const auto initial = ExpiringCredentials(std::chrono::minutes(1), "initial-key"); + ASSERT_THAT(credentialed->InitializeStorageCredentials(initial, nullptr), IsOk()); + auto invalid = initial; + invalid.front().prefix.clear(); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()).Times(0); + EXPECT_THAT(credentialed->InitializeStorageCredentials(invalid, provider), + IsError(ErrorKind::kValidationFailed)); + + EXPECT_THAT(io->NewInputFile("s3://bucket/key"), IsOk()); + EXPECT_EQ(credentialed->credentials(), initial); +} + TEST_F(ArrowS3FileIOTest, RefreshesOnceWhenOperationsRaceForIt) { auto result = MakeS3FileIO({}); ASSERT_THAT(result, IsOk()); @@ -317,12 +364,14 @@ TEST_F(ArrowS3FileIOTest, RefreshesOnceWhenOperationsRaceForIt) { ASSERT_NE(credentialed, nullptr); std::atomic refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); - }); - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key"), provider), IsOk()); constexpr int kThreads = 8; @@ -358,16 +407,18 @@ TEST_F(ArrowS3FileIOTest, RefreshesOnceWhenCredentialsHaveExpired) { bool release_refresh = false; std::atomic refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - std::unique_lock lock(mutex); - refresh_started = true; - cv.notify_all(); - cv.wait(lock, [&] { return release_refresh; }); - return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); - }); - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(-std::chrono::minutes(1), "expired-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + std::unique_lock lock(mutex); + refresh_started = true; + cv.notify_all(); + cv.wait(lock, [&] { return release_refresh; }); + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(-std::chrono::minutes(1), "expired-key"), provider), IsOk()); std::thread winner( @@ -410,15 +461,17 @@ TEST_F(ArrowS3FileIOTest, RefreshDoesNotUndoCredentialsInstalledWhileItRan) { bool refresh_started = false; bool release_refresh = false; - credentialed->SetCredentialRefresher([&]() -> Result> { - std::unique_lock lock(mutex); - refresh_started = true; - cv.notify_all(); - cv.wait(lock, [&] { return release_refresh; }); - return ExpiringCredentials(std::chrono::hours(1), "fetched-by-refresh"); - }); - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + std::unique_lock lock(mutex); + refresh_started = true; + cv.notify_all(); + cv.wait(lock, [&] { return release_refresh; }); + return ExpiringCredentials(std::chrono::hours(1), "fetched-by-refresh"); + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key"), provider), IsOk()); std::thread operation( @@ -430,7 +483,7 @@ TEST_F(ArrowS3FileIOTest, RefreshDoesNotUndoCredentialsInstalledWhileItRan) { std::unique_lock lock(mutex); cv.wait(lock, [&] { return refresh_started; }); } - ASSERT_THAT(credentialed->SetStorageCredentials(installed), IsOk()); + auto install_status = credentialed->SetStorageCredentials(installed); { std::lock_guard lock(mutex); release_refresh = true; @@ -438,6 +491,7 @@ TEST_F(ArrowS3FileIOTest, RefreshDoesNotUndoCredentialsInstalledWhileItRan) { cv.notify_all(); operation.join(); + ASSERT_THAT(install_status, IsOk()); EXPECT_EQ(credentialed->credentials(), installed); } @@ -450,10 +504,12 @@ TEST_F(ArrowS3FileIOTest, RefreshesSessionCredentialsWithoutAUsableExpiry) { ASSERT_NE(credentialed, nullptr); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); - }); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::hours(1), "refreshed-key"); + }); auto logger = std::make_shared(); ScopedDefaultLogger scoped(logger); @@ -464,8 +520,8 @@ TEST_F(ArrowS3FileIOTest, RefreshesSessionCredentialsWithoutAUsableExpiry) { if (!expiry.empty()) { config[std::string(S3Properties::kSessionTokenExpiresAtMs)] = std::string(expiry); } - ASSERT_THAT(credentialed->SetStorageCredentials( - {{.prefix = "s3", .config = std::move(config)}}), + ASSERT_THAT(credentialed->InitializeStorageCredentials( + {{.prefix = "s3", .config = std::move(config)}}, provider), IsOk()); EXPECT_TRUE(HasWarning(*logger, "session token")); @@ -486,11 +542,13 @@ TEST_F(ArrowS3FileIOTest, BacksOffWhenReplacementsAlsoLackAnExpiry) { {std::string(S3Properties::kSecretAccessKey), "secret"}, {std::string(S3Properties::kSessionToken), "token"}}}}; int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return undatable; - }); - ASSERT_THAT(credentialed->SetStorageCredentials(undatable), IsOk()); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return undatable; + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials(undatable, provider), IsOk()); EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); ASSERT_EQ(refresh_calls, 1); @@ -507,10 +565,12 @@ TEST_F(ArrowS3FileIOTest, IgnoresUnparseableExpiry) { ASSERT_NE(credentialed, nullptr); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return std::vector{}; - }); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return std::vector{}; + }); auto logger = std::make_shared(); ScopedDefaultLogger scoped(logger); @@ -520,7 +580,7 @@ TEST_F(ArrowS3FileIOTest, IgnoresUnparseableExpiry) { {std::string(S3Properties::kSessionTokenExpiresAtMs), "not-a-number"}}; const std::vector credentials = { {.prefix = "s3", .config = std::move(config)}}; - ASSERT_THAT(credentialed->SetStorageCredentials(credentials), IsOk()); + ASSERT_THAT(credentialed->InitializeStorageCredentials(credentials, provider), IsOk()); EXPECT_EQ(credentialed->credentials(), credentials); EXPECT_THAT(result.value()->NewInputFile("s3://bucket/key"), IsOk()); @@ -534,12 +594,14 @@ TEST_F(ArrowS3FileIOTest, BacksOffWhenTheReplacementIsAlsoCloseToExpiry) { ASSERT_NE(credentialed, nullptr); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return ExpiringCredentials(std::chrono::minutes(1), "short-lived-key"); - }); - ASSERT_THAT(credentialed->SetStorageCredentials( - ExpiringCredentials(std::chrono::minutes(1), "expiring-key")), + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return ExpiringCredentials(std::chrono::minutes(1), "short-lived-key"); + }); + ASSERT_THAT(credentialed->InitializeStorageCredentials( + ExpiringCredentials(std::chrono::minutes(1), "expiring-key"), provider), IsOk()); for (int i = 0; i < 3; ++i) { @@ -555,12 +617,14 @@ TEST_F(ArrowS3FileIOTest, KeepsCredentialsWhenRefreshFails) { ASSERT_NE(credentialed, nullptr); int refresh_calls = 0; - credentialed->SetCredentialRefresher([&]() -> Result> { - ++refresh_calls; - return NotFound("catalog unreachable"); - }); + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()) + .WillRepeatedly([&]() -> Result> { + ++refresh_calls; + return NotFound("catalog unreachable"); + }); const auto expiring = ExpiringCredentials(std::chrono::minutes(1), "expiring-key"); - ASSERT_THAT(credentialed->SetStorageCredentials(expiring), IsOk()); + ASSERT_THAT(credentialed->InitializeStorageCredentials(expiring, provider), IsOk()); auto logger = std::make_shared(); ScopedDefaultLogger scoped(logger); diff --git a/src/iceberg/test/rest_arrow_file_io_test.cc b/src/iceberg/test/rest_arrow_file_io_test.cc index 83966dc96..8b893c3f2 100644 --- a/src/iceberg/test/rest_arrow_file_io_test.cc +++ b/src/iceberg/test/rest_arrow_file_io_test.cc @@ -67,6 +67,11 @@ TEST_F(RestArrowFileIOTest, ReadsBackWhatItWroteThroughRealLocalFileIO) { #if ICEBERG_S3_ENABLED +class MockStorageCredentialProvider : public StorageCredentialProvider { + public: + MOCK_METHOD((Result>), Load, (), (override)); +}; + std::optional GetEnvIfSet(const char* key) { const char* value = std::getenv(key); if (value == nullptr || std::string_view(value).empty()) { @@ -195,6 +200,53 @@ TEST_F(RestArrowFileIOTest, ReadsBackWhatItWroteThroughAnOssLocation) { EXPECT_THAT(io.value()->DeleteFile(object_uri), IsOk()); } +TEST_F(RestArrowFileIOTest, RefreshesCredentialsThroughTheRealS3FileIO) { + const auto base_uri = GetEnvIfSet("ICEBERG_TEST_S3_URI"); + const auto access_key = GetEnvIfSet("AWS_ACCESS_KEY_ID"); + const auto secret_key = GetEnvIfSet("AWS_SECRET_ACCESS_KEY"); + if (!base_uri || !access_key || !secret_key) { + GTEST_SKIP() + << "Set ICEBERG_TEST_S3_URI, AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY"; + } + + std::unordered_map config = { + {"s3.access-key-id", *access_key}, {"s3.secret-access-key", *secret_key}}; + if (const auto token = GetEnvIfSet("AWS_SESSION_TOKEN")) { + config["s3.session-token"] = *token; + } + if (const auto endpoint = GetEnvIfSet("ICEBERG_TEST_S3_ENDPOINT")) { + config["s3.endpoint"] = *endpoint; + } + if (const auto region = GetEnvIfSet("AWS_REGION")) { + config["client.region"] = *region; + } + const std::vector refreshed = {{.prefix = "s3", .config = config}}; + config["s3.access-key-id"] = "bad-access-key"; + config["s3.secret-access-key"] = "bad-secret-key"; + config["s3.session-token"] = "expired-token"; + config["s3.session-token-expires-at-ms"] = "0"; + + // Only the provider can supply credentials that authenticate this round trip. + ScopedScrubbedAwsCredentialEnv scrubbed; + auto provider = std::make_shared(); + EXPECT_CALL(*provider, Load()).WillOnce(::testing::Return(refreshed)); + auto io = MakeTableFileIO({{"warehouse", "logical_warehouse_name"}}, + /*table_config=*/{}, + {{.prefix = "s3", .config = std::move(config)}}, provider); + ASSERT_THAT(io, IsOk()); + + auto object_uri = *base_uri; + if (!object_uri.ends_with('/')) { + object_uri += '/'; + } + object_uri += "iceberg_refresh_provider.txt"; + constexpr std::string_view kContent = "written after refreshing vended credentials"; + ASSERT_THAT(io.value()->WriteFile(object_uri, kContent), IsOk()); + EXPECT_THAT(io.value()->ReadFile(object_uri, std::nullopt), + HasValue(::testing::Eq(std::string(kContent)))); + EXPECT_THAT(io.value()->DeleteFile(object_uri), IsOk()); +} + #endif // ICEBERG_S3_ENABLED } // namespace