Skip to content
Open
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
1 change: 1 addition & 0 deletions src/iceberg/catalog/rest/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ set(ICEBERG_REST_SOURCES
auth/auth_managers.cc
auth/auth_properties.cc
auth/auth_session.cc
auth/auth_session_cache.cc
auth/oauth2_util.cc
auth/sigv4_manager.cc
auth/token_refresh_scheduler.cc
Expand Down
119 changes: 97 additions & 22 deletions src/iceberg/catalog/rest/auth/auth_manager.cc
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

#include "iceberg/catalog/rest/auth/auth_manager.h"

#include <algorithm>
#include <array>
#include <chrono>
#include <optional>
Expand All @@ -28,8 +29,10 @@
#include "iceberg/catalog/rest/auth/auth_manager_internal.h"
#include "iceberg/catalog/rest/auth/auth_properties.h"
#include "iceberg/catalog/rest/auth/auth_session.h"
#include "iceberg/catalog/rest/auth/auth_session_cache_internal.h"
#include "iceberg/catalog/rest/auth/auth_session_internal.h"
#include "iceberg/catalog/rest/auth/oauth2_util.h"
#include "iceberg/catalog/rest/auth/token_refresh_scheduler.h"
#include "iceberg/catalog/session_context.h"
#include "iceberg/util/base64.h"
#include "iceberg/util/macros.h"
Expand Down Expand Up @@ -156,13 +159,13 @@ class OAuth2Manager : public AuthManager {
start_time_ = std::chrono::steady_clock::now();
ICEBERG_ASSIGN_OR_RAISE(
auth_response_, OAuth2Util::FetchToken(*init_client, *init_session, config));
// TODO(lishuxu): Match Java OAuth2Util.AuthSession.fromTokenResponse here.
// TODO(lishuxu): Build the initialization session from the full token response.
return AuthSession::MakeDefault(
OAuth2Util::AuthHeaders(auth_response_->access_token));
}

if (!config.token().empty()) {
// TODO(lishuxu): Match Java OAuth2Util.AuthSession.fromAccessToken here.
// TODO(lishuxu): Preserve configured access-token expiration in the init session.
return AuthSession::MakeDefault(OAuth2Util::AuthHeaders(config.token()));
}

Expand All @@ -176,6 +179,31 @@ class OAuth2Manager : public AuthManager {
ICEBERG_PRECHECK(shared_client != nullptr,
"OAuth2 catalog session HTTP client must not be null");
refresh_client_ = std::move(shared_client);

// Initialize catalog-level configuration
keep_refreshed_ = config.keep_refreshed();
exchange_enabled_ = config.exchange_enabled();
session_timeout_ =
std::chrono::milliseconds(config.Get(AuthProperties::kSessionTimeoutMs));

// Create session cache
ICEBERG_ASSIGN_OR_RAISE(
session_cache_,
internal::AuthSessionCache::Make(
session_timeout_, [](std::shared_ptr<AuthSession> session) {
if (auto oauth2 =
std::dynamic_pointer_cast<internal::OAuth2Session>(session)) {
oauth2->StopRefreshing();
}
}));

// Periodically reclaim idle sessions. Sweep at most once a minute and at least
// once a second, so a zero timeout still gets periodic cleanup.
auto sweep_interval = std::clamp(session_timeout_, std::chrono::milliseconds(1000),
std::chrono::milliseconds(60000));
ICEBERG_RETURN_UNEXPECTED(session_cache_->StartPeriodicSweep(
TokenRefreshScheduler::Instance(), sweep_interval));

// Reuse the token response and start time from the init phase.
if (auth_response_.has_value()) {
return internal::MakeOAuth2Session(
Expand All @@ -184,8 +212,8 @@ class OAuth2Manager : public AuthManager {
config.optional_oauth_params(), refresh_client_, start_time_);
}

// TODO(lishuxu): Honor token-refresh-enabled for catalog bearer tokens, matching
// Java. If token is provided, use it directly.
// TODO(lishuxu): Honor token-refresh-enabled for catalog bearer tokens.
// If token is provided, use it directly.
if (!config.token().empty()) {
OAuthTokenResponse token_response{
.access_token = config.token(),
Expand Down Expand Up @@ -217,21 +245,31 @@ class OAuth2Manager : public AuthManager {

Result<std::shared_ptr<AuthSession>> ContextualSession(
const SessionContext& context, std::shared_ptr<AuthSession> parent) override {
// TODO(lishuxu): Add child-session caching and refresh, matching Java
// AuthSessionCache.
// Use session_id as cache key for contextual sessions
std::string cache_key = context.session_id.empty() ? "" : "ctx:" + context.session_id;
return MaybeCreateChildSession(context.credentials, /*allow_credential=*/true,
std::move(parent));
std::move(parent), cache_key);
}

Result<std::shared_ptr<AuthSession>> TableSession(
[[maybe_unused]] const TableIdentifier& table,
const std::unordered_map<std::string, std::string>& properties,
std::shared_ptr<AuthSession> parent) override {
// Use token value as cache key for table sessions
auto token_it = properties.find(AuthProperties::kToken.key());
std::string cache_key;
if (token_it != properties.end() && !token_it->second.empty()) {
cache_key = "tbl:" + token_it->second;
}
return MaybeCreateChildSession(FilterTableSessionProperties(properties),
/*allow_credential=*/false, std::move(parent));
/*allow_credential=*/false, std::move(parent),
cache_key);
}

Status Close() override {
if (session_cache_) {
session_cache_->Close();
}
refresh_client_.reset();
return {};
}
Expand Down Expand Up @@ -267,7 +305,8 @@ class OAuth2Manager : public AuthManager {

Result<std::shared_ptr<AuthSession>> MaybeCreateChildSession(
const std::unordered_map<std::string, std::string>& credentials,
bool allow_credential, std::shared_ptr<AuthSession> parent) {
bool allow_credential, std::shared_ptr<AuthSession> parent,
const std::string& cache_key) {
auto token_it = credentials.find(AuthProperties::kToken.key());
auto credential_it = credentials.find(AuthProperties::kCredential.key());
auto typed_token = FindPreferredTypedToken(credentials);
Expand All @@ -283,42 +322,78 @@ class OAuth2Manager : public AuthManager {
ICEBERG_PRECHECK(parent_info.has_value(),
"OAuth2 child session requires OAuth2 parent metadata");

// Use cache if we have a valid cache_key
if (!cache_key.empty() && session_cache_) {
return session_cache_->Get(cache_key,
[&]() -> Result<std::shared_ptr<AuthSession>> {
return CreateChildSessionUncached(
credentials, allow_credential, *parent_info);
});
}

// No cache, create session directly
return CreateChildSessionUncached(credentials, allow_credential, *parent_info);
}

Result<std::shared_ptr<AuthSession>> CreateChildSessionUncached(
const std::unordered_map<std::string, std::string>& credentials,
bool allow_credential, const OAuth2SessionInfo& parent_info) {
auto token_it = credentials.find(AuthProperties::kToken.key());
auto credential_it = credentials.find(AuthProperties::kCredential.key());
auto typed_token = FindPreferredTypedToken(credentials);

if (token_it != credentials.end()) {
ICEBERG_ASSIGN_OR_RAISE(auto config,
ChildConfig(*parent_info, parent_info->credential));
// Token child session: don't include parent credential in config.
// Token-based sessions are not refreshed; when token exchange is enabled,
// refresh will use the token itself rather than the parent credential.
auto properties = parent_info.optional_oauth_params;
properties[AuthProperties::kScope.key()] = parent_info.scope;
properties[AuthProperties::kOAuth2ServerUri.key()] = parent_info.oauth2_server_uri;
ICEBERG_ASSIGN_OR_RAISE(auto config, AuthProperties::FromProperties(properties));
return MakeSession(AccessTokenResponse(token_it->second), config,
/*keep_refreshed=*/false);
}

if (allow_credential && credential_it != credentials.end()) {
ICEBERG_ASSIGN_OR_RAISE(auto config,
ChildConfig(*parent_info, credential_it->second));
ICEBERG_ASSIGN_OR_RAISE(auto response,
OAuth2Util::FetchToken(*refresh_client_, *parent, config));
return MakeSession(response, config, /*keep_refreshed=*/false);
ChildConfig(parent_info, credential_it->second));
// Use parent session headers for FetchToken authentication
auto temp_parent = AuthSession::MakeDefault(parent_info.headers);
ICEBERG_ASSIGN_OR_RAISE(
auto response, OAuth2Util::FetchToken(*refresh_client_, *temp_parent, config));
// Credential child sessions use keep_refreshed from catalog level
return MakeSession(response, config, keep_refreshed_);
}

std::optional<std::string> actor_token;
std::optional<std::string> actor_token_type;
if (!parent_info->token.empty()) {
actor_token = parent_info->token;
actor_token_type = parent_info->issued_token_type;
if (!parent_info.token.empty()) {
actor_token = parent_info.token;
actor_token_type = parent_info.issued_token_type;
}
// Use parent session headers for ExchangeToken authentication
auto temp_parent = AuthSession::MakeDefault(parent_info.headers);
ICEBERG_ASSIGN_OR_RAISE(
auto response,
OAuth2Util::ExchangeToken(*refresh_client_, *parent, {}, typed_token->second,
OAuth2Util::ExchangeToken(*refresh_client_, *temp_parent, {}, typed_token->second,
typed_token->first, actor_token, actor_token_type,
parent_info->scope, parent_info->oauth2_server_uri,
parent_info->optional_oauth_params));
parent_info.scope, parent_info.oauth2_server_uri,
parent_info.optional_oauth_params));
ICEBERG_ASSIGN_OR_RAISE(auto config,
ChildConfig(*parent_info, parent_info->credential));
ChildConfig(parent_info, parent_info.credential));
return MakeSession(response, config, /*keep_refreshed=*/false);
}

/// Token response and start time captured by InitSession.
std::optional<OAuthTokenResponse> auth_response_;
std::optional<std::chrono::steady_clock::time_point> start_time_;
std::shared_ptr<HttpClient> refresh_client_;

// Catalog-level configuration initialized in CatalogSession
bool keep_refreshed_ = true;
bool exchange_enabled_ = true;
std::chrono::milliseconds session_timeout_{3'600'000};
std::shared_ptr<internal::AuthSessionCache> session_cache_;
};

Result<std::unique_ptr<AuthManager>> MakeOAuth2Manager(
Expand Down
4 changes: 4 additions & 0 deletions src/iceberg/catalog/rest/auth/auth_properties.h
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,10 @@ class ICEBERG_REST_EXPORT AuthProperties : public ConfigBase<AuthProperties> {
inline static Entry<std::string> kAudience{"audience", ""};
inline static Entry<std::string> kResource{"resource", ""};

/// Session cache timeout in milliseconds. Sessions will be eligible for eviction
/// after this duration of inactivity. Default is 1 hour (3,600,000 ms).
inline static Entry<int64_t> kSessionTimeoutMs{"auth.session-timeout-ms", 3'600'000};

// ---- OAuth2 token type constants ----

inline static const std::string kAccessTokenType =
Expand Down
2 changes: 2 additions & 0 deletions src/iceberg/catalog/rest/auth/auth_session.h
Original file line number Diff line number Diff line change
Expand Up @@ -37,11 +37,13 @@ namespace iceberg::rest::auth {
/// \brief OAuth2 metadata used to derive child authentication sessions.
struct ICEBERG_REST_EXPORT OAuth2SessionInfo {
std::string token;
std::string token_type;
std::string issued_token_type;
std::string credential;
std::string scope;
std::string oauth2_server_uri;
std::unordered_map<std::string, std::string> optional_oauth_params;
std::unordered_map<std::string, std::string> headers;
};

/// \brief An authentication session that can authenticate outgoing HTTP requests.
Expand Down
Loading
Loading