Skip to content

Commit f348bf3

Browse files
lucasfangJingsongLi
authored andcommitted
fix
1 parent 6c16150 commit f348bf3

9 files changed

Lines changed: 283 additions & 88 deletions

‎docs/source/user_guide/catalog.rst‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -404,6 +404,14 @@ the factory is called with are those of the accesses to authenticate, which is
404404
where an implementation reads its own configuration, such as the address of the
405405
token service, from.
406406

407+
Merge the credentials into your file system's options with
408+
``CredentialProvider::MergeOptionsWithCredentials`` rather than by hand, so they are
409+
shaped the same way the built-in data token file system shapes them. The default
410+
overlays the credentials key by key over the base options, mirroring the Java
411+
client. A provider that knows the file system its credentials are for overrides
412+
``MergeOptionsWithCredentials`` to normalize option aliases or clear the stale
413+
bucket-scoped variants the fresh credentials replace.
414+
407415
All factories share one identifier space, so an identifier a file system factory
408416
already takes — ``oss``, ``s3``, ``local``, ``jindo`` — would replace it; name the
409417
provider after where its credentials come from instead.

‎include/paimon/fs/credential_provider.h‎

Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,9 +46,29 @@ class PAIMON_EXPORT CredentialProvider {
4646
virtual ~CredentialProvider() = default;
4747

4848
/// Returns the credentials to sign an access with, reloading them when they are about
49-
/// to expire. The keys are file system options, e.g. "fs.oss.accessKeyId" or
50-
/// "fs.oss.securityToken", so a caller merges them over its own file system options.
49+
/// to expire. The keys are file system options, so a caller merges them over its own
50+
/// file system options with `MergeOptionsWithCredentials`.
5151
virtual Result<std::map<std::string, std::string>> GetCredentials() const = 0;
52+
53+
/// Merges the issued credentials into the file system options a delegate is built from.
54+
/// This is the canonical way to shape credentials into the options of an access: a
55+
/// caller that brings its own file system calls it so the credentials are applied the
56+
/// same way the built-in data token file system applies them.
57+
///
58+
/// The default overlays the credentials key by key over `base_options`, so they win
59+
/// wherever they overlap, mirroring the Java client and staying scheme-agnostic. A
60+
/// provider that knows the file system its credentials are for overrides this to
61+
/// normalize option aliases or clear stale bucket-scoped variants the credentials
62+
/// replace.
63+
virtual std::map<std::string, std::string> MergeOptionsWithCredentials(
64+
const std::map<std::string, std::string>& base_options,
65+
const std::map<std::string, std::string>& credentials) const {
66+
std::map<std::string, std::string> merged = base_options;
67+
for (const auto& [key, value] : credentials) {
68+
merged[key] = value;
69+
}
70+
return merged;
71+
}
5272
};
5373

5474
} // namespace paimon

‎src/paimon/rest/rest_catalog.cpp‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -483,7 +483,7 @@ Result<std::shared_ptr<FileSystem>> RestCatalog::GetTableFileSystem(
483483
// from the catalog options, which every table agrees on, so a rotation rebuilds them.
484484
const std::map<std::string, std::string>& catalog_options = api_->GetMergedOptions();
485485
std::shared_ptr<RestCredentialProvider> provider =
486-
std::make_shared<RestCredentialProvider>(api_, catalog_options, load_identifier);
486+
std::make_shared<RestCredentialProvider>(api_, load_identifier);
487487
return std::make_shared<RestTokenFileSystem>(std::move(provider), catalog_options,
488488
token_fs_cache_, fs_scheme_to_identifier_map_);
489489
}

‎src/paimon/rest/rest_credential_provider.cpp‎

Lines changed: 74 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -18,19 +18,40 @@
1818

1919
#include "paimon/rest/rest_credential_provider.h"
2020

21+
#include <algorithm>
2122
#include <chrono>
2223
#include <functional>
2324
#include <mutex>
25+
#include <set>
2426
#include <utility>
2527

2628
#include "paimon/catalog_options.h"
29+
#include "paimon/common/utils/string_utils.h"
2730

2831
namespace paimon {
2932

3033
namespace {
3134
/// Declared here rather than taken from the OSS file system, which is an optional build
3235
/// component this module must not depend on.
3336
constexpr const char kOssEndpointOption[] = "fs.oss.endpoint";
37+
constexpr const char kOssOptionPrefix[] = "fs.oss.";
38+
constexpr const char kOssBucketPrefix[] = "fs.oss.bucket.";
39+
// The two OSS names for the STS token: the backend reads securityToken first and only
40+
// consults sessionToken when it is empty, so a token under either name has to clear the
41+
// catalog value under both.
42+
constexpr const char kOssSecurityTokenSuffix[] = "securityToken";
43+
constexpr const char kOssSessionTokenSuffix[] = "sessionToken";
44+
45+
/// The suffix of a flat OSS option, i.e. what follows "fs.oss." for a key like
46+
/// "fs.oss.accessKeyId"; empty for a bucket-scoped ("fs.oss.bucket.<bucket>.<suffix>") or a
47+
/// non-OSS key.
48+
std::string FlatOssOptionSuffix(const std::string& key) {
49+
if (!StringUtils::StartsWith(key, kOssOptionPrefix) ||
50+
StringUtils::StartsWith(key, kOssBucketPrefix)) {
51+
return "";
52+
}
53+
return key.substr(std::string(kOssOptionPrefix).size());
54+
}
3455
} // namespace
3556

3657
size_t RestToken::Hash::operator()(const RestToken& rest_token) const {
@@ -42,11 +63,9 @@ size_t RestToken::Hash::operator()(const RestToken& rest_token) const {
4263
return result;
4364
}
4465

45-
RestCredentialProvider::RestCredentialProvider(
46-
const std::shared_ptr<RestApi>& api, const std::map<std::string, std::string>& catalog_options,
47-
const Identifier& identifier, Clock clock)
66+
RestCredentialProvider::RestCredentialProvider(const std::shared_ptr<RestApi>& api,
67+
const Identifier& identifier, Clock clock)
4868
: api_(api),
49-
catalog_options_(catalog_options),
5069
identifier_(identifier),
5170
clock_(std::move(clock)),
5271
logger_(Logger::GetLogger("RestCredentialProvider")) {}
@@ -60,13 +79,58 @@ bool RestCredentialProvider::ShouldRefresh() const {
6079
return token_->expires_at_millis - now_millis < RestApi::kTokenExpirationSafeTimeMillis;
6180
}
6281

63-
std::map<std::string, std::string> RestCredentialProvider::ApplyDlfEndpointOverride(
64-
const std::map<std::string, std::string>& token) const {
65-
std::map<std::string, std::string> merged = token;
82+
std::map<std::string, std::string> RestCredentialProvider::MergeOptionsWithCredentials(
83+
const std::map<std::string, std::string>& base_options,
84+
const std::map<std::string, std::string>& credentials) const {
85+
std::map<std::string, std::string> merged = base_options;
86+
for (const auto& [key, value] : credentials) {
87+
merged[key] = value;
88+
}
89+
// The OSS backend resolves a bucket-scoped option ("fs.oss.bucket.<b>.<suffix>") ahead of the
90+
// flat one and reads "fs.oss.securityToken" ahead of its "fs.oss.sessionToken" alias, so a
91+
// stale catalog value can shadow a credential the token just refreshed. Drop the stale
92+
// variants of every credential the token supplies -- including both security-token aliases
93+
// when it supplies either -- so the issued credential is the one that wins.
94+
std::set<std::string> refreshed_suffixes;
95+
bool refreshed_security_token = false;
96+
for (const auto& [key, value] : credentials) {
97+
std::string suffix = FlatOssOptionSuffix(key);
98+
if (suffix.empty()) {
99+
continue;
100+
}
101+
refreshed_suffixes.insert(suffix);
102+
refreshed_security_token = refreshed_security_token || suffix == kOssSecurityTokenSuffix ||
103+
suffix == kOssSessionTokenSuffix;
104+
}
105+
if (refreshed_security_token) {
106+
refreshed_suffixes.insert(kOssSecurityTokenSuffix);
107+
refreshed_suffixes.insert(kOssSessionTokenSuffix);
108+
}
109+
for (auto it = merged.begin(); it != merged.end();) {
110+
const std::string& key = it->first;
111+
// A credential the token itself supplies is authoritative and is never dropped, even when
112+
// it is bucket-scoped or the alias of another key the token carries.
113+
bool token_supplied = credentials.find(key) != credentials.end();
114+
bool stale_bucket_scoped = !token_supplied &&
115+
StringUtils::StartsWith(key, kOssBucketPrefix) &&
116+
std::any_of(refreshed_suffixes.begin(), refreshed_suffixes.end(),
117+
[&](const std::string& suffix) {
118+
return StringUtils::EndsWith(key, "." + suffix);
119+
});
120+
bool stale_security_token_alias =
121+
!token_supplied && refreshed_security_token &&
122+
(key == std::string(kOssOptionPrefix) + kOssSecurityTokenSuffix ||
123+
key == std::string(kOssOptionPrefix) + kOssSessionTokenSuffix);
124+
if (stale_bucket_scoped || stale_security_token_alias) {
125+
it = merged.erase(it);
126+
} else {
127+
++it;
128+
}
129+
}
66130
// The DLF OSS endpoint overrides the standard one, since the credentials are issued
67131
// for the DLF endpoint rather than for the endpoint the catalog was configured with.
68-
auto dlf_oss_endpoint = catalog_options_.find(CatalogOptions::DLF_OSS_ENDPOINT);
69-
if (dlf_oss_endpoint != catalog_options_.end() && !dlf_oss_endpoint->second.empty()) {
132+
auto dlf_oss_endpoint = base_options.find(CatalogOptions::DLF_OSS_ENDPOINT);
133+
if (dlf_oss_endpoint != base_options.end() && !dlf_oss_endpoint->second.empty()) {
70134
merged[kOssEndpointOption] = dlf_oss_endpoint->second;
71135
}
72136
return merged;
@@ -80,8 +144,7 @@ Status RestCredentialProvider::RefreshToken() const {
80144
identifier_.ToString().c_str(),
81145
static_cast<int64_t>(response.GetExpiresAtMillis()));
82146

83-
token_ =
84-
RestToken{ApplyDlfEndpointOverride(response.GetToken()), response.GetExpiresAtMillis()};
147+
token_ = RestToken{response.GetToken(), response.GetExpiresAtMillis()};
85148
return Status::OK();
86149
}
87150

‎src/paimon/rest/rest_credential_provider.h‎

Lines changed: 10 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -68,12 +68,9 @@ class RestCredentialProvider : public CredentialProvider {
6868

6969
/// @param api Client of the catalog that issues the credentials. Shared because a
7070
/// provider commonly outlives the catalog it was obtained from.
71-
/// @param catalog_options Options the credentials are merged over.
7271
/// @param identifier The table the credentials are requested for.
7372
/// @param clock Source of the current time, overridable for tests.
74-
RestCredentialProvider(const std::shared_ptr<RestApi>& api,
75-
const std::map<std::string, std::string>& catalog_options,
76-
const Identifier& identifier,
73+
RestCredentialProvider(const std::shared_ptr<RestApi>& api, const Identifier& identifier,
7774
Clock clock = std::chrono::system_clock::now);
7875

7976
~RestCredentialProvider() override = default;
@@ -89,6 +86,15 @@ class RestCredentialProvider : public CredentialProvider {
8986
/// source.
9087
Result<std::map<std::string, std::string>> GetCredentials() const override;
9188

89+
/// Merges the issued credentials over `base_options`, then corrects the OSS endpoint: the
90+
/// credentials are issued for the catalog's DLF OSS endpoint, which overrides both the
91+
/// endpoint the catalog was configured with and the one the server reported. The correction
92+
/// lives here rather than in the token so the token stays a minimal cache key that carries
93+
/// only the issued credentials.
94+
std::map<std::string, std::string> MergeOptionsWithCredentials(
95+
const std::map<std::string, std::string>& base_options,
96+
const std::map<std::string, std::string>& credentials) const override;
97+
9298
private:
9399
/// Reloads the credentials from the server. Called with the write lock of `mutex_`
94100
/// held.
@@ -97,17 +103,7 @@ class RestCredentialProvider : public CredentialProvider {
97103
/// Whether `token_` is absent or expires within the safe time.
98104
bool ShouldRefresh() const;
99105

100-
/// The issued `token` with the catalog's DLF OSS endpoint, when set, overriding the endpoint
101-
/// the server reported: the credentials are issued for the DLF endpoint, not the one the
102-
/// catalog was configured with. The token is otherwise left exactly as issued -- it is a file
103-
/// system cache key and what `ValidToken()` serves, so it carries only the credentials, never
104-
/// the whole catalog options; merging those over the token to build a delegate is the file
105-
/// system's `MergeTokenOptions`.
106-
std::map<std::string, std::string> ApplyDlfEndpointOverride(
107-
const std::map<std::string, std::string>& token) const;
108-
109106
std::shared_ptr<RestApi> api_;
110-
std::map<std::string, std::string> catalog_options_;
111107
Identifier identifier_;
112108
Clock clock_;
113109
std::shared_ptr<Logger> logger_;

‎src/paimon/rest/rest_credential_provider_test.cpp‎

Lines changed: 96 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,7 @@ class RestCredentialProviderTest : public ::testing::Test {
117117
}
118118
std::shared_ptr<RestApi> shared_api(std::move(api).value());
119119
return std::make_shared<RestCredentialProvider>(
120-
shared_api, catalog_options_, Identifier("db1", "t1"), [this] {
120+
shared_api, Identifier("db1", "t1"), [this] {
121121
return std::chrono::system_clock::time_point(
122122
std::chrono::milliseconds(now_millis_.load()));
123123
});
@@ -194,7 +194,7 @@ TEST_F(RestCredentialProviderTest, ExpiredTokenReloadsOnEveryCall) {
194194
ASSERT_EQ(2, state_->request_count.load());
195195
}
196196

197-
TEST_F(RestCredentialProviderTest, DlfEndpointOverridesTheServerEndpoint) {
197+
TEST_F(RestCredentialProviderTest, DlfEndpointOverridesTheServerEndpointOnMerge) {
198198
catalog_options_[CatalogOptions::DLF_OSS_ENDPOINT] = "dlf-endpoint";
199199
{
200200
std::lock_guard<std::mutex> lock(state_->mutex);
@@ -203,14 +203,20 @@ TEST_F(RestCredentialProviderTest, DlfEndpointOverridesTheServerEndpoint) {
203203
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
204204
ASSERT_NE(nullptr, provider);
205205

206+
// the token is the cache key and stays exactly as issued, carrying no catalog secrets and
207+
// not the endpoint correction, which is a merge concern
206208
ASSERT_OK_AND_ASSIGN(RestToken token, provider->ValidToken());
207209
ASSERT_EQ("ak-1", token.token.at("fs.oss.accessKeyId"));
208-
// the endpoint the credentials were issued for wins over the one the server reported
209-
ASSERT_EQ("dlf-endpoint", token.token.at(kOssEndpointOption));
210+
ASSERT_EQ("server-endpoint", token.token.at(kOssEndpointOption));
210211
ASSERT_EQ(kExpiresAtMillis, token.expires_at_millis);
211-
// the catalog options are not part of the token, so its secrets stay private
212212
ASSERT_EQ(0u, token.token.count(CatalogOptions::TOKEN));
213213
ASSERT_EQ(2u, token.token.size());
214+
215+
// merging shapes the credentials into the file system options: the endpoint the credentials
216+
// were issued for wins over the one the server reported
217+
Credentials merged = provider->MergeOptionsWithCredentials(catalog_options_, token.token);
218+
ASSERT_EQ("ak-1", merged.at("fs.oss.accessKeyId"));
219+
ASSERT_EQ("dlf-endpoint", merged.at(kOssEndpointOption));
214220
}
215221

216222
TEST_F(RestCredentialProviderTest, EmptyDlfOssEndpointIsNotApplied) {
@@ -222,10 +228,92 @@ TEST_F(RestCredentialProviderTest, EmptyDlfOssEndpointIsNotApplied) {
222228
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
223229
ASSERT_NE(nullptr, provider);
224230

225-
// an unset dlf endpoint leaves the endpoint the server reported alone
231+
// an unset dlf endpoint leaves the endpoint the credentials carry alone through the merge
226232
ASSERT_OK_AND_ASSIGN(RestToken token, provider->ValidToken());
227-
ASSERT_EQ("server-endpoint", token.token.at(kOssEndpointOption));
228-
ASSERT_EQ(1u, token.token.size());
233+
Credentials merged = provider->MergeOptionsWithCredentials(catalog_options_, token.token);
234+
ASSERT_EQ("server-endpoint", merged.at(kOssEndpointOption));
235+
}
236+
237+
TEST_F(RestCredentialProviderTest, IssuedCredentialsClearStaleBucketScopedCatalogVariants) {
238+
// The OSS backend resolves a bucket-scoped option ahead of the flat one, so a stale
239+
// bucket-scoped catalog value would shadow the flat credential the token just refreshed.
240+
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
241+
ASSERT_NE(nullptr, provider);
242+
243+
Credentials base = {
244+
{"fs.oss.accessKeyId", "catalog-ak"},
245+
{"fs.oss.bucket.b.accessKeyId", "catalog-bucket-ak"},
246+
{"unrelated", "kept"},
247+
};
248+
Credentials credentials = {{"fs.oss.accessKeyId", "token-ak"}};
249+
250+
Credentials merged = provider->MergeOptionsWithCredentials(base, credentials);
251+
ASSERT_EQ("token-ak", merged.at("fs.oss.accessKeyId"));
252+
ASSERT_EQ(0u, merged.count("fs.oss.bucket.b.accessKeyId"));
253+
ASSERT_EQ("kept", merged.at("unrelated"));
254+
}
255+
256+
TEST_F(RestCredentialProviderTest, SessionTokenCredentialClearsStaleSecurityTokenAliases) {
257+
// A token carrying a fresh session token must clear a stale catalog security token -- the
258+
// OSS backend reads securityToken first, so leaving it in place would pair the refreshed
259+
// key pair with the old STS token and fail authentication. Both the flat and the
260+
// bucket-scoped stale securityToken have to go.
261+
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
262+
ASSERT_NE(nullptr, provider);
263+
264+
Credentials base = {
265+
{"fs.oss.securityToken", "stale-sts"},
266+
{"fs.oss.bucket.b.securityToken", "stale-bucket-sts"},
267+
};
268+
Credentials credentials = {{"fs.oss.accessKeyId", "token-ak"},
269+
{"fs.oss.accessKeySecret", "token-sk"},
270+
{"fs.oss.sessionToken", "fresh-sts"}};
271+
272+
Credentials merged = provider->MergeOptionsWithCredentials(base, credentials);
273+
ASSERT_EQ("token-ak", merged.at("fs.oss.accessKeyId"));
274+
ASSERT_EQ("token-sk", merged.at("fs.oss.accessKeySecret"));
275+
ASSERT_EQ("fresh-sts", merged.at("fs.oss.sessionToken"));
276+
ASSERT_EQ(0u, merged.count("fs.oss.securityToken"));
277+
ASSERT_EQ(0u, merged.count("fs.oss.bucket.b.securityToken"));
278+
}
279+
280+
TEST_F(RestCredentialProviderTest, SecurityTokenCredentialClearsStaleSessionTokenAlias) {
281+
// The alias goes the other way too: a fresh securityToken clears a stale sessionToken so the
282+
// backend cannot fall back to it.
283+
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
284+
ASSERT_NE(nullptr, provider);
285+
286+
Credentials base = {{"fs.oss.sessionToken", "stale-sts"}};
287+
Credentials credentials = {{"fs.oss.securityToken", "fresh-sts"}};
288+
289+
Credentials merged = provider->MergeOptionsWithCredentials(base, credentials);
290+
ASSERT_EQ("fresh-sts", merged.at("fs.oss.securityToken"));
291+
ASSERT_EQ(0u, merged.count("fs.oss.sessionToken"));
292+
}
293+
294+
TEST_F(RestCredentialProviderTest, TokenSuppliedBucketScopedCredentialsArePreserved) {
295+
// When the token itself carries a complete bucket-scoped credential set, those values are
296+
// authoritative: the cleanup that a global credential of the same suffix would otherwise
297+
// trigger must not erase the token's own bucket-scoped keys.
298+
std::shared_ptr<RestCredentialProvider> provider = CreateProvider();
299+
ASSERT_NE(nullptr, provider);
300+
301+
Credentials base = {};
302+
Credentials credentials = {
303+
{"fs.oss.accessKeyId", "token-ak"},
304+
{"fs.oss.accessKeySecret", "token-sk"},
305+
{"fs.oss.securityToken", "token-sts"},
306+
{"fs.oss.bucket.b.accessKeyId", "token-bucket-ak"},
307+
{"fs.oss.bucket.b.accessKeySecret", "token-bucket-sk"},
308+
{"fs.oss.bucket.b.securityToken", "token-bucket-sts"},
309+
};
310+
311+
Credentials merged = provider->MergeOptionsWithCredentials(base, credentials);
312+
ASSERT_EQ("token-ak", merged.at("fs.oss.accessKeyId"));
313+
ASSERT_EQ("token-sts", merged.at("fs.oss.securityToken"));
314+
ASSERT_EQ("token-bucket-ak", merged.at("fs.oss.bucket.b.accessKeyId"));
315+
ASSERT_EQ("token-bucket-sk", merged.at("fs.oss.bucket.b.accessKeySecret"));
316+
ASSERT_EQ("token-bucket-sts", merged.at("fs.oss.bucket.b.securityToken"));
229317
}
230318

231319
TEST_F(RestCredentialProviderTest, ForbiddenIsReportedToTheCaller) {

0 commit comments

Comments
 (0)