Skip to content

Commit 668aaa4

Browse files
committed
feat(rest): vend storage credentials into per-table FileIO
The REST catalog declared the `.../credentials` endpoint path but `LoadTableResult` discarded the `storage-credentials` the server returns, so credential-vending catalogs could not access table data end to end. - Add a `StorageCredential` type and parse the `storage-credentials` list on `LoadTableResult` (load/create/register/stage responses), with JSON round-trip support. - Add `MakeTableFileIO`, which builds a table-scoped FileIO by layering the most-specific (longest-prefix) vended credential and the table's `config` on top of the catalog configuration, reusing the shared catalog FileIO when there are no overrides. This also resolves the per-table FileIO FIXME in `RestCatalog::LoadTable`. - Wire it through LoadTable / CreateTable / RegisterTable / StageCreateTable. - Unit tests for serde and for credential selection / FileIO layering.
1 parent c0c6b01 commit 668aaa4

8 files changed

Lines changed: 244 additions & 10 deletions

File tree

‎src/iceberg/catalog/rest/json_serde.cc‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,8 @@ constexpr std::string_view kSource = "source";
7171
constexpr std::string_view kDestination = "destination";
7272
constexpr std::string_view kMetadata = "metadata";
7373
constexpr std::string_view kConfig = "config";
74+
constexpr std::string_view kStorageCredentials = "storage-credentials";
75+
constexpr std::string_view kPrefix = "prefix";
7476
constexpr std::string_view kIdentifiers = "identifiers";
7577
constexpr std::string_view kOverrides = "overrides";
7678
constexpr std::string_view kDefaults = "defaults";
@@ -689,12 +691,32 @@ Result<RenameTableRequest> RenameTableRequestFromJson(const nlohmann::json& json
689691
return request;
690692
}
691693

694+
// StorageCredential (used by LoadTableResult)
695+
nlohmann::json ToJson(const StorageCredential& credential) {
696+
nlohmann::json json;
697+
json[kPrefix] = credential.prefix;
698+
SetContainerField(json, kConfig, credential.config);
699+
return json;
700+
}
701+
702+
Result<StorageCredential> StorageCredentialFromJson(const nlohmann::json& json) {
703+
StorageCredential credential;
704+
ICEBERG_ASSIGN_OR_RAISE(credential.prefix, GetJsonValue<std::string>(json, kPrefix));
705+
ICEBERG_ASSIGN_OR_RAISE(
706+
credential.config,
707+
GetJsonValueOrDefault<decltype(credential.config)>(json, kConfig));
708+
return credential;
709+
}
710+
692711
// LoadTableResult (used by CreateTableResponse, LoadTableResponse)
693712
nlohmann::json ToJson(const LoadTableResult& result) {
694713
nlohmann::json json;
695714
SetOptionalStringField(json, kMetadataLocation, result.metadata_location);
696715
json[kMetadata] = ToJson(*result.metadata);
697716
SetContainerField(json, kConfig, result.config);
717+
for (const auto& credential : result.storage_credentials) {
718+
json[kStorageCredentials].emplace_back(ToJson(credential));
719+
}
698720
return json;
699721
}
700722

@@ -707,6 +729,9 @@ Result<LoadTableResult> LoadTableResultFromJson(const nlohmann::json& json) {
707729
ICEBERG_ASSIGN_OR_RAISE(result.metadata, TableMetadataFromJson(metadata_json));
708730
ICEBERG_ASSIGN_OR_RAISE(result.config,
709731
GetJsonValueOrDefault<decltype(result.config)>(json, kConfig));
732+
ICEBERG_ASSIGN_OR_RAISE(result.storage_credentials,
733+
FromJsonList<StorageCredential>(json, kStorageCredentials,
734+
StorageCredentialFromJson));
710735
ICEBERG_RETURN_UNEXPECTED(result.Validate());
711736
return result;
712737
}

‎src/iceberg/catalog/rest/rest_catalog.cc‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -366,8 +366,12 @@ Result<std::shared_ptr<Table>> RestCatalog::CreateTable(
366366
ICEBERG_ASSIGN_OR_RAISE(auto result,
367367
CreateTableInternal(identifier, schema, spec, order, location,
368368
properties, /*stage_create=*/false));
369+
ICEBERG_ASSIGN_OR_RAISE(auto file_io,
370+
MakeTableFileIO(config_, file_io_, result.metadata->location,
371+
result.config, result.storage_credentials));
369372
return Table::Make(identifier, std::move(result.metadata),
370-
std::move(result.metadata_location), file_io_, shared_from_this());
373+
std::move(result.metadata_location), std::move(file_io),
374+
shared_from_this());
371375
}
372376

373377
Result<std::shared_ptr<Table>> RestCatalog::UpdateTable(
@@ -409,10 +413,13 @@ Result<std::shared_ptr<Transaction>> RestCatalog::StageCreateTable(
409413
ICEBERG_ASSIGN_OR_RAISE(auto result,
410414
CreateTableInternal(identifier, schema, spec, order, location,
411415
properties, /*stage_create=*/true));
416+
ICEBERG_ASSIGN_OR_RAISE(auto file_io,
417+
MakeTableFileIO(config_, file_io_, result.metadata->location,
418+
result.config, result.storage_credentials));
412419
ICEBERG_ASSIGN_OR_RAISE(auto staged_table,
413420
StagedTable::Make(identifier, std::move(result.metadata),
414-
std::move(result.metadata_location), file_io_,
415-
shared_from_this()));
421+
std::move(result.metadata_location),
422+
std::move(file_io), shared_from_this()));
416423
return Transaction::Make(std::move(staged_table), TransactionKind::kCreate);
417424
}
418425

@@ -479,9 +486,11 @@ Result<std::shared_ptr<Table>> RestCatalog::LoadTable(const TableIdentifier& ide
479486
ICEBERG_ASSIGN_OR_RAISE(const auto body, LoadTableInternal(identifier));
480487
ICEBERG_ASSIGN_OR_RAISE(auto json, FromJsonString(body));
481488
ICEBERG_ASSIGN_OR_RAISE(auto load_result, LoadTableResultFromJson(json));
482-
/// FIXME: support per-table FileIO creation
489+
ICEBERG_ASSIGN_OR_RAISE(
490+
auto file_io, MakeTableFileIO(config_, file_io_, load_result.metadata->location,
491+
load_result.config, load_result.storage_credentials));
483492
return Table::Make(identifier, std::move(load_result.metadata),
484-
std::move(load_result.metadata_location), file_io_,
493+
std::move(load_result.metadata_location), std::move(file_io),
485494
shared_from_this());
486495
}
487496

@@ -503,8 +512,11 @@ Result<std::shared_ptr<Table>> RestCatalog::RegisterTable(
503512

504513
ICEBERG_ASSIGN_OR_RAISE(auto json, FromJsonString(response.body()));
505514
ICEBERG_ASSIGN_OR_RAISE(auto load_result, LoadTableResultFromJson(json));
515+
ICEBERG_ASSIGN_OR_RAISE(
516+
auto file_io, MakeTableFileIO(config_, file_io_, load_result.metadata->location,
517+
load_result.config, load_result.storage_credentials));
506518
return Table::Make(identifier, std::move(load_result.metadata),
507-
std::move(load_result.metadata_location), file_io_,
519+
std::move(load_result.metadata_location), std::move(file_io),
508520
shared_from_this());
509521
}
510522

‎src/iceberg/catalog/rest/rest_file_io.cc‎

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,13 @@
1919

2020
#include "iceberg/catalog/rest/rest_file_io.h"
2121

22+
#include <memory>
2223
#include <string>
24+
#include <unordered_map>
25+
#include <utility>
26+
#include <vector>
2327

28+
#include "iceberg/catalog/rest/rest_util.h"
2429
#include "iceberg/file_io_registry.h"
2530
#include "iceberg/util/macros.h"
2631

@@ -92,4 +97,59 @@ Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(const RestCatalogProperties& c
9297
return FileIORegistry::Load(io_impl, config.configs());
9398
}
9499

100+
const StorageCredential* MatchStorageCredential(
101+
std::string_view location, const std::vector<StorageCredential>& credentials) {
102+
const StorageCredential* best = nullptr;
103+
for (const auto& credential : credentials) {
104+
if (!location.starts_with(credential.prefix)) {
105+
continue;
106+
}
107+
if (best == nullptr || credential.prefix.size() > best->prefix.size()) {
108+
best = &credential;
109+
}
110+
}
111+
return best;
112+
}
113+
114+
Result<std::shared_ptr<FileIO>> MakeTableFileIO(
115+
const RestCatalogProperties& catalog_config,
116+
const std::shared_ptr<FileIO>& catalog_file_io, std::string_view location,
117+
const std::unordered_map<std::string, std::string>& table_config,
118+
const std::vector<StorageCredential>& storage_credentials) {
119+
const StorageCredential* credential =
120+
MatchStorageCredential(location, storage_credentials);
121+
122+
// Without table-specific overrides, reuse the shared catalog FileIO.
123+
if (table_config.empty() && credential == nullptr) {
124+
return catalog_file_io;
125+
}
126+
127+
// Layer table config, then vended credentials, on top of the catalog properties.
128+
// Vended credentials are the most specific and therefore take precedence.
129+
static const std::unordered_map<std::string, std::string> kEmptyConfig;
130+
auto properties =
131+
MergeConfigs(catalog_config.configs(), table_config,
132+
credential != nullptr ? credential->config : kEmptyConfig);
133+
134+
std::string io_impl;
135+
if (auto it = properties.find(RestCatalogProperties::kIOImpl.key());
136+
it != properties.end()) {
137+
io_impl = it->second;
138+
}
139+
if (io_impl.empty()) {
140+
ICEBERG_ASSIGN_OR_RAISE(const auto detected_kind, DetectBuiltinFileIO(location));
141+
io_impl = std::string(BuiltinFileIOName(detected_kind));
142+
} else if (!location.empty() && IsBuiltinImpl(io_impl)) {
143+
ICEBERG_ASSIGN_OR_RAISE(const auto detected_kind, DetectBuiltinFileIO(location));
144+
if (io_impl != BuiltinFileIOName(detected_kind)) {
145+
return InvalidArgument(
146+
R"("io-impl" value '{}' is incompatible with table location '{}')", io_impl,
147+
location);
148+
}
149+
}
150+
151+
ICEBERG_ASSIGN_OR_RAISE(auto file_io, FileIORegistry::Load(io_impl, properties));
152+
return std::shared_ptr<FileIO>(std::move(file_io));
153+
}
154+
95155
} // namespace iceberg::rest

‎src/iceberg/catalog/rest/rest_file_io.h‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,12 @@
2222
#include <cstdint>
2323
#include <memory>
2424
#include <string_view>
25+
#include <unordered_map>
26+
#include <vector>
2527

2628
#include "iceberg/catalog/rest/catalog_properties.h"
2729
#include "iceberg/catalog/rest/iceberg_rest_export.h"
30+
#include "iceberg/catalog/rest/types.h"
2831
#include "iceberg/file_io.h"
2932
#include "iceberg/file_io_registry.h"
3033
#include "iceberg/result.h"
@@ -44,4 +47,27 @@ ICEBERG_REST_EXPORT std::string_view BuiltinFileIOName(BuiltinFileIOKind kind);
4447
ICEBERG_REST_EXPORT Result<std::unique_ptr<FileIO>> MakeCatalogFileIO(
4548
const RestCatalogProperties& config);
4649

50+
/// \brief Select the storage credential whose prefix is the most specific (longest)
51+
/// match for `location`, per the REST spec guidance.
52+
///
53+
/// \return A pointer into `credentials`, or nullptr if none match. The pointer is
54+
/// valid only as long as `credentials` is alive and unmodified.
55+
ICEBERG_REST_EXPORT const StorageCredential* MatchStorageCredential(
56+
std::string_view location, const std::vector<StorageCredential>& credentials);
57+
58+
/// \brief Build a FileIO for a single table from a load/create response.
59+
///
60+
/// Layers the table's vended storage credentials (the most specific prefix match for
61+
/// `location`) and table-specific `table_config` on top of the catalog configuration,
62+
/// then resolves a FileIO implementation. Vended credentials take precedence over
63+
/// `table_config`, which takes precedence over the catalog configuration.
64+
///
65+
/// When the table supplies neither vended credentials nor overriding config,
66+
/// `catalog_file_io` is returned unchanged so the shared instance is reused.
67+
ICEBERG_REST_EXPORT Result<std::shared_ptr<FileIO>> MakeTableFileIO(
68+
const RestCatalogProperties& catalog_config,
69+
const std::shared_ptr<FileIO>& catalog_file_io, std::string_view location,
70+
const std::unordered_map<std::string, std::string>& table_config,
71+
const std::vector<StorageCredential>& storage_credentials);
72+
4773
} // namespace iceberg::rest

‎src/iceberg/catalog/rest/types.cc‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,8 @@ bool CreateTableRequest::operator==(const CreateTableRequest& other) const {
8686
}
8787

8888
bool LoadTableResult::operator==(const LoadTableResult& other) const {
89-
if (metadata_location != other.metadata_location || config != other.config) {
89+
if (metadata_location != other.metadata_location || config != other.config ||
90+
storage_credentials != other.storage_credentials) {
9091
return false;
9192
}
9293

‎src/iceberg/catalog/rest/types.h‎

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -180,12 +180,29 @@ struct ICEBERG_REST_EXPORT CreateTableRequest {
180180
/// \brief An opaque token that allows clients to make use of pagination for list APIs.
181181
using PageToken = std::string;
182182

183+
/// \brief A storage credential vended by the REST catalog, scoped to a location prefix.
184+
///
185+
/// The REST catalog returns these so clients can access table data without holding
186+
/// long-lived storage credentials. When several credentials are available, clients
187+
/// should pick the one with the most specific (longest) matching prefix.
188+
struct ICEBERG_REST_EXPORT StorageCredential {
189+
/// Storage location prefix this credential applies to (e.g. "s3://bucket/db/table").
190+
std::string prefix; // required
191+
/// Credential properties to layer onto the FileIO configuration (e.g. access key,
192+
/// secret key, session token).
193+
std::unordered_map<std::string, std::string> config;
194+
195+
bool operator==(const StorageCredential&) const = default;
196+
};
197+
183198
/// \brief Result body for table create/load/register APIs.
184199
struct ICEBERG_REST_EXPORT LoadTableResult {
185200
std::string metadata_location;
186201
std::shared_ptr<TableMetadata> metadata; // required
187202
std::unordered_map<std::string, std::string> config;
188-
// TODO(Li Feiyang): Add std::shared_ptr<StorageCredential> storage_credential;
203+
/// Storage credentials for accessing the table's data. Clients should prefer these
204+
/// over any credentials embedded in `config`.
205+
std::vector<StorageCredential> storage_credentials;
189206

190207
/// \brief Validates the LoadTableResult.
191208
Status Validate() const {

‎src/iceberg/test/rest_file_io_test.cc‎

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,15 @@
1919

2020
#include "iceberg/catalog/rest/rest_file_io.h"
2121

22+
#include <memory>
23+
#include <string>
24+
#include <unordered_map>
25+
#include <vector>
26+
2227
#include <gmock/gmock.h>
2328
#include <gtest/gtest.h>
2429

30+
#include "iceberg/catalog/rest/types.h"
2531
#include "iceberg/file_io_registry.h"
2632
#include "iceberg/test/matchers.h"
2733

@@ -143,4 +149,68 @@ TEST(RestFileIOTest, MakeCatalogFileIOSkipsCheckWhenWarehouseAbsent) {
143149
ASSERT_THAT(result, IsOk());
144150
}
145151

152+
TEST(RestFileIOTest, MatchStorageCredentialPicksLongestPrefix) {
153+
std::vector<StorageCredential> credentials = {
154+
{.prefix = "s3://bucket", .config = {{"k", "broad"}}},
155+
{.prefix = "s3://bucket/db/table", .config = {{"k", "specific"}}},
156+
{.prefix = "s3://other", .config = {{"k", "other"}}},
157+
};
158+
159+
const auto* match =
160+
MatchStorageCredential("s3://bucket/db/table/data/f.parquet", credentials);
161+
ASSERT_NE(match, nullptr);
162+
EXPECT_EQ(match->prefix, "s3://bucket/db/table");
163+
164+
EXPECT_EQ(MatchStorageCredential("gs://nope/x", credentials), nullptr);
165+
}
166+
167+
TEST(RestFileIOTest, MatchStorageCredentialEmptyReturnsNull) {
168+
EXPECT_EQ(MatchStorageCredential("s3://bucket/x", {}), nullptr);
169+
}
170+
171+
TEST(RestFileIOTest, MakeTableFileIOReusesCatalogIOWhenNoOverrides) {
172+
auto catalog_io = std::make_shared<MockFileIO>();
173+
auto config = RestCatalogProperties::FromMap(
174+
{{"io-impl", std::string(FileIORegistry::kArrowS3FileIO)}});
175+
176+
auto result = MakeTableFileIO(config, catalog_io, "s3://bucket/test",
177+
/*table_config=*/{}, /*storage_credentials=*/{});
178+
ASSERT_THAT(result, IsOk());
179+
EXPECT_EQ(result.value(), catalog_io); // shared catalog instance reused
180+
}
181+
182+
TEST(RestFileIOTest, MakeTableFileIOAppliesVendedCredentials) {
183+
auto captured = std::make_shared<std::unordered_map<std::string, std::string>>();
184+
FileIORegistry::Register(
185+
std::string(FileIORegistry::kArrowS3FileIO),
186+
[captured](const std::unordered_map<std::string, std::string>& properties)
187+
-> Result<std::unique_ptr<FileIO>> {
188+
*captured = properties;
189+
return std::make_unique<MockFileIO>();
190+
});
191+
192+
auto catalog_io = std::make_shared<MockFileIO>();
193+
auto config = RestCatalogProperties::FromMap(
194+
{{"io-impl", std::string(FileIORegistry::kArrowS3FileIO)},
195+
{"s3.access-key-id", "catalog-key"}});
196+
197+
std::vector<StorageCredential> credentials = {
198+
{.prefix = "s3://bucket", .config = {{"s3.access-key-id", "broad"}}},
199+
{.prefix = "s3://bucket/test",
200+
.config = {{"s3.access-key-id", "vended"}, {"s3.session-token", "tok"}}},
201+
};
202+
203+
auto result = MakeTableFileIO(config, catalog_io, "s3://bucket/test/data/f.parquet",
204+
/*table_config=*/{{"write.parquet.compression", "zstd"}},
205+
credentials);
206+
ASSERT_THAT(result, IsOk());
207+
EXPECT_NE(result.value(), catalog_io); // a new, table-scoped FileIO
208+
209+
// The most specific vended credential wins over the catalog value.
210+
EXPECT_EQ((*captured)["s3.access-key-id"], "vended");
211+
EXPECT_EQ((*captured)["s3.session-token"], "tok");
212+
// Table-specific config is layered in as well.
213+
EXPECT_EQ((*captured)["write.parquet.compression"], "zstd");
214+
}
215+
146216
} // namespace iceberg::rest

‎src/iceberg/test/rest_json_serde_test.cc‎

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1116,7 +1116,17 @@ INSTANTIATE_TEST_SUITE_P(
11161116
.model = {.metadata_location = "s3://bucket/metadata/v1.json",
11171117
.metadata = MakeSimpleTableMetadata(),
11181118
.config = {{"warehouse", "s3://bucket/warehouse"},
1119-
{"foo", "bar"}}}}),
1119+
{"foo", "bar"}}}},
1120+
// With vended storage credentials
1121+
LoadTableResultParam{
1122+
.test_name = "WithStorageCredentials",
1123+
.expected_json_str =
1124+
R"({"metadata":{"current-schema-id":1,"current-snapshot-id":null,"default-sort-order-id":0,"default-spec-id":0,"format-version":2,"last-column-id":1,"last-partition-id":0,"last-sequence-number":0,"last-updated-ms":0,"location":"s3://bucket/test","metadata-log":[],"partition-specs":[{"fields":[],"spec-id":0}],"partition-statistics":[],"properties":{},"refs":{},"schemas":[{"fields":[{"id":1,"name":"id","required":true,"type":"int"}],"schema-id":1,"type":"struct"}],"snapshot-log":[],"snapshots":[],"sort-orders":[{"fields":[],"order-id":0}],"statistics":[],"table-uuid":"test-uuid-1234"},"storage-credentials":[{"config":{"s3.access-key-id":"AKIA","s3.secret-access-key":"secret"},"prefix":"s3://bucket/test"}]})",
1125+
.model = {.metadata = MakeSimpleTableMetadata(),
1126+
.storage_credentials = {StorageCredential{
1127+
.prefix = "s3://bucket/test",
1128+
.config = {{"s3.access-key-id", "AKIA"},
1129+
{"s3.secret-access-key", "secret"}}}}}}),
11201130
[](const ::testing::TestParamInfo<LoadTableResultParam>& info) {
11211131
return info.param.test_name;
11221132
});
@@ -1145,7 +1155,20 @@ INSTANTIATE_TEST_SUITE_P(
11451155
.json_str =
11461156
R"({"metadata":{"format-version":2,"table-uuid":"test-uuid-1234","location":"s3://bucket/test","last-sequence-number":0,"last-updated-ms":0,"last-column-id":1,"schemas":[{"type":"struct","schema-id":1,"fields":[{"id":1,"name":"id","type":"int","required":true}]}],"current-schema-id":1,"partition-specs":[{"spec-id":0,"fields":[]}],"default-spec-id":0,"last-partition-id":0,"sort-orders":[{"order-id":0,"fields":[]}],"default-sort-order-id":0,"properties":{}},"config":{"warehouse":"s3://bucket/warehouse"}})",
11471157
.expected_model = {.metadata = MakeSimpleTableMetadata(),
1148-
.config = {{"warehouse", "s3://bucket/warehouse"}}}}),
1158+
.config = {{"warehouse", "s3://bucket/warehouse"}}}},
1159+
// With multiple vended storage credentials
1160+
LoadTableResultDeserializeParam{
1161+
.test_name = "WithStorageCredentials",
1162+
.json_str =
1163+
R"({"metadata":{"format-version":2,"table-uuid":"test-uuid-1234","location":"s3://bucket/test","last-sequence-number":0,"last-updated-ms":0,"last-column-id":1,"schemas":[{"type":"struct","schema-id":1,"fields":[{"id":1,"name":"id","type":"int","required":true}]}],"current-schema-id":1,"partition-specs":[{"spec-id":0,"fields":[]}],"default-spec-id":0,"last-partition-id":0,"sort-orders":[{"order-id":0,"fields":[]}],"default-sort-order-id":0,"properties":{}},"storage-credentials":[{"prefix":"s3://bucket","config":{"s3.access-key-id":"BROAD"}},{"prefix":"s3://bucket/test","config":{"s3.access-key-id":"AKIA","s3.session-token":"tok"}}]})",
1164+
.expected_model =
1165+
{.metadata = MakeSimpleTableMetadata(),
1166+
.storage_credentials =
1167+
{StorageCredential{.prefix = "s3://bucket",
1168+
.config = {{"s3.access-key-id", "BROAD"}}},
1169+
StorageCredential{.prefix = "s3://bucket/test",
1170+
.config = {{"s3.access-key-id", "AKIA"},
1171+
{"s3.session-token", "tok"}}}}}}),
11491172
[](const ::testing::TestParamInfo<LoadTableResultDeserializeParam>& info) {
11501173
return info.param.test_name;
11511174
});

0 commit comments

Comments
 (0)