wgtmac commented on code in PR #898:
URL: https://github.com/apache/iceberg-cpp/pull/898#discussion_r4091927830


##########
src/iceberg/arrow/s3/arrow_s3_file_io.cc:
##########
@@ -209,27 +212,67 @@ class ArrowS3FileIO final : public FileIO, public 
SupportsStorageCredentials {
   Status SetStorageCredentials(
       const std::vector<StorageCredential>& storage_credentials) override;
 
-  const std::vector<StorageCredential>& credentials() const override {
+  std::vector<StorageCredential> credentials() const override {
+    std::shared_lock lock(mutex_);
     return storage_credentials_;
   }
 
   SupportsStorageCredentials* AsSupportsStorageCredentials() override { return 
this; }
 
  private:
-  ArrowFileSystemFileIO& FileIOForPath(std::string_view location);
-
-  ArrowFileSystemFileIO default_file_io_;
+  /// \brief Delegate serving `location`, pinned by the caller against a
+  /// concurrent credential install.
+  std::shared_ptr<ArrowFileSystemFileIO> FileIOForPath(std::string_view 
location);
+
+  using DelegatesByPrefix =
+      std::vector<std::pair<std::string, 
std::shared_ptr<ArrowFileSystemFileIO>>>;
+
+  /// \brief Longest-prefix match against one consistent view of the delegates.
+  static std::shared_ptr<ArrowFileSystemFileIO> MatchDelegate(
+      const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
+      const DelegatesByPrefix& by_prefix, std::string_view location);
+
+  /// \brief Build a delegate for each credential this FileIO can serve.
+  ///
+  /// Lock-free on purpose: building an S3 client can reach out to discover a
+  /// bucket region, which would stall every concurrent operation. Reads no
+  /// mutable member state.
+  Result<DelegatesByPrefix> BuildDelegates(
+      const std::vector<StorageCredential>& storage_credentials) const;
+
+  /// \brief Swap in credentials and delegates, handing back the retired ones.
+  ///
+  /// Callers must hold `mutex_` exclusively and let the returned generation
+  /// destruct only after releasing it: tearing down an S3 client can block on
+  /// in-flight requests, which would stall every operation.
+  void InstallCredentials(std::vector<StorageCredential>& storage_credentials,
+                          DelegatesByPrefix& delegates);
+
+  std::shared_ptr<ArrowFileSystemFileIO> default_file_io_;
   std::unordered_map<std::string, std::string> default_properties_;
+  // Guards everything below; shared because reads happen per file operation.
+  mutable std::shared_mutex mutex_;
   std::vector<StorageCredential> storage_credentials_;
-  std::vector<std::pair<std::string, std::unique_ptr<ArrowFileSystemFileIO>>>
-      file_io_by_prefix_;
+  DelegatesByPrefix file_io_by_prefix_;
 };
 
 Status ArrowS3FileIO::SetStorageCredentials(
     const std::vector<StorageCredential>& storage_credentials) {
-  std::vector<std::pair<std::string, std::unique_ptr<ArrowFileSystemFileIO>>>
-      file_io_by_prefix;
-  file_io_by_prefix.reserve(storage_credentials.size());
+  ICEBERG_ASSIGN_OR_RAISE(auto delegates, BuildDelegates(storage_credentials));

Review Comment:
   REST installs credentials only at FileIO creation. C++ has no refresh path. 
What production path reinstalls them on a live S3 FileIO?



##########
src/iceberg/arrow/s3/arrow_s3_file_io.cc:
##########
@@ -209,27 +212,67 @@ class ArrowS3FileIO final : public FileIO, public 
SupportsStorageCredentials {
   Status SetStorageCredentials(
       const std::vector<StorageCredential>& storage_credentials) override;
 
-  const std::vector<StorageCredential>& credentials() const override {
+  std::vector<StorageCredential> credentials() const override {
+    std::shared_lock lock(mutex_);
     return storage_credentials_;
   }
 
   SupportsStorageCredentials* AsSupportsStorageCredentials() override { return 
this; }
 
  private:
-  ArrowFileSystemFileIO& FileIOForPath(std::string_view location);
-
-  ArrowFileSystemFileIO default_file_io_;
+  /// \brief Delegate serving `location`, pinned by the caller against a
+  /// concurrent credential install.
+  std::shared_ptr<ArrowFileSystemFileIO> FileIOForPath(std::string_view 
location);
+
+  using DelegatesByPrefix =
+      std::vector<std::pair<std::string, 
std::shared_ptr<ArrowFileSystemFileIO>>>;
+
+  /// \brief Longest-prefix match against one consistent view of the delegates.
+  static std::shared_ptr<ArrowFileSystemFileIO> MatchDelegate(
+      const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
+      const DelegatesByPrefix& by_prefix, std::string_view location);
+
+  /// \brief Build a delegate for each credential this FileIO can serve.
+  ///
+  /// Lock-free on purpose: building an S3 client can reach out to discover a

Review Comment:
   ResolvingFileIO still builds the first S3 client under its lock. Region 
discovery can block other operations. Can we move it out?



##########
src/iceberg/arrow/s3/arrow_s3_file_io.cc:
##########
@@ -244,62 +287,91 @@ Status ArrowS3FileIO::SetStorageCredentials(
       properties[key] = value;
     }
     ICEBERG_ASSIGN_OR_RAISE(auto fs, BuildArrowS3FileSystem(properties));
-    file_io_by_prefix.emplace_back(
-        CanonicalizeS3Scheme(credential.prefix),
-        std::make_unique<ArrowFileSystemFileIO>(std::move(fs)));
+    delegates.emplace_back(CanonicalizeS3Scheme(credential.prefix),
+                           
std::make_shared<ArrowFileSystemFileIO>(std::move(fs)));
   }
-  if (file_io_by_prefix.empty() && !storage_credentials.empty()) {
+  if (delegates.empty() && !storage_credentials.empty()) {
     // Silent skipping of every vended credential is hard to diagnose: S3 
access
     // would proceed with the default credentials and fail only at IO time.
     ICEBERG_LOG_WARN(
         "None of the {} vended storage credential(s) has an S3-compatible 
prefix; "
         "S3 access will use the default credentials",
         storage_credentials.size());
   }
-  file_io_by_prefix_ = std::move(file_io_by_prefix);
-  storage_credentials_ = storage_credentials;
-  return {};
+  return delegates;
 }
 
-ArrowFileSystemFileIO& ArrowS3FileIO::FileIOForPath(std::string_view location) 
{
-  if (file_io_by_prefix_.empty()) {
-    return default_file_io_;
+void ArrowS3FileIO::InstallCredentials(
+    std::vector<StorageCredential>& storage_credentials, DelegatesByPrefix& 
delegates) {
+  file_io_by_prefix_.swap(delegates);
+  storage_credentials_.swap(storage_credentials);
+}
+
+std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::MatchDelegate(
+    const std::shared_ptr<ArrowFileSystemFileIO>& fallback,
+    const DelegatesByPrefix& by_prefix, std::string_view location) {
+  if (by_prefix.empty()) {
+    return fallback;
   }
   const std::string canonical = CanonicalizeS3Scheme(location);
-  ArrowFileSystemFileIO* best = &default_file_io_;
+  auto best = fallback;
   size_t best_len = 0;
-  for (const auto& [prefix, file_io] : file_io_by_prefix_) {
+  for (const auto& [prefix, file_io] : by_prefix) {
     if (prefix.size() > best_len && canonical.starts_with(prefix)) {
-      best = file_io.get();
+      best = file_io;
       best_len = prefix.size();
     }
   }
-  return *best;
+  return best;
+}
+
+std::shared_ptr<ArrowFileSystemFileIO> ArrowS3FileIO::FileIOForPath(
+    std::string_view location) {
+  std::shared_lock lock(mutex_);
+  return MatchDelegate(default_file_io_, file_io_by_prefix_, location);
 }
 
 Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(
     std::string file_location) {
-  return FileIOForPath(file_location).NewInputFile(std::move(file_location));
+  return FileIOForPath(file_location)->NewInputFile(std::move(file_location));
 }
 
 Result<std::unique_ptr<InputFile>> ArrowS3FileIO::NewInputFile(std::string 
file_location,
                                                                size_t length) {
-  return FileIOForPath(file_location).NewInputFile(std::move(file_location), 
length);
+  return FileIOForPath(file_location)->NewInputFile(std::move(file_location), 
length);
 }
 
 Result<std::unique_ptr<OutputFile>> ArrowS3FileIO::NewOutputFile(
     std::string file_location) {
-  return FileIOForPath(file_location).NewOutputFile(std::move(file_location));
+  return FileIOForPath(file_location)->NewOutputFile(std::move(file_location));
 }
 
 Status ArrowS3FileIO::DeleteFile(const std::string& file_location) {
-  return FileIOForPath(file_location).DeleteFile(file_location);
+  return FileIOForPath(file_location)->DeleteFile(file_location);
 }
 
 Status ArrowS3FileIO::DeleteFiles(const std::vector<std::string>& 
file_locations) {
-  std::unordered_map<ArrowFileSystemFileIO*, std::vector<std::string>> 
locations_by_io;
+  // 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<ArrowFileSystemFileIO> fallback;
+  DelegatesByPrefix by_prefix;
+  {
+    std::shared_lock lock(mutex_);
+    fallback = default_file_io_;
+    by_prefix = file_io_by_prefix_;
+  }
+  std::vector<std::pair<std::shared_ptr<ArrowFileSystemFileIO>, 
std::vector<std::string>>>

Review Comment:
   This retains the pre-existing fail-fast behavior. Java attempts all batches 
and counts failures. Is matching that behavior in scope?



##########
src/iceberg/test/arrow_s3_file_io_test.cc:
##########
@@ -237,6 +239,69 @@ TEST_F(ArrowS3FileIOTest, WarnsWhenNoCredentialApplies) {
   EXPECT_TRUE(HasWarning(*logger));
 }
 
+TEST_F(ArrowS3FileIOTest, DeleteFilesDispatchesAcrossCredentialPrefixes) {
+  auto result = MakeS3FileIO({});
+  ASSERT_THAT(result, IsOk());
+  auto* credentialed = result.value()->AsSupportsStorageCredentials();
+  ASSERT_NE(credentialed, nullptr);
+
+  auto credential = [](std::string_view prefix, std::string_view access_key) {
+    return StorageCredential{
+        .prefix = std::string(prefix),
+        .config = {{std::string(S3Properties::kAccessKeyId), 
std::string(access_key)},
+                   {std::string(S3Properties::kSecretAccessKey), "secret"}}};
+  };
+  ASSERT_THAT(credentialed->SetStorageCredentials({credential("s3://bucket-a", 
"key-a"),
+                                                   credential("s3://bucket-b", 
"key-b")}),
+              IsOk());
+
+  auto status = result.value()->DeleteFiles({"s3://bucket-a/%ZZ.parquet",

Review Comment:
   The first URI fails parsing, so bucket-b is never reached. Can this test 
verify that both prefixes receive deletes?



##########
src/iceberg/test/arrow_s3_file_io_test.cc:
##########
@@ -237,6 +239,69 @@ TEST_F(ArrowS3FileIOTest, WarnsWhenNoCredentialApplies) {
   EXPECT_TRUE(HasWarning(*logger));
 }
 
+TEST_F(ArrowS3FileIOTest, DeleteFilesDispatchesAcrossCredentialPrefixes) {
+  auto result = MakeS3FileIO({});
+  ASSERT_THAT(result, IsOk());
+  auto* credentialed = result.value()->AsSupportsStorageCredentials();
+  ASSERT_NE(credentialed, nullptr);
+
+  auto credential = [](std::string_view prefix, std::string_view access_key) {
+    return StorageCredential{
+        .prefix = std::string(prefix),
+        .config = {{std::string(S3Properties::kAccessKeyId), 
std::string(access_key)},
+                   {std::string(S3Properties::kSecretAccessKey), "secret"}}};
+  };
+  ASSERT_THAT(credentialed->SetStorageCredentials({credential("s3://bucket-a", 
"key-a"),
+                                                   credential("s3://bucket-b", 
"key-b")}),
+              IsOk());
+
+  auto status = result.value()->DeleteFiles({"s3://bucket-a/%ZZ.parquet",
+                                             "s3://bucket-a/second.parquet",
+                                             "s3://bucket-b/other.parquet"});
+  EXPECT_THAT(status, HasErrorMessage("Cannot parse URI"));
+}
+
+TEST_F(ArrowS3FileIOTest, OperationsSurviveConcurrentCredentialInstalls) {
+  auto result = MakeS3FileIO({});
+  ASSERT_THAT(result, IsOk());
+  auto* credentialed = result.value()->AsSupportsStorageCredentials();
+  ASSERT_NE(credentialed, nullptr);
+
+  auto credential = [](std::string_view access_key) {
+    return StorageCredential{
+        .prefix = "s3://bucket",
+        .config = {{std::string(S3Properties::kAccessKeyId), 
std::string(access_key)},
+                   {std::string(S3Properties::kSecretAccessKey), "secret"}}};
+  };
+  ASSERT_THAT(credentialed->SetStorageCredentials({credential("first")}), 
IsOk());
+
+  std::atomic<bool> stop = false;
+  std::atomic<int> failures = 0;
+  std::vector<std::thread> operations;
+  operations.reserve(4);
+  for (int i = 0; i < 4; ++i) {
+    operations.emplace_back([&] {
+      while (!stop.load()) {
+        if (!result.value()->NewInputFile("s3://bucket/key").has_value()) {

Review Comment:
   NewInputFile only creates a wrapper. Could this test hold a file across an 
install and perform an actual read?



##########
src/iceberg/arrow/s3/arrow_s3_file_io.cc:
##########
@@ -192,7 +195,7 @@ class ArrowS3FileIO final : public FileIO, public 
SupportsStorageCredentials {
  public:
   ArrowS3FileIO(std::shared_ptr<::arrow::fs::FileSystem> arrow_fs,
                 std::unordered_map<std::string, std::string> 
default_properties)
-      : default_file_io_(std::move(arrow_fs)),
+      : 
default_file_io_(std::make_shared<ArrowFileSystemFileIO>(std::move(arrow_fs))),

Review Comment:
   Why do we need to change this?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to