Skip to content

Commit fb8da85

Browse files
Sandeep GottimukkalaSandeep Gottimukkala
authored andcommitted
Address review comments and add tests for scan plan endpoints
- Cancel server-side plan when ResolveScanTasks fails partway through - Propagate use_snapshot_schema from scan context to PlanTableScanRequest: true for UseSnapshot/AsOfTime/tag refs and incremental scans, false for branch refs and default scans - Gate RestTable creation on effective scan-planning-mode config (table config overrides client config, default is client); error if server mode is requested but endpoint is not advertised - Add ScanPlanningMode enum and ScanPlanningModeFrom() parser to RestCatalogProperties - Make HttpClient methods virtual and add HttpResponse::MakeForTesting() to support unit test mocking - Add tests: use_snapshot_schema in table_scan_test, ScanPlanningModeFrom parsing in catalog_properties_test, and RestTableScan HTTP flow tests in rest_table_scan_test
1 parent 7a3bd35 commit fb8da85

13 files changed

Lines changed: 624 additions & 24 deletions

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

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
#include "iceberg/catalog/rest/catalog_properties.h"
2121

2222
#include <algorithm>
23+
#include <optional>
2324
#include <string>
2425
#include <string_view>
2526

@@ -61,4 +62,14 @@ Result<SnapshotMode> RestCatalogProperties::SnapshotLoadingMode() const {
6162
}
6263
}
6364

65+
Result<std::optional<ScanPlanningMode>> RestCatalogProperties::ScanPlanningModeFrom(
66+
const std::unordered_map<std::string, std::string>& config) {
67+
auto it = config.find(kScanPlanningMode.key());
68+
if (it == config.end()) return std::nullopt;
69+
std::string lower = StringUtils::ToLower(it->second);
70+
if (lower == "client") return ScanPlanningMode::kClient;
71+
if (lower == "server") return ScanPlanningMode::kServer;
72+
return InvalidArgument("Invalid scan planning mode: '{}'.", it->second);
73+
}
74+
6475
} // namespace iceberg::rest

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

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,9 @@ namespace iceberg::rest {
3434
/// \brief Snapshot loading mode for REST catalog.
3535
enum class SnapshotMode : uint8_t { kAll, kRefs };
3636

37+
/// \brief Scan planning mode for REST catalog.
38+
enum class ScanPlanningMode : uint8_t { kClient, kServer };
39+
3740
/// \brief Configuration class for a REST Catalog.
3841
class ICEBERG_REST_EXPORT RestCatalogProperties
3942
: public ConfigBase<RestCatalogProperties> {
@@ -55,6 +58,8 @@ class ICEBERG_REST_EXPORT RestCatalogProperties
5558
inline static Entry<std::string> kNamespaceSeparator{"namespace-separator", "%1F"};
5659
/// \brief The snapshot loading mode (ALL or REFS).
5760
inline static Entry<std::string> kSnapshotLoadingMode{"snapshot-loading-mode", "ALL"};
61+
/// \brief The scan planning mode (client or server).
62+
inline static Entry<std::string> kScanPlanningMode{"scan-planning-mode", "client"};
5863
/// \brief The prefix for HTTP headers.
5964
inline static constexpr std::string_view kHeaderPrefix = "header.";
6065

@@ -77,6 +82,11 @@ class ICEBERG_REST_EXPORT RestCatalogProperties
7782
/// "REFS", or an error if the value is invalid. Parsing is
7883
/// case-insensitive to match Java behavior.
7984
Result<SnapshotMode> SnapshotLoadingMode() const;
85+
86+
/// \brief Get the scan planning mode from the given config map, returning
87+
/// std::nullopt if the key is absent.
88+
static Result<std::optional<ScanPlanningMode>> ScanPlanningModeFrom(
89+
const std::unordered_map<std::string, std::string>& config);
8090
};
8191

8292
} // namespace iceberg::rest

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,15 @@ std::unordered_map<std::string, std::string> HttpResponse::headers() const {
6262
return impl_->headers();
6363
}
6464

65+
HttpResponse HttpResponse::MakeForTesting(int32_t status_code, std::string body) {
66+
cpr::Response cpr_response;
67+
cpr_response.status_code = status_code;
68+
cpr_response.text = std::move(body);
69+
HttpResponse response;
70+
response.impl_ = std::make_unique<HttpResponse::Impl>(std::move(cpr_response));
71+
return response;
72+
}
73+
6574
namespace {
6675

6776
/// \brief Default error type for unparseable REST responses.

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

Lines changed: 23 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,9 @@ class ICEBERG_REST_EXPORT HttpResponse {
6161
/// \brief Get the headers of the response as a map.
6262
std::unordered_map<std::string, std::string> headers() const;
6363

64+
/// \brief Create a response for use in unit tests.
65+
static HttpResponse MakeForTesting(int32_t status_code, std::string body);
66+
6467
private:
6568
friend class HttpClient;
6669
class Impl;
@@ -71,44 +74,45 @@ class ICEBERG_REST_EXPORT HttpResponse {
7174
class ICEBERG_REST_EXPORT HttpClient {
7275
public:
7376
explicit HttpClient(std::unordered_map<std::string, std::string> default_headers = {});
74-
~HttpClient();
77+
virtual ~HttpClient();
7578

7679
HttpClient(const HttpClient&) = delete;
7780
HttpClient& operator=(const HttpClient&) = delete;
7881
HttpClient(HttpClient&&) = delete;
7982
HttpClient& operator=(HttpClient&&) = delete;
8083

8184
/// \brief Sends a GET request.
82-
Result<HttpResponse> Get(const std::string& path,
83-
const std::unordered_map<std::string, std::string>& params,
84-
const std::unordered_map<std::string, std::string>& headers,
85-
const ErrorHandler& error_handler, auth::AuthSession& session);
85+
virtual Result<HttpResponse> Get(
86+
const std::string& path,
87+
const std::unordered_map<std::string, std::string>& params,
88+
const std::unordered_map<std::string, std::string>& headers,
89+
const ErrorHandler& error_handler, auth::AuthSession& session);
8690

8791
/// \brief Sends a POST request.
88-
Result<HttpResponse> Post(const std::string& path, const std::string& body,
89-
const std::unordered_map<std::string, std::string>& headers,
90-
const ErrorHandler& error_handler,
91-
auth::AuthSession& session);
92+
virtual Result<HttpResponse> Post(const std::string& path, const std::string& body,
93+
const std::unordered_map<std::string, std::string>& headers,
94+
const ErrorHandler& error_handler,
95+
auth::AuthSession& session);
9296

9397
/// \brief Sends a POST request with form data.
94-
Result<HttpResponse> PostForm(
98+
virtual Result<HttpResponse> PostForm(
9599
const std::string& path,
96100
const std::unordered_map<std::string, std::string>& form_data,
97101
const std::unordered_map<std::string, std::string>& headers,
98102
const ErrorHandler& error_handler, auth::AuthSession& session);
99103

100104
/// \brief Sends a HEAD request.
101-
Result<HttpResponse> Head(const std::string& path,
102-
const std::unordered_map<std::string, std::string>& headers,
103-
const ErrorHandler& error_handler,
104-
auth::AuthSession& session);
105+
virtual Result<HttpResponse> Head(
106+
const std::string& path,
107+
const std::unordered_map<std::string, std::string>& headers,
108+
const ErrorHandler& error_handler, auth::AuthSession& session);
105109

106110
/// \brief Sends a DELETE request.
107-
Result<HttpResponse> Delete(const std::string& path,
108-
const std::unordered_map<std::string, std::string>& params,
109-
const std::unordered_map<std::string, std::string>& headers,
110-
const ErrorHandler& error_handler,
111-
auth::AuthSession& session);
111+
virtual Result<HttpResponse> Delete(
112+
const std::string& path,
113+
const std::unordered_map<std::string, std::string>& params,
114+
const std::unordered_map<std::string, std::string>& headers,
115+
const ErrorHandler& error_handler, auth::AuthSession& session);
112116

113117
private:
114118
std::unordered_map<std::string, std::string> default_headers_;

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

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,8 @@
4343
#include "iceberg/catalog/rest/rest_util.h"
4444
#include "iceberg/catalog/rest/types.h"
4545
#include "iceberg/json_serde_internal.h"
46+
#include "iceberg/logging/logger.h"
47+
#include "iceberg/logging/log_level.h"
4648
#include "iceberg/partition_spec.h"
4749
#include "iceberg/result.h"
4850
#include "iceberg/schema.h"
@@ -875,7 +877,32 @@ Result<std::shared_ptr<Table>> RestCatalog::MakeTableFromLoadResult(
875877
auto table_catalog = std::make_shared<TableScopedCatalog>(
876878
shared_from_this(), context, identifier, table_config, table_session);
877879

878-
if (supported_endpoints_.contains(Endpoint::PlanTableScan())) {
880+
// Determine effective scan planning mode: table config overrides client config.
881+
ICEBERG_ASSIGN_OR_RAISE(auto client_mode,
882+
RestCatalogProperties::ScanPlanningModeFrom(config_.configs()));
883+
ICEBERG_ASSIGN_OR_RAISE(auto server_mode,
884+
RestCatalogProperties::ScanPlanningModeFrom(table_config));
885+
886+
if (client_mode.has_value() && server_mode.has_value() &&
887+
*client_mode != *server_mode) {
888+
Log(LogLevel::kWarn,
889+
"Scan planning mode mismatch for table {}: client config={}, server config={}. "
890+
"Server config will take precedence.",
891+
identifier.ToString(),
892+
*client_mode == ScanPlanningMode::kClient ? "client" : "server",
893+
*server_mode == ScanPlanningMode::kClient ? "client" : "server");
894+
}
895+
896+
ScanPlanningMode effective_mode =
897+
server_mode.value_or(client_mode.value_or(ScanPlanningMode::kClient));
898+
899+
if (effective_mode == ScanPlanningMode::kServer) {
900+
if (!supported_endpoints_.contains(Endpoint::PlanTableScan())) {
901+
return NotSupported(
902+
"Server requires server-side scan planning for table {} but does not support "
903+
"the PlanTableScan endpoint.",
904+
identifier.ToString());
905+
}
879906
RestScanContext rest_ctx{
880907
.client = client_,
881908
.paths = paths_,

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

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -100,8 +100,10 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> RestTableScan::PlanTableScan(
100100
if (context_.from_snapshot_id.has_value() && context_.to_snapshot_id.has_value()) {
101101
request.start_snapshot_id = context_.from_snapshot_id;
102102
request.end_snapshot_id = context_.to_snapshot_id;
103+
request.use_snapshot_schema = true;
103104
} else if (context_.snapshot_id.has_value()) {
104105
request.snapshot_id = context_.snapshot_id;
106+
request.use_snapshot_schema = context_.use_snapshot_schema;
105107
}
106108

107109
if (!context_.columns_to_keep_stats.empty()) {
@@ -129,8 +131,11 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> RestTableScan::PlanTableScan(
129131
plan_id = result.plan_id;
130132

131133
switch (result.plan_status) {
132-
case PlanStatus::kCompleted:
133-
return ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs);
134+
case PlanStatus::kCompleted: {
135+
auto tasks = ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs);
136+
if (!tasks.has_value()) CancelPlanning(plan_id);
137+
return tasks;
138+
}
134139
case PlanStatus::kSubmitted:
135140
return FetchPlanningResult(plan_id, specs);
136141
case PlanStatus::kFailed:
@@ -165,8 +170,11 @@ Result<std::vector<std::shared_ptr<FileScanTask>>> RestTableScan::FetchPlanningR
165170
ICEBERG_RETURN_UNEXPECTED(result.Validate());
166171

167172
switch (result.plan_status) {
168-
case PlanStatus::kCompleted:
169-
return ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs);
173+
case PlanStatus::kCompleted: {
174+
auto tasks = ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs);
175+
if (!tasks.has_value()) CancelPlanning(plan_id);
176+
return tasks;
177+
}
170178
case PlanStatus::kSubmitted: {
171179
auto elapsed_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
172180
std::chrono::steady_clock::now() - start)

‎src/iceberg/table_scan.cc‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -301,6 +301,7 @@ TableScanBuilder<ScanType>& TableScanBuilder<ScanType>::UseSnapshot(int64_t snap
301301
context_.snapshot_id.value());
302302
ICEBERG_BUILDER_ASSIGN_OR_RETURN(std::ignore, metadata_->SnapshotById(snapshot_id));
303303
context_.snapshot_id = snapshot_id;
304+
context_.use_snapshot_schema = true;
304305
return *this;
305306
}
306307

@@ -309,6 +310,7 @@ TableScanBuilder<ScanType>& TableScanBuilder<ScanType>::UseRef(const std::string
309310
if (ref == SnapshotRef::kMainBranch) {
310311
snapshot_schema_ = nullptr;
311312
context_.snapshot_id.reset();
313+
context_.use_snapshot_schema = false;
312314
return *this;
313315
}
314316

@@ -321,6 +323,7 @@ TableScanBuilder<ScanType>& TableScanBuilder<ScanType>::UseRef(const std::string
321323
const int64_t snapshot_id = iter->second->snapshot_id;
322324
ICEBERG_BUILDER_ASSIGN_OR_RETURN(std::ignore, metadata_->SnapshotById(snapshot_id));
323325
context_.snapshot_id = snapshot_id;
326+
context_.use_snapshot_schema = (iter->second->type() == SnapshotRefType::kTag);
324327

325328
return *this;
326329
}

‎src/iceberg/table_scan.h‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -229,6 +229,7 @@ struct TableScanContext {
229229
std::optional<int64_t> to_snapshot_id;
230230
std::string branch{};
231231
std::optional<int64_t> min_rows_requested;
232+
bool use_snapshot_schema{false};
232233
OptionalExecutor plan_executor;
233234

234235
// Validate the context parameters to see if they have conflicts.

‎src/iceberg/test/CMakeLists.txt‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -297,10 +297,12 @@ if(ICEBERG_BUILD_REST)
297297
add_rest_iceberg_test(rest_catalog_test
298298
SOURCES
299299
auth_manager_test.cc
300+
catalog_properties_test.cc
300301
error_handlers_test.cc
301302
endpoint_test.cc
302303
rest_file_io_test.cc
303304
rest_json_serde_test.cc
305+
rest_table_scan_test.cc
304306
rest_util_test.cc)
305307

306308
if(ICEBERG_SIGV4)
Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
20+
#include "iceberg/catalog/rest/catalog_properties.h"
21+
22+
#include <gtest/gtest.h>
23+
24+
#include "iceberg/test/matchers.h"
25+
26+
namespace iceberg::rest {
27+
28+
TEST(ScanPlanningModeTest, MissingKeyReturnsNullopt) {
29+
std::unordered_map<std::string, std::string> config;
30+
auto result = RestCatalogProperties::ScanPlanningModeFrom(config);
31+
ASSERT_THAT(result, IsOk());
32+
EXPECT_FALSE(result->has_value());
33+
}
34+
35+
TEST(ScanPlanningModeTest, ClientLowercaseReturnsKClient) {
36+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
37+
{{"scan-planning-mode", "client"}});
38+
ASSERT_THAT(result, IsOk());
39+
ASSERT_TRUE(result->has_value());
40+
EXPECT_EQ(**result, ScanPlanningMode::kClient);
41+
}
42+
43+
TEST(ScanPlanningModeTest, ServerLowercaseReturnsKServer) {
44+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
45+
{{"scan-planning-mode", "server"}});
46+
ASSERT_THAT(result, IsOk());
47+
ASSERT_TRUE(result->has_value());
48+
EXPECT_EQ(**result, ScanPlanningMode::kServer);
49+
}
50+
51+
TEST(ScanPlanningModeTest, ClientUppercaseReturnsKClient) {
52+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
53+
{{"scan-planning-mode", "CLIENT"}});
54+
ASSERT_THAT(result, IsOk());
55+
ASSERT_TRUE(result->has_value());
56+
EXPECT_EQ(**result, ScanPlanningMode::kClient);
57+
}
58+
59+
TEST(ScanPlanningModeTest, ServerUppercaseReturnsKServer) {
60+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
61+
{{"scan-planning-mode", "SERVER"}});
62+
ASSERT_THAT(result, IsOk());
63+
ASSERT_TRUE(result->has_value());
64+
EXPECT_EQ(**result, ScanPlanningMode::kServer);
65+
}
66+
67+
TEST(ScanPlanningModeTest, InvalidValueReturnsError) {
68+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
69+
{{"scan-planning-mode", "invalid"}});
70+
EXPECT_THAT(result, IsError(ErrorKind::kInvalidArgument));
71+
}
72+
73+
TEST(ScanPlanningModeTest, OtherKeysAreIgnored) {
74+
auto result = RestCatalogProperties::ScanPlanningModeFrom(
75+
{{"other-key", "server"}, {"scan-planning-mode", "client"}});
76+
ASSERT_THAT(result, IsOk());
77+
ASSERT_TRUE(result->has_value());
78+
EXPECT_EQ(**result, ScanPlanningMode::kClient);
79+
}
80+
81+
} // namespace iceberg::rest

0 commit comments

Comments
 (0)