Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 33 additions & 9 deletions src/iceberg/catalog/rest/json_serde.cc
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,25 @@ Result<StorageCredential> StorageCredentialFromJson(const nlohmann::json& json)
return credential;
}

/// \brief Reads the optional `storage-credentials` array shared by the
/// LoadTable and LoadCredentials responses.
Result<std::vector<StorageCredential>> StorageCredentialsFromJson(
const nlohmann::json& json) {
std::vector<StorageCredential> credentials;
auto it = json.find(kStorageCredentials);
if (it == json.end() || it->is_null()) {
return credentials;
}
if (!it->is_array()) {
return JsonParseError("Cannot parse storage credentials from non-array");
}
for (const auto& entry : *it) {
ICEBERG_ASSIGN_OR_RAISE(auto credential, StorageCredentialFromJson(entry));
credentials.push_back(std::move(credential));
}
return credentials;
}

template <typename Value>
Result<std::map<int32_t, Value>> KeyValueMapFromJson(const nlohmann::json& json,
std::string_view key) {
Expand Down Expand Up @@ -738,19 +757,24 @@ Result<LoadTableResult> LoadTableResultFromJson(const nlohmann::json& json) {
ICEBERG_ASSIGN_OR_RAISE(result.metadata, TableMetadataFromJson(metadata_json));
ICEBERG_ASSIGN_OR_RAISE(result.config,
GetJsonValueOrDefault<decltype(result.config)>(json, kConfig));
if (auto it = json.find(kStorageCredentials); it != json.end() && !it->is_null()) {
if (!it->is_array()) {
return JsonParseError("Cannot parse storage credentials from non-array");
}
for (const auto& entry : *it) {
ICEBERG_ASSIGN_OR_RAISE(auto cred, StorageCredentialFromJson(entry));
result.storage_credentials.push_back(std::move(cred));
}
}
ICEBERG_ASSIGN_OR_RAISE(result.storage_credentials, StorageCredentialsFromJson(json));
ICEBERG_RETURN_UNEXPECTED(result.Validate());
return result;
}

Result<LoadCredentialsResponse> LoadCredentialsResponseFromJson(
const nlohmann::json& json) {
// Required here, unlike in LoadTable: reading a malformed response as "no
// credentials" would look like a refresh that succeeded and dropped them.
if (auto it = json.find(kStorageCredentials); it == json.end() || it->is_null()) {
return JsonParseError("Missing '{}'", kStorageCredentials);
}
LoadCredentialsResponse response;
ICEBERG_ASSIGN_OR_RAISE(response.storage_credentials, StorageCredentialsFromJson(json));
ICEBERG_RETURN_UNEXPECTED(response.Validate());
return response;
}

nlohmann::json ToJson(const ListNamespacesResponse& response) {
nlohmann::json json;
SetOptionalStringField(json, kNextPageToken, response.next_page_token);
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/catalog/rest/json_serde_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,10 @@ template <>
ICEBERG_REST_EXPORT Result<LoadTableResult> FromJson(const nlohmann::json& json);
ICEBERG_REST_EXPORT Result<nlohmann::json> ToJson(const LoadTableResult& model);

// Response-only model: a client never serializes it, so no ToJson.
ICEBERG_REST_EXPORT Result<LoadCredentialsResponse> LoadCredentialsResponseFromJson(
const nlohmann::json& json);

ICEBERG_REST_EXPORT Result<CreateTableRequest> CreateTableRequestFromJson(
const nlohmann::json& json);
template <>
Expand Down
87 changes: 80 additions & 7 deletions src/iceberg/catalog/rest/rest_catalog.cc
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
#include "iceberg/catalog/rest/rest_catalog.h"

#include <memory>
#include <mutex>
#include <string>
#include <string_view>
#include <tuple>
#include <unordered_map>
Expand All @@ -42,6 +44,7 @@
#include "iceberg/catalog/rest/rest_util.h"
#include "iceberg/catalog/rest/types.h"
#include "iceberg/json_serde_internal.h"
#include "iceberg/logging/log_macros.h"
#include "iceberg/metrics/metrics_reporters.h"
#include "iceberg/partition_spec.h"
#include "iceberg/result.h"
Expand All @@ -60,6 +63,41 @@ namespace iceberg::rest {

namespace {

class RestStorageCredentialProvider final : public StorageCredentialProvider {
public:
RestStorageCredentialProvider(std::shared_ptr<HttpClient> client,
std::shared_ptr<auth::AuthSession> session,
std::string path)
: client_(std::move(client)),
session_(std::move(session)),
path_(std::move(path)) {}

Result<std::vector<StorageCredential>> Load() override {
// The provider can be shared by several FileIOs; keep fetches serialized.
std::lock_guard lock(mutex_);
ICEBERG_ASSIGN_OR_RAISE(const auto response,
client_->Get(path_, /*params=*/{}, /*headers=*/{},
*TableErrorHandler::Instance(), *session_));
// Parse errors can contain credential data; return a fixed message.
auto json = FromJsonString(response.body());
if (!json.has_value()) {
return JsonParseError("Malformed LoadCredentials response");
}
auto result = LoadCredentialsResponseFromJson(*json);
if (!result.has_value()) {
return std::unexpected<Error>(
{.kind = result.error().kind, .message = "Malformed LoadCredentials response"});
}
return std::move(result->storage_credentials);
}

private:
std::shared_ptr<HttpClient> client_;
std::shared_ptr<auth::AuthSession> session_;
std::string path_;
std::mutex mutex_;
};

/// \brief Get the default set of endpoints for backwards compatibility according to the
/// iceberg rest spec.
std::unordered_set<Endpoint> GetDefaultEndpoints() {
Expand Down Expand Up @@ -508,12 +546,45 @@ Result<std::shared_ptr<auth::AuthSession>> RestCatalog::TableAuthSession(
std::move(contextual_session));
}

std::shared_ptr<StorageCredentialProvider> RestCatalog::MakeStorageCredentialProvider(
const TableIdentifier& identifier,
std::shared_ptr<auth::AuthSession> table_session) const {
if (!supported_endpoints_.contains(Endpoint::TableCredentials())) {
// Not an error, but it surfaces much later as credentials expiring.
ICEBERG_LOG_DEBUG(
"Catalog does not advertise {}; vended credentials for '{}' will not be "
"refreshed",
Endpoint::TableCredentials().ToString(), ToString(identifier));
return nullptr;
}
auto path = paths_->Credentials(identifier);
if (!path.has_value()) {
ICEBERG_LOG_WARN(
"Cannot build the credentials path for '{}' ({}); its vended credentials "
"will not be refreshed",
ToString(identifier), path.error().message);
return nullptr;
}
auto client = client_;
auto credentials_path = std::move(path.value());
auto session = std::move(table_session);
return std::make_shared<RestStorageCredentialProvider>(
std::move(client), std::move(session), std::move(credentials_path));
}

Result<std::shared_ptr<FileIO>> RestCatalog::TableFileIO(
const SessionContext& /*context*/,
const SessionContext& /*context*/, const TableIdentifier& identifier,
const std::unordered_map<std::string, std::string>& table_config,
const std::vector<StorageCredential>& storage_credentials) const {
const std::vector<StorageCredential>& storage_credentials,
std::shared_ptr<auth::AuthSession> table_session) const {
if (!table_config.empty() || !storage_credentials.empty()) {
return MakeTableFileIO(config_.configs(), table_config, storage_credentials);
// Only vended credentials expire, so only they need a provider.
std::shared_ptr<StorageCredentialProvider> provider;
if (!storage_credentials.empty()) {
provider = MakeStorageCredentialProvider(identifier, std::move(table_session));
}
return MakeTableFileIO(config_.configs(), table_config, storage_credentials,
std::move(provider));
}

return file_io_;
Expand Down Expand Up @@ -772,11 +843,12 @@ Result<std::shared_ptr<Transaction>> RestCatalog::StageCreateTable(
/*stage_create=*/true, *contextual_session));
auto table_config = std::move(result.config);
auto storage_credentials = std::move(result.storage_credentials);
ICEBERG_ASSIGN_OR_RAISE(auto table_io,
TableFileIO(context, table_config, storage_credentials));
// Before the FileIO: refreshing its credentials reuses the table session.
ICEBERG_ASSIGN_OR_RAISE(
auto table_session,
TableAuthSession(identifier, table_config, std::move(contextual_session)));
ICEBERG_ASSIGN_OR_RAISE(auto table_io, TableFileIO(context, identifier, table_config,
storage_credentials, table_session));
ICEBERG_ASSIGN_OR_RAISE(auto reporter, MakeTableReporter(identifier, table_session));
auto table_catalog = std::make_shared<TableScopedCatalog>(
shared_from_this(), context, identifier, table_config, std::move(table_session),
Expand Down Expand Up @@ -890,11 +962,12 @@ Result<std::shared_ptr<Table>> RestCatalog::MakeTableFromLoadResult(
std::shared_ptr<auth::AuthSession> contextual_session) {
auto table_config = std::move(result.config);
auto storage_credentials = std::move(result.storage_credentials);
ICEBERG_ASSIGN_OR_RAISE(auto table_io,
TableFileIO(context, table_config, storage_credentials));
// Before the FileIO: refreshing its credentials reuses the table session.
ICEBERG_ASSIGN_OR_RAISE(
auto table_session,
TableAuthSession(identifier, table_config, std::move(contextual_session)));
ICEBERG_ASSIGN_OR_RAISE(auto table_io, TableFileIO(context, identifier, table_config,
storage_credentials, table_session));
ICEBERG_ASSIGN_OR_RAISE(auto reporter, MakeTableReporter(identifier, table_session));
auto table_catalog = std::make_shared<TableScopedCatalog>(
shared_from_this(), context, identifier, table_config, table_session, table_io);
Expand Down
12 changes: 10 additions & 2 deletions src/iceberg/catalog/rest/rest_catalog.h
Original file line number Diff line number Diff line change
Expand Up @@ -85,9 +85,17 @@ class ICEBERG_REST_EXPORT RestCatalog final
std::shared_ptr<auth::AuthSession> contextual_session);

Result<std::shared_ptr<FileIO>> TableFileIO(
const SessionContext& context,
const SessionContext& context, const TableIdentifier& identifier,
const std::unordered_map<std::string, std::string>& table_config,
const std::vector<StorageCredential>& storage_credentials) const;
const std::vector<StorageCredential>& storage_credentials,
std::shared_ptr<auth::AuthSession> table_session) const;

/// \brief Build a provider for this table's vended credentials.
///
/// Returns nullptr when the catalog does not serve LoadCredentials.
std::shared_ptr<StorageCredentialProvider> MakeStorageCredentialProvider(
const TableIdentifier& identifier,
std::shared_ptr<auth::AuthSession> table_session) const;

Result<std::vector<Namespace>> ListNamespaces(const Namespace& ns,
auth::AuthSession& session) const;
Expand Down
16 changes: 14 additions & 2 deletions src/iceberg/catalog/rest/rest_file_io.cc
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,13 @@

#include <string>
#include <unordered_map>
#include <utility>
#include <vector>

#include "iceberg/catalog/rest/types.h"
#include "iceberg/file_io.h"
#include "iceberg/file_io_registry.h"
#include "iceberg/logging/log_macros.h"
#include "iceberg/resolving_file_io.h"
#include "iceberg/util/macros.h"

Expand Down Expand Up @@ -57,15 +59,25 @@ Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(const RestCatalogProperties& c
Result<std::unique_ptr<FileIO>> MakeTableFileIO(
const std::unordered_map<std::string, std::string>& catalog_config,
const std::unordered_map<std::string, std::string>& table_config,
const std::vector<StorageCredential>& storage_credentials) {
const std::vector<StorageCredential>& storage_credentials,
std::shared_ptr<StorageCredentialProvider> provider) {
const auto default_properties = MergeFileIOProperties(catalog_config, table_config);
ICEBERG_ASSIGN_OR_RAISE(
auto io, MakeCatalogFileIO(RestCatalogProperties::FromMap(default_properties)));

if (storage_credentials.empty()) {
return io;
} else if (auto* credentialed = io->AsSupportsStorageCredentials()) {
ICEBERG_RETURN_UNEXPECTED(credentialed->SetStorageCredentials(storage_credentials));
// First, so the FileIO never briefly holds credentials it cannot replace.
auto status = credentialed->InitializeStorageCredentials(storage_credentials,
std::move(provider));
if (!status && !storage_credentials.empty() &&
status.error().kind == ErrorKind::kNotSupported) {
ICEBERG_LOG_WARN("Configured FileIO cannot refresh vended storage credentials: {}",
status.error().message);
status = credentialed->SetStorageCredentials(storage_credentials);
}
ICEBERG_RETURN_UNEXPECTED(status);
} else {
return NotSupported("Configured FileIO does not support vended storage credentials");
}
Expand Down
3 changes: 2 additions & 1 deletion src/iceberg/catalog/rest/rest_file_io.h
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ ICEBERG_REST_EXPORT Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(
ICEBERG_REST_EXPORT Result<std::unique_ptr<FileIO>> MakeTableFileIO(
const std::unordered_map<std::string, std::string>& catalog_config,
const std::unordered_map<std::string, std::string>& table_config,
const std::vector<StorageCredential>& storage_credentials);
const std::vector<StorageCredential>& storage_credentials,
std::shared_ptr<StorageCredentialProvider> provider = nullptr);

} // namespace iceberg::rest
15 changes: 15 additions & 0 deletions src/iceberg/catalog/rest/types.h
Original file line number Diff line number Diff line change
Expand Up @@ -209,6 +209,21 @@ using CreateTableResponse = LoadTableResult;
/// \brief Alias of LoadTableResult used as the body of LoadTableResponse
using LoadTableResponse = LoadTableResult;

/// \brief Response body of the LoadCredentials API.
struct ICEBERG_REST_EXPORT LoadCredentialsResponse {
std::vector<StorageCredential> storage_credentials;

/// \brief Validates the LoadCredentialsResponse.
Status Validate() const {
for (const auto& credential : storage_credentials) {
ICEBERG_RETURN_UNEXPECTED(credential.Validate());
}
return {};
}

bool operator==(const LoadCredentialsResponse& other) const = default;
};

/// \brief Response body for listing namespaces.
struct ICEBERG_REST_EXPORT ListNamespacesResponse {
PageToken next_page_token;
Expand Down
13 changes: 13 additions & 0 deletions src/iceberg/file_io.h
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
#include <span>
#include <string>
#include <string_view>
#include <utility>
#include <vector>

#include "iceberg/iceberg_export.h"
Expand Down Expand Up @@ -198,6 +199,18 @@ class ICEBERG_EXPORT SupportsStorageCredentials {
/// By value because a concurrent install may replace them. An implementation
/// that delegates may report what was installed on it.
virtual std::vector<StorageCredential> credentials() const = 0;

/// \brief Install initial credentials and their refresh source before first use.
///
/// The default implementation supports static credentials only.
virtual Status InitializeStorageCredentials(
const std::vector<StorageCredential>& storage_credentials,
std::shared_ptr<StorageCredentialProvider> provider) {
if (provider) {
return NotSupported("Credential refresh is not supported");
}
return SetStorageCredentials(storage_credentials);
}
};

} // namespace iceberg
Loading
Loading