Skip to content

Commit bc689cc

Browse files
refactor(io): load resolver delegates off the lock; test real IO across installs
ResolvingFileIO::FileIOForPath held its mutex across FileIORegistry::Load and the delegate's SetStorageCredentials, so building the first S3 client (which can look up a bucket region) stalled every other operation on the resolver. It now loads without the lock and caches the result only if no credential install happened meanwhile; otherwise it loads again with the new credentials. Delegates retired by an install are torn down outside the lock. Tests now exercise real IO where the old ones stopped short: the multi-prefix DeleteFiles test failed URI parsing before reaching any delegate, and the concurrency test only built input-file wrappers. The replacements write, delete and read through the test object store, and the stress test waits for every worker to run before installing.
1 parent 0eb0820 commit bc689cc

5 files changed

Lines changed: 207 additions & 45 deletions

File tree

‎src/iceberg/arrow/s3/arrow_s3_file_io.cc‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -221,7 +221,7 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials {
221221

222222
/// \brief Build a delegate for each credential this FileIO can serve.
223223
///
224-
/// Lock-free on purpose: building an S3 client can reach out to discover a
224+
/// Runs without holding `mutex_`: building an S3 client can reach out to discover a
225225
/// bucket region, which would stall every concurrent operation. Reads no
226226
/// mutable member state.
227227
Result<DelegatesByPrefix> BuildDelegates(

‎src/iceberg/resolving_file_io.cc‎

Lines changed: 42 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -40,28 +40,44 @@ Result<std::shared_ptr<FileIO>> ResolvingFileIO::FileIOForPath(
4040
const auto scheme = StringUtils::ToLower(LocationUtil::ParseScheme(location));
4141
ICEBERG_ASSIGN_OR_RAISE(const auto name, FileIORegistry::Resolve(scheme));
4242

43-
{
44-
std::shared_lock lock(mutex_);
45-
if (const auto cached = io_by_name_.find(name); cached != io_by_name_.end()) {
46-
return cached->second;
47-
}
48-
}
49-
50-
std::unique_lock lock(mutex_);
51-
auto it = io_by_name_.find(name);
52-
if (it == io_by_name_.end()) {
53-
ICEBERG_ASSIGN_OR_RAISE(auto io, FileIORegistry::Load(name, properties_));
54-
// Forward all credentials; each implementation applies the prefixes it
55-
// understands.
56-
if (!storage_credentials_.empty()) {
43+
// Loads without holding `mutex_`: building a client can reach the network (an
44+
// S3 client may look up its bucket region), which would stall every other
45+
// operation. Forwards all credentials; each implementation applies the
46+
// prefixes it understands.
47+
auto load = [&](const std::vector<StorageCredential>& credentials)
48+
-> Result<std::shared_ptr<FileIO>> {
49+
ICEBERG_ASSIGN_OR_RAISE(std::shared_ptr<FileIO> io,
50+
FileIORegistry::Load(name, properties_));
51+
if (!credentials.empty()) {
5752
if (auto* credentialed = io->AsSupportsStorageCredentials()) {
58-
ICEBERG_RETURN_UNEXPECTED(
59-
credentialed->SetStorageCredentials(storage_credentials_));
53+
ICEBERG_RETURN_UNEXPECTED(credentialed->SetStorageCredentials(credentials));
6054
}
6155
}
62-
it = io_by_name_.emplace(std::string(name), std::move(io)).first;
56+
return io;
57+
};
58+
59+
while (true) {
60+
uint64_t generation = 0;
61+
std::vector<StorageCredential> credentials;
62+
{
63+
std::shared_lock lock(mutex_);
64+
if (const auto cached = io_by_name_.find(name); cached != io_by_name_.end()) {
65+
return cached->second;
66+
}
67+
generation = credential_generation_;
68+
credentials = storage_credentials_;
69+
}
70+
// Declared before the lock, so a delegate that is not cached is torn down
71+
// only after the lock is released.
72+
auto loaded = load(credentials);
73+
std::unique_lock lock(mutex_);
74+
if (generation != credential_generation_) {
75+
continue; // Credentials were replaced mid-load; load again with them.
76+
}
77+
ICEBERG_RETURN_UNEXPECTED(loaded);
78+
// A concurrent first access may have cached one already; that one wins.
79+
return io_by_name_.try_emplace(name, *loaded).first->second;
6380
}
64-
return it->second;
6581
}
6682

6783
Result<std::unique_ptr<InputFile>> ResolvingFileIO::NewInputFile(
@@ -103,9 +119,14 @@ Status ResolvingFileIO::SetStorageCredentials(
103119
const std::vector<StorageCredential>& storage_credentials) {
104120
// Rebuild delegates lazily with the new credentials. Updating live delegates
105121
// instead would leave the resolver inconsistent if one of them rejected them.
106-
std::unique_lock lock(mutex_);
107-
storage_credentials_ = storage_credentials;
108-
io_by_name_.clear();
122+
// Retired outside the lock: tearing down a delegate can block.
123+
decltype(io_by_name_) retired;
124+
{
125+
std::unique_lock lock(mutex_);
126+
storage_credentials_ = storage_credentials;
127+
++credential_generation_;
128+
retired.swap(io_by_name_);
129+
}
109130
return {};
110131
}
111132

‎src/iceberg/resolving_file_io.h‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
/// \file iceberg/resolving_file_io.h
2323
/// \brief FileIO that resolves the concrete implementation per file-path scheme.
2424

25+
#include <cstdint>
2526
#include <memory>
2627
#include <shared_mutex>
2728
#include <string>
@@ -73,6 +74,9 @@ class ICEBERG_EXPORT ResolvingFileIO final : public FileIO,
7374
// Guards lazy resolution and credential state.
7475
mutable std::shared_mutex mutex_;
7576
std::vector<StorageCredential> storage_credentials_;
77+
// Bumped by every credential install, so a delegate loaded from an older set
78+
// never reaches the cache.
79+
uint64_t credential_generation_ = 0;
7680
std::unordered_map<std::string, std::shared_ptr<FileIO>, StringHash, StringEqual>
7781
io_by_name_;
7882
};

‎src/iceberg/test/arrow_s3_file_io_test.cc‎

Lines changed: 96 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
#include <iostream>
2424
#include <memory>
2525
#include <optional>
26+
#include <span>
2627
#include <string>
2728
#include <string_view>
2829
#include <thread>
@@ -239,28 +240,6 @@ TEST_F(ArrowS3FileIOTest, WarnsWhenNoCredentialApplies) {
239240
EXPECT_TRUE(HasWarning(*logger));
240241
}
241242

242-
TEST_F(ArrowS3FileIOTest, DeleteFilesDispatchesAcrossCredentialPrefixes) {
243-
auto result = MakeS3FileIO({});
244-
ASSERT_THAT(result, IsOk());
245-
auto* credentialed = result.value()->AsSupportsStorageCredentials();
246-
ASSERT_NE(credentialed, nullptr);
247-
248-
auto credential = [](std::string_view prefix, std::string_view access_key) {
249-
return StorageCredential{
250-
.prefix = std::string(prefix),
251-
.config = {{std::string(S3Properties::kAccessKeyId), std::string(access_key)},
252-
{std::string(S3Properties::kSecretAccessKey), "secret"}}};
253-
};
254-
ASSERT_THAT(credentialed->SetStorageCredentials({credential("s3://bucket-a", "key-a"),
255-
credential("s3://bucket-b", "key-b")}),
256-
IsOk());
257-
258-
auto status = result.value()->DeleteFiles({"s3://bucket-a/%ZZ.parquet",
259-
"s3://bucket-a/second.parquet",
260-
"s3://bucket-b/other.parquet"});
261-
EXPECT_THAT(status, HasErrorMessage("Cannot parse URI"));
262-
}
263-
264243
TEST_F(ArrowS3FileIOTest, OperationsSurviveConcurrentCredentialInstalls) {
265244
auto result = MakeS3FileIO({});
266245
ASSERT_THAT(result, IsOk());
@@ -276,18 +255,28 @@ TEST_F(ArrowS3FileIOTest, OperationsSurviveConcurrentCredentialInstalls) {
276255
ASSERT_THAT(credentialed->SetStorageCredentials({credential("first")}), IsOk());
277256

278257
std::atomic<bool> stop = false;
258+
std::atomic<int> started = 0;
279259
std::atomic<int> failures = 0;
280260
std::vector<std::thread> operations;
281261
operations.reserve(4);
282262
for (int i = 0; i < 4; ++i) {
283263
operations.emplace_back([&] {
284-
while (!stop.load()) {
264+
auto open = [&] {
285265
if (!result.value()->NewInputFile("s3://bucket/key").has_value()) {
286266
++failures;
287267
}
268+
};
269+
open();
270+
++started;
271+
while (!stop.load()) {
272+
open();
288273
}
289274
});
290275
}
276+
// Every worker has run and is still looping before the first install.
277+
while (started.load() < 4) {
278+
std::this_thread::yield();
279+
}
291280
// No assertions until the threads are joined: a fatal assertion here would
292281
// destroy joinable threads and terminate the binary, masking the failure.
293282
Status install_status = {};
@@ -406,6 +395,90 @@ TEST_F(ArrowS3FileIOTest, AppliesOssCredentialInRealRoundTrip) {
406395
EXPECT_THAT(CheckReadWrite(*io, s3_uri, "hello oss with vended credentials"), IsOk());
407396
}
408397

398+
TEST_F(ArrowS3FileIOTest, DeleteFilesReachesEveryCredentialPrefix) {
399+
if (!HasIntegrationEnv()) {
400+
GTEST_SKIP() << "Set ICEBERG_TEST_S3_URI to enable S3 IO test";
401+
}
402+
403+
auto properties = PropertiesFromEnv();
404+
if (!properties.contains(std::string(S3Properties::kAccessKeyId)) ||
405+
!properties.contains(std::string(S3Properties::kSecretAccessKey))) {
406+
GTEST_SKIP() << "Set AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY to enable "
407+
"credential routing test";
408+
}
409+
410+
// Only the prefix delegates can authenticate, so a location that falls to
411+
// the default one fails the batch instead of passing unnoticed.
412+
auto bad_defaults = properties;
413+
for (const auto& [key, value] : BadS3Credentials()) {
414+
bad_defaults.insert_or_assign(key, value);
415+
}
416+
ICEBERG_UNWRAP_OR_FAIL(auto io, MakeS3FileIO(std::move(bad_defaults)));
417+
auto* credentialed = io->AsSupportsStorageCredentials();
418+
ASSERT_NE(credentialed, nullptr);
419+
420+
const auto a = ObjectUri("delete_files_a/");
421+
const auto b = ObjectUri("delete_files_b/");
422+
ASSERT_THAT(credentialed->SetStorageCredentials({{.prefix = a, .config = properties},
423+
{.prefix = b, .config = properties}}),
424+
IsOk());
425+
426+
const std::vector<std::string> paths = {a + "first", b + "only", a + "second"};
427+
for (const auto& path : paths) {
428+
ASSERT_THAT(io->WriteFile(path, "payload"), IsOk());
429+
}
430+
ASSERT_THAT(io->DeleteFiles(paths), IsOk());
431+
// Deleting a missing key succeeds on S3, so check the objects are gone. The
432+
// writes above authenticated on these paths, so a failed read means absent.
433+
for (const auto& path : paths) {
434+
EXPECT_FALSE(io->ReadFile(path, std::nullopt).has_value()) << path;
435+
}
436+
}
437+
438+
TEST_F(ArrowS3FileIOTest, InputFileOutlivesCredentialInstall) {
439+
if (!HasIntegrationEnv()) {
440+
GTEST_SKIP() << "Set ICEBERG_TEST_S3_URI to enable S3 IO test";
441+
}
442+
443+
auto properties = PropertiesFromEnv();
444+
if (!properties.contains(std::string(S3Properties::kAccessKeyId)) ||
445+
!properties.contains(std::string(S3Properties::kSecretAccessKey))) {
446+
GTEST_SKIP() << "Set AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY to enable "
447+
"credential routing test";
448+
}
449+
450+
auto bad_defaults = properties;
451+
for (const auto& [key, value] : BadS3Credentials()) {
452+
bad_defaults.insert_or_assign(key, value);
453+
}
454+
ICEBERG_UNWRAP_OR_FAIL(auto io, MakeS3FileIO(std::move(bad_defaults)));
455+
auto* credentialed = io->AsSupportsStorageCredentials();
456+
ASSERT_NE(credentialed, nullptr);
457+
458+
const auto prefix = ObjectUri("retained_input/");
459+
const std::vector<StorageCredential> credentials = {
460+
{.prefix = prefix, .config = properties}};
461+
ASSERT_THAT(credentialed->SetStorageCredentials(credentials), IsOk());
462+
463+
const auto uri = prefix + "object";
464+
constexpr std::string_view kContent = "written before the credential install";
465+
ASSERT_THAT(io->WriteFile(uri, kContent), IsOk());
466+
ICEBERG_UNWRAP_OR_FAIL(auto file, io->NewInputFile(uri));
467+
ICEBERG_UNWRAP_OR_FAIL(auto opened_before, file->Open());
468+
469+
// Retires the delegate `file` came from; both handles must still read.
470+
ASSERT_THAT(credentialed->SetStorageCredentials(credentials), IsOk());
471+
472+
ICEBERG_UNWRAP_OR_FAIL(auto opened_after, file->Open());
473+
for (auto* stream : {opened_before.get(), opened_after.get()}) {
474+
std::string read(kContent.size(), '\0');
475+
EXPECT_THAT(stream->ReadFully(0, std::as_writable_bytes(std::span(read))), IsOk());
476+
EXPECT_EQ(read, kContent);
477+
EXPECT_THAT(stream->Close(), IsOk());
478+
}
479+
EXPECT_THAT(io->DeleteFile(uri), IsOk());
480+
}
481+
409482
#if ICEBERG_S3_ENABLED
410483
TEST_F(ArrowS3FileIOTest, ClientRegion) {
411484
auto result =

‎src/iceberg/test/resolving_file_io_test.cc‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,13 @@
1919

2020
#include "iceberg/resolving_file_io.h"
2121

22+
#include <chrono>
23+
#include <condition_variable>
24+
#include <future>
2225
#include <memory>
26+
#include <mutex>
2327
#include <string>
28+
#include <thread>
2429
#include <unordered_map>
2530
#include <vector>
2631

@@ -254,4 +259,63 @@ TEST(ResolvingFileIOTest, ForwardsAllCredentialsToResolvedImplementations) {
254259
ASSERT_NE(last_local_io, nullptr);
255260
}
256261

262+
TEST(ResolvingFileIOTest, LoadsWithoutTheLockAndDropsStaleDelegates) {
263+
// The first load parks until new credentials are installed. Loading outside
264+
// the lock lets that install proceed; what the parked load built then carries
265+
// stale credentials and must be redone rather than cached.
266+
// Static because registry factories are process-global; reset for reruns.
267+
static std::mutex gate;
268+
static std::condition_variable cv;
269+
static bool parked;
270+
static bool resume;
271+
static int calls;
272+
static RecordingCredentialedFileIO* last;
273+
parked = resume = false;
274+
calls = 0;
275+
last = nullptr;
276+
FileIORegistry::Register(
277+
"test.file-io.slow",
278+
{.create =
279+
[](const FileIORegistry::Properties&) -> Result<std::unique_ptr<FileIO>> {
280+
if (++calls == 1) {
281+
std::unique_lock lock(gate);
282+
parked = true;
283+
cv.notify_all();
284+
cv.wait(lock, [] { return resume; });
285+
}
286+
auto io = std::make_unique<RecordingCredentialedFileIO>();
287+
last = io.get();
288+
return io;
289+
},
290+
.accepts = [](std::string_view scheme) { return scheme == "slow"; }});
291+
292+
ResolvingFileIO io({});
293+
const std::vector<StorageCredential> stale = {
294+
{.prefix = "slow", .config = {{"k", "1"}}}};
295+
const std::vector<StorageCredential> fresh = {
296+
{.prefix = "slow", .config = {{"k", "2"}}}};
297+
ASSERT_THAT(io.SetStorageCredentials(stale), IsOk());
298+
299+
std::thread reader([&] { (void)io.NewInputFile("slow://bucket/file"); });
300+
{
301+
std::unique_lock lock(gate);
302+
cv.wait(lock, [] { return parked; });
303+
}
304+
auto install =
305+
std::async(std::launch::async, [&] { return io.SetStorageCredentials(fresh); });
306+
const bool install_waited =
307+
install.wait_for(std::chrono::seconds(10)) != std::future_status::ready;
308+
{
309+
std::lock_guard lock(gate);
310+
resume = true;
311+
}
312+
cv.notify_all();
313+
reader.join();
314+
315+
EXPECT_FALSE(install_waited) << "a credential install waited on an in-flight load";
316+
EXPECT_THAT(install.get(), IsOk());
317+
ASSERT_EQ(calls, 2);
318+
EXPECT_EQ(last->credentials(), fresh);
319+
}
320+
257321
} // namespace iceberg

0 commit comments

Comments
 (0)