Skip to content

Commit 896b3ce

Browse files
refactor(io): hand out S3 delegates by shared_ptr and build them off the lock
1 parent 307adb1 commit 896b3ce

7 files changed

Lines changed: 182 additions & 40 deletions

File tree

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

Lines changed: 102 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,12 @@
1717
* under the License.
1818
*/
1919

20+
#include <algorithm>
2021
#include <cstdlib>
2122
#include <memory>
23+
#include <mutex>
2224
#include <optional>
25+
#include <shared_mutex>
2326
#include <string>
2427
#include <string_view>
2528
#include <unordered_map>
@@ -192,7 +195,7 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials {
192195
public:
193196
ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs,
194197
std::unordered_map<std::string, std::string> default_properties)
195-
: default_file_io_(std::move(arrow_fs)),
198+
: default_file_io_(std::make_shared<ArrowFileSystemFileIO>(std::move(arrow_fs))),
196199
default_properties_(std::move(default_properties)) {}
197200

198201
Result<std::unique_ptr<InputFile>> NewInputFile(std::string file_location) override;
@@ -209,27 +212,67 @@ class ArrowS3FileIO final : public FileIO, public SupportsStorageCredentials {
209212
Status SetStorageCredentials(
210213
const std::vector<StorageCredential>& storage_credentials) override;
211214

212-
const std::vector<StorageCredential>& credentials() const override {
215+
std::vector<StorageCredential> credentials() const override {
216+
std::shared_lock lock(mutex_);
213217
return storage_credentials_;
214218
}
215219

216220
SupportsStorageCredentials* AsSupportsStorageCredentials() override { return this; }
217221

218222
private:
219-
ArrowFileSystemFileIO& FileIOForPath(std::string_view location);
220-
221-
ArrowFileSystemFileIO default_file_io_;
223+
/// \brief Delegate serving `location`, pinned by the caller against a
224+
/// concurrent credential install.
225+
std::shared_ptr<ArrowFileSystemFileIO> FileIOForPath(std::string_view location);
226+
227+
using DelegatesByPrefix =
228+
std::vector<std::pair<std::string, std::shared_ptr<ArrowFileSystemFileIO>>>;
229+
230+
/// \brief Longest-prefix match against one consistent view of the delegates.
231+
static std::shared_ptr<ArrowFileSystemFileIO> MatchDelegate(
232+
const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
233+
const DelegatesByPrefix& by_prefix, std::string_view location);
234+
235+
/// \brief Build a delegate for each credential this FileIO can serve.
236+
///
237+
/// Lock-free on purpose: building an S3 client can reach out to discover a
238+
/// bucket region, which would stall every concurrent operation. Reads no
239+
/// mutable member state.
240+
Result<DelegatesByPrefix> BuildDelegates(
241+
const std::vector<StorageCredential>& storage_credentials) const;
242+
243+
/// \brief Swap in credentials and delegates, handing back the retired ones.
244+
///
245+
/// Callers must hold `mutex_` exclusively and let the returned generation
246+
/// destruct only after releasing it: tearing down an S3 client can block on
247+
/// in-flight requests, which would stall every operation.
248+
void InstallCredentials(std::vector<StorageCredential>& storage_credentials,
249+
DelegatesByPrefix& delegates);
250+
251+
std::shared_ptr<ArrowFileSystemFileIO> default_file_io_;
222252
std::unordered_map<std::string, std::string> default_properties_;
253+
// Guards everything below; shared because reads happen per file operation.
254+
mutable std::shared_mutex mutex_;
223255
std::vector<StorageCredential> storage_credentials_;
224-
std::vector<std::pair<std::string, std::unique_ptr<ArrowFileSystemFileIO>>>
225-
file_io_by_prefix_;
256+
DelegatesByPrefix file_io_by_prefix_;
226257
};
227258

228259
Status ArrowS3FileIO::SetStorageCredentials(
229260
const std::vector<StorageCredential>& storage_credentials) {
230-
std::vector<std::pair<std::string, std::unique_ptr<ArrowFileSystemFileIO>>>
231-
file_io_by_prefix;
232-
file_io_by_prefix.reserve(storage_credentials.size());
261+
ICEBERG_ASSIGN_OR_RAISE(auto delegates, BuildDelegates(storage_credentials));
262+
auto credentials = storage_credentials;
263+
{
264+
std::unique_lock lock(mutex_);
265+
InstallCredentials(credentials, delegates);
266+
}
267+
// `credentials` and `delegates` now hold the retired generation and destruct
268+
// here, outside the lock.
269+
return {};
270+
}
271+
272+
Result<ArrowS3FileIO::DelegatesByPrefix> ArrowS3FileIO::BuildDelegates(
273+
const std::vector<StorageCredential>& storage_credentials) const {
274+
DelegatesByPrefix delegates;
275+
delegates.reserve(storage_credentials.size());
233276
// TODO(gangwu): Refresh vended credentials via credentials.uri before tokens expire.
234277
for (const auto& credential : storage_credentials) {
235278
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
@@ -244,62 +287,91 @@ Status ArrowS3FileIO::SetStorageCredentials(
244287
properties[key] = value;
245288
}
246289
ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties));
247-
file_io_by_prefix.emplace_back(
248-
CanonicalizeS3Scheme(credential.prefix),
249-
std::make_unique<ArrowFileSystemFileIO>(std::move(fs)));
290+
delegates.emplace_back(CanonicalizeS3Scheme(credential.prefix),
291+
std::make_shared<ArrowFileSystemFileIO>(std::move(fs)));
250292
}
251-
if (file_io_by_prefix.empty() && !storage_credentials.empty()) {
293+
if (delegates.empty() && !storage_credentials.empty()) {
252294
// Silent skipping of every vended credential is hard to diagnose: S3 access
253295
// would proceed with the default credentials and fail only at IO time.
254296
ICEBERG_LOG_WARN(
255297
"None of the {} vended storage credential(s) has an S3-compatible prefix; "
256298
"S3 access will use the default credentials",
257299
storage_credentials.size());
258300
}
259-
file_io_by_prefix_ = std::move(file_io_by_prefix);
260-
storage_credentials_ = storage_credentials;
261-
return {};
301+
return delegates;
262302
}
263303

264-
ArrowFileSystemFileIO& ArrowS3FileIO::FileIOForPath(std::string_view location) {
265-
if (file_io_by_prefix_.empty()) {
266-
return default_file_io_;
304+
void ArrowS3FileIO::InstallCredentials(
305+
std::vector<StorageCredential>& storage_credentials, DelegatesByPrefix& delegates) {
306+
file_io_by_prefix_.swap(delegates);
307+
storage_credentials_.swap(storage_credentials);
308+
}
309+
310+
std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::MatchDelegate(
311+
const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
312+
const DelegatesByPrefix& by_prefix, std::string_view location) {
313+
if (by_prefix.empty()) {
314+
return fallback;
267315
}
268316
const std::string canonical = CanonicalizeS3Scheme(location);
269-
ArrowFileSystemFileIO* best = &default_file_io_;
317+
auto best = fallback;
270318
size_t best_len = 0;
271-
for (const auto& [prefix, file_io] : file_io_by_prefix_) {
319+
for (const auto& [prefix, file_io] : by_prefix) {
272320
if (prefix.size() > best_len && canonical.starts_with(prefix)) {
273-
best = file_io.get();
321+
best = file_io;
274322
best_len = prefix.size();
275323
}
276324
}
277-
return *best;
325+
return best;
326+
}
327+
328+
std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::FileIOForPath(
329+
std::string_view location) {
330+
std::shared_lock lock(mutex_);
331+
return MatchDelegate(default_file_io_, file_io_by_prefix_, location);
278332
}
279333

280334
Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(
281335
std::string file_location) {
282-
return FileIOForPath(file_location).NewInputFile(std::move(file_location));
336+
return FileIOForPath(file_location)->NewInputFile(std::move(file_location));
283337
}
284338

285339
Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(std::string file_location,
286340
size_t length) {
287-
return FileIOForPath(file_location).NewInputFile(std::move(file_location), length);
341+
return FileIOForPath(file_location)->NewInputFile(std::move(file_location), length);
288342
}
289343

290344
Result<std::unique_ptr<OutputFile>> ArrowS3FileIO::NewOutputFile(
291345
std::string file_location) {
292-
return FileIOForPath(file_location).NewOutputFile(std::move(file_location));
346+
return FileIOForPath(file_location)->NewOutputFile(std::move(file_location));
293347
}
294348

295349
Status ArrowS3FileIO::DeleteFile(const std::string& file_location) {
296-
return FileIOForPath(file_location).DeleteFile(file_location);
350+
return FileIOForPath(file_location)->DeleteFile(file_location);
297351
}
298352

299353
Status ArrowS3FileIO::DeleteFiles(const std::vector<std::string>& file_locations) {
300-
std::unordered_map<ArrowFileSystemFileIO*, std::vector<std::string>> locations_by_io;
354+
// One snapshot so the whole batch matches the same delegate generation; only
355+
// ever a handful of delegates, so a linear scan beats hashing.
356+
std::shared_ptr<ArrowFileSystemFileIO> fallback;
357+
DelegatesByPrefix by_prefix;
358+
{
359+
std::shared_lock lock(mutex_);
360+
fallback = default_file_io_;
361+
by_prefix = file_io_by_prefix_;
362+
}
363+
std::vector<std::pair<std::shared_ptr<ArrowFileSystemFileIO>, std::vector<std::string>>>
364+
locations_by_io;
301365
for (const auto& file_location : file_locations) {
302-
locations_by_io[&FileIOForPath(file_location)].push_back(file_location);
366+
auto file_io = MatchDelegate(fallback, by_prefix, file_location);
367+
auto it = std::ranges::find_if(
368+
locations_by_io, [&](const auto& entry) { return entry.first == file_io; });
369+
if (it == locations_by_io.end()) {
370+
locations_by_io.emplace_back(std::move(file_io),
371+
std::vector<std::string>{file_location});
372+
} else {
373+
it->second.push_back(file_location);
374+
}
303375
}
304376
for (auto& [file_io, locations] : locations_by_io) {
305377
ICEBERG_RETURN_UNEXPECTED(file_io->DeleteFiles(locations));

‎src/iceberg/file_io.h‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -193,8 +193,11 @@ class ICEBERG_EXPORT SupportsStorageCredentials {
193193
virtual Status SetStorageCredentials(
194194
const std::vector<StorageCredential>& storage_credentials) = 0;
195195

196-
/// \brief Return currently installed storage credentials.
197-
virtual const std::vector<StorageCredential>& credentials() const = 0;
196+
/// \brief Return the storage credentials this FileIO holds.
197+
///
198+
/// By value because a concurrent install may replace them. An implementation
199+
/// that delegates may report what was installed on it.
200+
virtual std::vector<StorageCredential> credentials() const = 0;
198201
};
199202

200203
} // namespace iceberg

‎src/iceberg/resolving_file_io.cc‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,7 +109,8 @@ Status ResolvingFileIO::SetStorageCredentials(
109109
return {};
110110
}
111111

112-
const std::vector<StorageCredential>& ResolvingFileIO::credentials() const {
112+
std::vector<StorageCredential> ResolvingFileIO::credentials() const {
113+
std::shared_lock lock(mutex_);
113114
return storage_credentials_;
114115
}
115116

‎src/iceberg/resolving_file_io.h‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,9 @@
3838
namespace iceberg {
3939

4040
/// \brief FileIO that resolves and caches implementations by registry name.
41+
///
42+
/// Vended credentials are forwarded to every resolved implementation that
43+
/// supports them; each applies what it understands.
4144
class ICEBERG_EXPORT ResolvingFileIO final : public FileIO,
4245
public SupportsStorageCredentials {
4346
public:
@@ -58,7 +61,7 @@ class ICEBERG_EXPORT ResolvingFileIO final : public FileIO,
5861
Status SetStorageCredentials(
5962
const std::vector<StorageCredential>& storage_credentials) override;
6063

61-
const std::vector<StorageCredential>& credentials() const override;
64+
std::vector<StorageCredential> credentials() const override;
6265

6366
SupportsStorageCredentials* AsSupportsStorageCredentials() override { return this; }
6467

@@ -67,8 +70,8 @@ class ICEBERG_EXPORT ResolvingFileIO final : public FileIO,
6770
Result<std::shared_ptr<FileIO>> FileIOForPath(std::string_view location);
6871

6972
std::unordered_map<std::string, std::string> properties_;
70-
// Guards lazy resolution and credential refresh.
71-
std::shared_mutex mutex_;
73+
// Guards lazy resolution and credential state.
74+
mutable std::shared_mutex mutex_;
7275
std::vector<StorageCredential> storage_credentials_;
7376
std::unordered_map<std::string, std::shared_ptr<FileIO>, StringHash, StringEqual>
7477
io_by_name_;

‎src/iceberg/test/arrow_s3_file_io_test.cc‎

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,14 @@
1818
*/
1919

2020
#include <algorithm>
21+
#include <atomic>
2122
#include <cstdlib>
2223
#include <iostream>
2324
#include <memory>
2425
#include <optional>
2526
#include <string>
2627
#include <string_view>
28+
#include <thread>
2729
#include <unordered_map>
2830
#include <utility>
2931
#include <vector>
@@ -237,6 +239,69 @@ TEST_F(ArrowS3FileIOTest, WarnsWhenNoCredentialApplies) {
237239
EXPECT_TRUE(HasWarning(*logger));
238240
}
239241

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+
264+
TEST_F(ArrowS3FileIOTest, OperationsSurviveConcurrentCredentialInstalls) {
265+
auto result = MakeS3FileIO({});
266+
ASSERT_THAT(result, IsOk());
267+
auto* credentialed = result.value()->AsSupportsStorageCredentials();
268+
ASSERT_NE(credentialed, nullptr);
269+
270+
auto credential = [](std::string_view access_key) {
271+
return StorageCredential{
272+
.prefix = "s3://bucket",
273+
.config = {{std::string(S3Properties::kAccessKeyId), std::string(access_key)},
274+
{std::string(S3Properties::kSecretAccessKey), "secret"}}};
275+
};
276+
ASSERT_THAT(credentialed->SetStorageCredentials({credential("first")}), IsOk());
277+
278+
std::atomic<bool> stop = false;
279+
std::atomic<int> failures = 0;
280+
std::vector<std::thread> operations;
281+
operations.reserve(4);
282+
for (int i = 0; i < 4; ++i) {
283+
operations.emplace_back([&] {
284+
while (!stop.load()) {
285+
if (!result.value()->NewInputFile("s3://bucket/key").has_value()) {
286+
++failures;
287+
}
288+
}
289+
});
290+
}
291+
// No assertions until the threads are joined: a fatal assertion here would
292+
// destroy joinable threads and terminate the binary, masking the failure.
293+
Status install_status = {};
294+
for (int round = 0; round < 3 && install_status.has_value(); ++round) {
295+
install_status = credentialed->SetStorageCredentials({credential("replacement")});
296+
}
297+
stop = true;
298+
for (auto& operation : operations) {
299+
operation.join();
300+
}
301+
ASSERT_THAT(install_status, IsOk());
302+
EXPECT_EQ(failures, 0);
303+
}
304+
240305
TEST_F(ArrowS3FileIOTest, RejectsIncompleteStaticCredentials) {
241306
auto result =
242307
MakeS3FileIO({{std::string(S3Properties::kAccessKeyId), "access-key-only"}});

‎src/iceberg/test/resolving_file_io_test.cc‎

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -60,9 +60,7 @@ class RecordingCredentialedFileIO : public RecordingFileIO,
6060
return {};
6161
}
6262

63-
const std::vector<StorageCredential>& credentials() const override {
64-
return credentials_;
65-
}
63+
std::vector<StorageCredential> credentials() const override { return credentials_; }
6664

6765
SupportsStorageCredentials* AsSupportsStorageCredentials() override { return this; }
6866

‎src/iceberg/test/rest_file_io_test.cc‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ class MockCredentialedFileIO : public MockFileIO, public SupportsStorageCredenti
6161
return {};
6262
}
6363

64-
const std::vector<StorageCredential>& credentials() const override {
64+
std::vector<StorageCredential> credentials() const override {
6565
return captured_storage_credentials;
6666
}
6767

0 commit comments

Comments
 (0)